- 1
//! Distributed event fabric bridge (docs/design/53-distributed-bus.md). - 2
//! - 3
//! Bridges the server's local `SystemEvent` model to vak-bus's distributed - 4
//! messaging fabric. When NATS is configured (`[server.bus]` section), the - 5
//! `ServerBus` publishes system events to JetStream subjects with Merkle causal - 6
//! chaining, AES-256-GCM envelope encryption, and role-based ACL enforcement. - 7
//! When NATS is absent, the `InMemoryBus` backend provides the same semantic - 8
//! contract for local operation and CI. - 9
//! - 10
//! The bus is **additive**: the existing local `tokio::broadcast` hub keeps - 11
//! working exactly as before. `ServerBus` is only consulted when present, and - 12
//! a missing/unavailable backend fails closed (local broadcast continues). - 13
- 14
#![allow(dead_code)] - 15
- 16
use std::sync::Arc; - 17
- 18
use tokio::sync::Mutex; - 19
use vak_bus::BusMetricsSnapshot; - 20
use vak_bus::bus::{EventPublisher, EventSubscriber, InMemoryBus, NatsBus, NatsConfig}; - 21
use vak_bus::envelope::MessageEnvelope; - 22
use vak_bus::subjects::{AclPolicy, Subject}; - 23
- 24
use crate::events::SystemEvent; - 25
- 26
/// The backend engine backing `ServerBus`. - 27
#[allow(clippy::large_enum_variant)] - 28
enum Backend { - 29
InMemory(InMemoryBus), - 30
Distributed(NatsBus), - 31
} - 32
- 33
impl std::fmt::Debug for Backend { - 34
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - 35
match self { - 36
Backend::InMemory(_) => write!(f, "Backend::InMemory"), - 37
Backend::Distributed(_) => write!(f, "Backend::Distributed"), - 38
} - 39
} - 40
} - 41
- 42
/// The NATS secrets, resolved through the secrets chain at the moment a - 43
/// bus is built (`Core::bus_secret`), never read from TOML. - 44
#[derive(Default)] - 45
pub struct BusCredentials { - 46
pub jwt: Option<String>, - 47
pub nkey_seed: Option<String>, - 48
} - 49
- 50
impl BusCredentials { - 51
pub fn of(core: &vak_core::Core) -> Self { - 52
Self { - 53
jwt: core.bus_secret(vak_config::BUS_NATS_JWT_VAR), - 54
nkey_seed: core.bus_secret(vak_config::BUS_NATS_NKEY_SEED_VAR), - 55
} - 56
} - 57
} - 58
- 59
/// Distributed event bus for the vak server. - 60
/// - 61
/// Wraps a vak-bus backend (`InMemoryBus` for local operation, `NatsBus` for - 62
/// distributed fan-out) and translates between the server's `SystemEvent` - 63
/// model and vak-bus's `MessageEnvelope` wire format. - 64
pub struct ServerBus { - 65
backend: Backend, - 66
workspace_id: String, - 67
/// Optional per-workspace encryption key. When set, payloads are - 68
/// AES-256-GCM encrypted before publishing. - 69
workspace_secret: Option<Arc<[u8]>>, - 70
/// Monotonic sequence number for causal chaining on the event plane. - 71
seq: Arc<Mutex<u64>>, - 72
} - 73
- 74
impl std::fmt::Debug for ServerBus { - 75
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - 76
f.debug_struct("ServerBus") - 77
.field("backend", &self.backend) - 78
.field("workspace_id", &self.workspace_id) - 79
.finish_non_exhaustive() - 80
} - 81
} - 82
- 83
impl ServerBus { - 84
/// Create a local-only bus backed by `InMemoryBus`. - 85
pub fn local(workspace_id: impl Into<String>) -> Self { - 86
Self { - 87
backend: Backend::InMemory(InMemoryBus::new()), - 88
workspace_id: workspace_id.into(), - 89
workspace_secret: None, - 90
seq: Arc::new(Mutex::new(0)), - 91
} - 92
} - 93
- 94
/// Create a distributed bus backed by NATS Core + JetStream. - 95
pub async fn distributed( - 96
workspace_id: impl Into<String>, - 97
config: &vak_config::BusResolved, - 98
credentials: &BusCredentials, - 99
) -> Result<Self, vak_bus::BusError> { - 100
let nats_config = NatsConfig { - 101
url: config - 102
.nats_url - 103
.clone() - 104
.unwrap_or_else(|| NatsConfig::default().url), - 105
credentials_jwt: credentials.jwt.clone(), - 106
nkey_seed: credentials.nkey_seed.clone(), - 107
..NatsConfig::default() - 108
}; - 109
let bus = NatsBus::connect(nats_config).await?; - 110
Ok(Self { - 111
backend: Backend::Distributed(bus), - 112
workspace_id: workspace_id.into(), - 113
workspace_secret: config.workspace_secret.as_deref().map(Arc::from), - 114
seq: Arc::new(Mutex::new(0)), - 115
}) - 116
} - 117
- 118
/// Build a `ServerBus` from resolved config: distributed if NATS is - 119
/// configured, local otherwise. Never fails for a missing configuration — - 120
/// returns a local bus so the server always has an event fabric. - 121
/// - 122
/// A NATS connection failure degrades to `InMemoryBus` with a warning, - 123
/// preserving local observability. - 124
pub async fn from_resolved( - 125
workspace_id: impl Into<String>, - 126
config: &vak_config::BusResolved, - 127
credentials: &BusCredentials, - 128
) -> Self { - 129
let ws = workspace_id.into(); - 130
if config.nats_url.is_some() { - 131
match Self::distributed(&ws, config, credentials).await { - 132
Ok(bus) => bus, - 133
Err(e) => { - 134
eprintln!( - 135
"vak-server: NATS bus failed to connect ({e}); \ - 136
falling back to local InMemoryBus" - 137
); - 138
Self::local(&ws) - 139
} - 140
} - 141
} else { - 142
Self::local(&ws) - 143
} - 144
} - 145
- 146
/// Publish a `SystemEvent` to the distributed fabric as a `MessageEnvelope`. - 147
/// - 148
/// Maps each event variant to an appropriate vak-bus `Subject`, attaches - 149
/// the workspace/session context, computes the Merkle causal link from - 150
/// the previous event's hash, and optionally encrypts the payload. - 151
pub async fn emit( - 152
&self, - 153
event: &SystemEvent, - 154
session_id: Option<&str>, - 155
) -> Result<(), vak_bus::BusError> { - 156
let seq = { - 157
let mut s = self.seq.lock().await; - 158
*s += 1; - 159
*s - 160
}; - 161
- 162
let subject = self.subject_for(event, session_id); - 163
let prev_hash = if seq == 1 { - 164
MessageEnvelope::GENESIS_HASH.to_string() - 165
} else { - 166
format!("{}", seq - 1) - 167
}; - 168
- 169
let payload = serde_json::to_vec(event).unwrap_or_else(|_| b"{}".to_vec()); - 170
- 171
let mut envelope = MessageEnvelope::new( - 172
format!("vak://server/{}", self.workspace_id), - 173
subject_event_type(&subject), - 174
&self.workspace_id, - 175
"server", - 176
seq, - 177
&prev_hash, - 178
payload, - 179
None, - 180
); - 181
- 182
if let Some(sid) = session_id { - 183
envelope = envelope.with_session(sid); - 184
} - 185
- 186
if let Some(secret) = &self.workspace_secret { - 187
let _ = vak_bus::encrypt_envelope(&mut envelope, secret, "v1"); - 188
} - 189
- 190
let subject_str = subject.to_subject_string(); - 191
match &self.backend { - 192
Backend::InMemory(b) => EventPublisher::publish(b, &subject_str, envelope).await, - 193
Backend::Distributed(b) => EventPublisher::publish(b, &subject_str, envelope).await, - 194
} - 195
} - 196
- 197
/// Subscribe to events matching a subject pattern. - 198
pub async fn subscribe( - 199
&self, - 200
pattern: &str, - 201
) -> Result<tokio::sync::mpsc::Receiver<MessageEnvelope>, vak_bus::BusError> { - 202
match &self.backend { - 203
Backend::InMemory(b) => EventSubscriber::subscribe(b, pattern).await, - 204
Backend::Distributed(b) => EventSubscriber::subscribe(b, pattern).await, - 205
} - 206
} - 207
- 208
/// Subscribe to all system events for this workspace. - 209
pub async fn subscribe_workspace_events( - 210
&self, - 211
) -> Result<tokio::sync::mpsc::Receiver<MessageEnvelope>, vak_bus::BusError> { - 212
self.subscribe(&format!("vak.events.{}.*", self.workspace_id)) - 213
.await - 214
} - 215
- 216
/// Current bus operational metrics. - 217
pub fn metrics(&self) -> BusMetricsSnapshot { - 218
match &self.backend { - 219
Backend::InMemory(b) => b.metrics().snapshot(), - 220
Backend::Distributed(b) => b.metrics().snapshot(), - 221
} - 222
} - 223
- 224
/// Whether the bus is backed by a distributed NATS connection. - 225
pub fn is_distributed(&self) -> bool { - 226
matches!(self.backend, Backend::Distributed(_)) - 227
} - 228
- 229
/// The workspace ID this bus is scoped to. - 230
pub fn workspace_id(&self) -> &str { - 231
&self.workspace_id - 232
} - 233
- 234
/// Operational status snapshot for the admin console / ops center. - 235
pub fn status(&self) -> serde_json::Value { - 236
let m = self.metrics(); - 237
serde_json::json!({ - 238
"workspace_id": self.workspace_id, - 239
"backend": if self.is_distributed() { "nats" } else { "memory" }, - 240
"connected": self.is_distributed(), - 241
"encrypted": self.workspace_secret.is_some(), - 242
"metrics": { - 243
"published_count": m.published_count, - 244
"received_count": m.received_count, - 245
"bytes_published": m.bytes_published, - 246
"bytes_received": m.bytes_received, - 247
"dead_letter_count": m.dead_letter_count, - 248
"active_queue_lag": m.active_queue_lag, - 249
}, - 250
}) - 251
} - 252
- 253
/// ACL policy for the workspace's event plane. - 254
pub fn acl_for(&self, session_id: &str, agent_id: &str, role: &str) -> AclPolicy { - 255
AclPolicy::for_worker(&self.workspace_id, session_id, agent_id, role) - 256
} - 257
- 258
fn subject_for(&self, event: &SystemEvent, session_id: Option<&str>) -> Subject { - 259
let sess = session_id.unwrap_or(""); - 260
match event { - 261
SystemEvent::SessionCreated { .. } => Subject::Custom(format!( - 262
"vak.events.{}.{}.session.created", - 263
self.workspace_id, sess - 264
)), - 265
SystemEvent::SessionEntryAppended { - 266
session_id: sid, .. - 267
} => Subject::Custom(format!("vak.events.{}.{}.entry", self.workspace_id, sid)), - 268
SystemEvent::ConfigChanged { .. } => Subject::Custom(format!( - 269
"vak.events.{}.{}.config.changed", - 270
self.workspace_id, sess - 271
)), - 272
SystemEvent::GatewayInbound { .. } => Subject::Custom(format!( - 273
"vak.events.{}.{}.gateway.inbound", - 274
self.workspace_id, sess - 275
)), - 276
SystemEvent::ApprovalRequested { session_id, .. } => Subject::Approvals { - 277
workspace_id: self.workspace_id.clone(), - 278
session_id: session_id.clone(), - 279
}, - 280
SystemEvent::ApprovalGranted { .. } | SystemEvent::ApprovalDenied { .. } => { - 281
Subject::Custom(format!( - 282
"vak.events.{}.{}.approval.{}", - 283
self.workspace_id, - 284
sess, - 285
if matches!(event, SystemEvent::ApprovalGranted { .. }) { - 286
"granted" - 287
} else { - 288
"denied" - 289
} - 290
)) - 291
} - 292
SystemEvent::SecurityEvent { .. } => Subject::Custom(format!( - 293
"vak.events.{}.{}.security", - 294
self.workspace_id, sess - 295
)), - 296
SystemEvent::ProviderError { .. } => Subject::Custom(format!( - 297
"vak.events.{}.{}.provider.error", - 298
self.workspace_id, sess - 299
)), - 300
SystemEvent::RateLimit { .. } => Subject::Custom(format!( - 301
"vak.events.{}.{}.rate_limit", - 302
self.workspace_id, sess - 303
)), - 304
SystemEvent::Heartbeat => Subject::Custom(format!( - 305
"vak.events.{}.{}.heartbeat", - 306
self.workspace_id, sess - 307
)), - 308
SystemEvent::Agent(_) => { - 309
Subject::Custom(format!("vak.events.{}.{}.agent", self.workspace_id, sess)) - 310
} - 311
} - 312
} - 313
} - 314
- 315
/// Human-readable event type string for the bus subject. - 316
fn subject_event_type(subject: &Subject) -> &str { - 317
match subject { - 318
Subject::EventsTokens { .. } => "vak.events.tokens", - 319
Subject::EventsTerminal { .. } => "vak.events.terminal", - 320
Subject::EventsTelemetry { .. } => "vak.events.telemetry", - 321
Subject::WorkTask { .. } => "vak.work.task", - 322
Subject::AgentInbox { .. } => "vak.agent.inbox", - 323
Subject::Receipts { .. } => "vak.receipts.completed", - 324
Subject::Approvals { .. } => "vak.approvals.request", - 325
Subject::DeadLetter { .. } => "vak.dlq.failed", - 326
Subject::Custom(s) => s, - 327
} - 328
} - 329
- 330
/// Decrypt a received envelope if it carries encryption metadata. - 331
pub fn try_decrypt( - 332
envelope: &mut MessageEnvelope, - 333
workspace_secret: &[u8], - 334
) -> Result<(), vak_bus::BusError> { - 335
if !envelope.encrypted { - 336
return Ok(()); - 337
} - 338
vak_bus::decrypt_envelope(envelope, workspace_secret) - 339
.map_err(|e| vak_bus::BusError::Serialization(e.to_string())) - 340
} - 341
- 342
// ---- tests ---- - 343
- 344
#[cfg(test)] - 345
mod tests { - 346
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 347
- 348
use super::*; - 349
- 350
#[tokio::test] - 351
async fn local_bus_publishes_and_receives_system_event() { - 352
let bus = ServerBus::local("ws_test"); - 353
- 354
// Subscribe before emitting: broadcast channels don't deliver to - 355
// late joiners. - 356
let mut rx = bus - 357
.subscribe("vak.events.ws_test.>") - 358
.await - 359
.expect("subscribe"); - 360
- 361
// Yield to let the forwarding task start before we publish. - 362
tokio::task::yield_now().await; - 363
- 364
let event = SystemEvent::Heartbeat; - 365
bus.emit(&event, Some("sess1")).await.expect("publish"); - 366
- 367
let received = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) - 368
.await - 369
.expect("timeout") - 370
.expect("received envelope"); - 371
- 372
let payload: serde_json::Value = - 373
serde_json::from_slice(&received.payload).expect("decode payload"); - 374
assert_eq!(payload["type"], "Heartbeat"); - 375
} - 376
- 377
#[tokio::test] - 378
async fn local_bus_chains_causal_sequence() { - 379
let bus = ServerBus::local("ws_seq"); - 380
- 381
// Subscribe before emitting (broadcast channels don't backfill - 382
// to late joiners). - 383
let mut rx = bus.subscribe("vak.events.ws_seq.>").await.expect("sub"); - 384
- 385
// Yield to let the forwarding task start before we publish. - 386
tokio::task::yield_now().await; - 387
- 388
bus.emit(&SystemEvent::Heartbeat, Some("sess1")) - 389
.await - 390
.expect("emit 1"); - 391
bus.emit(&SystemEvent::Heartbeat, Some("sess1")) - 392
.await - 393
.expect("emit 2"); - 394
- 395
let env1 = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) - 396
.await - 397
.expect("timeout 1") - 398
.expect("recv 1"); - 399
let env2 = tokio::time::timeout(std::time::Duration::from_millis(500), rx.recv()) - 400
.await - 401
.expect("timeout 2") - 402
.expect("recv 2"); - 403
- 404
assert!(env2.lineage.seq_num == env1.lineage.seq_num + 1); - 405
} - 406
- 407
#[tokio::test] - 408
async fn local_bus_metrics_track_publish() { - 409
let bus = ServerBus::local("ws_metrics"); - 410
bus.emit(&SystemEvent::Heartbeat, None).await.expect("emit"); - 411
let snap = bus.metrics(); - 412
assert!(snap.published_count >= 1); - 413
} - 414
- 415
#[tokio::test] - 416
async fn distributed_bus_falls_back_when_nats_unreachable() { - 417
let cfg = vak_config::BusResolved { - 418
nats_url: Some("nats://127.0.0.1:1".to_string()), - 419
..Default::default() - 420
}; - 421
let bus = ServerBus::from_resolved("ws_fb", &cfg, &BusCredentials::default()).await; - 422
bus.emit(&SystemEvent::Heartbeat, None) - 423
.await - 424
.expect("emit on fallback"); - 425
} - 426
- 427
#[test] - 428
fn approval_subject_uses_typed_variant() { - 429
let bus = ServerBus::local("ws_approvals"); - 430
let event = SystemEvent::ApprovalRequested { - 431
id: "appr1".to_string(), - 432
session_id: "sess_x".to_string(), - 433
tool: "bash".to_string(), - 434
reason: "test".to_string(), - 435
}; - 436
let subject = bus.subject_for(&event, Some("sess_x")); - 437
match subject { - 438
Subject::Approvals { - 439
workspace_id, - 440
session_id, - 441
} => { - 442
assert_eq!(workspace_id, "ws_approvals"); - 443
assert_eq!(session_id, "sess_x"); - 444
} - 445
other => panic!("expected Approvals subject, got {:?}", other), - 446
} - 447
} - 448
- 449
#[test] - 450
fn acl_policy_is_workspace_scoped() { - 451
let bus = ServerBus::local("ws_acl"); - 452
let acl = bus.acl_for("sess1", "agent_a", "coder"); - 453
assert!(acl.can_publish("vak.events.ws_acl.sess1.tokens")); - 454
assert!(!acl.can_publish("vak.events.ws_other.sess1.tokens")); - 455
} - 456
- 457
#[test] - 458
fn resolved_bus_config_defaults_to_local() { - 459
let cfg = vak_config::BusResolved::default(); - 460
assert!(cfg.nats_url.is_none()); - 461
} - 462
} - 463
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.