Skip to main content

abyss_plugin_protocol/
message.rs

1//! Handshake and broker control wire payloads for plugin protocol version 1.
2
3use serde::{Deserialize, Serialize};
4use thiserror::Error;
5
6/// Broker plugin protocol versions supported by this SDK release.
7#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)]
8#[non_exhaustive]
9#[serde(try_from = "u16", into = "u16")]
10pub enum PluginProtocolVersion {
11    /// Initial handshake, live stream, and Agent event contract.
12    V1,
13}
14
15impl PluginProtocolVersion {
16    /// Returns the integer carried on the wire.
17    #[must_use]
18    pub const fn wire_value(self) -> u16 {
19        match self {
20            Self::V1 => 1,
21        }
22    }
23}
24
25impl From<PluginProtocolVersion> for u16 {
26    fn from(version: PluginProtocolVersion) -> Self {
27        version.wire_value()
28    }
29}
30
31impl TryFrom<u16> for PluginProtocolVersion {
32    type Error = UnsupportedPluginProtocolVersion;
33
34    fn try_from(version: u16) -> Result<Self, Self::Error> {
35        match version {
36            1 => Ok(Self::V1),
37            version => Err(UnsupportedPluginProtocolVersion { version }),
38        }
39    }
40}
41
42/// An unsupported broker plugin protocol version received from the wire.
43#[derive(Debug, Error)]
44#[error("unsupported broker plugin protocol version {version}")]
45pub struct UnsupportedPluginProtocolVersion {
46    version: u16,
47}
48
49impl UnsupportedPluginProtocolVersion {
50    /// Returns the unsupported integer received from the peer.
51    #[must_use]
52    pub const fn version(&self) -> u16 {
53        self.version
54    }
55}
56
57/// Initial handshake sent from a plugin process to `abyss-broker`.
58#[derive(Debug, Deserialize, Serialize)]
59#[serde(deny_unknown_fields)]
60pub struct PluginHello {
61    /// Plugin protocol version requested by the plugin.
62    pub protocol_version: PluginProtocolVersion,
63    /// Stable identity used to describe this plugin connection.
64    pub plugin_id: String,
65}
66
67impl PluginHello {
68    /// Creates a version 1 plugin handshake.
69    #[must_use]
70    // Valid plugin identifiers are non-empty runtime strings, so exposing this
71    // as const would suggest a construction mode that the wire contract rejects.
72    #[expect(
73        clippy::missing_const_for_fn,
74        reason = "valid plugin identifiers are non-empty runtime strings"
75    )]
76    pub fn new(plugin_id: String) -> Self {
77        Self {
78            protocol_version: PluginProtocolVersion::V1,
79            plugin_id,
80        }
81    }
82}
83
84/// Handshake response sent from `abyss-broker` to one plugin process.
85#[derive(Debug, Deserialize, Serialize)]
86#[serde(deny_unknown_fields)]
87pub struct BrokerHello {
88    /// Plugin protocol version confirmed by the broker.
89    pub protocol_version: PluginProtocolVersion,
90}
91
92impl BrokerHello {
93    /// Creates a version 1 broker handshake response.
94    #[must_use]
95    pub const fn v1() -> Self {
96        Self {
97            protocol_version: PluginProtocolVersion::V1,
98        }
99    }
100}
101
102/// Handshake rejection sent as the broker's first response frame.
103#[derive(Debug, Deserialize, Serialize)]
104#[serde(deny_unknown_fields)]
105pub struct BrokerError {
106    /// Stable rejection code interpreted in the handshake phase.
107    pub code: u32,
108    /// Human-readable rejection reason suitable for diagnostics.
109    pub reason: String,
110}
111
112/// Deliberate final frame sent before the broker closes an accepted session.
113#[derive(Debug, Deserialize, Serialize)]
114#[serde(deny_unknown_fields)]
115pub struct BrokerClose {
116    /// Stable close code interpreted after a successful handshake.
117    pub code: u32,
118    /// Human-readable close reason suitable for diagnostics.
119    pub reason: String,
120}
121
122/// Stable handshake rejection meanings defined by protocol version 1.
123#[derive(Clone, Copy, Debug)]
124#[non_exhaustive]
125pub enum BrokerErrorCode {
126    /// The requested protocol version is not supported.
127    UnsupportedProtocolVersion,
128    /// The first frame is malformed or contains an invalid plugin identifier.
129    InvalidHandshake,
130    /// The broker cannot accept another plugin session.
131    ResourceLimit,
132}
133
134impl BrokerErrorCode {
135    /// Returns the integer carried in a `BrokerError` frame.
136    #[must_use]
137    pub const fn wire_value(self) -> u32 {
138        match self {
139            Self::UnsupportedProtocolVersion => 1,
140            Self::InvalidHandshake => 2,
141            Self::ResourceLimit => 3,
142        }
143    }
144}
145
146/// Stable deliberate-close meanings defined by protocol version 1.
147#[derive(Clone, Copy, Debug)]
148#[non_exhaustive]
149pub enum BrokerCloseCode {
150    /// The broker is shutting down normally.
151    BrokerShutdown,
152    /// The plugin fell behind the bounded live event stream.
153    EventStreamTooSlow,
154}
155
156impl BrokerCloseCode {
157    /// Returns the integer carried in a `BrokerClose` frame.
158    #[must_use]
159    pub const fn wire_value(self) -> u32 {
160        match self {
161            Self::BrokerShutdown => 100,
162            Self::EventStreamTooSlow => 101,
163        }
164    }
165}
166
167impl BrokerError {
168    /// Creates a typed version 1 handshake rejection payload.
169    #[must_use]
170    pub fn new<T>(code: BrokerErrorCode, reason: T) -> Self
171    where
172        T: Into<String>,
173    {
174        Self {
175            code: code.wire_value(),
176            reason: reason.into(),
177        }
178    }
179}
180
181impl BrokerClose {
182    /// Creates a typed version 1 deliberate-close payload.
183    #[must_use]
184    pub fn new<T>(code: BrokerCloseCode, reason: T) -> Self
185    where
186        T: Into<String>,
187    {
188        Self {
189            code: code.wire_value(),
190            reason: reason.into(),
191        }
192    }
193}
194
195#[cfg(test)]
196mod tests {
197    use serde_json::json;
198
199    use super::{BrokerClose, BrokerError, BrokerHello, PluginHello, PluginProtocolVersion};
200
201    #[test]
202    fn unsupported_protocol_version_is_rejected() {
203        let error = serde_json::from_value::<PluginProtocolVersion>(json!(7_u16))
204            .expect_err("protocol version 7 must be rejected");
205
206        assert!(
207            error
208                .to_string()
209                .contains("unsupported broker plugin protocol version 7"),
210            "error should identify the unsupported protocol version"
211        );
212    }
213
214    #[test]
215    fn constructors_emit_direct_v1_handshake_payloads() {
216        let plugin_hello = serde_json::to_value(PluginHello::new("sample-plugin".to_owned()))
217            .expect("plugin hello should serialize");
218        let broker_hello =
219            serde_json::to_value(BrokerHello::v1()).expect("broker hello should serialize");
220
221        assert!(
222            plugin_hello.get("type").is_none(),
223            "session phase should identify PluginHello without a type discriminator"
224        );
225        assert!(
226            broker_hello.get("type").is_none(),
227            "session phase should identify BrokerHello without a type discriminator"
228        );
229        assert_eq!(
230            plugin_hello["protocol_version"], 1_u16,
231            "hello must request v1"
232        );
233        assert_eq!(
234            plugin_hello["plugin_id"], "sample-plugin",
235            "plugin id must be stable"
236        );
237        assert_eq!(
238            broker_hello["protocol_version"], 1_u16,
239            "broker hello must confirm plugin protocol v1"
240        );
241    }
242
243    #[test]
244    fn broker_control_payloads_serialize_without_wrapper_enums() {
245        let broker_error = serde_json::to_value(BrokerError {
246            code: 1,
247            reason: "unsupported protocol version".to_owned(),
248        })
249        .expect("broker error should serialize");
250        let broker_close = serde_json::to_value(BrokerClose {
251            code: 101,
252            reason: "plugin event stream is too slow".to_owned(),
253        })
254        .expect("broker close should serialize");
255
256        assert!(broker_error.get("type").is_none());
257        assert_eq!(broker_error["code"], 1_u32);
258        assert_eq!(broker_error["reason"], "unsupported protocol version");
259        assert!(broker_close.get("type").is_none());
260        assert_eq!(broker_close["code"], 101_u32);
261        assert_eq!(broker_close["reason"], "plugin event stream is too slow");
262    }
263}