Skip to main content

abyss_plugin_protocol/
event.rs

1//! Language-neutral Agent event wire types shared by the broker and Rust SDK.
2//!
3//! These types are independent of the current backend ingest request. A
4//! delivery plugin may translate an event into its destination's API without
5//! making that remote API part of the broker contract.
6//! Agent event compatibility is part of the broker plugin protocol version, so
7//! an individual event does not carry a separate schema version.
8
9use std::num::NonZeroU32;
10
11use chrono::{DateTime, Utc};
12use serde::{Deserialize, Deserializer, Serialize, Serializer};
13/// One event decoded with the contract defined by the active plugin protocol.
14#[derive(Debug, Deserialize, Serialize)]
15#[serde(deny_unknown_fields)]
16pub struct AgentEvent {
17    /// Stable event identifier assigned by the event producer.
18    pub event_id: String,
19    /// Time at which the broker observed the logical event.
20    pub occurred_at: DateTime<Utc>,
21    /// Device context observed by the broker.
22    pub device: DeviceContext,
23    /// Agent product context inferred from traffic.
24    pub agent: AgentContext,
25    /// Provider or collector session identifier.
26    pub session_id: String,
27    /// One-based collector turn index within the session.
28    pub turn_index: NonZeroU32,
29    /// Provider and model context.
30    pub llm: LlmContext,
31    /// Request or response side represented by this event.
32    pub side: AgentEventSide,
33    /// Policy-retained request or response text.
34    #[serde(default, skip_serializing_if = "Option::is_none")]
35    pub text: Option<String>,
36    /// Token counters attributed to this event side.
37    pub token_usage: TokenUsage,
38    /// Normalized tool calls associated with this event side.
39    #[serde(default, skip_serializing_if = "Vec::is_empty")]
40    pub tool_calls: Vec<ToolCall>,
41    /// Normalized tool results associated with this event side.
42    #[serde(default, skip_serializing_if = "Vec::is_empty")]
43    pub tool_results: Vec<ToolResult>,
44    /// Policy-retained image attachments.
45    #[serde(default, skip_serializing_if = "Vec::is_empty")]
46    pub attachments: Vec<ImageAttachment>,
47}
48
49/// Device context attached to one Agent event.
50#[derive(Debug, Deserialize, Serialize)]
51#[serde(deny_unknown_fields)]
52pub struct DeviceContext {
53    /// Human-readable host name.
54    pub host_name: String,
55    /// Platform namespace such as `linux`, `macos`, or `windows`.
56    pub platform: String,
57    /// Operating-system version when available.
58    #[serde(default, skip_serializing_if = "Option::is_none")]
59    pub os_version: Option<String>,
60}
61
62/// Agent product context attached to one Agent event.
63#[derive(Debug, Deserialize, Serialize)]
64#[serde(deny_unknown_fields)]
65pub struct AgentContext {
66    /// Product name such as `codex` or `claude-code`.
67    pub name: String,
68    /// Agent version when it can be inferred from traffic.
69    #[serde(default, skip_serializing_if = "Option::is_none")]
70    pub version: Option<String>,
71}
72
73/// LLM provider and model context attached to one Agent event.
74#[derive(Debug, Deserialize, Serialize)]
75#[serde(deny_unknown_fields)]
76pub struct LlmContext {
77    /// Provider namespace with typed variants for built-in providers.
78    pub provider: LlmProvider,
79    /// Provider model name.
80    pub model: String,
81}
82
83/// Known and extension LLM provider namespaces.
84#[derive(Debug, Eq, PartialEq)]
85#[non_exhaustive]
86pub enum LlmProvider {
87    /// OpenAI-compatible first-party API traffic.
88    OpenAi,
89    /// Anthropic first-party API traffic.
90    Anthropic,
91    /// Provider namespace not yet represented by a built-in variant.
92    Other(String),
93}
94
95impl LlmProvider {
96    /// Creates a provider from its wire namespace.
97    #[must_use]
98    pub fn from_wire_name(provider: String) -> Self {
99        match provider.as_str() {
100            "openai" => Self::OpenAi,
101            "anthropic" => Self::Anthropic,
102            _ => Self::Other(provider),
103        }
104    }
105
106    /// Returns the provider namespace carried on the wire.
107    #[must_use]
108    pub fn wire_name(&self) -> &str {
109        match self {
110            Self::OpenAi => "openai",
111            Self::Anthropic => "anthropic",
112            Self::Other(provider) => provider,
113        }
114    }
115}
116
117impl Serialize for LlmProvider {
118    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
119    where
120        S: Serializer,
121    {
122        serializer.serialize_str(self.wire_name())
123    }
124}
125
126impl<'de> Deserialize<'de> for LlmProvider {
127    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
128    where
129        D: Deserializer<'de>,
130    {
131        String::deserialize(deserializer).map(Self::from_wire_name)
132    }
133}
134
135/// Request or response side represented by one usage event.
136#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
137#[non_exhaustive]
138#[serde(rename_all = "snake_case")]
139pub enum AgentEventSide {
140    /// Agent request content and input-side token usage.
141    Request,
142    /// Provider response content and output-side token usage.
143    Response,
144}
145
146impl AgentEventSide {
147    /// Returns the request/response name carried by backend translations.
148    #[must_use]
149    pub const fn wire_name(self) -> &'static str {
150        match self {
151            Self::Request => "request",
152            Self::Response => "response",
153        }
154    }
155}
156
157/// One normalized model-requested tool invocation.
158#[derive(Debug, Deserialize, Serialize)]
159#[serde(deny_unknown_fields)]
160pub struct ToolCall {
161    /// Broker-normalized identifier shared with the corresponding tool result.
162    pub call_id: String,
163    /// Broker-normalized tool name, such as `Bash`, `Read`, or `exec`.
164    pub name: String,
165    /// Complete provider-visible tool input normalized as text by the broker.
166    pub input: String,
167    /// Lowercase SHA-256 digest of `input`.
168    pub input_sha256: String,
169}
170
171/// One normalized result submitted back to the model after a tool call.
172#[derive(Debug, Deserialize, Serialize)]
173#[serde(deny_unknown_fields)]
174pub struct ToolResult {
175    /// Broker-normalized identifier shared with the originating tool call.
176    pub call_id: String,
177    /// Provider-visible result normalized as text by the broker.
178    pub output: String,
179    /// Lowercase SHA-256 digest of `output`.
180    pub output_sha256: String,
181}
182
183/// Normalized non-negative token counters.
184#[derive(Debug, Deserialize, Serialize)]
185#[serde(deny_unknown_fields)]
186pub struct TokenUsage {
187    /// Tokens consumed by prompt or input content.
188    pub input_tokens: u64,
189    /// Tokens produced by model output.
190    pub output_tokens: u64,
191    /// Input tokens read from a provider cache.
192    pub cache_read_tokens: u64,
193    /// Input tokens written into a provider cache.
194    pub cache_write_tokens: u64,
195    /// Tokens spent in model reasoning traces.
196    pub reasoning_tokens: u64,
197    /// Provider total or a total derived by the producer.
198    pub total_tokens: u64,
199}
200
201/// One policy-retained image attachment.
202#[derive(Debug, Deserialize, Serialize)]
203#[serde(deny_unknown_fields)]
204pub struct ImageAttachment {
205    /// Stable zero-based presentation position within the event.
206    pub position: u32,
207    /// Validated image media type.
208    pub media_type: ImageMediaType,
209    /// Decoded byte count before base64 transport encoding.
210    pub byte_size: u64,
211    /// Lowercase SHA-256 digest of the decoded image bytes.
212    pub sha256: String,
213    /// Image content when policy permits full image retention.
214    #[serde(default, skip_serializing_if = "Option::is_none")]
215    pub content_base64: Option<String>,
216}
217
218/// Image media types supported by broker plugin protocol version 1.
219#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
220#[non_exhaustive]
221pub enum ImageMediaType {
222    /// Portable Network Graphics.
223    #[serde(rename = "image/png")]
224    Png,
225    /// Joint Photographic Experts Group image.
226    #[serde(rename = "image/jpeg")]
227    Jpeg,
228    /// WebP image.
229    #[serde(rename = "image/webp")]
230    Webp,
231    /// Graphics Interchange Format image.
232    #[serde(rename = "image/gif")]
233    Gif,
234}
235
236impl ImageMediaType {
237    /// Returns the MIME type carried on the wire.
238    #[must_use]
239    pub const fn wire_name(self) -> &'static str {
240        match self {
241            Self::Png => "image/png",
242            Self::Jpeg => "image/jpeg",
243            Self::Webp => "image/webp",
244            Self::Gif => "image/gif",
245        }
246    }
247}
248
249#[cfg(test)]
250mod tests {
251    use std::num::NonZeroU32;
252
253    use super::{
254        AgentContext, AgentEvent, AgentEventSide, DeviceContext, LlmContext, LlmProvider,
255        TokenUsage, ToolCall, ToolResult,
256    };
257    use chrono::{TimeZone as _, Utc};
258
259    #[test]
260    fn agent_event_serializes_as_one_flat_typed_payload() {
261        let event = sample_event(LlmProvider::OpenAi);
262        let serialized = serde_json::to_value(event).expect("Agent event should serialize");
263
264        assert!(
265            serialized.get("schema_version").is_none(),
266            "Agent events should be versioned by the plugin protocol"
267        );
268        assert_eq!(
269            serialized["side"], "request",
270            "the event side should be a direct typed field"
271        );
272        assert!(
273            serialized.get("event_type").is_none() && serialized.get("payload").is_none(),
274            "a single event kind should not add a speculative discriminator or payload wrapper"
275        );
276        assert!(
277            serialized.get("metadata").is_none(),
278            "the public event must not expose an unstructured metadata object"
279        );
280        assert!(
281            serialized.get("text").is_none(),
282            "absent optional content should be omitted"
283        );
284    }
285
286    #[test]
287    fn tool_activity_serializes_as_named_structures() {
288        let mut event = sample_event(LlmProvider::OpenAi);
289        event.tool_calls.push(ToolCall {
290            call_id: "call-1".to_owned(),
291            name: "exec".to_owned(),
292            input: "pwd".to_owned(),
293            input_sha256: "call-hash".to_owned(),
294        });
295        event.tool_results.push(ToolResult {
296            call_id: "call-1".to_owned(),
297            output: "/workspace".to_owned(),
298            output_sha256: "result-hash".to_owned(),
299        });
300
301        let serialized = serde_json::to_value(event).expect("Agent event should serialize");
302
303        assert_eq!(serialized["tool_calls"][0]["name"], "exec");
304        assert_eq!(serialized["tool_calls"][0]["input"], "pwd");
305        assert_eq!(serialized["tool_results"][0]["call_id"], "call-1");
306        assert_eq!(serialized["tool_results"][0]["output"], "/workspace");
307    }
308
309    #[test]
310    fn agent_event_rejects_an_unstructured_metadata_object() {
311        let mut serialized = serde_json::to_value(sample_event(LlmProvider::OpenAi))
312            .expect("event should serialize");
313        serialized["metadata"] = serde_json::json!({"provider_private_field": true});
314
315        let error = serde_json::from_value::<AgentEvent>(serialized)
316            .expect_err("undeclared metadata should be rejected");
317
318        assert!(
319            error.to_string().contains("unknown field `metadata`"),
320            "error should identify metadata as outside the public contract"
321        );
322    }
323
324    #[test]
325    fn extension_provider_namespace_round_trips() {
326        let serialized = serde_json::to_value(sample_event(LlmProvider::Other(
327            "customer-private-llm".to_owned(),
328        )))
329        .expect("Agent event should serialize");
330        let decoded: AgentEvent =
331            serde_json::from_value(serialized).expect("Agent event should deserialize");
332
333        assert_eq!(
334            decoded.llm.provider,
335            LlmProvider::Other("customer-private-llm".to_owned()),
336            "unknown provider namespaces should be preserved"
337        );
338    }
339
340    fn sample_event(provider: LlmProvider) -> AgentEvent {
341        AgentEvent {
342            event_id: "evt-test".to_owned(),
343            occurred_at: Utc
344                .with_ymd_and_hms(2026, 8, 19, 10, 0, 0)
345                .single()
346                .expect("sample timestamp should be valid"),
347            device: DeviceContext {
348                host_name: "test-host".to_owned(),
349                platform: "macos".to_owned(),
350                os_version: None,
351            },
352            agent: AgentContext {
353                name: "codex".to_owned(),
354                version: None,
355            },
356            session_id: "session-1".to_owned(),
357            turn_index: NonZeroU32::new(1).expect("one should be non-zero"),
358            llm: LlmContext {
359                provider,
360                model: "gpt-test".to_owned(),
361            },
362            side: AgentEventSide::Request,
363            text: None,
364            token_usage: TokenUsage {
365                input_tokens: 10,
366                output_tokens: 0,
367                cache_read_tokens: 0,
368                cache_write_tokens: 0,
369                reasoning_tokens: 0,
370                total_tokens: 10,
371            },
372            tool_calls: Vec::new(),
373            tool_results: Vec::new(),
374            attachments: Vec::new(),
375        }
376    }
377}