conversation_api/execution/event.rs
1//! Durable Agent events and their delivery boundary.
2
3use async_trait::async_trait;
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6
7use crate::execution::{
8 EventId, ExternalError, InvocationContext, MessageId, PreparedAction, RuntimeSnapshot,
9 UsageSummary,
10};
11
12/// Event emitted by the Runtime.
13#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
14#[serde(tag = "type", rename_all = "snake_case")]
15#[serde(deny_unknown_fields)]
16pub enum AgentEvent {
17 /// A run has started or resumed.
18 #[serde(deserialize_with = "crate::execution::deserialize_empty_variant")]
19 Started,
20 /// The run requires user input before it can continue.
21 InteractionRequired {
22 /// Batch identifier answered by `SubmitInteraction`.
23 batch_id: MessageId,
24 /// Prepared actions awaiting one decision each.
25 actions: Vec<PreparedAction>,
26 /// Usage accumulated before Runtime suspended the run.
27 usage: UsageSummary,
28 },
29 /// The run completed successfully.
30 Completed {
31 /// Final structured result.
32 result: Value,
33 /// Aggregated usage for the completed run.
34 usage: UsageSummary,
35 },
36 /// The run failed definitively.
37 Failed {
38 /// Catalog member from [`crate::RunErrorCode`].
39 code: String,
40 /// Safe user-facing or diagnostic message.
41 message: String,
42 /// Usage settled before the failure, cancellation, or interruption.
43 usage: UsageSummary,
44 },
45}
46
47/// Best-effort, request-scoped observation that is not part of the durable outbox.
48#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
49#[serde(tag = "type", rename_all = "snake_case")]
50#[serde(deny_unknown_fields)]
51pub enum AgentObservation {
52 /// Coarse progress for live presentation.
53 Progress {
54 /// Stable progress label.
55 status: String,
56 },
57 /// Temporary assistant text emitted before a tool batch.
58 CompanionText {
59 /// Text shown while work continues.
60 content: String,
61 },
62 /// A tool attempt is starting.
63 ToolStarted {
64 /// Model tool-call identifier.
65 call_id: String,
66 /// Stable tool name.
67 tool_name: String,
68 },
69 /// A tool attempt produced a result.
70 ToolResult {
71 /// Model tool-call identifier.
72 call_id: String,
73 /// Stable tool name.
74 tool_name: String,
75 /// Structured result.
76 result: Value,
77 /// Whether execution failed.
78 is_error: bool,
79 },
80 /// A validated model tool-call batch has been checkpointed for execution.
81 ToolCallsScheduled {
82 /// Tool names in model order.
83 tool_names: Vec<String>,
84 },
85 /// A tool batch completed and its results are available to the next model step.
86 ToolCallsCompleted {
87 /// Number of completed tool calls.
88 count: u32,
89 },
90 /// Optional diagnostic payload requested by the caller.
91 DebugData {
92 /// Diagnostic namespace.
93 scope: String,
94 /// Structured diagnostic value.
95 payload: Value,
96 },
97}
98
99/// Sequenced event persisted as part of a thread transition.
100#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
101pub struct DurableEvent {
102 /// Globally unique event identifier used for deduplication.
103 pub id: EventId,
104 /// Monotonic sequence within a thread.
105 pub sequence: u64,
106 /// Frozen routing, ownership, and run identity for this event.
107 pub context: InvocationContext,
108 /// Event payload.
109 pub event: AgentEvent,
110}
111
112/// Delivery boundary for events that have already been durably committed.
113#[async_trait]
114pub trait EventPublisher: Send + Sync {
115 /// Publishes one committed stable-boundary event with at-least-once delivery semantics.
116 ///
117 /// `snapshot` is the exact Runtime generation committed with `event`. Implementations must
118 /// never reload a newer snapshot while encoding the event.
119 async fn publish(
120 &self,
121 snapshot: &RuntimeSnapshot,
122 event: &DurableEvent,
123 ) -> Result<(), ExternalError>;
124}
125
126/// Best-effort observer for non-durable UI progress.
127#[async_trait]
128pub trait Observer: Send + Sync {
129 /// Delivers one request-scoped observation.
130 async fn observe(&self, context: &InvocationContext, observation: &AgentObservation);
131}