khive_wire_protocol/frame.rs
1//! The closed set of application frame kinds (ADR-137, "Decision" and
2//! "Protocol contract completeness").
3
4use serde::de::{Error as DeError, MapAccess, Visitor};
5use serde::ser::{Error as SerError, SerializeMap};
6use serde::{Deserialize, Deserializer, Serialize, Serializer};
7
8use crate::version::ProtocolVersion;
9
10/// A caller-generated operation id.
11///
12/// Unique across `request`, `subscribe`, and `unsubscribe` frames for the
13/// lifetime of one connection. The server echoes it on the operation's
14/// single terminal frame.
15///
16/// The codec rejects an EMPTY operation id in BOTH directions: an empty
17/// string can never be a unique caller-generated id, so it is a
18/// frame-grammar violation rather than a value the protocol has to give
19/// meaning to. Decode rejects it in `OperationId`'s [`Deserialize`] impl;
20/// encode rejects it in [`crate::codec::encode_frame_with_max`], so a
21/// locally constructed frame carrying one can never leave this crate.
22/// In-memory construction (`From<String>`, `From<&str>`) remains
23/// unrestricted; only the wire form is validated.
24#[derive(Debug, Clone, PartialEq, Eq, Hash)]
25pub struct OperationId(pub String);
26
27impl From<String> for OperationId {
28 fn from(value: String) -> Self {
29 Self(value)
30 }
31}
32
33impl From<&str> for OperationId {
34 fn from(value: &str) -> Self {
35 Self(value.to_string())
36 }
37}
38
39impl std::fmt::Display for OperationId {
40 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
41 f.write_str(&self.0)
42 }
43}
44
45impl Serialize for OperationId {
46 fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
47 serializer.serialize_str(&self.0)
48 }
49}
50
51impl<'de> Deserialize<'de> for OperationId {
52 fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
53 let value = String::deserialize(deserializer)?;
54 if value.is_empty() {
55 return Err(D::Error::custom("operation id must be a non-empty string"));
56 }
57 Ok(Self(value))
58 }
59}
60
61/// A per-topic, server-assigned, strictly increasing resumption cursor.
62pub type Cursor = u64;
63
64/// The names of every frame kind, in the order ADR-137 lists them. Used to
65/// drive the closed-set unknown-kind check in the codec, and exposed so a
66/// caller can enumerate the protocol's frame vocabulary without matching on
67/// [`Frame`] itself.
68pub const FRAME_KINDS: &[&str] = &[
69 "handshake",
70 "handshake_ack",
71 "request",
72 "response",
73 "error",
74 "cancel",
75 "subscribe",
76 "subscribe_ack",
77 "unsubscribe",
78 "unsubscribe_ack",
79 "event",
80];
81
82/// The frame kinds a client may send to the server, in the direction each
83/// is legal. The complement of this set within [`FRAME_KINDS`] —
84/// `handshake_ack`, `response`, `error`, `subscribe_ack`, `unsubscribe_ack`,
85/// and `event` — is server→client only; a server-side inbound gate rejects
86/// them as grammar violations ([`crate::handshake::HandshakeGate`]).
87pub const CLIENT_TO_SERVER_KINDS: &[&str] =
88 &["handshake", "request", "cancel", "subscribe", "unsubscribe"];
89
90/// The field set of a [`Frame::Handshake`] payload.
91///
92/// Strict within a protocol version: [`deny_unknown_fields`](serde) — a
93/// payload carrying any field this struct does not declare is rejected at
94/// decode (ADR-137's closed grammar; see the crate documentation's "Strict
95/// field rejection" section).
96#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
97#[serde(deny_unknown_fields)]
98pub struct HandshakePayload {
99 /// The protocol version the client wants to speak.
100 pub version: ProtocolVersion,
101}
102
103/// The field set of a [`Frame::HandshakeAck`] payload.
104#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
105#[serde(deny_unknown_fields)]
106pub struct HandshakeAckPayload {
107 /// The protocol version the connection now speaks.
108 pub version: ProtocolVersion,
109}
110
111/// The field set of a [`Frame::Request`] payload.
112#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
113#[serde(deny_unknown_fields)]
114pub struct RequestPayload {
115 /// Caller-generated, connection-unique operation id.
116 pub id: OperationId,
117 /// The request DSL string (ADR-016's function-call or JSON form).
118 pub ops: String,
119 /// Optional deadline in milliseconds, measured from server receipt
120 /// of this frame against the server's monotonic clock. Scopes the
121 /// entire request frame (the whole DSL batch or chain).
122 #[serde(default, skip_serializing_if = "Option::is_none")]
123 pub deadline_ms: Option<u64>,
124 /// Frame-level namespace override. Legal only on transports that
125 /// accept caller-supplied identity context; a mapped transport
126 /// (ADR-137's TCP transport) rejects any request carrying this with
127 /// [`crate::error::WireErrorCode::ContextRejected`].
128 #[serde(default, skip_serializing_if = "Option::is_none")]
129 pub namespace: Option<String>,
130 /// Frame-level actor override; see `namespace` above.
131 #[serde(default, skip_serializing_if = "Option::is_none")]
132 pub actor_id: Option<String>,
133 /// Frame-level visible-namespace-set override; see `namespace`
134 /// above.
135 #[serde(default, skip_serializing_if = "Option::is_none")]
136 pub visible_namespaces: Option<Vec<String>>,
137}
138
139/// The field set of a [`Frame::Response`] payload.
140#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
141#[serde(deny_unknown_fields)]
142pub struct ResponsePayload {
143 /// Echoes the originating `request`'s operation id.
144 pub id: OperationId,
145 /// The verb-dispatch result, exactly as ADR-016's `request` verb
146 /// surface returns it (an aggregate `{ok, tool, result}` /
147 /// `{ok, summary, ...}` payload). Opaque to this crate.
148 pub result: serde_json::Value,
149}
150
151/// The field set of a [`Frame::Error`] payload.
152#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
153#[serde(deny_unknown_fields)]
154pub struct ErrorPayload {
155 /// The operation id this error terminates, for a request-scoped
156 /// error. `None` for a connection-terminal error, which carries no
157 /// operation id (ADR-137, "Operation correlation").
158 #[serde(default, skip_serializing_if = "Option::is_none")]
159 pub id: Option<OperationId>,
160 /// The wire error code.
161 pub code: crate::error::WireErrorCode,
162 /// A human-readable detail message. Not part of the closed
163 /// contract — callers must branch on `code`, never on this string.
164 pub message: String,
165}
166
167/// The field set of a [`Frame::Cancel`] payload.
168#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
169#[serde(deny_unknown_fields)]
170pub struct CancelPayload {
171 /// The `request` operation id to cancel. A `cancel` naming a
172 /// subscribe/unsubscribe id, or an unknown or already-terminal
173 /// request id, is a no-op.
174 pub id: OperationId,
175}
176
177/// The field set of a [`Frame::Subscribe`] payload.
178#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
179#[serde(deny_unknown_fields)]
180pub struct SubscribePayload {
181 /// Caller-generated, connection-unique operation id.
182 pub id: OperationId,
183 /// The topic to subscribe to, `<domain>.<event>`.
184 pub topic: String,
185 /// Resume position. Absent starts delivery at new events only;
186 /// present replays every retained event with a cursor greater than
187 /// this value before delivering new events.
188 #[serde(default, skip_serializing_if = "Option::is_none")]
189 pub resume_cursor: Option<Cursor>,
190}
191
192/// The field set of a [`Frame::SubscribeAck`] payload.
193#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
194#[serde(deny_unknown_fields)]
195pub struct SubscribeAckPayload {
196 /// Echoes the originating `subscribe`'s operation id.
197 pub id: OperationId,
198 /// The subscribed topic.
199 pub topic: String,
200 /// The cursor position delivery begins after.
201 pub start_cursor: Cursor,
202}
203
204/// The field set of a [`Frame::Unsubscribe`] payload.
205#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
206#[serde(deny_unknown_fields)]
207pub struct UnsubscribePayload {
208 /// Caller-generated, connection-unique operation id.
209 pub id: OperationId,
210 /// The topic to unsubscribe from. Naming a topic with no active
211 /// subscription is an idempotent no-op.
212 pub topic: String,
213}
214
215/// The field set of a [`Frame::UnsubscribeAck`] payload.
216#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
217#[serde(deny_unknown_fields)]
218pub struct UnsubscribeAckPayload {
219 /// Echoes the originating `unsubscribe`'s operation id.
220 pub id: OperationId,
221 /// The unsubscribed topic.
222 pub topic: String,
223}
224
225/// The field set of a [`Frame::Event`] payload.
226#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
227#[serde(deny_unknown_fields)]
228pub struct EventPayload {
229 /// The topic this event belongs to.
230 pub topic: String,
231 /// Server-assigned, per-topic, strictly increasing resumption
232 /// cursor.
233 pub cursor: Cursor,
234 /// Server-assigned event time, RFC 3339.
235 pub occurred_at: String,
236 /// Topic-specific payload. Field-by-field shape is owned by the
237 /// per-topic catalog (ADR-137, "Implementation-phase deliverables"),
238 /// not by this crate.
239 pub payload: serde_json::Value,
240}
241
242/// One application frame.
243///
244/// Every variant corresponds to exactly one entry in [`FRAME_KINDS`] and is
245/// carried on the wire as a JSON object with a `"kind"` discriminant field
246/// holding that variant's `snake_case` name, followed by the variant's own
247/// fields flattened into the same object. This is the wire framing's
248/// internally tagged encoding; see the crate documentation for a worked
249/// example of the exact bytes.
250///
251/// The set is closed within a protocol version (ADR-137, "Decision"): a
252/// decoder that encounters a `"kind"` value outside [`FRAME_KINDS`] must
253/// reject the frame ([`crate::codec::CodecError::UnknownFrameKind`]), never
254/// skip or ignore it.
255///
256/// Serde is implemented by hand (rather than `#[serde(tag = "kind")]`) so
257/// the decode path can enforce the closed grammar ADR-137 requires: the
258/// codec checks the `"kind"` discriminant against [`FRAME_KINDS`] itself,
259/// then hands the payload to the matching kind's payload struct
260/// ([`HandshakePayload`], [`RequestPayload`], ...), every one of which
261/// carries `#[serde(deny_unknown_fields)]`. A payload with any field its
262/// kind does not declare is therefore rejected — never silently ignored —
263/// and a missing or non-string `"kind"` is rejected before any kind is
264/// matched. The visitor likewise enforces the ADR-137 id/scope rule for
265/// `error` frames and captures the unknown-code fallback's raw string
266/// ([`Frame::Error`](Frame)'s `unrecognized_code`), so a DIRECT serde
267/// decode of `Frame` — not just the codec — can never represent an
268/// inconsistent error frame. Encoding writes `"kind"` first, then the
269/// kind's fields in
270/// declaration order, skipping absent optional fields; this is the exact
271/// byte layout the golden fixtures pin.
272#[derive(Debug, Clone, PartialEq)]
273pub enum Frame {
274 /// The first application frame on every connection. Names the protocol
275 /// version the client supports.
276 Handshake {
277 /// The protocol version the client wants to speak.
278 version: ProtocolVersion,
279 },
280
281 /// The server's acceptance of a [`Frame::Handshake`], naming the
282 /// accepted protocol version.
283 HandshakeAck {
284 /// The protocol version the connection now speaks.
285 version: ProtocolVersion,
286 },
287
288 /// A caller-issued operation: a DSL batch or chain (ADR-016) to execute.
289 Request {
290 /// Caller-generated, connection-unique operation id.
291 id: OperationId,
292 /// The request DSL string (ADR-016's function-call or JSON form).
293 ops: String,
294 /// Optional deadline in milliseconds, measured from server receipt
295 /// of this frame against the server's monotonic clock. Scopes the
296 /// entire request frame (the whole DSL batch or chain).
297 deadline_ms: Option<u64>,
298 /// Frame-level namespace override. Legal only on transports that
299 /// accept caller-supplied identity context; a mapped transport
300 /// (ADR-137's TCP transport) rejects any request carrying this with
301 /// [`crate::error::WireErrorCode::ContextRejected`].
302 namespace: Option<String>,
303 /// Frame-level actor override; see `namespace` above.
304 actor_id: Option<String>,
305 /// Frame-level visible-namespace-set override; see `namespace`
306 /// above.
307 visible_namespaces: Option<Vec<String>>,
308 },
309
310 /// The successful terminal frame for a `request`.
311 Response {
312 /// Echoes the originating `request`'s operation id.
313 id: OperationId,
314 /// The verb-dispatch result, exactly as ADR-016's `request` verb
315 /// surface returns it (an aggregate `{ok, tool, result}` /
316 /// `{ok, summary, ...}` payload). Opaque to this crate.
317 result: serde_json::Value,
318 },
319
320 /// A wire-level failure terminal frame.
321 Error {
322 /// The operation id this error terminates, for a request-scoped
323 /// error. `None` for a connection-terminal error, which carries no
324 /// operation id (ADR-137, "Operation correlation").
325 id: Option<OperationId>,
326 /// The wire error code.
327 code: crate::error::WireErrorCode,
328 /// A human-readable detail message. Not part of the closed
329 /// contract — callers must branch on `code`, never on this string.
330 message: String,
331 /// The raw code string when the wire carried a code OUTSIDE the
332 /// closed set ([`crate::error::WIRE_ERROR_CODES`]) and serde's
333 /// `#[serde(other)]` fallback mapped it to
334 /// [`crate::error::WireErrorCode::Internal`]; `None` for every
335 /// recognized code, and for every frame this crate produced by
336 /// encoding or by in-memory construction. Diagnostic only: it is
337 /// never serialized, and only DECODE paths fill it — this type's
338 /// serde visitor, whether driven by the codec's `decode_payload`
339 /// or by a direct serde decode. A frame carrying it is a decoded
340 /// fallback and cannot be re-encoded by this relay: the encode path
341 /// rejects it ([`crate::codec::CodecError::FallbackFrameNotEncodable`])
342 /// rather than emit `internal` and silently discard the newer code.
343 /// The connection remains healthy; only this relay attempt failed.
344 /// If it has no operation id, a consumer should surface it as a
345 /// connection-level diagnostic rather than guess which request to
346 /// fail.
347 unrecognized_code: Option<String>,
348 },
349
350 /// Asks the server to terminate an in-flight `request`.
351 Cancel {
352 /// The `request` operation id to cancel. A `cancel` naming a
353 /// subscribe/unsubscribe id, or an unknown or already-terminal
354 /// request id, is a no-op.
355 id: OperationId,
356 },
357
358 /// Opens delivery for one topic on the connection.
359 Subscribe {
360 /// Caller-generated, connection-unique operation id.
361 id: OperationId,
362 /// The topic to subscribe to, `<domain>.<event>`.
363 topic: String,
364 /// Resume position. Absent starts delivery at new events only;
365 /// present replays every retained event with a cursor greater than
366 /// this value before delivering new events.
367 resume_cursor: Option<Cursor>,
368 },
369
370 /// The successful terminal frame for a `subscribe`.
371 SubscribeAck {
372 /// Echoes the originating `subscribe`'s operation id.
373 id: OperationId,
374 /// The subscribed topic.
375 topic: String,
376 /// The cursor position delivery begins after.
377 start_cursor: Cursor,
378 },
379
380 /// Ends delivery for one topic on the connection.
381 Unsubscribe {
382 /// Caller-generated, connection-unique operation id.
383 id: OperationId,
384 /// The topic to unsubscribe from. Naming a topic with no active
385 /// subscription is an idempotent no-op.
386 topic: String,
387 },
388
389 /// The terminal frame for an `unsubscribe`.
390 UnsubscribeAck {
391 /// Echoes the originating `unsubscribe`'s operation id.
392 id: OperationId,
393 /// The unsubscribed topic.
394 topic: String,
395 },
396
397 /// A server-pushed state-change delivery for a subscribed topic.
398 ///
399 /// Carries no operation id; correlated by topic and ordered by cursor
400 /// instead.
401 Event {
402 /// The topic this event belongs to.
403 topic: String,
404 /// Server-assigned, per-topic, strictly increasing resumption
405 /// cursor.
406 cursor: Cursor,
407 /// Server-assigned event time, RFC 3339.
408 occurred_at: String,
409 /// Topic-specific payload. Field-by-field shape is owned by the
410 /// per-topic catalog (ADR-137, "Implementation-phase deliverables"),
411 /// not by this crate.
412 payload: serde_json::Value,
413 },
414}
415
416impl Frame {
417 /// The `snake_case` frame-kind name of this frame, matching its `"kind"`
418 /// discriminant on the wire.
419 pub const fn kind(&self) -> &'static str {
420 match self {
421 Frame::Handshake { .. } => "handshake",
422 Frame::HandshakeAck { .. } => "handshake_ack",
423 Frame::Request { .. } => "request",
424 Frame::Response { .. } => "response",
425 Frame::Error { .. } => "error",
426 Frame::Cancel { .. } => "cancel",
427 Frame::Subscribe { .. } => "subscribe",
428 Frame::SubscribeAck { .. } => "subscribe_ack",
429 Frame::Unsubscribe { .. } => "unsubscribe",
430 Frame::UnsubscribeAck { .. } => "unsubscribe_ack",
431 Frame::Event { .. } => "event",
432 }
433 }
434}
435
436/// Validate the invariants shared by the codec and direct serde encoding.
437fn validate_frame_for_serialize(frame: &Frame) -> Result<(), String> {
438 use crate::error::TerminalScope;
439
440 fn check_id(kind: &str, id: &OperationId) -> Result<(), String> {
441 if id.0.is_empty() {
442 return Err(format!(
443 "frame kind {kind:?}: operation id must be a non-empty string"
444 ));
445 }
446 Ok(())
447 }
448
449 fn check_version(kind: &str, version: ProtocolVersion) -> Result<(), String> {
450 if version.get() == 0 {
451 return Err(format!(
452 "frame kind {kind:?}: protocol version 0 does not exist"
453 ));
454 }
455 Ok(())
456 }
457
458 match frame {
459 Frame::Handshake { version } => check_version("handshake", *version)?,
460 Frame::HandshakeAck { version } => check_version("handshake_ack", *version)?,
461 Frame::Request { id, .. } => check_id("request", id)?,
462 Frame::Response { id, .. } => check_id("response", id)?,
463 Frame::Cancel { id } => check_id("cancel", id)?,
464 Frame::Subscribe { id, .. } => check_id("subscribe", id)?,
465 Frame::SubscribeAck { id, .. } => check_id("subscribe_ack", id)?,
466 Frame::Unsubscribe { id, .. } => check_id("unsubscribe", id)?,
467 Frame::UnsubscribeAck { id, .. } => check_id("unsubscribe_ack", id)?,
468 Frame::Event { .. } => {}
469 Frame::Error {
470 id,
471 code,
472 unrecognized_code,
473 ..
474 } => {
475 if let Some(raw_code) = unrecognized_code {
476 return Err(format!(
477 "fallback error frame (unrecognized code {raw_code:?}) is not re-encodable: \
478 re-encoding would emit \"internal\" and discard the newer wire code"
479 ));
480 }
481 if let Some(id) = id {
482 check_id("error", id)?;
483 }
484 match (code.terminal_scope(), id) {
485 (TerminalScope::Connection, Some(id)) => {
486 return Err(format!(
487 "error frame violates the id/scope rule: connection-terminal code {code} \
488 must not carry an operation id, got {id}"
489 ));
490 }
491 (TerminalScope::Request, None) => {
492 return Err(format!(
493 "error frame violates the id/scope rule: request-terminal code {code} \
494 must echo the operation id it terminates"
495 ));
496 }
497 (TerminalScope::Connection, None) | (TerminalScope::Request, Some(_)) => {}
498 }
499 }
500 }
501 Ok(())
502}
503
504/// Serialize one frame as its wire object: `"kind"` first, then the kind's
505/// fields in declaration order, absent optional fields skipped. The byte
506/// layout is pinned by the golden fixtures in `tests/fixtures/*.hex`.
507impl Serialize for Frame {
508 fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
509 if let Err(detail) = validate_frame_for_serialize(self) {
510 return Err(S::Error::custom(detail));
511 }
512 let size_hint = match self {
513 Frame::Handshake { .. } | Frame::HandshakeAck { .. } | Frame::Cancel { .. } => 2,
514 Frame::Subscribe { resume_cursor, .. } => {
515 if resume_cursor.is_some() {
516 4
517 } else {
518 3
519 }
520 }
521 Frame::Request {
522 deadline_ms,
523 namespace,
524 actor_id,
525 visible_namespaces,
526 ..
527 } => {
528 3 + usize::from(deadline_ms.is_some())
529 + usize::from(namespace.is_some())
530 + usize::from(actor_id.is_some())
531 + usize::from(visible_namespaces.is_some())
532 }
533 Frame::Error { id, .. } => 3 + usize::from(id.is_some()),
534 Frame::Response { .. } | Frame::Unsubscribe { .. } | Frame::UnsubscribeAck { .. } => 3,
535 Frame::SubscribeAck { .. } => 4,
536 Frame::Event { .. } => 5,
537 };
538 let mut map = serializer.serialize_map(Some(size_hint))?;
539 map.serialize_entry("kind", self.kind())?;
540 match self {
541 Frame::Handshake { version } => {
542 map.serialize_entry("version", version)?;
543 }
544 Frame::HandshakeAck { version } => {
545 map.serialize_entry("version", version)?;
546 }
547 Frame::Request {
548 id,
549 ops,
550 deadline_ms,
551 namespace,
552 actor_id,
553 visible_namespaces,
554 } => {
555 map.serialize_entry("id", id)?;
556 map.serialize_entry("ops", ops)?;
557 if let Some(deadline_ms) = deadline_ms {
558 map.serialize_entry("deadline_ms", deadline_ms)?;
559 }
560 if let Some(namespace) = namespace {
561 map.serialize_entry("namespace", namespace)?;
562 }
563 if let Some(actor_id) = actor_id {
564 map.serialize_entry("actor_id", actor_id)?;
565 }
566 if let Some(visible_namespaces) = visible_namespaces {
567 map.serialize_entry("visible_namespaces", visible_namespaces)?;
568 }
569 }
570 Frame::Response { id, result } => {
571 map.serialize_entry("id", id)?;
572 map.serialize_entry("result", result)?;
573 }
574 // `unrecognized_code` is a decode-only diagnostic; validation
575 // above rejects it rather than silently dropping the wire code.
576 Frame::Error {
577 id, code, message, ..
578 } => {
579 if let Some(id) = id {
580 map.serialize_entry("id", id)?;
581 }
582 map.serialize_entry("code", code)?;
583 map.serialize_entry("message", message)?;
584 }
585 Frame::Cancel { id } => {
586 map.serialize_entry("id", id)?;
587 }
588 Frame::Subscribe {
589 id,
590 topic,
591 resume_cursor,
592 } => {
593 map.serialize_entry("id", id)?;
594 map.serialize_entry("topic", topic)?;
595 if let Some(resume_cursor) = resume_cursor {
596 map.serialize_entry("resume_cursor", resume_cursor)?;
597 }
598 }
599 Frame::SubscribeAck {
600 id,
601 topic,
602 start_cursor,
603 } => {
604 map.serialize_entry("id", id)?;
605 map.serialize_entry("topic", topic)?;
606 map.serialize_entry("start_cursor", start_cursor)?;
607 }
608 Frame::Unsubscribe { id, topic } => {
609 map.serialize_entry("id", id)?;
610 map.serialize_entry("topic", topic)?;
611 }
612 Frame::UnsubscribeAck { id, topic } => {
613 map.serialize_entry("id", id)?;
614 map.serialize_entry("topic", topic)?;
615 }
616 Frame::Event {
617 topic,
618 cursor,
619 occurred_at,
620 payload,
621 } => {
622 map.serialize_entry("topic", topic)?;
623 map.serialize_entry("cursor", cursor)?;
624 map.serialize_entry("occurred_at", occurred_at)?;
625 map.serialize_entry("payload", payload)?;
626 }
627 }
628 map.end()
629 }
630}
631
632/// The decode half of the closed grammar; see the type-level docs and the
633/// crate documentation's "Strict field rejection" section. The codec's
634/// `decode_payload` drives this via `serde_json::from_value::<Frame>`
635/// after its own closed-set `"kind"` check (which produces the
636/// finer-grained [`crate::codec::CodecError::UnknownFrameKind`]).
637///
638/// The visitor enforces every decode-time rule that only needs the fields
639/// of one frame, so a DIRECT serde decode of `Frame` (e.g.
640/// `serde_json::from_str::<Frame>`) agrees with the codec path:
641///
642/// - the ADR-137 id/scope pairing for `error` frames — a
643/// connection-terminal code must carry no operation id, and a
644/// request-terminal code must echo the one it terminates
645/// ([`crate::codec::CodecError::InconsistentErrorScope`] is the codec's
646/// typed form of this rejection; a direct serde decode reports the same
647/// rule through its deserializer's error type); and
648/// - the unknown-code fallback diagnostic: an `error` frame whose wire
649/// code is outside the closed set ([`crate::error::WIRE_ERROR_CODES`])
650/// carries the raw string in `unrecognized_code`.
651///
652/// The one guarantee only the codec path provides is the finer-grained
653/// [`crate::codec::CodecError::UnknownFrameKind`] classification for a
654/// `"kind"` outside [`FRAME_KINDS`]; this visitor rejects unknown kinds
655/// too, with the deserializer's own error type.
656impl<'de> Deserialize<'de> for Frame {
657 fn deserialize<D: Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
658 struct FrameVisitor;
659
660 impl<'de> Visitor<'de> for FrameVisitor {
661 type Value = Frame;
662
663 fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
664 formatter.write_str("a JSON object with a string \"kind\" discriminant field")
665 }
666
667 fn visit_map<A: MapAccess<'de>>(self, map: A) -> Result<Frame, A::Error> {
668 // Buffer the whole object, then inspect `"kind"` before any
669 // per-kind parsing: the closed-set check needs the
670 // discriminant up front, and the per-kind payload structs
671 // (all `deny_unknown_fields`) must never see the `"kind"`
672 // key itself.
673 let mut object: serde_json::Map<String, serde_json::Value> =
674 Deserialize::deserialize(serde::de::value::MapAccessDeserializer::new(map))?;
675 let kind_value = object
676 .get("kind")
677 .ok_or_else(|| A::Error::custom("missing field `kind`"))?;
678 let kind = kind_value
679 .as_str()
680 .ok_or_else(|| A::Error::custom("field `kind` must be a string"))?;
681 let kind = kind.to_string();
682 object.remove("kind");
683
684 fn parse<T: serde::de::DeserializeOwned>(
685 object: &serde_json::Map<String, serde_json::Value>,
686 ) -> Result<T, serde_json::Error> {
687 serde_json::from_value(serde_json::Value::Object(object.clone()))
688 }
689
690 match kind.as_str() {
691 "handshake" => {
692 let payload: HandshakePayload = parse(&object).map_err(A::Error::custom)?;
693 Ok(Frame::Handshake {
694 version: payload.version,
695 })
696 }
697 "handshake_ack" => {
698 let payload: HandshakeAckPayload =
699 parse(&object).map_err(A::Error::custom)?;
700 Ok(Frame::HandshakeAck {
701 version: payload.version,
702 })
703 }
704 "request" => {
705 let payload: RequestPayload = parse(&object).map_err(A::Error::custom)?;
706 Ok(Frame::Request {
707 id: payload.id,
708 ops: payload.ops,
709 deadline_ms: payload.deadline_ms,
710 namespace: payload.namespace,
711 actor_id: payload.actor_id,
712 visible_namespaces: payload.visible_namespaces,
713 })
714 }
715 "response" => {
716 let payload: ResponsePayload = parse(&object).map_err(A::Error::custom)?;
717 Ok(Frame::Response {
718 id: payload.id,
719 result: payload.result,
720 })
721 }
722 "error" => {
723 // Capture the raw code string BEFORE payload
724 // parsing: serde's `#[serde(other)]` fallback maps
725 // every code outside the closed set
726 // ([`crate::error::WIRE_ERROR_CODES`]) to
727 // [`crate::error::WireErrorCode::Internal`] and
728 // erases which of the two the wire carried. Both
729 // the id/scope check and the fallback diagnostic
730 // below need the pre-fallback string.
731 let raw_code = object
732 .get("code")
733 .and_then(|c| c.as_str())
734 .map(str::to_string);
735 let payload: ErrorPayload = parse(&object).map_err(A::Error::custom)?;
736
737 // ADR-137, "Operation correlation": a
738 // connection-terminal code carries no operation id,
739 // and a request-terminal code echoes the one it
740 // terminates. Enforced HERE — inside the visitor —
741 // so every decode path agrees: a direct
742 // `serde_json::from_str::<Frame>` can never
743 // represent an inconsistent error frame, and
744 // [`crate::codec::decode_payload`] (which drives
745 // this same visitor) re-classifies the rejection
746 // into its typed
747 // [`crate::codec::CodecError::InconsistentErrorScope`].
748 // The check covers closed-set codes only: an unknown
749 // code's true scope is unknowable to this protocol
750 // version, and the ADR directs the fallback-to-
751 // `internal` treatment rather than rejection.
752 let code_in_closed_set = raw_code
753 .as_deref()
754 .is_some_and(|c| crate::error::WIRE_ERROR_CODES.contains(&c));
755 if code_in_closed_set {
756 match (payload.code.terminal_scope(), payload.id.as_ref()) {
757 (crate::error::TerminalScope::Connection, Some(id)) => {
758 return Err(A::Error::custom(format!(
759 "{}connection-terminal code {code} must not carry an operation id, got {id}",
760 crate::codec::INCONSISTENT_SCOPE_ERROR_PREFIX,
761 code = payload.code
762 )));
763 }
764 (crate::error::TerminalScope::Request, None) => {
765 return Err(A::Error::custom(format!(
766 "{}request-terminal code {code} must echo the operation id it terminates",
767 crate::codec::INCONSISTENT_SCOPE_ERROR_PREFIX,
768 code = payload.code
769 )));
770 }
771 (crate::error::TerminalScope::Connection, None)
772 | (crate::error::TerminalScope::Request, Some(_)) => {}
773 }
774 }
775
776 Ok(Frame::Error {
777 id: payload.id,
778 code: payload.code,
779 message: payload.message,
780 // The decoded-fallback marker, filled on EVERY
781 // decode path (codec or direct serde): the raw
782 // string when the code fell back to `Internal`,
783 // `None` for every closed-set code — including a
784 // literal `"internal"`, which is a closed-set
785 // member, not a fallback.
786 unrecognized_code: if code_in_closed_set { None } else { raw_code },
787 })
788 }
789 "cancel" => {
790 let payload: CancelPayload = parse(&object).map_err(A::Error::custom)?;
791 Ok(Frame::Cancel { id: payload.id })
792 }
793 "subscribe" => {
794 let payload: SubscribePayload = parse(&object).map_err(A::Error::custom)?;
795 Ok(Frame::Subscribe {
796 id: payload.id,
797 topic: payload.topic,
798 resume_cursor: payload.resume_cursor,
799 })
800 }
801 "subscribe_ack" => {
802 let payload: SubscribeAckPayload =
803 parse(&object).map_err(A::Error::custom)?;
804 Ok(Frame::SubscribeAck {
805 id: payload.id,
806 topic: payload.topic,
807 start_cursor: payload.start_cursor,
808 })
809 }
810 "unsubscribe" => {
811 let payload: UnsubscribePayload =
812 parse(&object).map_err(A::Error::custom)?;
813 Ok(Frame::Unsubscribe {
814 id: payload.id,
815 topic: payload.topic,
816 })
817 }
818 "unsubscribe_ack" => {
819 let payload: UnsubscribeAckPayload =
820 parse(&object).map_err(A::Error::custom)?;
821 Ok(Frame::UnsubscribeAck {
822 id: payload.id,
823 topic: payload.topic,
824 })
825 }
826 "event" => {
827 let payload: EventPayload = parse(&object).map_err(A::Error::custom)?;
828 Ok(Frame::Event {
829 topic: payload.topic,
830 cursor: payload.cursor,
831 occurred_at: payload.occurred_at,
832 payload: payload.payload,
833 })
834 }
835 // The codec's closed-set check against `FRAME_KINDS`
836 // rejects unknown kinds with `UnknownFrameKind` before
837 // this point; this arm only fires for a direct
838 // `serde_json::from_value::<Frame>` call.
839 other => Err(A::Error::custom(format!("unknown frame kind: {other:?}"))),
840 }
841 }
842 }
843
844 deserializer.deserialize_map(FrameVisitor)
845 }
846}
847
848#[cfg(test)]
849mod tests {
850 use super::*;
851
852 #[test]
853 fn empty_operation_id_is_rejected_at_deserialization() {
854 // An empty operation id can never be a unique caller-
855 // generated id, so the wire form rejects it. In-memory
856 // construction stays unrestricted (checked below).
857 let err = serde_json::from_str::<OperationId>(r#""""#).unwrap_err();
858 assert!(
859 err.to_string().contains("non-empty"),
860 "unexpected error: {err}"
861 );
862 assert_eq!(OperationId::from(""), OperationId("".to_string()));
863 }
864
865 #[test]
866 fn non_empty_operation_id_round_trips() {
867 let id: OperationId = serde_json::from_str(r#""op-1""#).unwrap();
868 assert_eq!(id, OperationId::from("op-1"));
869 assert_eq!(serde_json::to_string(&id).unwrap(), r#""op-1""#);
870 }
871
872 #[test]
873 fn direct_serde_rejects_every_invalid_wire_shape() {
874 let frames = [
875 Frame::Handshake {
876 version: ProtocolVersion::new(0),
877 },
878 Frame::HandshakeAck {
879 version: ProtocolVersion::new(0),
880 },
881 Frame::Request {
882 id: OperationId::from(""),
883 ops: "stats()".to_string(),
884 deadline_ms: None,
885 namespace: None,
886 actor_id: None,
887 visible_namespaces: None,
888 },
889 Frame::Response {
890 id: OperationId::from(""),
891 result: serde_json::json!({}),
892 },
893 Frame::Error {
894 id: Some(OperationId::from("")),
895 code: crate::error::WireErrorCode::Internal,
896 message: "failure".to_string(),
897 unrecognized_code: None,
898 },
899 Frame::Cancel {
900 id: OperationId::from(""),
901 },
902 Frame::Subscribe {
903 id: OperationId::from(""),
904 topic: "a.b".to_string(),
905 resume_cursor: None,
906 },
907 Frame::SubscribeAck {
908 id: OperationId::from(""),
909 topic: "a.b".to_string(),
910 start_cursor: 0,
911 },
912 Frame::Unsubscribe {
913 id: OperationId::from(""),
914 topic: "a.b".to_string(),
915 },
916 Frame::UnsubscribeAck {
917 id: OperationId::from(""),
918 topic: "a.b".to_string(),
919 },
920 Frame::Error {
921 id: Some(OperationId::from("op-1")),
922 code: crate::error::WireErrorCode::FrameTooLarge,
923 message: "too big".to_string(),
924 unrecognized_code: None,
925 },
926 Frame::Error {
927 id: None,
928 code: crate::error::WireErrorCode::Cancelled,
929 message: "cancelled".to_string(),
930 unrecognized_code: None,
931 },
932 Frame::Error {
933 id: None,
934 code: crate::error::WireErrorCode::Internal,
935 message: "future".to_string(),
936 unrecognized_code: Some("future_code_xyz".to_string()),
937 },
938 ];
939
940 for frame in frames {
941 assert!(
942 serde_json::to_vec(&frame).is_err(),
943 "invalid frame {:?} serialized successfully",
944 frame.kind()
945 );
946 }
947 }
948
949 #[test]
950 fn direct_serde_matches_encode_payload_for_a_valid_frame() {
951 let frame = Frame::Request {
952 id: OperationId::from("op-1"),
953 ops: "stats()".to_string(),
954 deadline_ms: Some(5000),
955 namespace: Some("default".to_string()),
956 actor_id: Some("actor".to_string()),
957 visible_namespaces: Some(vec!["default".to_string()]),
958 };
959 let direct = serde_json::to_vec(&frame).unwrap();
960 let encoded = crate::codec::encode_frame(&frame).unwrap();
961 assert_eq!(direct, &encoded[crate::codec::LENGTH_PREFIX_BYTES..]);
962 assert_eq!(serde_json::from_slice::<Frame>(&direct).unwrap(), frame);
963 }
964
965 #[test]
966 fn unknown_kind_through_direct_serde_is_rejected() {
967 let err = serde_json::from_str::<Frame>(r#"{"kind":"ping"}"#).unwrap_err();
968 assert!(err.to_string().contains("unknown frame kind"));
969 }
970
971 #[test]
972 fn client_to_server_kinds_are_all_frame_kinds() {
973 for kind in CLIENT_TO_SERVER_KINDS {
974 assert!(
975 FRAME_KINDS.contains(kind),
976 "client-to-server kind {kind:?} is missing from FRAME_KINDS"
977 );
978 }
979 }
980
981 #[test]
982 fn direct_serde_rejects_connection_terminal_error_carrying_an_id() {
983 // The id/scope rule is enforced inside the visitor itself, so a
984 // DIRECT serde decode — bypassing `decode_payload` entirely —
985 // cannot represent an inconsistent error frame either.
986 let err = serde_json::from_str::<Frame>(
987 r#"{"kind":"error","id":"op-1","code":"frame_too_large","message":"too big"}"#,
988 )
989 .unwrap_err();
990 let message = err.to_string();
991 assert!(
992 message.contains("id/scope rule"),
993 "unexpected error: {message}"
994 );
995 assert!(message.contains("frame_too_large"), "error: {message}");
996 assert!(message.contains("op-1"), "error: {message}");
997 }
998
999 #[test]
1000 fn direct_serde_rejects_request_terminal_error_without_an_id() {
1001 let err = serde_json::from_str::<Frame>(
1002 r#"{"kind":"error","code":"cancelled","message":"cancelled"}"#,
1003 )
1004 .unwrap_err();
1005 let message = err.to_string();
1006 assert!(
1007 message.contains("id/scope rule"),
1008 "unexpected error: {message}"
1009 );
1010 assert!(message.contains("cancelled"), "error: {message}");
1011 }
1012
1013 #[test]
1014 fn direct_serde_accepts_both_consistent_error_scopes() {
1015 let connection_terminal: Frame = serde_json::from_str(
1016 r#"{"kind":"error","code":"unsupported_version","message":"no common version"}"#,
1017 )
1018 .unwrap();
1019 assert!(matches!(connection_terminal, Frame::Error { id: None, .. }));
1020
1021 let request_terminal: Frame = serde_json::from_str(
1022 r#"{"kind":"error","id":"op-9","code":"deadline_exceeded","message":"too slow"}"#,
1023 )
1024 .unwrap();
1025 match request_terminal {
1026 Frame::Error {
1027 id,
1028 code,
1029 unrecognized_code,
1030 ..
1031 } => {
1032 assert_eq!(id, Some(OperationId::from("op-9")));
1033 assert_eq!(code, crate::error::WireErrorCode::DeadlineExceeded);
1034 assert!(unrecognized_code.is_none());
1035 }
1036 other => panic!("expected an error frame, got {other:?}"),
1037 }
1038 }
1039
1040 #[test]
1041 fn direct_serde_fills_unrecognized_code_for_an_unknown_wire_code() {
1042 // The fallback diagnostic is filled by the visitor, on every
1043 // decode path — not only by the codec. The id/scope pairing is NOT
1044 // enforced for an unknown code (its true scope is unknown), so
1045 // both id shapes decode.
1046 for json in [
1047 r#"{"kind":"error","id":"op-7","code":"future_code_xyz","message":"from newer peer"}"#,
1048 r#"{"kind":"error","code":"future_code_xyz","message":"from newer peer"}"#,
1049 ] {
1050 let frame: Frame = serde_json::from_str(json).unwrap();
1051 match frame {
1052 Frame::Error {
1053 code,
1054 unrecognized_code,
1055 ..
1056 } => {
1057 assert_eq!(code, crate::error::WireErrorCode::Internal);
1058 assert_eq!(unrecognized_code.as_deref(), Some("future_code_xyz"));
1059 }
1060 other => panic!("expected an error frame, got {other:?}"),
1061 }
1062 }
1063 }
1064
1065 #[test]
1066 fn direct_serde_leaves_unrecognized_code_none_for_a_literal_internal() {
1067 // `"internal"` is a closed-set member, not a fallback: the
1068 // diagnostic must stay `None` so the encode side does not mistake
1069 // an honest `internal` for a decoded-fallback frame.
1070 let frame: Frame = serde_json::from_str(
1071 r#"{"kind":"error","id":"op-1","code":"internal","message":"boom"}"#,
1072 )
1073 .unwrap();
1074 match frame {
1075 Frame::Error {
1076 unrecognized_code, ..
1077 } => assert!(unrecognized_code.is_none()),
1078 other => panic!("expected an error frame, got {other:?}"),
1079 }
1080 }
1081}