1mod codec;
7mod dynamic_subscriber;
8mod frame;
9mod nats_transport;
10mod traits;
11mod transport;
12pub mod zmq_transport;
13
14pub use codec::{Codec, MsgpackCodec};
15pub use dynamic_subscriber::DynamicSubscriber;
16pub use frame::{FRAME_HEADER_SIZE, FRAME_VERSION, Frame, FrameError, FrameHeader};
17pub use traits::{EventEnvelope, EventStream, TypedEventStream};
18pub use transport::{EventTransportRx, EventTransportTx, WireStream};
19pub use zmq_transport::{
20 ValidatedEnvelope, ValidatedZmqSource, ValidatedZmqSourceError, ZmqPubTransport,
21 ZmqSubTransport,
22};
23
24pub use crate::discovery::{EventScope, EventTransportKind};
26
27use std::num::NonZeroUsize;
28use std::sync::Arc;
29use std::sync::atomic::{AtomicU64, Ordering};
30use std::time::{SystemTime, UNIX_EPOCH};
31
32use anyhow::Result;
33use bytes::Bytes;
34use futures::{Stream, StreamExt};
35use lru::LruCache;
36use rand::TryRngCore;
37use serde::Serialize;
38use serde::de::DeserializeOwned;
39use std::pin::Pin;
40use std::task::{Context, Poll};
41
42use crate::DistributedRuntime;
43use crate::component::{Component, Endpoint, Namespace};
44use crate::discovery::{
45 Discovery, DiscoveryInstance, DiscoveryQuery, DiscoverySpec, EventChannelQuery, EventTransport,
46 MAX_JSON_SAFE_PUBLISHER_ID,
47};
48use crate::protocols::EndpointId;
49use crate::traits::DistributedRuntimeProvider;
50use crate::utils::local_ip_for_advertise;
51
52#[derive(Debug, Clone)]
58struct BrokerEndpoints {
59 xsub_endpoints: Vec<String>,
60 xpub_endpoints: Vec<String>,
61}
62
63pub fn uses_direct_zmq(transport_kind: EventTransportKind) -> bool {
69 uses_direct_zmq_from_lookup(transport_kind, |key| std::env::var_os(key))
70}
71
72fn uses_direct_zmq_from_lookup(
73 transport_kind: EventTransportKind,
74 mut get_env: impl FnMut(&str) -> Option<std::ffi::OsString>,
75) -> bool {
76 if transport_kind != EventTransportKind::Zmq {
77 return false;
78 }
79
80 if get_env(crate::config::environment_names::zmq_broker::DYN_ZMQ_BROKER_URL).is_some() {
81 return false;
82 }
83
84 !get_env(crate::config::environment_names::zmq_broker::DYN_ZMQ_BROKER_ENABLED)
85 .is_some_and(|value| crate::config::is_truthy(&value.to_string_lossy()))
86}
87
88async fn resolve_zmq_broker(
91 drt: &DistributedRuntime,
92 scope: &EventScope,
93) -> Result<Option<BrokerEndpoints>> {
94 if let Ok(broker_url) =
96 std::env::var(crate::config::environment_names::zmq_broker::DYN_ZMQ_BROKER_URL)
97 {
98 let (xsub_endpoints, xpub_endpoints) = parse_broker_url(&broker_url)?;
99 tracing::info!(
100 num_xsub = xsub_endpoints.len(),
101 num_xpub = xpub_endpoints.len(),
102 "Using explicit ZMQ broker URL"
103 );
104 return Ok(Some(BrokerEndpoints {
105 xsub_endpoints,
106 xpub_endpoints,
107 }));
108 }
109
110 if crate::config::env_is_truthy(
112 crate::config::environment_names::zmq_broker::DYN_ZMQ_BROKER_ENABLED,
113 ) {
114 let query = DiscoveryQuery::EventChannels(EventChannelQuery::component(
115 scope.namespace().to_string(),
116 "zmq_broker".to_string(),
117 ));
118
119 let instances = drt.discovery().list(query).await?;
120
121 let mut xsub_endpoints = Vec::new();
123 let mut xpub_endpoints = Vec::new();
124
125 for instance in instances {
126 if let DiscoveryInstance::EventChannel { transport, .. } = instance
127 && let EventTransport::ZmqBroker {
128 xsub_endpoints: xsubs,
129 xpub_endpoints: xpubs,
130 } = transport
131 {
132 xsub_endpoints.extend(xsubs);
133 xpub_endpoints.extend(xpubs);
134 }
135 }
136
137 if xsub_endpoints.is_empty() {
138 anyhow::bail!(
139 "DYN_ZMQ_BROKER_ENABLED is set but no broker found in discovery for namespace '{}'",
140 scope.namespace()
141 );
142 }
143
144 tracing::info!(
145 num_brokers = xsub_endpoints.len(),
146 "Discovered ZMQ brokers from discovery plane"
147 );
148
149 return Ok(Some(BrokerEndpoints {
150 xsub_endpoints,
151 xpub_endpoints,
152 }));
153 }
154
155 Ok(None)
157}
158
159fn parse_broker_url(url: &str) -> Result<(Vec<String>, Vec<String>)> {
161 let parts: Vec<&str> = url.split(',').map(|s| s.trim()).collect();
162 if parts.len() != 2 {
163 anyhow::bail!(
164 "Invalid broker URL format. Expected 'xsub=<urls> , xpub=<urls>', got: {}",
165 url
166 );
167 }
168
169 let mut xsub_endpoints = Vec::new();
170 let mut xpub_endpoints = Vec::new();
171
172 for part in parts {
173 if let Some(urls_str) = part.strip_prefix("xsub=") {
174 xsub_endpoints = urls_str
175 .split(';')
176 .map(|s| s.trim().to_string())
177 .filter(|s| !s.is_empty())
178 .collect();
179 } else if let Some(urls_str) = part.strip_prefix("xpub=") {
180 xpub_endpoints = urls_str
181 .split(';')
182 .map(|s| s.trim().to_string())
183 .filter(|s| !s.is_empty())
184 .collect();
185 } else {
186 anyhow::bail!(
187 "Invalid broker URL part. Expected 'xsub=' or 'xpub=' prefix, got: {}",
188 part
189 );
190 }
191 }
192
193 if xsub_endpoints.is_empty() || xpub_endpoints.is_empty() {
194 anyhow::bail!(
195 "Broker URL must contain at least one xsub and one xpub endpoint. Got xsub={:?}, xpub={:?}",
196 xsub_endpoints,
197 xpub_endpoints
198 );
199 }
200
201 Ok((xsub_endpoints, xpub_endpoints))
202}
203
204struct DeduplicatingStream {
207 inner: WireStream,
208 codec: Arc<Codec>,
209 seen_events: LruCache<(u64, u64), ()>, }
211
212impl DeduplicatingStream {
213 fn new(inner: WireStream, codec: Arc<Codec>, cache_size: usize) -> Self {
214 Self {
215 inner,
216 codec,
217 seen_events: LruCache::new(
218 NonZeroUsize::new(cache_size).expect("cache_size must be non-zero"),
219 ),
220 }
221 }
222}
223
224impl Stream for DeduplicatingStream {
225 type Item = Result<Bytes>;
226
227 fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
228 loop {
229 match Pin::new(&mut self.inner).poll_next(cx) {
230 Poll::Ready(Some(Ok(bytes))) => {
231 match self.codec.decode_envelope_identity(&bytes) {
233 Ok(key) => {
234 if self.seen_events.contains(&key) {
236 tracing::debug!(
238 publisher_id = key.0,
239 sequence = key.1,
240 "Filtered duplicate event from multi-broker setup"
241 );
242 continue;
243 }
244
245 self.seen_events.put(key, ());
247 return Poll::Ready(Some(Ok(bytes)));
248 }
249 Err(e) => {
250 tracing::warn!(error = %e, "Failed to decode envelope for deduplication");
251 return Poll::Ready(Some(Err(e)));
252 }
253 }
254 }
255 Poll::Ready(Some(Err(e))) => return Poll::Ready(Some(Err(e))),
256 Poll::Ready(None) => return Poll::Ready(None),
257 Poll::Pending => return Poll::Pending,
258 }
259 }
260 }
261}
262
263fn discovery_safe_publisher_id(random_id: u64) -> u64 {
265 random_id & MAX_JSON_SAFE_PUBLISHER_ID
266}
267
268pub struct EventPublisher {
270 transport_kind: EventTransportKind,
271 topic: String,
272 subject: String,
273 publisher_id: u64,
274 sequence: AtomicU64,
275 tx: Arc<dyn EventTransportTx>,
276 codec: Arc<Codec>,
277 runtime_handle: tokio::runtime::Handle,
278 graceful_shutdown_tracker: Arc<crate::utils::GracefulShutdownTracker>,
280 discovery_client: Option<Arc<dyn Discovery>>,
282 discovery_instance: Option<crate::discovery::DiscoveryInstance>,
283}
284
285impl EventPublisher {
286 pub async fn for_endpoint(endpoint: &Endpoint, topic: impl Into<String>) -> Result<Self> {
288 let transport_kind = endpoint.drt().default_event_transport_kind();
289 Self::for_endpoint_with_transport(endpoint, topic, transport_kind).await
290 }
291
292 pub async fn for_endpoint_with_transport(
294 endpoint: &Endpoint,
295 topic: impl Into<String>,
296 transport_kind: EventTransportKind,
297 ) -> Result<Self> {
298 Self::new_internal(
299 endpoint.drt(),
300 EventScope::Endpoint {
301 endpoint: endpoint.id(),
302 },
303 topic.into(),
304 transport_kind,
305 )
306 .await
307 }
308
309 pub async fn for_endpoint_id(
311 drt: &DistributedRuntime,
312 endpoint: &EndpointId,
313 topic: impl Into<String>,
314 ) -> Result<Self> {
315 let transport_kind = drt.default_event_transport_kind();
316 Self::for_endpoint_id_with_transport(drt, endpoint, topic, transport_kind).await
317 }
318
319 pub async fn for_endpoint_id_with_transport(
321 drt: &DistributedRuntime,
322 endpoint: &EndpointId,
323 topic: impl Into<String>,
324 transport_kind: EventTransportKind,
325 ) -> Result<Self> {
326 Self::new_internal(
327 drt,
328 EventScope::Endpoint {
329 endpoint: endpoint.clone(),
330 },
331 topic.into(),
332 transport_kind,
333 )
334 .await
335 }
336
337 pub async fn for_component(comp: &Component, topic: impl Into<String>) -> Result<Self> {
345 let transport_kind = comp.drt().default_event_transport_kind();
346 Self::for_component_with_transport(comp, topic, transport_kind).await
347 }
348
349 pub async fn for_component_with_transport(
351 comp: &Component,
352 topic: impl Into<String>,
353 transport_kind: EventTransportKind,
354 ) -> Result<Self> {
355 let drt = comp.drt();
356 let scope = EventScope::Component {
357 namespace: comp.namespace().name(),
358 component: comp.name().to_string(),
359 };
360 Self::new_internal(drt, scope, topic.into(), transport_kind).await
361 }
362
363 pub async fn for_namespace(ns: &Namespace, topic: impl Into<String>) -> Result<Self> {
371 let transport_kind = ns.drt().default_event_transport_kind();
372 Self::for_namespace_with_transport(ns, topic, transport_kind).await
373 }
374
375 pub async fn for_namespace_with_transport(
377 ns: &Namespace,
378 topic: impl Into<String>,
379 transport_kind: EventTransportKind,
380 ) -> Result<Self> {
381 let drt = ns.drt();
382 let scope = EventScope::Namespace { name: ns.name() };
383 Self::new_internal(drt, scope, topic.into(), transport_kind).await
384 }
385
386 async fn new_internal(
387 drt: &DistributedRuntime,
388 scope: EventScope,
389 topic: String,
390 transport_kind: EventTransportKind,
391 ) -> Result<Self> {
392 let publisher_id = discovery_safe_publisher_id(
397 rand::rngs::OsRng
398 .try_next_u64()
399 .map_err(|error| anyhow::anyhow!("failed to generate publisher ID: {error}"))?,
400 );
401 let discovery = Some(drt.discovery());
402 let runtime_handle = drt.runtime().secondary();
403 let subject = scope.subject(&topic);
404 let graceful_shutdown_tracker = drt.graceful_shutdown_tracker();
405
406 enum TransportSetup {
408 Nats(Arc<dyn EventTransportTx>, Arc<Codec>),
409 ZmqDirect(Arc<dyn EventTransportTx>, Arc<Codec>, String), ZmqBroker(Arc<dyn EventTransportTx>, Arc<Codec>),
411 }
412
413 let transport_setup = match transport_kind {
414 EventTransportKind::Nats => {
415 let transport = Arc::new(nats_transport::NatsTransport::new_publisher(
416 drt.clone(),
417 subject.clone(),
418 ));
419 let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
420 TransportSetup::Nats(transport as Arc<dyn EventTransportTx>, codec)
421 }
422 EventTransportKind::Zmq => {
423 if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
425 let pub_transport = if broker.xsub_endpoints.len() == 1 {
427 zmq_transport::ZmqPubTransport::connect(&broker.xsub_endpoints[0], &subject)
428 .await?
429 } else {
430 zmq_transport::ZmqPubTransport::connect_multiple(
431 &broker.xsub_endpoints,
432 &subject,
433 )
434 .await?
435 };
436
437 let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
438 TransportSetup::ZmqBroker(
439 Arc::new(pub_transport) as Arc<dyn EventTransportTx>,
440 codec,
441 )
442 } else {
443 let (pub_transport, actual_bind_endpoint) = std::thread::spawn({
445 let topic = topic.clone();
446 move || {
447 let rt = tokio::runtime::Builder::new_current_thread()
448 .enable_all()
449 .build()
450 .expect("Failed to create Tokio runtime for ZMQ");
451
452 rt.block_on(async move {
453 zmq_transport::ZmqPubTransport::bind("tcp://0.0.0.0:0", &topic)
454 .await
455 .expect("Failed to bind ZMQ publisher")
456 })
457 }
458 })
459 .join()
460 .expect("Failed to join ZMQ initialization thread");
461
462 let actual_port: u16 = actual_bind_endpoint
464 .rsplit(':')
465 .next()
466 .and_then(|s| s.parse().ok())
467 .expect("Failed to parse port from bind endpoint");
468 let local_ip = local_ip_for_advertise();
469 let public_endpoint = format!("tcp://{}:{}", local_ip, actual_port);
470
471 let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
472 TransportSetup::ZmqDirect(
473 Arc::new(pub_transport) as Arc<dyn EventTransportTx>,
474 codec,
475 public_endpoint,
476 )
477 }
478 }
479 };
480
481 let (tx, codec, discovery_instance) = match transport_setup {
483 TransportSetup::Nats(tx, codec) => {
484 let transport_config = EventTransport::nats(scope.subject_prefix());
485 let spec = DiscoverySpec::EventChannel {
486 scope: scope.clone(),
487 topic: topic.clone(),
488 publisher_id,
489 transport: transport_config,
490 };
491
492 let discovery_instance = drt.discovery().register(spec).await?;
493 tracing::info!(
494 topic = %topic,
495 transport = ?transport_kind,
496 publisher_id = %publisher_id,
497 "EventPublisher registered with discovery"
498 );
499 (tx, codec, Some(discovery_instance))
500 }
501 TransportSetup::ZmqDirect(tx, codec, public_endpoint) => {
502 let transport_config = EventTransport::zmq(public_endpoint);
503 let spec = DiscoverySpec::EventChannel {
504 scope: scope.clone(),
505 topic: topic.clone(),
506 publisher_id,
507 transport: transport_config,
508 };
509
510 let discovery_instance = drt.discovery().register(spec).await?;
511 tracing::info!(
512 topic = %topic,
513 transport = ?transport_kind,
514 publisher_id = %publisher_id,
515 "EventPublisher registered with discovery (direct mode)"
516 );
517 (tx, codec, Some(discovery_instance))
518 }
519 TransportSetup::ZmqBroker(tx, codec) => {
520 tracing::info!(
521 topic = %topic,
522 transport = ?transport_kind,
523 "EventPublisher in broker mode - skipping discovery registration"
524 );
525 (tx, codec, None)
526 }
527 };
528
529 Ok(Self {
530 transport_kind,
531 topic,
532 subject,
533 publisher_id,
534 sequence: AtomicU64::new(0),
535 tx,
536 codec,
537 runtime_handle,
538 graceful_shutdown_tracker,
539 discovery_client: discovery,
540 discovery_instance,
541 })
542 }
543
544 pub async fn publish<T: Serialize + Send + Sync>(&self, event: &T) -> Result<()> {
546 let payload = self.codec.encode_payload(event)?;
547 self.publish_bytes_ref(payload.as_ref()).await
548 }
549
550 pub async fn publish_bytes(&self, bytes: Vec<u8>) -> Result<()> {
552 self.publish_bytes_ref(&bytes).await
553 }
554
555 pub async fn publish_bytes_ref(&self, bytes: &[u8]) -> Result<()> {
557 let envelope_bytes = self.codec.encode_envelope_parts(
558 self.publisher_id,
559 self.sequence.fetch_add(1, Ordering::SeqCst),
560 current_timestamp_ms(),
561 &self.topic,
562 bytes,
563 )?;
564
565 self.tx.publish(&self.subject, envelope_bytes).await
566 }
567
568 pub fn publisher_id(&self) -> u64 {
570 self.publisher_id
571 }
572
573 pub fn topic(&self) -> &str {
575 &self.topic
576 }
577
578 pub fn transport_kind(&self) -> EventTransportKind {
580 self.transport_kind
581 }
582}
583
584impl Drop for EventPublisher {
585 fn drop(&mut self) {
586 if let (Some(discovery), Some(instance)) =
588 (self.discovery_client.take(), self.discovery_instance.take())
589 {
590 let topic = self.topic.clone();
591 let publisher_id = instance.instance_id();
592 let runtime_handle = self.runtime_handle.clone();
593 let shutdown_guard = self.graceful_shutdown_tracker.register_task();
594
595 let spawn_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
598 runtime_handle.spawn(async move {
599 let _shutdown_guard = shutdown_guard;
600 match discovery.unregister(instance).await {
601 Ok(()) => {
602 tracing::info!(
603 topic = %topic,
604 publisher_id = %publisher_id,
605 "EventPublisher unregistered from discovery"
606 );
607 }
608 Err(e) => {
609 tracing::warn!(
610 topic = %topic,
611 publisher_id = %publisher_id,
612 error = %e,
613 "Failed to unregister EventPublisher from discovery"
614 );
615 }
616 }
617 });
618 }));
619
620 if spawn_result.is_err() {
621 tracing::warn!(
622 topic = %self.topic,
623 publisher_id = %publisher_id,
624 "Skipping EventPublisher unregister during drop because the runtime is unavailable"
625 );
626 }
627 }
628 }
629}
630
631pub struct EventSubscriber {
633 stream: EventStream,
634 #[allow(dead_code)]
635 scope: EventScope,
636 #[allow(dead_code)]
637 topic: String,
638 codec: Arc<Codec>,
639}
640
641impl EventSubscriber {
642 pub async fn for_endpoint(endpoint: &Endpoint, topic: impl Into<String>) -> Result<Self> {
644 let transport_kind = endpoint.drt().default_event_transport_kind();
645 Self::for_endpoint_with_transport(endpoint, topic, transport_kind).await
646 }
647
648 pub async fn for_endpoint_with_transport(
650 endpoint: &Endpoint,
651 topic: impl Into<String>,
652 transport_kind: EventTransportKind,
653 ) -> Result<Self> {
654 Self::for_endpoint_id_with_transport(endpoint.drt(), &endpoint.id(), topic, transport_kind)
655 .await
656 }
657
658 pub async fn for_endpoint_id(
661 drt: &DistributedRuntime,
662 endpoint: &EndpointId,
663 topic: impl Into<String>,
664 ) -> Result<Self> {
665 let transport_kind = drt.default_event_transport_kind();
666 Self::for_endpoint_id_with_transport(drt, endpoint, topic, transport_kind).await
667 }
668
669 pub async fn for_endpoint_id_with_transport(
671 drt: &DistributedRuntime,
672 endpoint: &EndpointId,
673 topic: impl Into<String>,
674 transport_kind: EventTransportKind,
675 ) -> Result<Self> {
676 Self::new_internal(
677 drt,
678 EventScope::Endpoint {
679 endpoint: endpoint.clone(),
680 },
681 topic.into(),
682 transport_kind,
683 )
684 .await
685 }
686
687 pub async fn for_component(comp: &Component, topic: impl Into<String>) -> Result<Self> {
695 let transport_kind = comp.drt().default_event_transport_kind();
696 Self::for_component_with_transport(comp, topic, transport_kind).await
697 }
698
699 pub async fn for_component_with_transport(
701 comp: &Component,
702 topic: impl Into<String>,
703 transport_kind: EventTransportKind,
704 ) -> Result<Self> {
705 let drt = comp.drt();
706 let scope = EventScope::Component {
707 namespace: comp.namespace().name(),
708 component: comp.name().to_string(),
709 };
710 Self::new_internal(drt, scope, topic.into(), transport_kind).await
711 }
712
713 pub async fn for_namespace(ns: &Namespace, topic: impl Into<String>) -> Result<Self> {
721 let transport_kind = ns.drt().default_event_transport_kind();
722 Self::for_namespace_with_transport(ns, topic, transport_kind).await
723 }
724
725 pub async fn for_namespace_with_transport(
727 ns: &Namespace,
728 topic: impl Into<String>,
729 transport_kind: EventTransportKind,
730 ) -> Result<Self> {
731 let drt = ns.drt();
732 let scope = EventScope::Namespace { name: ns.name() };
733 Self::new_internal(drt, scope, topic.into(), transport_kind).await
734 }
735
736 async fn new_internal(
737 drt: &DistributedRuntime,
738 scope: EventScope,
739 topic: String,
740 transport_kind: EventTransportKind,
741 ) -> Result<Self> {
742 let discovery = drt.discovery();
743 let routing_key = scope.subject(&topic);
744
745 let (wire_stream, codec): (WireStream, Arc<Codec>) = match transport_kind {
747 EventTransportKind::Nats => {
748 let transport = nats_transport::NatsTransport::new(drt.clone());
749 let stream = transport.subscribe(&routing_key).await?;
750 let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
751 (stream, codec)
752 }
753 EventTransportKind::Zmq => {
754 if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
756 let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
758
759 let stream: WireStream = if broker.xpub_endpoints.len() == 1 {
760 let stream = zmq_transport::ZmqSubTransport::connect_single_consumer(
763 &broker.xpub_endpoints[0],
764 &routing_key,
765 )
766 .await?;
767 Box::pin(stream.map(|result| result.map(|message| message.payload)))
768 } else {
769 let inner_stream =
771 zmq_transport::ZmqSubTransport::connect_single_consumer_multiple(
772 &broker.xpub_endpoints,
773 &routing_key,
774 )
775 .await?;
776 let inner_stream: WireStream = Box::pin(
777 inner_stream.map(|result| result.map(|message| message.payload)),
778 );
779
780 Box::pin(DeduplicatingStream::new(
782 inner_stream,
783 codec.clone(),
784 100_000,
785 ))
786 };
787
788 (stream, codec)
789 } else {
790 let query = match &scope {
792 EventScope::Namespace { name } => {
793 crate::discovery::DiscoveryQuery::EventChannels(
794 crate::discovery::EventChannelQuery::namespace_topic(
795 name.clone(),
796 topic.clone(),
797 ),
798 )
799 }
800 EventScope::Component {
801 namespace,
802 component,
803 } => crate::discovery::DiscoveryQuery::EventChannels(
804 crate::discovery::EventChannelQuery::topic(
805 namespace.clone(),
806 component.clone(),
807 topic.clone(),
808 ),
809 ),
810 EventScope::Endpoint { endpoint } => {
811 crate::discovery::DiscoveryQuery::EventChannels(
812 crate::discovery::EventChannelQuery::endpoint_topic(
813 endpoint.clone(),
814 topic.clone(),
815 ),
816 )
817 }
818 };
819
820 let subscriber = Arc::new(DynamicSubscriber::with_cancel_token(
821 discovery,
822 query,
823 topic.clone(),
824 drt.primary_token().child_token(),
825 ));
826
827 let stream = subscriber.start_zmq().await?;
828 let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
829 (stream, codec)
830 }
831 }
832 };
833
834 let topic_filter = topic.clone();
836 let codec_for_stream = codec.clone();
837 let stream = wire_stream.filter_map(move |result| {
838 let codec = codec_for_stream.clone();
839 let topic_filter = topic_filter.clone();
840 async move {
841 match result {
842 Ok(bytes) => match codec.decode_envelope(&bytes) {
843 Ok(envelope) => {
844 if envelope.topic == topic_filter {
846 Some(Ok(envelope))
847 } else {
848 None
849 }
850 }
851 Err(e) => Some(Err(e)),
852 },
853 Err(e) => Some(Err(e)),
854 }
855 }
856 });
857
858 tracing::info!(
859 topic = %topic,
860 transport = ?transport_kind,
861 "EventSubscriber created"
862 );
863
864 Ok(Self {
865 stream: Box::pin(stream),
866 scope,
867 topic,
868 codec,
869 })
870 }
871
872 pub async fn next(&mut self) -> Option<Result<EventEnvelope>> {
874 self.stream.next().await
875 }
876
877 pub fn typed<T: DeserializeOwned + Send + 'static>(self) -> TypedEventSubscriber<T> {
879 TypedEventSubscriber {
880 stream: self.stream,
881 codec: self.codec,
882 _marker: std::marker::PhantomData,
883 }
884 }
885}
886
887pub struct TypedEventSubscriber<T> {
889 stream: EventStream,
890 codec: Arc<Codec>,
891 _marker: std::marker::PhantomData<T>,
892}
893
894impl<T: DeserializeOwned + Send + 'static> TypedEventSubscriber<T> {
895 pub async fn next(&mut self) -> Option<Result<(EventEnvelope, T)>> {
897 std::future::poll_fn(|cx| self.poll_next(cx)).await
898 }
899
900 pub fn poll_next(&mut self, cx: &mut Context<'_>) -> Poll<Option<Result<(EventEnvelope, T)>>> {
902 match self.stream.as_mut().poll_next(cx) {
903 Poll::Ready(Some(envelope)) => Poll::Ready(Some(match envelope {
904 Ok(env) => match self.codec.decode_payload(&env.payload) {
905 Ok(typed) => Ok((env, typed)),
906 Err(e) => Err(e),
907 },
908 Err(e) => Err(e),
909 })),
910 Poll::Ready(None) => Poll::Ready(None),
911 Poll::Pending => Poll::Pending,
912 }
913 }
914}
915
916fn current_timestamp_ms() -> u64 {
918 SystemTime::now()
919 .duration_since(UNIX_EPOCH)
920 .map(|d| d.as_millis() as u64)
921 .unwrap_or(0)
922}
923
924#[cfg(test)]
925mod tests {
926 use super::*;
927 use crate::config::environment_names::zmq_broker as broker_env;
928
929 #[tokio::test]
930 async fn multi_broker_stream_deduplicates_by_publisher_and_sequence() {
931 let codec = Arc::new(Codec::default());
932 let first = codec
933 .encode_envelope_parts(7, 11, 1, "events", b"first")
934 .unwrap();
935 let second = codec
936 .encode_envelope_parts(7, 12, 2, "events", b"second")
937 .unwrap();
938 let other_publisher = codec
939 .encode_envelope_parts(8, 11, 3, "events", b"other")
940 .unwrap();
941 let expected_first = first.clone();
942 let expected_second = second.clone();
943 let expected_other = other_publisher.clone();
944 let inner: WireStream = Box::pin(futures::stream::iter(vec![
945 Ok(first.clone()),
946 Ok(first),
947 Ok(second),
948 Ok(other_publisher),
949 ]));
950 let mut stream = DeduplicatingStream::new(inner, codec, 16);
951
952 assert_eq!(stream.next().await.unwrap().unwrap(), expected_first);
953 assert_eq!(stream.next().await.unwrap().unwrap(), expected_second);
954 assert_eq!(stream.next().await.unwrap().unwrap(), expected_other);
955 assert!(stream.next().await.is_none());
956 }
957
958 #[test]
959 fn publisher_ids_survive_a_json_number_round_trip() {
960 let unsafe_id: u64 = 13_584_172_880_116_487_724;
963 assert!(unsafe_id > MAX_JSON_SAFE_PUBLISHER_ID);
964 assert_ne!(unsafe_id as f64 as u64, unsafe_id);
965
966 for random_id in [0, 1, u64::MAX, unsafe_id, 6_633_287_539_119_378] {
967 let publisher_id = discovery_safe_publisher_id(random_id);
968 assert!(
969 publisher_id <= MAX_JSON_SAFE_PUBLISHER_ID,
970 "publisher ID {publisher_id} exceeds the JSON-safe integer range"
971 );
972 assert_eq!(
973 publisher_id as f64 as u64, publisher_id,
974 "publisher ID {publisher_id} must survive an f64 round trip"
975 );
976 }
977 }
978
979 #[test]
980 fn direct_zmq_topology_selection_is_narrow() {
981 let lookup = |url: Option<&str>, enabled: Option<&str>| {
982 let url = url.map(std::ffi::OsString::from);
983 let enabled = enabled.map(std::ffi::OsString::from);
984 move |key: &str| match key {
985 broker_env::DYN_ZMQ_BROKER_URL => url.clone(),
986 broker_env::DYN_ZMQ_BROKER_ENABLED => enabled.clone(),
987 _ => None,
988 }
989 };
990
991 assert!(uses_direct_zmq_from_lookup(
992 EventTransportKind::Zmq,
993 lookup(None, None)
994 ));
995 assert!(!uses_direct_zmq_from_lookup(
996 EventTransportKind::Zmq,
997 lookup(Some("xsub=tcp://broker:5555,xpub=tcp://broker:5556"), None)
998 ));
999 assert!(!uses_direct_zmq_from_lookup(
1000 EventTransportKind::Zmq,
1001 lookup(None, Some("true"))
1002 ));
1003 assert!(uses_direct_zmq_from_lookup(
1004 EventTransportKind::Zmq,
1005 lookup(None, Some("false"))
1006 ));
1007 assert!(!uses_direct_zmq_from_lookup(
1008 EventTransportKind::Nats,
1009 lookup(None, None)
1010 ));
1011 }
1012
1013 #[tokio::test]
1014 async fn direct_zmq_endpoint_scopes_are_isolated() {
1015 temp_env::async_with_vars(
1016 [
1017 (broker_env::DYN_ZMQ_BROKER_URL, None::<&str>),
1018 (broker_env::DYN_ZMQ_BROKER_ENABLED, None::<&str>),
1019 ],
1020 async {
1021 let runtime = crate::Runtime::from_current().expect("create runtime handle");
1022 let drt = DistributedRuntime::new(
1023 runtime,
1024 crate::distributed::DistributedConfig::process_local(),
1025 )
1026 .await
1027 .expect("create distributed runtime");
1028 let component = drt
1029 .namespace("endpoint-event-isolation-test")
1030 .expect("create namespace")
1031 .component("worker")
1032 .expect("create component");
1033 let endpoint_a = component.endpoint("a");
1034 let endpoint_b = component.endpoint("b");
1035
1036 let publisher_a = EventPublisher::for_endpoint_with_transport(
1037 &endpoint_a,
1038 "events",
1039 EventTransportKind::Zmq,
1040 )
1041 .await
1042 .expect("create endpoint A publisher");
1043 let publisher_b = EventPublisher::for_endpoint_with_transport(
1044 &endpoint_b,
1045 "events",
1046 EventTransportKind::Zmq,
1047 )
1048 .await
1049 .expect("create endpoint B publisher");
1050 let mut subscriber_a = EventSubscriber::for_endpoint_with_transport(
1051 &endpoint_a,
1052 "events",
1053 EventTransportKind::Zmq,
1054 )
1055 .await
1056 .expect("create endpoint A subscriber");
1057 let mut subscriber_b = EventSubscriber::for_endpoint_with_transport(
1058 &endpoint_b,
1059 "events",
1060 EventTransportKind::Zmq,
1061 )
1062 .await
1063 .expect("create endpoint B subscriber");
1064
1065 let receive = async {
1066 loop {
1067 publisher_a
1068 .publish_bytes(vec![0xa1])
1069 .await
1070 .expect("publish endpoint A event");
1071 publisher_b
1072 .publish_bytes(vec![0xb2])
1073 .await
1074 .expect("publish endpoint B event");
1075
1076 let a = tokio::time::timeout(
1077 std::time::Duration::from_millis(100),
1078 subscriber_a.next(),
1079 )
1080 .await;
1081 let b = tokio::time::timeout(
1082 std::time::Duration::from_millis(100),
1083 subscriber_b.next(),
1084 )
1085 .await;
1086 if let (Ok(Some(Ok(a))), Ok(Some(Ok(b)))) = (a, b) {
1087 assert_eq!(a.publisher_id, publisher_a.publisher_id());
1088 assert_eq!(a.payload.as_ref(), &[0xa1]);
1089 assert_eq!(b.publisher_id, publisher_b.publisher_id());
1090 assert_eq!(b.payload.as_ref(), &[0xb2]);
1091 break;
1092 }
1093
1094 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1095 }
1096 };
1097
1098 tokio::time::timeout(std::time::Duration::from_secs(5), receive)
1099 .await
1100 .expect("each endpoint subscriber should receive only its own publisher");
1101 },
1102 )
1103 .await;
1104 }
1105
1106 #[tokio::test]
1107 async fn direct_zmq_publishers_in_one_endpoint_fan_into_one_subscriber() {
1108 temp_env::async_with_vars(
1109 [
1110 (broker_env::DYN_ZMQ_BROKER_URL, None::<&str>),
1111 (broker_env::DYN_ZMQ_BROKER_ENABLED, None::<&str>),
1112 ],
1113 async {
1114 let runtime = crate::Runtime::from_current().expect("create runtime handle");
1115 let drt = DistributedRuntime::new(
1116 runtime,
1117 crate::distributed::DistributedConfig::process_local(),
1118 )
1119 .await
1120 .expect("create distributed runtime");
1121 let endpoint = drt
1122 .namespace("endpoint-event-fan-in-test")
1123 .expect("create namespace")
1124 .component("worker")
1125 .expect("create component")
1126 .endpoint("generate");
1127 let publisher_a = EventPublisher::for_endpoint_with_transport(
1128 &endpoint,
1129 "events",
1130 EventTransportKind::Zmq,
1131 )
1132 .await
1133 .expect("create first publisher");
1134 let publisher_b = EventPublisher::for_endpoint_with_transport(
1135 &endpoint,
1136 "events",
1137 EventTransportKind::Zmq,
1138 )
1139 .await
1140 .expect("create second publisher");
1141 let mut subscriber = EventSubscriber::for_endpoint_with_transport(
1142 &endpoint,
1143 "events",
1144 EventTransportKind::Zmq,
1145 )
1146 .await
1147 .expect("create subscriber");
1148
1149 let receive = async {
1150 let mut publisher_ids = std::collections::HashSet::new();
1151 while publisher_ids.len() < 2 {
1152 publisher_a.publish_bytes(vec![0xa1]).await.unwrap();
1153 publisher_b.publish_bytes(vec![0xb2]).await.unwrap();
1154 if let Ok(Some(Ok(envelope))) = tokio::time::timeout(
1155 std::time::Duration::from_millis(100),
1156 subscriber.next(),
1157 )
1158 .await
1159 {
1160 publisher_ids.insert(envelope.publisher_id);
1161 }
1162 }
1163 assert_eq!(
1164 publisher_ids,
1165 std::collections::HashSet::from([
1166 publisher_a.publisher_id(),
1167 publisher_b.publisher_id(),
1168 ])
1169 );
1170 for publisher_id in [publisher_a.publisher_id(), publisher_b.publisher_id()] {
1171 assert!(
1172 publisher_id <= MAX_JSON_SAFE_PUBLISHER_ID,
1173 "publisher ID {publisher_id} exceeds the JSON-safe integer range"
1174 );
1175 }
1176 };
1177 tokio::time::timeout(std::time::Duration::from_secs(5), receive)
1178 .await
1179 .expect("subscriber should receive both endpoint publishers");
1180 },
1181 )
1182 .await;
1183 }
1184
1185 #[tokio::test]
1186 async fn same_topic_publishers_are_independent_across_recreation() {
1187 temp_env::async_with_vars(
1188 [
1189 (broker_env::DYN_ZMQ_BROKER_URL, None::<&str>),
1190 (broker_env::DYN_ZMQ_BROKER_ENABLED, None::<&str>),
1191 ],
1192 async {
1193 let runtime = crate::Runtime::from_current().expect("create runtime handle");
1194 let drt = DistributedRuntime::new(
1195 runtime,
1196 crate::distributed::DistributedConfig::process_local(),
1197 )
1198 .await
1199 .expect("create distributed runtime");
1200 let component = drt
1201 .namespace("event-publisher-test")
1202 .expect("create namespace")
1203 .component("worker")
1204 .expect("create component");
1205
1206 let publisher_a = EventPublisher::for_component_with_transport(
1207 &component,
1208 "events",
1209 EventTransportKind::Zmq,
1210 )
1211 .await
1212 .expect("create first publisher");
1213 let publisher_b = EventPublisher::for_component_with_transport(
1214 &component,
1215 "events",
1216 EventTransportKind::Zmq,
1217 )
1218 .await
1219 .expect("create second publisher");
1220 let publisher_a_id = publisher_a.publisher_id();
1221 let publisher_b_id = publisher_b.publisher_id();
1222
1223 assert_ne!(publisher_a_id, publisher_b_id);
1224
1225 let query = DiscoveryQuery::EventChannels(EventChannelQuery::topic(
1226 "event-publisher-test",
1227 "worker",
1228 "events",
1229 ));
1230 let instances = drt
1231 .discovery()
1232 .list(query.clone())
1233 .await
1234 .expect("list event publishers");
1235 assert_eq!(instances.len(), 2);
1236 assert!(
1237 instances
1238 .iter()
1239 .any(|instance| instance.instance_id() == publisher_a_id)
1240 );
1241 assert!(
1242 instances
1243 .iter()
1244 .any(|instance| instance.instance_id() == publisher_b_id)
1245 );
1246
1247 let mut subscriber = EventSubscriber::for_component_with_transport(
1248 &component,
1249 "events",
1250 EventTransportKind::Zmq,
1251 )
1252 .await
1253 .expect("create subscriber");
1254 let mut received_a = false;
1255 let mut received_b = false;
1256
1257 tokio::time::timeout(std::time::Duration::from_secs(5), async {
1258 while !received_a || !received_b {
1259 publisher_a
1260 .publish_bytes(vec![0xa1])
1261 .await
1262 .expect("publish from first publisher");
1263 publisher_b
1264 .publish_bytes(vec![0xb2])
1265 .await
1266 .expect("publish from second publisher");
1267
1268 if let Ok(Some(envelope)) = tokio::time::timeout(
1269 std::time::Duration::from_millis(100),
1270 subscriber.next(),
1271 )
1272 .await
1273 {
1274 let envelope = envelope.expect("receive event envelope");
1275 if envelope.publisher_id == publisher_a_id {
1276 assert_eq!(envelope.payload.as_ref(), &[0xa1]);
1277 received_a = true;
1278 } else if envelope.publisher_id == publisher_b_id {
1279 assert_eq!(envelope.payload.as_ref(), &[0xb2]);
1280 received_b = true;
1281 } else {
1282 panic!("event from unexpected publisher");
1283 }
1284 }
1285
1286 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1287 }
1288 })
1289 .await
1290 .expect("subscriber should receive events from both publishers");
1291
1292 drop(publisher_a);
1293 let publisher_a_recreated = EventPublisher::for_component_with_transport(
1294 &component,
1295 "events",
1296 EventTransportKind::Zmq,
1297 )
1298 .await
1299 .expect("recreate first publisher");
1300 let publisher_a_recreated_id = publisher_a_recreated.publisher_id();
1301
1302 assert_ne!(publisher_a_recreated_id, publisher_a_id);
1303 assert_ne!(publisher_a_recreated_id, publisher_b_id);
1304 assert_eq!(
1305 publisher_a_recreated.sequence.load(Ordering::SeqCst),
1306 0,
1307 "a recreated publisher starts a new sequence space"
1308 );
1309
1310 tokio::time::timeout(std::time::Duration::from_secs(1), async {
1311 loop {
1312 let instances = drt
1313 .discovery()
1314 .list(query.clone())
1315 .await
1316 .expect("list event publishers after recreation");
1317 if instances.len() == 2
1318 && instances
1319 .iter()
1320 .any(|instance| instance.instance_id() == publisher_b_id)
1321 && instances
1322 .iter()
1323 .any(|instance| instance.instance_id() == publisher_a_recreated_id)
1324 {
1325 break;
1326 }
1327 tokio::task::yield_now().await;
1328 }
1329 })
1330 .await
1331 .expect("old publisher should unregister without removing current publishers");
1332
1333 let mut received_b_after_recreation = false;
1334 let mut received_recreated_a = false;
1335
1336 tokio::time::timeout(std::time::Duration::from_secs(5), async {
1337 while !received_b_after_recreation || !received_recreated_a {
1338 publisher_b
1339 .publish_bytes(vec![0xb3])
1340 .await
1341 .expect("publish from second publisher after recreation");
1342 publisher_a_recreated
1343 .publish_bytes(vec![0xa2])
1344 .await
1345 .expect("publish from recreated publisher");
1346
1347 if let Ok(Some(envelope)) = tokio::time::timeout(
1348 std::time::Duration::from_millis(100),
1349 subscriber.next(),
1350 )
1351 .await
1352 {
1353 let envelope = envelope.expect("receive event envelope after drop");
1354 if envelope.publisher_id == publisher_b_id
1355 && envelope.payload.as_ref() == [0xb3]
1356 {
1357 received_b_after_recreation = true;
1358 } else if envelope.publisher_id == publisher_a_recreated_id
1359 && envelope.payload.as_ref() == [0xa2]
1360 {
1361 received_recreated_a = true;
1362 }
1363 }
1364
1365 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1366 }
1367 })
1368 .await
1369 .expect("subscriber should receive from surviving and recreated publishers");
1370 },
1371 )
1372 .await;
1373 }
1374
1375 #[tokio::test]
1376 async fn dropped_publisher_unregister_completes_within_graceful_shutdown() {
1377 temp_env::async_with_vars(
1378 [
1379 (broker_env::DYN_ZMQ_BROKER_URL, None::<&str>),
1380 (broker_env::DYN_ZMQ_BROKER_ENABLED, None::<&str>),
1381 ],
1382 async {
1383 let runtime = crate::Runtime::from_current().expect("create runtime handle");
1384 let drt = DistributedRuntime::new(
1385 runtime,
1386 crate::distributed::DistributedConfig::process_local(),
1387 )
1388 .await
1389 .expect("create distributed runtime");
1390 let component = drt
1391 .namespace("event-publisher-shutdown-test")
1392 .expect("create namespace")
1393 .component("worker")
1394 .expect("create component");
1395
1396 let publisher = EventPublisher::for_component_with_transport(
1397 &component,
1398 "events",
1399 EventTransportKind::Zmq,
1400 )
1401 .await
1402 .expect("create publisher");
1403 let publisher_id = publisher.publisher_id();
1404
1405 let query = DiscoveryQuery::EventChannels(EventChannelQuery::topic(
1406 "event-publisher-shutdown-test",
1407 "worker",
1408 "events",
1409 ));
1410 let instances = drt
1411 .discovery()
1412 .list(query.clone())
1413 .await
1414 .expect("list event publishers");
1415 assert_eq!(instances.len(), 1);
1416 assert_eq!(instances[0].instance_id(), publisher_id);
1417
1418 let tracker = drt.graceful_shutdown_tracker();
1419 assert_eq!(tracker.get_count(), 0);
1420
1421 let main_token = drt.runtime().primary_token();
1422 let endpoint_token = drt.runtime().child_token();
1423
1424 drop(publisher);
1428 assert_eq!(
1429 tracker.get_count(),
1430 1,
1431 "dropping a publisher must register its unregister work with the graceful-shutdown tracker"
1432 );
1433
1434 drt.runtime().shutdown();
1435
1436 tokio::time::timeout(std::time::Duration::from_secs(5), main_token.cancelled())
1437 .await
1438 .expect("graceful shutdown should complete once the unregister task finishes");
1439
1440 assert!(endpoint_token.is_cancelled());
1441 assert_eq!(
1442 tracker.get_count(),
1443 0,
1444 "the unregister task must release its graceful-shutdown guard"
1445 );
1446
1447 let instances = drt
1450 .discovery()
1451 .list(query)
1452 .await
1453 .expect("list event publishers after shutdown");
1454 assert!(
1455 instances.is_empty(),
1456 "unregister must complete within the graceful-shutdown window"
1457 );
1458 },
1459 )
1460 .await;
1461 }
1462
1463 #[tokio::test]
1464 async fn runtime_cancellation_stops_retained_direct_zmq_subscriber() {
1465 temp_env::async_with_vars(
1466 [
1467 (broker_env::DYN_ZMQ_BROKER_URL, None::<&str>),
1468 (broker_env::DYN_ZMQ_BROKER_ENABLED, None::<&str>),
1469 ],
1470 async {
1471 let runtime = crate::Runtime::from_current().expect("create runtime handle");
1472 let drt = DistributedRuntime::new(
1473 runtime,
1474 crate::distributed::DistributedConfig::process_local(),
1475 )
1476 .await
1477 .expect("create distributed runtime");
1478 let component = drt
1479 .namespace("event-subscriber-shutdown-test")
1480 .expect("create namespace")
1481 .component("worker")
1482 .expect("create component");
1483
1484 let mut subscriber = EventSubscriber::for_component_with_transport(
1485 &component,
1486 "events",
1487 EventTransportKind::Zmq,
1488 )
1489 .await
1490 .expect("create subscriber");
1491
1492 drt.primary_token().cancel();
1493
1494 let next =
1495 tokio::time::timeout(std::time::Duration::from_secs(1), subscriber.next())
1496 .await
1497 .expect("runtime cancellation should stop the subscriber stream");
1498 assert!(next.is_none());
1499 },
1500 )
1501 .await;
1502 }
1503
1504 #[test]
1505 fn test_event_scope_subject_prefix() {
1506 let scopes = [
1507 (
1508 EventScope::Namespace {
1509 name: "ns.one".to_string(),
1510 },
1511 "namespace.ns%2Eone",
1512 ),
1513 (
1514 EventScope::Component {
1515 namespace: "ns.one".to_string(),
1516 component: "worker/*".to_string(),
1517 },
1518 "namespace.ns%2Eone.component.worker%2F%2A",
1519 ),
1520 (
1521 EventScope::Endpoint {
1522 endpoint: EndpointId {
1523 namespace: "ns.one".to_string(),
1524 component: "worker/*".to_string(),
1525 name: "generate.>".to_string(),
1526 },
1527 },
1528 "namespace.ns%2Eone.component.worker%2F%2A.endpoint.generate%2E%3E",
1529 ),
1530 ];
1531
1532 for (scope, expected_prefix) in scopes {
1533 assert_eq!(scope.subject_prefix(), expected_prefix);
1534 assert_eq!(
1535 scope.subject("kv.events/*"),
1536 format!("{expected_prefix}.kv%2Eevents%2F%2A")
1537 );
1538 }
1539 }
1540
1541 #[test]
1542 fn test_event_scope_accessors() {
1543 let ns_scope = EventScope::Namespace {
1544 name: "my-ns".to_string(),
1545 };
1546 assert_eq!(ns_scope.namespace(), "my-ns");
1547 assert_eq!(ns_scope.component(), None);
1548
1549 let comp_scope = EventScope::Component {
1550 namespace: "my-ns".to_string(),
1551 component: "my-comp".to_string(),
1552 };
1553 assert_eq!(comp_scope.namespace(), "my-ns");
1554 assert_eq!(comp_scope.component(), Some("my-comp"));
1555 }
1556
1557 #[test]
1558 fn test_event_envelope_serde() {
1559 let envelope = EventEnvelope {
1560 publisher_id: 42,
1561 sequence: 10,
1562 published_at: 1700000000000,
1563 topic: "test-topic".to_string(),
1564 payload: Bytes::from("test data"),
1565 };
1566
1567 let json = serde_json::to_string(&envelope).expect("serialize");
1568 let deserialized: EventEnvelope = serde_json::from_str(&json).expect("deserialize");
1569
1570 assert_eq!(deserialized.publisher_id, 42);
1571 assert_eq!(deserialized.sequence, 10);
1572 assert_eq!(deserialized.published_at, 1700000000000);
1573 assert_eq!(deserialized.topic, "test-topic");
1574 assert_eq!(deserialized.payload, Bytes::from("test data"));
1575 }
1576}