abyss_plugin_protocol/
message.rs1use serde::{Deserialize, Serialize};
4use thiserror::Error;
5
6#[derive(Clone, Copy, Debug, Eq, PartialEq, Deserialize, Serialize)]
8#[non_exhaustive]
9#[serde(try_from = "u16", into = "u16")]
10pub enum PluginProtocolVersion {
11 V1,
13}
14
15impl PluginProtocolVersion {
16 #[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#[derive(Debug, Error)]
44#[error("unsupported broker plugin protocol version {version}")]
45pub struct UnsupportedPluginProtocolVersion {
46 version: u16,
47}
48
49impl UnsupportedPluginProtocolVersion {
50 #[must_use]
52 pub const fn version(&self) -> u16 {
53 self.version
54 }
55}
56
57#[derive(Debug, Deserialize, Serialize)]
59#[serde(deny_unknown_fields)]
60pub struct PluginHello {
61 pub protocol_version: PluginProtocolVersion,
63 pub plugin_id: String,
65}
66
67impl PluginHello {
68 #[must_use]
70 #[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#[derive(Debug, Deserialize, Serialize)]
86#[serde(deny_unknown_fields)]
87pub struct BrokerHello {
88 pub protocol_version: PluginProtocolVersion,
90}
91
92impl BrokerHello {
93 #[must_use]
95 pub const fn v1() -> Self {
96 Self {
97 protocol_version: PluginProtocolVersion::V1,
98 }
99 }
100}
101
102#[derive(Debug, Deserialize, Serialize)]
104#[serde(deny_unknown_fields)]
105pub struct BrokerError {
106 pub code: u32,
108 pub reason: String,
110}
111
112#[derive(Debug, Deserialize, Serialize)]
114#[serde(deny_unknown_fields)]
115pub struct BrokerClose {
116 pub code: u32,
118 pub reason: String,
120}
121
122#[derive(Clone, Copy, Debug)]
124#[non_exhaustive]
125pub enum BrokerErrorCode {
126 UnsupportedProtocolVersion,
128 InvalidHandshake,
130 ResourceLimit,
132}
133
134impl BrokerErrorCode {
135 #[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#[derive(Clone, Copy, Debug)]
148#[non_exhaustive]
149pub enum BrokerCloseCode {
150 BrokerShutdown,
152 EventStreamTooSlow,
154}
155
156impl BrokerCloseCode {
157 #[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 #[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 #[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}