Skip to main content

tatara_engine/nats/
mod.rs

1//! NATS event bus for tatara — inter-task communication, log aggregation,
2//! and catalog change notifications.
3//!
4//! Optional: if NATS is not configured, all publish operations are no-ops.
5//! Uses JetStream for guaranteed delivery of critical operations and
6//! core NATS for ephemeral data (logs, metrics).
7
8use anyhow::Result;
9use serde::{Deserialize, Serialize};
10use tracing::{debug, info, warn};
11
12use tatara_core::catalog::ServiceEntry;
13use tatara_core::domain::event::Event;
14
15/// Configuration for the NATS event bus.
16#[derive(Debug, Clone, Serialize, Deserialize)]
17pub struct NatsConfig {
18    /// NATS server URL.
19    #[serde(default = "default_nats_url")]
20    pub url: String,
21
22    /// Whether NATS integration is enabled.
23    #[serde(default)]
24    pub enabled: bool,
25
26    /// Reserved for future JetStream support (not yet implemented).
27    /// Currently all NATS operations use core NATS publish/subscribe.
28    #[serde(default)]
29    pub _jetstream_reserved: bool,
30
31    /// Subject prefix for all tatara NATS subjects.
32    #[serde(default = "default_subject_prefix")]
33    pub subject_prefix: String,
34}
35
36fn default_nats_url() -> String {
37    "nats://127.0.0.1:4222".to_string()
38}
39fn default_true() -> bool {
40    true
41}
42fn default_subject_prefix() -> String {
43    "tatara".to_string()
44}
45
46impl Default for NatsConfig {
47    fn default() -> Self {
48        Self {
49            url: default_nats_url(),
50            enabled: false,
51            _jetstream_reserved: false,
52            subject_prefix: default_subject_prefix(),
53        }
54    }
55}
56
57/// NATS event bus for publishing events, logs, and catalog changes.
58///
59/// Subject hierarchy:
60/// - `{prefix}.events.{kind}` — cluster events
61/// - `{prefix}.logs.{alloc_id}.{task_name}` — task logs
62/// - `{prefix}.health.{service_name}` — health probe results
63/// - `{prefix}.catalog.changes` — service catalog mutations
64pub struct NatsEventBus {
65    config: NatsConfig,
66    client: Option<async_nats::Client>,
67}
68
69impl NatsEventBus {
70    /// Connect to NATS. If disabled or connection fails, creates a no-op bus.
71    pub async fn connect(config: NatsConfig) -> Self {
72        if !config.enabled {
73            debug!("NATS event bus disabled");
74            return Self {
75                config,
76                client: None,
77            };
78        }
79
80        match async_nats::connect(&config.url).await {
81            Ok(client) => {
82                info!(url = %config.url, "connected to NATS");
83                Self {
84                    config,
85                    client: Some(client),
86                }
87            }
88            Err(e) => {
89                warn!(url = %config.url, error = %e, "failed to connect to NATS, running without event bus");
90                Self {
91                    config,
92                    client: None,
93                }
94            }
95        }
96    }
97
98    /// Create a disconnected (no-op) event bus.
99    pub fn disconnected() -> Self {
100        Self {
101            config: NatsConfig::default(),
102            client: None,
103        }
104    }
105
106    /// Check if NATS is connected.
107    pub fn is_connected(&self) -> bool {
108        self.client.is_some()
109    }
110
111    /// Publish a cluster event.
112    pub async fn publish_event(&self, event: &Event) -> Result<()> {
113        let Some(client) = &self.client else {
114            return Ok(());
115        };
116        let subject = format!("{}.events.{}", self.config.subject_prefix, event.kind_str());
117        let payload = serde_json::to_vec(event)?;
118        client
119            .publish(subject, payload.into())
120            .await
121            .map_err(|e| anyhow::anyhow!("NATS publish failed: {e}"))?;
122        Ok(())
123    }
124
125    /// Publish a log entry for cross-node aggregation.
126    pub async fn publish_log(
127        &self,
128        alloc_id: &str,
129        task_name: &str,
130        message: &str,
131        stream: &str,
132    ) -> Result<()> {
133        let Some(client) = &self.client else {
134            return Ok(());
135        };
136        let subject = format!(
137            "{}.logs.{}.{}",
138            self.config.subject_prefix, alloc_id, task_name
139        );
140        let payload = serde_json::json!({
141            "alloc_id": alloc_id,
142            "task_name": task_name,
143            "message": message,
144            "stream": stream,
145            "timestamp": chrono::Utc::now().to_rfc3339(),
146        });
147        client
148            .publish(subject, serde_json::to_vec(&payload)?.into())
149            .await
150            .map_err(|e| anyhow::anyhow!("NATS publish failed: {e}"))?;
151        Ok(())
152    }
153
154    /// Publish a health probe result.
155    pub async fn publish_health(
156        &self,
157        service_name: &str,
158        service_id: &str,
159        healthy: bool,
160    ) -> Result<()> {
161        let Some(client) = &self.client else {
162            return Ok(());
163        };
164        let subject = format!("{}.health.{}", self.config.subject_prefix, service_name);
165        let payload = serde_json::json!({
166            "service_id": service_id,
167            "healthy": healthy,
168            "timestamp": chrono::Utc::now().to_rfc3339(),
169        });
170        client
171            .publish(subject, serde_json::to_vec(&payload)?.into())
172            .await
173            .map_err(|e| anyhow::anyhow!("NATS publish failed: {e}"))?;
174        Ok(())
175    }
176
177    /// Publish a catalog change (service registered/deregistered).
178    pub async fn publish_catalog_change(&self, action: &str, entry: &ServiceEntry) -> Result<()> {
179        let Some(client) = &self.client else {
180            return Ok(());
181        };
182        let subject = format!("{}.catalog.changes", self.config.subject_prefix);
183        let payload = serde_json::json!({
184            "action": action,
185            "service_name": entry.service_name,
186            "service_id": entry.service_id,
187            "address": entry.address,
188            "port": entry.port,
189            "timestamp": chrono::Utc::now().to_rfc3339(),
190        });
191        client
192            .publish(subject, serde_json::to_vec(&payload)?.into())
193            .await
194            .map_err(|e| anyhow::anyhow!("NATS publish failed: {e}"))?;
195        Ok(())
196    }
197
198    /// Subscribe to events matching a filter pattern.
199    pub async fn subscribe_events(
200        &self,
201        kind_filter: &str,
202    ) -> Result<Option<async_nats::Subscriber>> {
203        let Some(client) = &self.client else {
204            return Ok(None);
205        };
206        let subject = format!("{}.events.{}", self.config.subject_prefix, kind_filter);
207        let sub = client.subscribe(subject).await?;
208        Ok(Some(sub))
209    }
210
211    /// Subscribe to logs for a specific allocation.
212    pub async fn subscribe_logs(&self, alloc_id: &str) -> Result<Option<async_nats::Subscriber>> {
213        let Some(client) = &self.client else {
214            return Ok(None);
215        };
216        let subject = format!("{}.logs.{}.>", self.config.subject_prefix, alloc_id);
217        let sub = client.subscribe(subject).await?;
218        Ok(Some(sub))
219    }
220}
221
222/// Extension to use Event's Display-based kind string for NATS subjects.
223trait EventKindStr {
224    fn kind_str(&self) -> String;
225}
226
227impl EventKindStr for Event {
228    fn kind_str(&self) -> String {
229        self.kind.to_string()
230    }
231}