1use std::fmt;
9
10use crate::frame_v1 as backend;
11
12pub const DAEMON_FRAME_V1_VERSION: u8 = backend::ENVELOPE_VERSION;
14
15pub const DAEMON_FRAME_V1_MAX_BODY_BYTES: usize = backend::MAX_FRAME_BYTES;
17
18#[derive(Clone, Debug, PartialEq, Eq)]
26pub struct DaemonFrame {
27 envelope_version: u32,
28 kind: i32,
29 payload_protocol: u32,
30 payload: Vec<u8>,
31 request_id: u64,
32 payload_encoding: i32,
33 deadline_unix_ms: u64,
34 traceparent: String,
35 tracestate: String,
36}
37
38#[derive(Clone, Copy, Debug, PartialEq, Eq)]
43pub enum DaemonFrameKind {
44 Request,
46 Response,
48 Event,
50 Cancel,
52 Unknown(i32),
54}
55
56#[derive(Clone, Copy, Debug, PartialEq, Eq)]
61pub enum DaemonPayloadEncoding {
62 None,
64 Zstd,
66 Snappy,
68 Lz4,
70 Unknown(i32),
72}
73
74impl DaemonFrame {
75 #[must_use]
77 pub fn request(payload_protocol: u32, payload: Vec<u8>) -> Self {
78 Self::from_backend(backend::Frame::request(payload_protocol, payload))
79 }
80
81 #[must_use]
86 pub fn response_to(request: &Self, payload: Vec<u8>) -> Self {
87 Self::from_backend(backend::Frame::response_to(&request.to_backend(), payload))
88 }
89
90 #[must_use]
92 pub fn with_request_id(mut self, request_id: u64) -> Self {
93 self.request_id = request_id;
94 self
95 }
96
97 #[must_use]
99 pub fn with_raw_kind(mut self, kind: i32) -> Self {
100 self.kind = kind;
101 self
102 }
103
104 #[must_use]
106 pub fn with_raw_payload_encoding(mut self, payload_encoding: i32) -> Self {
107 self.payload_encoding = payload_encoding;
108 self
109 }
110
111 #[must_use]
113 pub fn with_deadline_unix_ms(mut self, deadline_unix_ms: u64) -> Self {
114 self.deadline_unix_ms = deadline_unix_ms;
115 self
116 }
117
118 #[must_use]
124 pub fn with_trace_context(
125 mut self,
126 trace_id: impl AsRef<str>,
127 span_id: impl AsRef<str>,
128 ) -> Self {
129 self.traceparent = format!("00-{}-{}-01", trace_id.as_ref(), span_id.as_ref());
130 self
131 }
132
133 #[must_use]
135 pub fn with_trace_state(mut self, trace_state: impl Into<String>) -> Self {
136 self.tracestate = trace_state.into();
137 self
138 }
139
140 #[must_use]
142 pub fn envelope_version(&self) -> u32 {
143 self.envelope_version
144 }
145
146 #[must_use]
148 pub fn kind(&self) -> i32 {
149 self.kind
150 }
151
152 #[must_use]
154 pub fn kind_classification(&self) -> DaemonFrameKind {
155 match self.kind {
156 value if value == backend::FrameKind::Request as i32 => DaemonFrameKind::Request,
157 value if value == backend::FrameKind::Response as i32 => DaemonFrameKind::Response,
158 value if value == backend::FrameKind::Event as i32 => DaemonFrameKind::Event,
159 value if value == backend::FrameKind::Cancel as i32 => DaemonFrameKind::Cancel,
160 value => DaemonFrameKind::Unknown(value),
161 }
162 }
163
164 #[must_use]
166 pub fn payload_protocol(&self) -> u32 {
167 self.payload_protocol
168 }
169
170 #[must_use]
172 pub fn payload(&self) -> &[u8] {
173 &self.payload
174 }
175
176 #[must_use]
178 pub fn request_id(&self) -> u64 {
179 self.request_id
180 }
181
182 #[must_use]
184 pub fn payload_encoding(&self) -> i32 {
185 self.payload_encoding
186 }
187
188 #[must_use]
190 pub fn payload_encoding_classification(&self) -> DaemonPayloadEncoding {
191 match self.payload_encoding {
192 value if value == backend::PayloadEncoding::None as i32 => DaemonPayloadEncoding::None,
193 value if value == backend::PayloadEncoding::Zstd as i32 => DaemonPayloadEncoding::Zstd,
194 value if value == backend::PayloadEncoding::Snappy as i32 => {
195 DaemonPayloadEncoding::Snappy
196 }
197 value if value == backend::PayloadEncoding::Lz4 as i32 => DaemonPayloadEncoding::Lz4,
198 value => DaemonPayloadEncoding::Unknown(value),
199 }
200 }
201
202 #[must_use]
204 pub fn deadline_unix_ms(&self) -> u64 {
205 self.deadline_unix_ms
206 }
207
208 #[must_use]
210 pub fn trace_id(&self) -> Option<&str> {
211 trace_ids(&self.traceparent).map(|(trace_id, _)| trace_id)
212 }
213
214 #[must_use]
216 pub fn span_id(&self) -> Option<&str> {
217 trace_ids(&self.traceparent).map(|(_, span_id)| span_id)
218 }
219
220 #[must_use]
222 pub fn trace_state(&self) -> &str {
223 &self.tracestate
224 }
225
226 fn from_backend(frame: backend::Frame) -> Self {
227 Self {
228 envelope_version: frame.envelope_version,
229 kind: frame.kind,
230 payload_protocol: frame.payload_protocol,
231 payload: frame.payload,
232 request_id: frame.request_id,
233 payload_encoding: frame.payload_encoding,
234 deadline_unix_ms: frame.deadline_unix_ms,
235 traceparent: frame.traceparent,
236 tracestate: frame.tracestate,
237 }
238 }
239
240 fn to_backend(&self) -> backend::Frame {
241 backend::Frame {
242 envelope_version: self.envelope_version,
243 kind: self.kind,
244 payload_protocol: self.payload_protocol,
245 payload: self.payload.clone(),
246 request_id: self.request_id,
247 payload_encoding: self.payload_encoding,
248 deadline_unix_ms: self.deadline_unix_ms,
249 traceparent: self.traceparent.clone(),
250 tracestate: self.tracestate.clone(),
251 }
252 }
253}
254
255#[derive(Clone, Copy, Debug, Default)]
257pub struct DaemonFrameCodec;
258
259impl DaemonFrameCodec {
260 pub fn encode(frame: &DaemonFrame) -> Result<Vec<u8>, DaemonFrameError> {
262 backend::encode_framed(&frame.to_backend()).map_err(DaemonFrameError::from_backend)
263 }
264
265 pub fn encode_request(
267 payload_protocol: u32,
268 payload: Vec<u8>,
269 request_id: u64,
270 ) -> Result<Vec<u8>, DaemonFrameError> {
271 Self::encode(&DaemonFrame::request(payload_protocol, payload).with_request_id(request_id))
272 }
273
274 pub fn encode_response_to(
276 request: &DaemonFrame,
277 payload: Vec<u8>,
278 ) -> Result<Vec<u8>, DaemonFrameError> {
279 Self::encode(&DaemonFrame::response_to(request, payload))
280 }
281
282 pub fn decode(buffer: &[u8]) -> Result<DaemonFrameDecode, DaemonFrameError> {
288 match backend::try_decode_framed(buffer).map_err(DaemonFrameError::from_backend)? {
289 Some(decoded) => Ok(DaemonFrameDecode::Frame {
290 frame: DaemonFrame::from_backend(decoded.frame),
291 consumed: decoded.consumed,
292 }),
293 None => Ok(DaemonFrameDecode::NeedMoreBytes),
294 }
295 }
296}
297
298#[derive(Clone, Debug, PartialEq, Eq)]
300pub enum DaemonFrameDecode {
301 NeedMoreBytes,
303 Frame {
305 frame: DaemonFrame,
307 consumed: usize,
309 },
310}
311
312#[derive(Clone, Debug, PartialEq, Eq)]
314pub enum DaemonFrameError {
315 UnsupportedFrameVersion {
317 received: u8,
319 expected: u8,
321 },
322 FrameTooLarge {
324 body_length: usize,
326 maximum: usize,
328 },
329 MalformedFrame,
331}
332
333impl DaemonFrameError {
334 fn from_backend(error: backend::FramingError) -> Self {
335 match error {
336 backend::FramingError::UnsupportedFramingVersion { got, expected } => {
337 Self::UnsupportedFrameVersion {
338 received: got,
339 expected,
340 }
341 }
342 backend::FramingError::FrameTooLarge { body_length, cap } => Self::FrameTooLarge {
343 body_length,
344 maximum: cap,
345 },
346 backend::FramingError::UnexpectedEof { .. }
347 | backend::FramingError::Io(_)
348 | backend::FramingError::Decode(_) => Self::MalformedFrame,
349 }
350 }
351}
352
353impl fmt::Display for DaemonFrameError {
354 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
355 match self {
356 Self::UnsupportedFrameVersion { received, expected } => {
357 write!(
358 formatter,
359 "unsupported daemon-frame version {received}; expected {expected}"
360 )
361 }
362 Self::FrameTooLarge {
363 body_length,
364 maximum,
365 } => write!(
366 formatter,
367 "daemon-frame body {body_length} exceeds maximum {maximum}"
368 ),
369 Self::MalformedFrame => formatter.write_str("malformed daemon-frame body"),
370 }
371 }
372}
373
374impl std::error::Error for DaemonFrameError {}
375
376#[doc(hidden)]
378pub const fn is_first_party_payload_protocol(payload_protocol: u32) -> bool {
379 backend::registry::is_first_party(payload_protocol)
380}
381
382#[doc(hidden)]
384pub const fn is_registered_payload_protocol(payload_protocol: u32) -> bool {
385 backend::registry::is_registered_consumer_id(payload_protocol)
386}
387
388#[doc(hidden)]
390pub const fn is_private_payload_protocol(payload_protocol: u32) -> bool {
391 backend::registry::is_private_use_id(payload_protocol)
392}
393
394#[macro_export]
400macro_rules! register_daemon_frame_payload_protocol {
401 ($(#[$meta:meta])* $vis:vis const $name:ident: u32 = $value:expr;) => {
402 $(#[$meta])*
403 $vis const $name: u32 = $value;
404
405 const _: () = {
406 assert!(
407 !$crate::daemon_frame_v1::is_first_party_payload_protocol($name),
408 concat!(
409 stringify!($name),
410 " collides with a first-party daemon-frame payload protocol",
411 ),
412 );
413 assert!(
414 $crate::daemon_frame_v1::is_registered_payload_protocol($name)
415 || $crate::daemon_frame_v1::is_private_payload_protocol($name),
416 concat!(
417 stringify!($name),
418 " must lie in the registered-consumer range (0x7000..=0x7EFF) ",
419 "or the private-use range (0xF000..=0xFFFF)",
420 ),
421 );
422 };
423 };
424}
425
426fn trace_ids(traceparent: &str) -> Option<(&str, &str)> {
427 let mut fields = traceparent.split('-');
428 let _version = fields.next()?;
429 let trace_id = fields.next()?;
430 let span_id = fields.next()?;
431 let _flags = fields.next()?;
432 if fields.next().is_some() || trace_id.is_empty() || span_id.is_empty() {
433 return None;
434 }
435 Some((trace_id, span_id))
436}
437
438#[cfg(test)]
439mod compatibility_tests {
440 use super::*;
441
442 #[test]
443 fn request_matches_literal_consumer_wire_fixture() {
444 let request =
445 DaemonFrame::request(0x7A63, b"ping".to_vec()).with_request_id(0x0102_0304_0506_0708);
446 assert_eq!(
447 DaemonFrameCodec::encode(&request).unwrap(),
448 [
449 0x01, 0x16, 0, 0, 0, 0x08, 0x01, 0x18, 0xE3, 0xF4, 0x01, 0x22, 0x04, b'p', b'i',
450 b'n', b'g', 0x28, 0x88, 0x8E, 0x98, 0xA8, 0xC0, 0xE0, 0x80, 0x81, 0x01,
451 ]
452 );
453 }
454
455 #[test]
456 fn decoded_trace_headers_are_not_reconstructed_from_semantic_ids() {
457 let mut raw = backend::Frame::request(0x7A63, vec![0, 255]);
458 raw.traceparent = "01-ABCDEF0123456789ABCDEF0123456789-ABCDEF0123456789-fe".into();
459 raw.tracestate = "vendor=One,another=Two".into();
460 let wire = backend::encode_framed(&raw).unwrap();
461 let DaemonFrameDecode::Frame { frame, .. } = DaemonFrameCodec::decode(&wire).unwrap()
462 else {
463 panic!("complete frame")
464 };
465 assert_eq!(DaemonFrameCodec::encode(&frame).unwrap(), wire);
466 let reply = DaemonFrame::response_to(&frame, Vec::new()).to_backend();
467 assert_eq!(reply.traceparent, raw.traceparent);
468 assert_eq!(reply.tracestate, raw.tracestate);
469 }
470
471 #[test]
472 fn raw_unknown_enums_trace_and_correlation_round_trip_in_frozen_bytes() {
473 let frame = DaemonFrame::request(0x7A63, b"payload".to_vec())
474 .with_request_id(77)
475 .with_raw_kind(37)
476 .with_raw_payload_encoding(91)
477 .with_trace_context("0123456789abcdef0123456789abcdef", "0123456789abcdef")
478 .with_trace_state("vendor=value");
479 let wire = DaemonFrameCodec::encode(&frame).expect("encode");
480 let DaemonFrameDecode::Frame {
481 frame: decoded,
482 consumed,
483 } = DaemonFrameCodec::decode(&wire).expect("decode")
484 else {
485 panic!("complete")
486 };
487 assert_eq!(consumed, wire.len());
488 assert_eq!(decoded.kind_classification(), DaemonFrameKind::Unknown(37));
489 assert_eq!(
490 decoded.payload_encoding_classification(),
491 DaemonPayloadEncoding::Unknown(91)
492 );
493 assert_eq!(decoded.request_id(), 77);
494 assert_eq!(decoded.trace_id(), Some("0123456789abcdef0123456789abcdef"));
495 assert_eq!(DaemonFrameCodec::encode(&decoded).expect("reencode"), wire);
496 }
497
498 #[test]
499 fn incremental_and_frozen_error_mapping_are_stable() {
500 assert_eq!(
501 DaemonFrameCodec::decode(&[]).expect("partial"),
502 DaemonFrameDecode::NeedMoreBytes
503 );
504 assert!(matches!(
505 DaemonFrameCodec::decode(&[2]),
506 Err(DaemonFrameError::UnsupportedFrameVersion { .. })
507 ));
508 assert!(matches!(
509 DaemonFrameCodec::decode(&[1, 1, 0, 0, 0, 0xff]),
510 Err(DaemonFrameError::MalformedFrame)
511 ));
512 }
513}