Skip to main content

core_api/node/contracts/
agent_gateway.rs

1//! Internal Agent Gateway ingress. Authentication supplies tenant, scope and Node identity.
2use serde::{Deserialize, Serialize};
3
4pub const SUBMIT_TARGET: &str = "/agent/submit";
5
6#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
7#[serde(rename_all = "camelCase", deny_unknown_fields)]
8pub struct AudioReference {
9    pub uri: String,
10    pub mime_type: String,
11    #[serde(default, skip_serializing_if = "Option::is_none")]
12    pub duration_ms: Option<u64>,
13}
14
15#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
16#[serde(
17    tag = "type",
18    rename_all = "camelCase",
19    rename_all_fields = "camelCase",
20    deny_unknown_fields
21)]
22pub enum SessionSelection {
23    New { operation_id: String },
24    Existing { thread_id: String },
25}
26
27#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
28#[serde(tag = "type", rename_all = "camelCase", deny_unknown_fields)]
29pub enum Input {
30    Text { text: String },
31    Audio { audio: AudioReference },
32}
33
34#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
35#[serde(rename_all = "camelCase")]
36pub enum OutputFormat {
37    #[default]
38    Text,
39    Audio,
40}
41
42#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
43#[serde(rename_all = "camelCase")]
44pub enum EventSelection {
45    #[default]
46    Final,
47    Conversation,
48}
49
50/// Optional live conversation events, in addition to the terminal reply. These are
51/// best-effort; the canonical event's versions support duplicate/stale detection.
52#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
53#[serde(rename_all = "camelCase", deny_unknown_fields)]
54pub struct ConversationEventRequest {
55    pub request_id: String,
56    pub endpoint_id: String,
57    pub event: serde_json::Value,
58}
59
60#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
61#[serde(rename_all = "camelCase", deny_unknown_fields)]
62pub struct OutputOptions {
63    #[serde(default)]
64    pub format: OutputFormat,
65    #[serde(default)]
66    pub events: EventSelection,
67}
68
69#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
70#[serde(rename_all = "camelCase", deny_unknown_fields)]
71pub struct SubmitRequest {
72    pub request_id: String,
73    pub endpoint_id: String,
74    pub session: SessionSelection,
75    pub input: Input,
76    #[serde(default)]
77    pub output: OutputOptions,
78}
79
80#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
81#[serde(rename_all = "camelCase", deny_unknown_fields)]
82pub struct SubmitResponse {
83    pub request_id: String,
84    pub accepted: bool,
85    pub expires_at_unix_ms: i64,
86}
87
88#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
89#[serde(rename_all = "camelCase")]
90pub enum ResultStatus {
91    Completed,
92    Failed,
93    Cancelled,
94    RequiresAction,
95}
96
97#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
98#[serde(rename_all = "camelCase", deny_unknown_fields)]
99pub struct ReplyRequest {
100    pub delivery_id: String,
101    pub request_id: String,
102    pub thread_id: String,
103    pub endpoint_id: String,
104    pub status: ResultStatus,
105    pub expires_at_unix_ms: i64,
106    #[serde(default, skip_serializing_if = "Option::is_none")]
107    pub text: Option<String>,
108    #[serde(default, skip_serializing_if = "Option::is_none")]
109    pub audio: Option<AudioReference>,
110}
111
112#[cfg(test)]
113mod tests {
114    use super::*;
115
116    #[test]
117    fn live_event_envelope_preserves_canonical_event_and_rejects_extra_identity() {
118        let value = serde_json::json!({
119            "requestId": "request", "endpointId": "endpoint",
120            "event": {"type": "messageDelta", "threadId": "thread", "version": 4}
121        });
122        let event: ConversationEventRequest = serde_json::from_value(value.clone()).unwrap();
123        assert_eq!(serde_json::to_value(event).unwrap(), value);
124        let mut forged = value;
125        forged["nodeId"] = serde_json::json!("other-node");
126        assert!(serde_json::from_value::<ConversationEventRequest>(forged).is_err());
127    }
128
129    #[test]
130    fn generic_input_defaults_to_final_text_and_has_no_spoofable_caller_identity() {
131        let value = serde_json::json!({"requestId":"r","endpointId":"speaker", "session":{"type":"new","operationId":"wake"},"input":{"type":"text","text":"hello"}});
132        let request: SubmitRequest = serde_json::from_value(value.clone()).unwrap();
133        assert_eq!(
134            request.output,
135            OutputOptions {
136                format: OutputFormat::Text,
137                events: EventSelection::Final
138            }
139        );
140        let mut forged = value;
141        forged["nodeId"] = serde_json::json!("other");
142        assert!(serde_json::from_value::<SubmitRequest>(forged).is_err());
143    }
144
145    #[test]
146    fn audio_follow_up_and_failure_without_audio_round_trip() {
147        let request = SubmitRequest {
148            request_id: "r".into(),
149            endpoint_id: "speaker".into(),
150            session: SessionSelection::Existing {
151                thread_id: "thread".into(),
152            },
153            input: Input::Audio {
154                audio: AudioReference {
155                    uri: "meow-artifact://recording".into(),
156                    mime_type: "audio/ogg".into(),
157                    duration_ms: Some(1000),
158                },
159            },
160            output: OutputOptions {
161                format: OutputFormat::Audio,
162                events: EventSelection::Final,
163            },
164        };
165        let value = serde_json::to_value(&request).unwrap();
166        assert_eq!(value["session"]["threadId"], "thread");
167        assert_eq!(
168            serde_json::from_value::<SubmitRequest>(value).unwrap(),
169            request
170        );
171        let reply = ReplyRequest {
172            delivery_id: "d".into(),
173            request_id: "r".into(),
174            thread_id: "thread".into(),
175            endpoint_id: "speaker".into(),
176            status: ResultStatus::Failed,
177            expires_at_unix_ms: 1,
178            text: None,
179            audio: None,
180        };
181        assert_eq!(
182            serde_json::from_value::<ReplyRequest>(serde_json::to_value(&reply).unwrap()).unwrap(),
183            reply
184        );
185    }
186}