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        /// 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}