Skip to main content

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}