Skip to main content

dynamo_runtime/transports/event_plane/
mod.rs

1// SPDX-FileCopyrightText: Copyright (c) 2024-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Generic Event Plane for transport-agnostic pub/sub communication.
5
6mod 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::{ZmqPubTransport, ZmqSubTransport};
20
21// Re-export transport kind from discovery for convenience
22pub use crate::discovery::EventTransportKind;
23
24use std::num::NonZeroUsize;
25use std::sync::Arc;
26use std::sync::atomic::{AtomicU64, Ordering};
27use std::time::{SystemTime, UNIX_EPOCH};
28
29use anyhow::Result;
30use bytes::Bytes;
31use futures::{Stream, StreamExt};
32use lru::LruCache;
33use rand::TryRngCore;
34use serde::Serialize;
35use serde::de::DeserializeOwned;
36use std::pin::Pin;
37use std::task::{Context, Poll};
38
39use crate::DistributedRuntime;
40use crate::component::{Component, Namespace};
41use crate::discovery::{
42    Discovery, DiscoveryInstance, DiscoveryQuery, DiscoverySpec, EventChannelQuery, EventTransport,
43};
44use crate::traits::DistributedRuntimeProvider;
45use crate::utils::local_ip_for_advertise;
46
47/// Scope of the event plane - determines the subject prefix for pub/sub.
48#[derive(Debug, Clone)]
49pub enum EventScope {
50    /// Namespace-level scope: `namespace.{name}`
51    Namespace { name: String },
52    /// Component-level scope: `namespace.{namespace}.component.{component}`
53    Component {
54        namespace: String,
55        component: String,
56    },
57}
58
59impl EventScope {
60    /// Returns the subject prefix for this scope.
61    pub fn subject_prefix(&self) -> String {
62        match self {
63            EventScope::Namespace { name } => format!("namespace.{}", name),
64            EventScope::Component {
65                namespace,
66                component,
67            } => {
68                format!("namespace.{}.component.{}", namespace, component)
69            }
70        }
71    }
72
73    /// Get the namespace name
74    pub fn namespace(&self) -> &str {
75        match self {
76            EventScope::Namespace { name } => name,
77            EventScope::Component { namespace, .. } => namespace,
78        }
79    }
80
81    /// Get the component name (if component-scoped)
82    pub fn component(&self) -> Option<&str> {
83        match self {
84            EventScope::Namespace { .. } => None,
85            EventScope::Component { component, .. } => Some(component),
86        }
87    }
88}
89
90// ============================================================================
91// Broker Resolution Logic
92// ============================================================================
93
94/// Broker endpoints for ZMQ broker mode
95#[derive(Debug, Clone)]
96struct BrokerEndpoints {
97    xsub_endpoints: Vec<String>,
98    xpub_endpoints: Vec<String>,
99}
100
101/// Resolve ZMQ broker endpoints from environment or discovery
102/// Returns None if broker mode is not configured (direct mode)
103async fn resolve_zmq_broker(
104    drt: &DistributedRuntime,
105    scope: &EventScope,
106) -> Result<Option<BrokerEndpoints>> {
107    // Priority 1: Explicit URL from DYN_ZMQ_BROKER_URL
108    if let Ok(broker_url) =
109        std::env::var(crate::config::environment_names::zmq_broker::DYN_ZMQ_BROKER_URL)
110    {
111        let (xsub_endpoints, xpub_endpoints) = parse_broker_url(&broker_url)?;
112        tracing::info!(
113            num_xsub = xsub_endpoints.len(),
114            num_xpub = xpub_endpoints.len(),
115            "Using explicit ZMQ broker URL"
116        );
117        return Ok(Some(BrokerEndpoints {
118            xsub_endpoints,
119            xpub_endpoints,
120        }));
121    }
122
123    // Priority 2: Discovery-based lookup if DYN_ZMQ_BROKER_ENABLED=true
124    if std::env::var(crate::config::environment_names::zmq_broker::DYN_ZMQ_BROKER_ENABLED)
125        .unwrap_or_default()
126        == "true"
127    {
128        let query = DiscoveryQuery::EventChannels(EventChannelQuery::component(
129            scope.namespace().to_string(),
130            "zmq_broker".to_string(),
131        ));
132
133        let instances = drt.discovery().list(query).await?;
134
135        // Collect all broker instances (multiple brokers for HA)
136        let mut xsub_endpoints = Vec::new();
137        let mut xpub_endpoints = Vec::new();
138
139        for instance in instances {
140            if let DiscoveryInstance::EventChannel { transport, .. } = instance
141                && let EventTransport::ZmqBroker {
142                    xsub_endpoints: xsubs,
143                    xpub_endpoints: xpubs,
144                } = transport
145            {
146                xsub_endpoints.extend(xsubs);
147                xpub_endpoints.extend(xpubs);
148            }
149        }
150
151        if xsub_endpoints.is_empty() {
152            anyhow::bail!(
153                "DYN_ZMQ_BROKER_ENABLED=true but no broker found in discovery for namespace '{}'",
154                scope.namespace()
155            );
156        }
157
158        tracing::info!(
159            num_brokers = xsub_endpoints.len(),
160            "Discovered ZMQ brokers from discovery plane"
161        );
162
163        return Ok(Some(BrokerEndpoints {
164            xsub_endpoints,
165            xpub_endpoints,
166        }));
167    }
168
169    // No broker configured - use direct mode
170    Ok(None)
171}
172
173/// Parse broker URL format: "xsub=tcp://host1:5555;tcp://host2:5555 , xpub=tcp://host1:5556;tcp://host2:5556"
174fn parse_broker_url(url: &str) -> Result<(Vec<String>, Vec<String>)> {
175    let parts: Vec<&str> = url.split(',').map(|s| s.trim()).collect();
176    if parts.len() != 2 {
177        anyhow::bail!(
178            "Invalid broker URL format. Expected 'xsub=<urls> , xpub=<urls>', got: {}",
179            url
180        );
181    }
182
183    let mut xsub_endpoints = Vec::new();
184    let mut xpub_endpoints = Vec::new();
185
186    for part in parts {
187        if let Some(urls_str) = part.strip_prefix("xsub=") {
188            xsub_endpoints = urls_str
189                .split(';')
190                .map(|s| s.trim().to_string())
191                .filter(|s| !s.is_empty())
192                .collect();
193        } else if let Some(urls_str) = part.strip_prefix("xpub=") {
194            xpub_endpoints = urls_str
195                .split(';')
196                .map(|s| s.trim().to_string())
197                .filter(|s| !s.is_empty())
198                .collect();
199        } else {
200            anyhow::bail!(
201                "Invalid broker URL part. Expected 'xsub=' or 'xpub=' prefix, got: {}",
202                part
203            );
204        }
205    }
206
207    if xsub_endpoints.is_empty() || xpub_endpoints.is_empty() {
208        anyhow::bail!(
209            "Broker URL must contain at least one xsub and one xpub endpoint. Got xsub={:?}, xpub={:?}",
210            xsub_endpoints,
211            xpub_endpoints
212        );
213    }
214
215    Ok((xsub_endpoints, xpub_endpoints))
216}
217
218/// Deduplicates events based on (publisher_id, sequence) tuple
219/// Required when connecting to multiple brokers in HA mode
220struct DeduplicatingStream {
221    inner: WireStream,
222    codec: Arc<Codec>,
223    seen_events: LruCache<(u64, u64), ()>, // (publisher_id, sequence) -> ()
224}
225
226impl DeduplicatingStream {
227    fn new(inner: WireStream, codec: Arc<Codec>, cache_size: usize) -> Self {
228        Self {
229            inner,
230            codec,
231            seen_events: LruCache::new(
232                NonZeroUsize::new(cache_size).expect("cache_size must be non-zero"),
233            ),
234        }
235    }
236}
237
238impl Stream for DeduplicatingStream {
239    type Item = Result<Bytes>;
240
241    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
242        loop {
243            match Pin::new(&mut self.inner).poll_next(cx) {
244                Poll::Ready(Some(Ok(bytes))) => {
245                    // Decode envelope to extract publisher_id and sequence
246                    match self.codec.decode_envelope(&bytes) {
247                        Ok(envelope) => {
248                            let key = (envelope.publisher_id, envelope.sequence);
249
250                            // Check if we've seen this event before
251                            if self.seen_events.contains(&key) {
252                                // Duplicate - skip and continue loop
253                                tracing::debug!(
254                                    publisher_id = envelope.publisher_id,
255                                    sequence = envelope.sequence,
256                                    "Filtered duplicate event from multi-broker setup"
257                                );
258                                continue;
259                            }
260
261                            // New event - record and return
262                            self.seen_events.put(key, ());
263                            return Poll::Ready(Some(Ok(bytes)));
264                        }
265                        Err(e) => {
266                            tracing::warn!(error = %e, "Failed to decode envelope for deduplication");
267                            return Poll::Ready(Some(Err(e)));
268                        }
269                    }
270                }
271                Poll::Ready(Some(Err(e))) => return Poll::Ready(Some(Err(e))),
272                Poll::Ready(None) => return Poll::Ready(None),
273                Poll::Pending => return Poll::Pending,
274            }
275        }
276    }
277}
278
279/// Event publisher for a specific topic.
280pub struct EventPublisher {
281    transport_kind: EventTransportKind,
282    topic: String,
283    subject: String,
284    publisher_id: u64,
285    sequence: AtomicU64,
286    tx: Arc<dyn EventTransportTx>,
287    codec: Arc<Codec>,
288    runtime_handle: tokio::runtime::Handle,
289    // Keeps unregister work in graceful-shutdown Phase 2 when dropped before shutdown.
290    graceful_shutdown_tracker: Arc<crate::utils::GracefulShutdownTracker>,
291    /// Discovery client and registered instance for unregistration on drop
292    discovery_client: Option<Arc<dyn Discovery>>,
293    discovery_instance: Option<crate::discovery::DiscoveryInstance>,
294}
295
296impl EventPublisher {
297    /// Create a publisher for a component-scoped topic.
298    ///
299    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
300    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
301    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
302    /// Use [`for_component_with_transport`](Self::for_component_with_transport) to
303    /// override explicitly.
304    pub async fn for_component(comp: &Component, topic: impl Into<String>) -> Result<Self> {
305        let transport_kind = comp.drt().default_event_transport_kind();
306        Self::for_component_with_transport(comp, topic, transport_kind).await
307    }
308
309    /// Create a publisher with explicit transport.
310    pub async fn for_component_with_transport(
311        comp: &Component,
312        topic: impl Into<String>,
313        transport_kind: EventTransportKind,
314    ) -> Result<Self> {
315        let drt = comp.drt();
316        let scope = EventScope::Component {
317            namespace: comp.namespace().name(),
318            component: comp.name().to_string(),
319        };
320        Self::new_internal(drt, scope, topic.into(), transport_kind).await
321    }
322
323    /// Create a publisher for a namespace-scoped topic.
324    ///
325    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
326    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
327    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
328    /// Use [`for_namespace_with_transport`](Self::for_namespace_with_transport) to
329    /// override explicitly.
330    pub async fn for_namespace(ns: &Namespace, topic: impl Into<String>) -> Result<Self> {
331        let transport_kind = ns.drt().default_event_transport_kind();
332        Self::for_namespace_with_transport(ns, topic, transport_kind).await
333    }
334
335    /// Create a namespace publisher with explicit transport.
336    pub async fn for_namespace_with_transport(
337        ns: &Namespace,
338        topic: impl Into<String>,
339        transport_kind: EventTransportKind,
340    ) -> Result<Self> {
341        let drt = ns.drt();
342        let scope = EventScope::Namespace { name: ns.name() };
343        Self::new_internal(drt, scope, topic.into(), transport_kind).await
344    }
345
346    async fn new_internal(
347        drt: &DistributedRuntime,
348        scope: EventScope,
349        topic: String,
350        transport_kind: EventTransportKind,
351    ) -> Result<Self> {
352        // Publishers are discovery objects in their own right. A single process
353        // can host multiple publishers for the same scope/topic, each with its
354        // own ZMQ endpoint and sequence space, so the process ID is not unique
355        // enough here.
356        let publisher_id = rand::rngs::OsRng
357            .try_next_u64()
358            .map_err(|error| anyhow::anyhow!("failed to generate publisher ID: {error}"))?;
359        let discovery = Some(drt.discovery());
360        let runtime_handle = drt.runtime().secondary();
361        let subject = format!("{}.{}", scope.subject_prefix(), topic);
362        let graceful_shutdown_tracker = drt.graceful_shutdown_tracker();
363
364        // Use Msgpack codec for all transports
365        enum TransportSetup {
366            Nats(Arc<dyn EventTransportTx>, Arc<Codec>),
367            ZmqDirect(Arc<dyn EventTransportTx>, Arc<Codec>, String), // includes public endpoint
368            ZmqBroker(Arc<dyn EventTransportTx>, Arc<Codec>),
369        }
370
371        let transport_setup = match transport_kind {
372            EventTransportKind::Nats => {
373                let transport = Arc::new(nats_transport::NatsTransport::new_publisher(
374                    drt.clone(),
375                    subject.clone(),
376                ));
377                let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
378                TransportSetup::Nats(transport as Arc<dyn EventTransportTx>, codec)
379            }
380            EventTransportKind::Zmq => {
381                // Check for broker mode
382                if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
383                    // BROKER MODE: Connect to broker (single or multiple endpoints)
384                    let pub_transport = if broker.xsub_endpoints.len() == 1 {
385                        zmq_transport::ZmqPubTransport::connect(&broker.xsub_endpoints[0], &topic)
386                            .await?
387                    } else {
388                        zmq_transport::ZmqPubTransport::connect_multiple(
389                            &broker.xsub_endpoints,
390                            &topic,
391                        )
392                        .await?
393                    };
394
395                    let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
396                    TransportSetup::ZmqBroker(
397                        Arc::new(pub_transport) as Arc<dyn EventTransportTx>,
398                        codec,
399                    )
400                } else {
401                    // DIRECT MODE: Bind PUB socket
402                    let (pub_transport, actual_bind_endpoint) = std::thread::spawn({
403                        let topic = topic.clone();
404                        move || {
405                            let rt = tokio::runtime::Builder::new_current_thread()
406                                .enable_all()
407                                .build()
408                                .expect("Failed to create Tokio runtime for ZMQ");
409
410                            rt.block_on(async move {
411                                zmq_transport::ZmqPubTransport::bind("tcp://0.0.0.0:0", &topic)
412                                    .await
413                                    .expect("Failed to bind ZMQ publisher")
414                            })
415                        }
416                    })
417                    .join()
418                    .expect("Failed to join ZMQ initialization thread");
419
420                    // Get local IP for public endpoint
421                    let actual_port: u16 = actual_bind_endpoint
422                        .rsplit(':')
423                        .next()
424                        .and_then(|s| s.parse().ok())
425                        .expect("Failed to parse port from bind endpoint");
426                    let local_ip = local_ip_for_advertise();
427                    let public_endpoint = format!("tcp://{}:{}", local_ip, actual_port);
428
429                    let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
430                    TransportSetup::ZmqDirect(
431                        Arc::new(pub_transport) as Arc<dyn EventTransportTx>,
432                        codec,
433                        public_endpoint,
434                    )
435                }
436            }
437        };
438
439        // Extract transport and codec, and register if needed
440        let (tx, codec, discovery_instance) = match transport_setup {
441            TransportSetup::Nats(tx, codec) => {
442                let transport_config = EventTransport::nats(scope.subject_prefix());
443                let spec = DiscoverySpec::EventChannel {
444                    namespace: scope.namespace().to_string(),
445                    component: scope.component().unwrap_or("").to_string(),
446                    topic: topic.clone(),
447                    publisher_id,
448                    transport: transport_config,
449                };
450
451                let discovery_instance = drt.discovery().register(spec).await?;
452                tracing::info!(
453                    topic = %topic,
454                    transport = ?transport_kind,
455                    publisher_id = %publisher_id,
456                    "EventPublisher registered with discovery"
457                );
458                (tx, codec, Some(discovery_instance))
459            }
460            TransportSetup::ZmqDirect(tx, codec, public_endpoint) => {
461                let transport_config = EventTransport::zmq(public_endpoint);
462                let spec = DiscoverySpec::EventChannel {
463                    namespace: scope.namespace().to_string(),
464                    component: scope.component().unwrap_or("").to_string(),
465                    topic: topic.clone(),
466                    publisher_id,
467                    transport: transport_config,
468                };
469
470                let discovery_instance = drt.discovery().register(spec).await?;
471                tracing::info!(
472                    topic = %topic,
473                    transport = ?transport_kind,
474                    publisher_id = %publisher_id,
475                    "EventPublisher registered with discovery (direct mode)"
476                );
477                (tx, codec, Some(discovery_instance))
478            }
479            TransportSetup::ZmqBroker(tx, codec) => {
480                tracing::info!(
481                    topic = %topic,
482                    transport = ?transport_kind,
483                    "EventPublisher in broker mode - skipping discovery registration"
484                );
485                (tx, codec, None)
486            }
487        };
488
489        Ok(Self {
490            transport_kind,
491            topic,
492            subject,
493            publisher_id,
494            sequence: AtomicU64::new(0),
495            tx,
496            codec,
497            runtime_handle,
498            graceful_shutdown_tracker,
499            discovery_client: discovery,
500            discovery_instance,
501        })
502    }
503
504    /// Publish a serializable event.
505    pub async fn publish<T: Serialize + Send + Sync>(&self, event: &T) -> Result<()> {
506        let payload = self.codec.encode_payload(event)?;
507        self.publish_bytes_ref(payload.as_ref()).await
508    }
509
510    /// Publish raw bytes.
511    pub async fn publish_bytes(&self, bytes: Vec<u8>) -> Result<()> {
512        self.publish_bytes_ref(&bytes).await
513    }
514
515    /// Publish raw bytes without taking ownership of the payload buffer.
516    pub async fn publish_bytes_ref(&self, bytes: &[u8]) -> Result<()> {
517        let envelope_bytes = self.codec.encode_envelope_parts(
518            self.publisher_id,
519            self.sequence.fetch_add(1, Ordering::SeqCst),
520            current_timestamp_ms(),
521            &self.topic,
522            bytes,
523        )?;
524
525        self.tx.publish(&self.subject, envelope_bytes).await
526    }
527
528    /// Get the publisher ID.
529    pub fn publisher_id(&self) -> u64 {
530        self.publisher_id
531    }
532
533    /// Get the topic.
534    pub fn topic(&self) -> &str {
535        &self.topic
536    }
537
538    /// Get the transport kind.
539    pub fn transport_kind(&self) -> EventTransportKind {
540        self.transport_kind
541    }
542}
543
544impl Drop for EventPublisher {
545    fn drop(&mut self) {
546        // Unregister from discovery on drop
547        if let (Some(discovery), Some(instance)) =
548            (self.discovery_client.take(), self.discovery_instance.take())
549        {
550            let topic = self.topic.clone();
551            let publisher_id = instance.instance_id();
552            let runtime_handle = self.runtime_handle.clone();
553            let shutdown_guard = self.graceful_shutdown_tracker.register_task();
554
555            // Drop can run outside any Tokio context (notably via PyO3 finalizers), so use
556            // the runtime that created the publisher rather than the ambient thread state.
557            let spawn_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(move || {
558                runtime_handle.spawn(async move {
559                    let _shutdown_guard = shutdown_guard;
560                    match discovery.unregister(instance).await {
561                        Ok(()) => {
562                            tracing::info!(
563                                topic = %topic,
564                                publisher_id = %publisher_id,
565                                "EventPublisher unregistered from discovery"
566                            );
567                        }
568                        Err(e) => {
569                            tracing::warn!(
570                                topic = %topic,
571                                publisher_id = %publisher_id,
572                                error = %e,
573                                "Failed to unregister EventPublisher from discovery"
574                            );
575                        }
576                    }
577                });
578            }));
579
580            if spawn_result.is_err() {
581                tracing::warn!(
582                    topic = %self.topic,
583                    publisher_id = %publisher_id,
584                    "Skipping EventPublisher unregister during drop because the runtime is unavailable"
585                );
586            }
587        }
588    }
589}
590
591/// Event subscriber for a specific topic.
592pub struct EventSubscriber {
593    stream: EventStream,
594    #[allow(dead_code)]
595    scope: EventScope,
596    #[allow(dead_code)]
597    topic: String,
598    codec: Arc<Codec>,
599}
600
601impl EventSubscriber {
602    /// Create a subscriber for a component-scoped topic.
603    ///
604    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
605    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
606    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
607    /// Use [`for_component_with_transport`](Self::for_component_with_transport) to
608    /// override explicitly.
609    pub async fn for_component(comp: &Component, topic: impl Into<String>) -> Result<Self> {
610        let transport_kind = comp.drt().default_event_transport_kind();
611        Self::for_component_with_transport(comp, topic, transport_kind).await
612    }
613
614    /// Create a subscriber with explicit transport.
615    pub async fn for_component_with_transport(
616        comp: &Component,
617        topic: impl Into<String>,
618        transport_kind: EventTransportKind,
619    ) -> Result<Self> {
620        let drt = comp.drt();
621        let scope = EventScope::Component {
622            namespace: comp.namespace().name(),
623            component: comp.name().to_string(),
624        };
625        Self::new_internal(drt, scope, topic.into(), transport_kind).await
626    }
627
628    /// Create a subscriber for a namespace-scoped topic.
629    ///
630    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
631    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
632    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
633    /// Use [`for_namespace_with_transport`](Self::for_namespace_with_transport) to
634    /// override explicitly.
635    pub async fn for_namespace(ns: &Namespace, topic: impl Into<String>) -> Result<Self> {
636        let transport_kind = ns.drt().default_event_transport_kind();
637        Self::for_namespace_with_transport(ns, topic, transport_kind).await
638    }
639
640    /// Create a namespace subscriber with explicit transport.
641    pub async fn for_namespace_with_transport(
642        ns: &Namespace,
643        topic: impl Into<String>,
644        transport_kind: EventTransportKind,
645    ) -> Result<Self> {
646        let drt = ns.drt();
647        let scope = EventScope::Namespace { name: ns.name() };
648        Self::new_internal(drt, scope, topic.into(), transport_kind).await
649    }
650
651    async fn new_internal(
652        drt: &DistributedRuntime,
653        scope: EventScope,
654        topic: String,
655        transport_kind: EventTransportKind,
656    ) -> Result<Self> {
657        let discovery = drt.discovery();
658
659        // Use Msgpack codec for all transports
660        let (wire_stream, codec): (WireStream, Arc<Codec>) = match transport_kind {
661            EventTransportKind::Nats => {
662                let transport = nats_transport::NatsTransport::new(drt.clone());
663                let subject = format!("{}.{}", scope.subject_prefix(), topic);
664                let stream = transport.subscribe(&subject).await?;
665                let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
666                (stream, codec)
667            }
668            EventTransportKind::Zmq => {
669                // Check for broker mode
670                if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
671                    // BROKER MODE: Connect to broker's XPUB (single or multiple endpoints)
672                    let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
673
674                    let stream: WireStream = if broker.xpub_endpoints.len() == 1 {
675                        // Single broker - no deduplication needed
676                        let sub_transport = zmq_transport::ZmqSubTransport::connect_broker(
677                            &broker.xpub_endpoints[0],
678                            &topic,
679                        )
680                        .await?;
681                        sub_transport.subscribe(&topic).await?
682                    } else {
683                        // Multiple brokers - need deduplication
684                        let sub_transport =
685                            zmq_transport::ZmqSubTransport::connect_broker_multiple(
686                                &broker.xpub_endpoints,
687                                &topic,
688                            )
689                            .await?;
690                        let inner_stream = sub_transport.subscribe(&topic).await?;
691
692                        // Wrap with deduplication (default cache size: 100,000 entries)
693                        Box::pin(DeduplicatingStream::new(
694                            inner_stream,
695                            codec.clone(),
696                            100_000,
697                        ))
698                    };
699
700                    (stream, codec)
701                } else {
702                    // DIRECT MODE: Use dynamic subscriber to discover and connect to publishers
703                    let query = match &scope {
704                        EventScope::Namespace { name } => {
705                            crate::discovery::DiscoveryQuery::EventChannels(
706                                crate::discovery::EventChannelQuery::namespace(name.clone()),
707                            )
708                        }
709                        EventScope::Component {
710                            namespace,
711                            component,
712                        } => crate::discovery::DiscoveryQuery::EventChannels(
713                            crate::discovery::EventChannelQuery::topic(
714                                namespace.clone(),
715                                component.clone(),
716                                topic.clone(),
717                            ),
718                        ),
719                    };
720
721                    let subscriber =
722                        Arc::new(DynamicSubscriber::new(discovery, query, topic.clone()));
723
724                    let stream = subscriber.start_zmq().await?;
725                    let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
726                    (stream, codec)
727                }
728            }
729        };
730
731        // Filter by topic and decode envelopes
732        let topic_filter = topic.clone();
733        let codec_for_stream = codec.clone();
734        let stream = wire_stream.filter_map(move |result| {
735            let codec = codec_for_stream.clone();
736            let topic_filter = topic_filter.clone();
737            async move {
738                match result {
739                    Ok(bytes) => match codec.decode_envelope(&bytes) {
740                        Ok(envelope) => {
741                            // Filter by topic for transports that don't support native filtering
742                            if envelope.topic == topic_filter {
743                                Some(Ok(envelope))
744                            } else {
745                                None
746                            }
747                        }
748                        Err(e) => Some(Err(e)),
749                    },
750                    Err(e) => Some(Err(e)),
751                }
752            }
753        });
754
755        tracing::info!(
756            topic = %topic,
757            transport = ?transport_kind,
758            "EventSubscriber created"
759        );
760
761        Ok(Self {
762            stream: Box::pin(stream),
763            scope,
764            topic,
765            codec,
766        })
767    }
768
769    /// Get the next event envelope.
770    pub async fn next(&mut self) -> Option<Result<EventEnvelope>> {
771        self.stream.next().await
772    }
773
774    /// Subscribe with automatic deserialization.
775    pub fn typed<T: DeserializeOwned + Send + 'static>(self) -> TypedEventSubscriber<T> {
776        TypedEventSubscriber {
777            stream: self.stream,
778            codec: self.codec,
779            _marker: std::marker::PhantomData,
780        }
781    }
782}
783
784/// Typed event subscriber that deserializes payloads.
785pub struct TypedEventSubscriber<T> {
786    stream: EventStream,
787    codec: Arc<Codec>,
788    _marker: std::marker::PhantomData<T>,
789}
790
791impl<T: DeserializeOwned + Send + 'static> TypedEventSubscriber<T> {
792    /// Get the next typed event with its envelope.
793    pub async fn next(&mut self) -> Option<Result<(EventEnvelope, T)>> {
794        std::future::poll_fn(|cx| self.poll_next(cx)).await
795    }
796
797    /// Poll for the next typed event.
798    pub fn poll_next(&mut self, cx: &mut Context<'_>) -> Poll<Option<Result<(EventEnvelope, T)>>> {
799        match self.stream.as_mut().poll_next(cx) {
800            Poll::Ready(Some(envelope)) => Poll::Ready(Some(match envelope {
801                Ok(env) => match self.codec.decode_payload(&env.payload) {
802                    Ok(typed) => Ok((env, typed)),
803                    Err(e) => Err(e),
804                },
805                Err(e) => Err(e),
806            })),
807            Poll::Ready(None) => Poll::Ready(None),
808            Poll::Pending => Poll::Pending,
809        }
810    }
811}
812
813/// Get current timestamp in milliseconds since Unix epoch.
814fn current_timestamp_ms() -> u64 {
815    SystemTime::now()
816        .duration_since(UNIX_EPOCH)
817        .map(|d| d.as_millis() as u64)
818        .unwrap_or(0)
819}
820
821#[cfg(test)]
822mod tests {
823    use super::*;
824    use crate::config::environment_names::zmq_broker as broker_env;
825
826    #[tokio::test]
827    async fn same_topic_publishers_are_independent_across_recreation() {
828        temp_env::async_with_vars(
829            [
830                (broker_env::DYN_ZMQ_BROKER_URL, None::<&str>),
831                (broker_env::DYN_ZMQ_BROKER_ENABLED, None::<&str>),
832            ],
833            async {
834                let runtime = crate::Runtime::from_current().expect("create runtime handle");
835                let drt = DistributedRuntime::new(
836                    runtime,
837                    crate::distributed::DistributedConfig::process_local(),
838                )
839                .await
840                .expect("create distributed runtime");
841                let component = drt
842                    .namespace("event-publisher-test")
843                    .expect("create namespace")
844                    .component("worker")
845                    .expect("create component");
846
847                let publisher_a = EventPublisher::for_component_with_transport(
848                    &component,
849                    "events",
850                    EventTransportKind::Zmq,
851                )
852                .await
853                .expect("create first publisher");
854                let publisher_b = EventPublisher::for_component_with_transport(
855                    &component,
856                    "events",
857                    EventTransportKind::Zmq,
858                )
859                .await
860                .expect("create second publisher");
861                let publisher_a_id = publisher_a.publisher_id();
862                let publisher_b_id = publisher_b.publisher_id();
863
864                assert_ne!(publisher_a_id, publisher_b_id);
865
866                let query = DiscoveryQuery::EventChannels(EventChannelQuery::topic(
867                    "event-publisher-test",
868                    "worker",
869                    "events",
870                ));
871                let instances = drt
872                    .discovery()
873                    .list(query.clone())
874                    .await
875                    .expect("list event publishers");
876                assert_eq!(instances.len(), 2);
877                assert!(
878                    instances
879                        .iter()
880                        .any(|instance| instance.instance_id() == publisher_a_id)
881                );
882                assert!(
883                    instances
884                        .iter()
885                        .any(|instance| instance.instance_id() == publisher_b_id)
886                );
887
888                let mut subscriber = EventSubscriber::for_component_with_transport(
889                    &component,
890                    "events",
891                    EventTransportKind::Zmq,
892                )
893                .await
894                .expect("create subscriber");
895                let mut received_a = false;
896                let mut received_b = false;
897
898                tokio::time::timeout(std::time::Duration::from_secs(5), async {
899                    while !received_a || !received_b {
900                        publisher_a
901                            .publish_bytes(vec![0xa1])
902                            .await
903                            .expect("publish from first publisher");
904                        publisher_b
905                            .publish_bytes(vec![0xb2])
906                            .await
907                            .expect("publish from second publisher");
908
909                        if let Ok(Some(envelope)) = tokio::time::timeout(
910                            std::time::Duration::from_millis(100),
911                            subscriber.next(),
912                        )
913                        .await
914                        {
915                            let envelope = envelope.expect("receive event envelope");
916                            if envelope.publisher_id == publisher_a_id {
917                                assert_eq!(envelope.payload.as_ref(), &[0xa1]);
918                                received_a = true;
919                            } else if envelope.publisher_id == publisher_b_id {
920                                assert_eq!(envelope.payload.as_ref(), &[0xb2]);
921                                received_b = true;
922                            } else {
923                                panic!("event from unexpected publisher");
924                            }
925                        }
926
927                        tokio::time::sleep(std::time::Duration::from_millis(10)).await;
928                    }
929                })
930                .await
931                .expect("subscriber should receive events from both publishers");
932
933                drop(publisher_a);
934                let publisher_a_recreated = EventPublisher::for_component_with_transport(
935                    &component,
936                    "events",
937                    EventTransportKind::Zmq,
938                )
939                .await
940                .expect("recreate first publisher");
941                let publisher_a_recreated_id = publisher_a_recreated.publisher_id();
942
943                assert_ne!(publisher_a_recreated_id, publisher_a_id);
944                assert_ne!(publisher_a_recreated_id, publisher_b_id);
945                assert_eq!(
946                    publisher_a_recreated.sequence.load(Ordering::SeqCst),
947                    0,
948                    "a recreated publisher starts a new sequence space"
949                );
950
951                tokio::time::timeout(std::time::Duration::from_secs(1), async {
952                    loop {
953                        let instances = drt
954                            .discovery()
955                            .list(query.clone())
956                            .await
957                            .expect("list event publishers after recreation");
958                        if instances.len() == 2
959                            && instances
960                                .iter()
961                                .any(|instance| instance.instance_id() == publisher_b_id)
962                            && instances
963                                .iter()
964                                .any(|instance| instance.instance_id() == publisher_a_recreated_id)
965                        {
966                            break;
967                        }
968                        tokio::task::yield_now().await;
969                    }
970                })
971                .await
972                .expect("old publisher should unregister without removing current publishers");
973
974                let mut received_b_after_recreation = false;
975                let mut received_recreated_a = false;
976
977                tokio::time::timeout(std::time::Duration::from_secs(5), async {
978                    while !received_b_after_recreation || !received_recreated_a {
979                        publisher_b
980                            .publish_bytes(vec![0xb3])
981                            .await
982                            .expect("publish from second publisher after recreation");
983                        publisher_a_recreated
984                            .publish_bytes(vec![0xa2])
985                            .await
986                            .expect("publish from recreated publisher");
987
988                        if let Ok(Some(envelope)) = tokio::time::timeout(
989                            std::time::Duration::from_millis(100),
990                            subscriber.next(),
991                        )
992                        .await
993                        {
994                            let envelope = envelope.expect("receive event envelope after drop");
995                            if envelope.publisher_id == publisher_b_id
996                                && envelope.payload.as_ref() == [0xb3]
997                            {
998                                received_b_after_recreation = true;
999                            } else if envelope.publisher_id == publisher_a_recreated_id
1000                                && envelope.payload.as_ref() == [0xa2]
1001                            {
1002                                received_recreated_a = true;
1003                            }
1004                        }
1005
1006                        tokio::time::sleep(std::time::Duration::from_millis(10)).await;
1007                    }
1008                })
1009                .await
1010                .expect("subscriber should receive from surviving and recreated publishers");
1011            },
1012        )
1013        .await;
1014    }
1015
1016    #[tokio::test]
1017    async fn dropped_publisher_unregister_completes_within_graceful_shutdown() {
1018        temp_env::async_with_vars(
1019            [
1020                (broker_env::DYN_ZMQ_BROKER_URL, None::<&str>),
1021                (broker_env::DYN_ZMQ_BROKER_ENABLED, None::<&str>),
1022            ],
1023            async {
1024                let runtime = crate::Runtime::from_current().expect("create runtime handle");
1025                let drt = DistributedRuntime::new(
1026                    runtime,
1027                    crate::distributed::DistributedConfig::process_local(),
1028                )
1029                .await
1030                .expect("create distributed runtime");
1031                let component = drt
1032                    .namespace("event-publisher-shutdown-test")
1033                    .expect("create namespace")
1034                    .component("worker")
1035                    .expect("create component");
1036
1037                let publisher = EventPublisher::for_component_with_transport(
1038                    &component,
1039                    "events",
1040                    EventTransportKind::Zmq,
1041                )
1042                .await
1043                .expect("create publisher");
1044                let publisher_id = publisher.publisher_id();
1045
1046                let query = DiscoveryQuery::EventChannels(EventChannelQuery::topic(
1047                    "event-publisher-shutdown-test",
1048                    "worker",
1049                    "events",
1050                ));
1051                let instances = drt
1052                    .discovery()
1053                    .list(query.clone())
1054                    .await
1055                    .expect("list event publishers");
1056                assert_eq!(instances.len(), 1);
1057                assert_eq!(instances[0].instance_id(), publisher_id);
1058
1059                let tracker = drt.graceful_shutdown_tracker();
1060                assert_eq!(tracker.get_count(), 0);
1061
1062                let main_token = drt.runtime().primary_token();
1063                let endpoint_token = drt.runtime().child_token();
1064
1065                // Dropping the publisher schedules the async discovery unregister.
1066                // The drop path must synchronously take a graceful-shutdown guard so
1067                // that `Runtime::shutdown` Phase 2 waits for the unregister task.
1068                drop(publisher);
1069                assert_eq!(
1070                    tracker.get_count(),
1071                    1,
1072                    "dropping a publisher must register its unregister work with the graceful-shutdown tracker"
1073                );
1074
1075                drt.runtime().shutdown();
1076
1077                tokio::time::timeout(std::time::Duration::from_secs(5), main_token.cancelled())
1078                    .await
1079                    .expect("graceful shutdown should complete once the unregister task finishes");
1080
1081                assert!(endpoint_token.is_cancelled());
1082                assert_eq!(
1083                    tracker.get_count(),
1084                    0,
1085                    "the unregister task must release its graceful-shutdown guard"
1086                );
1087
1088                // Phase 3 (main token cancellation) only runs after Phase 2 drained
1089                // the tracker, so the dropped publisher must already be unregistered.
1090                let instances = drt
1091                    .discovery()
1092                    .list(query)
1093                    .await
1094                    .expect("list event publishers after shutdown");
1095                assert!(
1096                    instances.is_empty(),
1097                    "unregister must complete within the graceful-shutdown window"
1098                );
1099            },
1100        )
1101        .await;
1102    }
1103
1104    #[test]
1105    fn test_event_scope_subject_prefix() {
1106        let ns_scope = EventScope::Namespace {
1107            name: "test-ns".to_string(),
1108        };
1109        assert_eq!(ns_scope.subject_prefix(), "namespace.test-ns");
1110
1111        let comp_scope = EventScope::Component {
1112            namespace: "test-ns".to_string(),
1113            component: "test-comp".to_string(),
1114        };
1115        assert_eq!(
1116            comp_scope.subject_prefix(),
1117            "namespace.test-ns.component.test-comp"
1118        );
1119    }
1120
1121    #[test]
1122    fn test_event_scope_accessors() {
1123        let ns_scope = EventScope::Namespace {
1124            name: "my-ns".to_string(),
1125        };
1126        assert_eq!(ns_scope.namespace(), "my-ns");
1127        assert_eq!(ns_scope.component(), None);
1128
1129        let comp_scope = EventScope::Component {
1130            namespace: "my-ns".to_string(),
1131            component: "my-comp".to_string(),
1132        };
1133        assert_eq!(comp_scope.namespace(), "my-ns");
1134        assert_eq!(comp_scope.component(), Some("my-comp"));
1135    }
1136
1137    #[test]
1138    fn test_timestamp_generation() {
1139        let ts = current_timestamp_ms();
1140
1141        // Should be after Jan 1, 2020 (1577836800000) and before Jan 1, 2100 (4102444800000)
1142        assert!(ts > 1577836800000, "Timestamp should be after 2020");
1143        assert!(ts < 4102444800000, "Timestamp should be before 2100");
1144    }
1145
1146    #[test]
1147    fn test_event_envelope_serde() {
1148        let envelope = EventEnvelope {
1149            publisher_id: 42,
1150            sequence: 10,
1151            published_at: 1700000000000,
1152            topic: "test-topic".to_string(),
1153            payload: Bytes::from("test data"),
1154        };
1155
1156        let json = serde_json::to_string(&envelope).expect("serialize");
1157        let deserialized: EventEnvelope = serde_json::from_str(&json).expect("deserialize");
1158
1159        assert_eq!(deserialized.publisher_id, 42);
1160        assert_eq!(deserialized.sequence, 10);
1161        assert_eq!(deserialized.published_at, 1700000000000);
1162        assert_eq!(deserialized.topic, "test-topic");
1163        assert_eq!(deserialized.payload, Bytes::from("test data"));
1164    }
1165}