namespace.rs 4.6 KB
Newer Older
Ryan Olson's avatar
Ryan Olson committed
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
// SPDX-FileCopyrightText: Copyright (c) 2024-2025 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
// SPDX-License-Identifier: Apache-2.0
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

16
use anyhow::Context;
Ryan Olson's avatar
Ryan Olson committed
17
use async_trait::async_trait;
18
19
use futures::stream::StreamExt;
use futures::{Stream, TryStreamExt};
Ryan Olson's avatar
Ryan Olson committed
20
21

use super::*;
22
use crate::metrics::MetricsRegistry;
23
use crate::traits::events::{EventPublisher, EventSubscriber};
Ryan Olson's avatar
Ryan Olson committed
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54

#[async_trait]
impl EventPublisher for Namespace {
    fn subject(&self) -> String {
        format!("namespace.{}", self.name)
    }

    async fn publish(
        &self,
        event_name: impl AsRef<str> + Send + Sync,
        event: &(impl Serialize + Send + Sync),
    ) -> Result<()> {
        let bytes = serde_json::to_vec(event)?;
        self.publish_bytes(event_name, bytes).await
    }

    async fn publish_bytes(
        &self,
        event_name: impl AsRef<str> + Send + Sync,
        bytes: Vec<u8>,
    ) -> Result<()> {
        let subject = format!("{}.{}", self.subject(), event_name.as_ref());
        Ok(self
            .drt()
            .nats_client()
            .client()
            .publish(subject, bytes.into())
            .await?)
    }
}

55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
#[async_trait]
impl EventSubscriber for Namespace {
    async fn subscribe(
        &self,
        event_name: impl AsRef<str> + Send + Sync,
    ) -> Result<async_nats::Subscriber> {
        let subject = format!("{}.{}", self.subject(), event_name.as_ref());
        Ok(self.drt().nats_client().client().subscribe(subject).await?)
    }

    async fn subscribe_with_type<T: for<'de> Deserialize<'de> + Send + 'static>(
        &self,
        event_name: impl AsRef<str> + Send + Sync,
    ) -> Result<impl Stream<Item = Result<T>> + Send> {
        let subscriber = self.subscribe(event_name).await?;

        // Transform the subscriber into a stream of deserialized events
        let stream = subscriber.map(move |msg| {
            serde_json::from_slice::<T>(&msg.payload)
74
                .with_context(|| format!("Failed to deserialize event payload: {:?}", msg.payload))
75
76
77
78
79
80
        });

        Ok(stream)
    }
}

81
82
83
84
85
86
impl MetricsRegistry for Namespace {
    fn basename(&self) -> String {
        self.name.clone()
    }

    fn parent_hierarchy(&self) -> Vec<String> {
87
88
        // Build as: [ "" (DRT), non-empty parent basenames from root -> leaf ]
        let mut names = vec![String::new()]; // Start with empty string for DRT
89

90
91
92
93
94
95
        // Collect parent basenames from root to leaf
        let parent_names: Vec<String> =
            std::iter::successors(self.parent.as_deref(), |ns| ns.parent.as_deref())
                .map(|ns| ns.basename())
                .filter(|name| !name.is_empty())
                .collect();
96

97
98
99
        // Append parent names in reverse order (root to leaf)
        names.extend(parent_names.into_iter().rev());
        names
100
    }
101
102
}

103
#[cfg(feature = "integration")]
Ryan Olson's avatar
Ryan Olson committed
104
105
106
107
108
109
110
111
#[cfg(test)]
mod tests {
    use super::*;

    // todo - make a distributed runtime fixture
    // todo - two options - fully mocked or integration test
    #[tokio::test]
    async fn test_publish() {
112
113
        let rt = Runtime::from_current().unwrap();
        let dtr = DistributedRuntime::from_settings(rt.clone()).await.unwrap();
114
115
        let ns = dtr.namespace("test_namespace_publish".to_string()).unwrap();
        ns.publish("test_event", &"test".to_string()).await.unwrap();
116
        rt.shutdown();
Ryan Olson's avatar
Ryan Olson committed
117
    }
118
119
120
121
122

    #[tokio::test]
    async fn test_subscribe() {
        let rt = Runtime::from_current().unwrap();
        let dtr = DistributedRuntime::from_settings(rt.clone()).await.unwrap();
123
124
125
        let ns = dtr
            .namespace("test_namespace_subscribe".to_string())
            .unwrap();
126
127

        // Create a subscriber
128
        let mut subscriber = ns.subscribe("test_event").await.unwrap();
129
130

        // Publish a message
131
        ns.publish("test_event", &"test_message".to_string())
132
133
134
135
136
137
138
139
140
141
142
            .await
            .unwrap();

        // Receive the message
        if let Some(msg) = subscriber.next().await {
            let received = String::from_utf8(msg.payload.to_vec()).unwrap();
            assert_eq!(received, "\"test_message\"");
        }

        rt.shutdown();
    }
Ryan Olson's avatar
Ryan Olson committed
143
}