1mod codec;
7mod dynamic_subscriber;
8mod frame;
9mod nats_transport;
10mod traits;
11mod transport;
12pub mod zmq_transport;
13
14pub use codec::{Codec, MsgpackCodec};
15pub use dynamic_subscriber::DynamicSubscriber;
16pub use frame::{FRAME_HEADER_SIZE, FRAME_VERSION, Frame, FrameError, FrameHeader};
17pub use traits::{EventEnvelope, EventStream, TypedEventStream};
18pub use transport::{EventTransportRx, EventTransportTx, WireStream};
19pub use zmq_transport::{ZmqPubTransport, ZmqSubTransport};
20
21pub 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#[derive(Debug, Clone)]
49pub enum EventScope {
50 Namespace { name: String },
52 Component {
54 namespace: String,
55 component: String,
56 },
57}
58
59impl EventScope {
60 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 pub fn namespace(&self) -> &str {
75 match self {
76 EventScope::Namespace { name } => name,
77 EventScope::Component { namespace, .. } => namespace,
78 }
79 }
80
81 pub fn component(&self) -> Option<&str> {
83 match self {
84 EventScope::Namespace { .. } => None,
85 EventScope::Component { component, .. } => Some(component),
86 }
87 }
88}
89
90#[derive(Debug, Clone)]
96struct BrokerEndpoints {
97 xsub_endpoints: Vec<String>,
98 xpub_endpoints: Vec<String>,
99}
100
101async fn resolve_zmq_broker(
104 drt: &DistributedRuntime,
105 scope: &EventScope,
106) -> Result<Option<BrokerEndpoints>> {
107 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 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 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 Ok(None)
171}
172
173fn 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
218struct DeduplicatingStream {
221 inner: WireStream,
222 codec: Arc<Codec>,
223 seen_events: LruCache<(u64, u64), ()>, }
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 match self.codec.decode_envelope(&bytes) {
247 Ok(envelope) => {
248 let key = (envelope.publisher_id, envelope.sequence);
249
250 if self.seen_events.contains(&key) {
252 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 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
279pub 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 graceful_shutdown_tracker: Arc<crate::utils::GracefulShutdownTracker>,
291 discovery_client: Option<Arc<dyn Discovery>>,
293 discovery_instance: Option<crate::discovery::DiscoveryInstance>,
294}
295
296impl EventPublisher {
297 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 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 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 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 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 enum TransportSetup {
366 Nats(Arc<dyn EventTransportTx>, Arc<Codec>),
367 ZmqDirect(Arc<dyn EventTransportTx>, Arc<Codec>, String), 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 if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
383 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 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 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 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 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 pub async fn publish_bytes(&self, bytes: Vec<u8>) -> Result<()> {
512 self.publish_bytes_ref(&bytes).await
513 }
514
515 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 pub fn publisher_id(&self) -> u64 {
530 self.publisher_id
531 }
532
533 pub fn topic(&self) -> &str {
535 &self.topic
536 }
537
538 pub fn transport_kind(&self) -> EventTransportKind {
540 self.transport_kind
541 }
542}
543
544impl Drop for EventPublisher {
545 fn drop(&mut self) {
546 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 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
591pub 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 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 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 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 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 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 if let Some(broker) = resolve_zmq_broker(drt, &scope).await? {
671 let codec = Arc::new(Codec::Msgpack(MsgpackCodec));
673
674 let stream: WireStream = if broker.xpub_endpoints.len() == 1 {
675 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 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 Box::pin(DeduplicatingStream::new(
694 inner_stream,
695 codec.clone(),
696 100_000,
697 ))
698 };
699
700 (stream, codec)
701 } else {
702 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 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 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 pub async fn next(&mut self) -> Option<Result<EventEnvelope>> {
771 self.stream.next().await
772 }
773
774 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
784pub 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 pub async fn next(&mut self) -> Option<Result<(EventEnvelope, T)>> {
794 std::future::poll_fn(|cx| self.poll_next(cx)).await
795 }
796
797 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
813fn 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 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 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 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}