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::{
20    ValidatedEnvelope, ValidatedZmqSource, ValidatedZmqSourceError, ZmqPubTransport,
21    ZmqSubTransport,
22};
23
24// Re-export transport kind from discovery for convenience
25pub 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// ============================================================================
53// Broker Resolution Logic
54// ============================================================================
55
56/// Broker endpoints for ZMQ broker mode
57#[derive(Debug, Clone)]
58struct BrokerEndpoints {
59    xsub_endpoints: Vec<String>,
60    xpub_endpoints: Vec<String>,
61}
62
63/// Whether the configured event plane selects direct per-publisher ZMQ sockets.
64///
65/// This intentionally answers only the topology question needed by specialized
66/// consumers. Broker discovery and connection errors remain owned by the normal
67/// [`EventSubscriber`] construction path.
68pub 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
88/// Resolve ZMQ broker endpoints from environment or discovery
89/// Returns None if broker mode is not configured (direct mode)
90async fn resolve_zmq_broker(
91    drt: &DistributedRuntime,
92    scope: &EventScope,
93) -> Result<Option<BrokerEndpoints>> {
94    // Priority 1: Explicit URL from DYN_ZMQ_BROKER_URL
95    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    // Priority 2: Discovery-based lookup if DYN_ZMQ_BROKER_ENABLED is truthy
111    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        // Collect all broker instances (multiple brokers for HA)
122        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    // No broker configured - use direct mode
156    Ok(None)
157}
158
159/// Parse broker URL format: "xsub=tcp://host1:5555;tcp://host2:5555 , xpub=tcp://host1:5556;tcp://host2:5556"
160fn 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
204/// Deduplicates events based on (publisher_id, sequence) tuple
205/// Required when connecting to multiple brokers in HA mode
206struct DeduplicatingStream {
207    inner: WireStream,
208    codec: Arc<Codec>,
209    seen_events: LruCache<(u64, u64), ()>, // (publisher_id, sequence) -> ()
210}
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                    // Decode envelope to extract publisher_id and sequence
232                    match self.codec.decode_envelope_identity(&bytes) {
233                        Ok(key) => {
234                            // Check if we've seen this event before
235                            if self.seen_events.contains(&key) {
236                                // Duplicate - skip and continue loop
237                                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                            // New event - record and return
246                            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
263/// Keep publisher IDs exactly representable in float64-backed JSON metadata.
264fn discovery_safe_publisher_id(random_id: u64) -> u64 {
265    random_id & MAX_JSON_SAFE_PUBLISHER_ID
266}
267
268/// Event publisher for a specific topic.
269pub 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    // Keeps unregister work in graceful-shutdown Phase 2 when dropped before shutdown.
279    graceful_shutdown_tracker: Arc<crate::utils::GracefulShutdownTracker>,
280    /// Discovery client and registered instance for unregistration on drop
281    discovery_client: Option<Arc<dyn Discovery>>,
282    discovery_instance: Option<crate::discovery::DiscoveryInstance>,
283}
284
285impl EventPublisher {
286    /// Create a publisher for an endpoint-scoped topic.
287    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    /// Create an endpoint-scoped publisher with explicit transport.
293    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    /// Create a publisher for an endpoint identity without constructing a local endpoint handle.
310    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    /// Create an endpoint-identity publisher with an explicit transport.
320    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    /// Create a publisher for a component-scoped topic.
338    ///
339    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
340    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
341    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
342    /// Use [`for_component_with_transport`](Self::for_component_with_transport) to
343    /// override explicitly.
344    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    /// Create a publisher with explicit transport.
350    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    /// Create a publisher for a namespace-scoped topic.
364    ///
365    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
366    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
367    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
368    /// Use [`for_namespace_with_transport`](Self::for_namespace_with_transport) to
369    /// override explicitly.
370    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    /// Create a namespace publisher with explicit transport.
376    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        // Publishers are discovery objects in their own right. A single process
393        // can host multiple publishers for the same scope/topic, each with its
394        // own ZMQ endpoint and sequence space, so the process ID is not unique
395        // enough here.
396        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        // Use Msgpack codec for all transports
407        enum TransportSetup {
408            Nats(Arc<dyn EventTransportTx>, Arc<Codec>),
409            ZmqDirect(Arc<dyn EventTransportTx>, Arc<Codec>, String), // includes public endpoint
410            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                // Check for broker mode
424                if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
425                    // BROKER MODE: Connect to broker (single or multiple endpoints)
426                    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                    // DIRECT MODE: Bind PUB socket
444                    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                    // Get local IP for public endpoint
463                    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        // Extract transport and codec, and register if needed
482        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    /// Publish a serializable event.
545    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    /// Publish raw bytes.
551    pub async fn publish_bytes(&self, bytes: Vec<u8>) -> Result<()> {
552        self.publish_bytes_ref(&bytes).await
553    }
554
555    /// Publish raw bytes without taking ownership of the payload buffer.
556    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    /// Get the publisher ID.
569    pub fn publisher_id(&self) -> u64 {
570        self.publisher_id
571    }
572
573    /// Get the topic.
574    pub fn topic(&self) -> &str {
575        &self.topic
576    }
577
578    /// Get the transport kind.
579    pub fn transport_kind(&self) -> EventTransportKind {
580        self.transport_kind
581    }
582}
583
584impl Drop for EventPublisher {
585    fn drop(&mut self) {
586        // Unregister from discovery on drop
587        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            // Drop can run outside any Tokio context (notably via PyO3 finalizers), so use
596            // the runtime that created the publisher rather than the ambient thread state.
597            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
631/// Event subscriber for a specific topic.
632pub 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    /// Create a subscriber for an endpoint-scoped topic.
643    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    /// Create an endpoint-scoped subscriber with explicit transport.
649    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    /// Create a subscriber for an endpoint identity without constructing a
659    /// local [`Endpoint`] handle.
660    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    /// Create an endpoint-identity subscriber with explicit transport.
670    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    /// Create a subscriber for a component-scoped topic.
688    ///
689    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
690    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
691    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
692    /// Use [`for_component_with_transport`](Self::for_component_with_transport) to
693    /// override explicitly.
694    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    /// Create a subscriber with explicit transport.
700    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    /// Create a subscriber for a namespace-scoped topic.
714    ///
715    /// The event transport is chosen automatically: if `DYN_EVENT_PLANE` is set that
716    /// value is used; otherwise the runtime's default is used (ZMQ for local backends
717    /// such as `file`/`mem`, NATS for distributed backends such as `etcd`/`kubernetes`).
718    /// Use [`for_namespace_with_transport`](Self::for_namespace_with_transport) to
719    /// override explicitly.
720    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    /// Create a namespace subscriber with explicit transport.
726    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        // Use Msgpack codec for all transports
746        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                // Check for broker mode
755                if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
756                    // BROKER MODE: Connect to broker's XPUB (single or multiple endpoints)
757                    let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
758
759                    let stream: WireStream = if broker.xpub_endpoints.len() == 1 {
760                        // One EventSubscriber has one consumer. Poll the ZMQ socket
761                        // directly instead of forwarding through a lossy broadcast channel.
762                        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                        // Multiple brokers - need deduplication
770                        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                        // Wrap with deduplication (default cache size: 100,000 entries)
781                        Box::pin(DeduplicatingStream::new(
782                            inner_stream,
783                            codec.clone(),
784                            100_000,
785                        ))
786                    };
787
788                    (stream, codec)
789                } else {
790                    // DIRECT MODE: Use dynamic subscriber to discover and connect to publishers
791                    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        // Filter by topic and decode envelopes
835        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                            // Filter by topic for transports that don't support native filtering
845                            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    /// Get the next event envelope.
873    pub async fn next(&mut self) -> Option<Result<EventEnvelope>> {
874        self.stream.next().await
875    }
876
877    /// Subscribe with automatic deserialization.
878    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
887/// Typed event subscriber that deserializes payloads.
888pub 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    /// Get the next typed event with its envelope.
896    pub async fn next(&mut self) -> Option<Result<(EventEnvelope, T)>> {
897        std::future::poll_fn(|cx| self.poll_next(cx)).await
898    }
899
900    /// Poll for the next typed event.
901    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
916/// Get current timestamp in milliseconds since Unix epoch.
917fn 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        // This historical full-range ID is not JSON-safe. It was observed
961        // rounding to 13584172880116488000 after a discovery round trip.
962        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                // Dropping the publisher schedules the async discovery unregister.
1425                // The drop path must synchronously take a graceful-shutdown guard so
1426                // that `Runtime::shutdown` Phase 2 waits for the unregister task.
1427                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                // Phase 3 (main token cancellation) only runs after Phase 2 drained
1448                // the tracker, so the dropped publisher must already be unregistered.
1449                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}