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 /// Whether a client should append rather than replace.
62 append: bool,
63 },
64 /// A tool attempt is starting.
65 ToolStarted {
66 /// Model tool-call identifier.
67 call_id: String,
68 /// Stable tool name.
69 tool_name: String,
70 },
71 /// A tool attempt produced a result.
72 ToolResult {
73 /// Model tool-call identifier.
74 call_id: String,
75 /// Stable tool name.
76 tool_name: String,
77 /// Structured result.
78 result: Value,
79 /// Whether execution failed.
80 is_error: bool,
81 },
82 /// A validated model tool-call batch has been checkpointed for execution.
83 ToolCallsScheduled {
84 /// Tool names in model order.
85 tool_names: Vec<String>,
86 },
87 /// A tool batch completed and its results are available to the next model step.
88 ToolCallsCompleted {
89 /// Number of completed tool calls.
90 count: u32,
91 },
92 /// Optional diagnostic payload requested by the caller.
93 DebugData {
94 /// Diagnostic namespace.
95 scope: String,
96 /// Structured diagnostic value.
97 payload: Value,
98 },
99}
100
101/// Sequenced event persisted as part of a thread transition.
102#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
103pub struct DurableEvent {
104 /// Globally unique event identifier used for deduplication.
105 pub id: EventId,
106 /// Monotonic sequence within a thread.
107 pub sequence: u64,
108 /// Frozen routing, ownership, and run identity for this event.
109 pub context: InvocationContext,
110 /// Event payload.
111 pub event: AgentEvent,
112}
113
114/// Delivery boundary for events that have already been durably committed.
115#[async_trait]
116pub trait EventPublisher: Send + Sync {
117 /// Publishes one committed stable-boundary event with at-least-once delivery semantics.
118 ///
119 /// `snapshot` is the exact Runtime generation committed with `event`. Implementations must
120 /// never reload a newer snapshot while encoding the event.
121 async fn publish(
122 &self,
123 snapshot: &RuntimeSnapshot,
124 event: &DurableEvent,
125 ) -> Result<(), ExternalError>;
126}
127
128/// Best-effort observer for non-durable UI progress.
129#[async_trait]
130pub trait Observer: Send + Sync {
131 /// Delivers one request-scoped observation.
132 async fn observe(&self, context: &InvocationContext, observation: &AgentObservation);
133}