tatara_engine/nats/
mod.rs1use 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#[derive(Debug, Clone, Serialize, Deserialize)]
17pub struct NatsConfig {
18 #[serde(default = "default_nats_url")]
20 pub url: String,
21
22 #[serde(default)]
24 pub enabled: bool,
25
26 #[serde(default)]
29 pub _jetstream_reserved: bool,
30
31 #[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
57pub struct NatsEventBus {
65 config: NatsConfig,
66 client: Option<async_nats::Client>,
67}
68
69impl NatsEventBus {
70 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 pub fn disconnected() -> Self {
100 Self {
101 config: NatsConfig::default(),
102 client: None,
103 }
104 }
105
106 pub fn is_connected(&self) -> bool {
108 self.client.is_some()
109 }
110
111 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 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 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 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 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 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
222trait 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}