- 1
//! Integration tests for the vak-bus distributed event fabric - 2
//! (docs/design/53-distributed-bus.md). - 3
//! - 4
//! InMemory tests run unconditionally. NatsBus tests are `#[ignore]`'d - 5
//! and require a live NATS server: - 6
//! - 7
//! ```sh - 8
//! nats-server -p 4222 # or: docker run -p 4222:4222 nats - 9
//! cargo test -p vak-bus --test nats -- --ignored - 10
//! ``` - 11
- 12
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 13
- 14
use std::time::Duration; - 15
- 16
use vak_bus::{ - 17
BusError, EventPublisher, EventSubscriber, InMemoryBus, NatsBus, NatsConfig, WorkQueue, - 18
}; - 19
use vak_bus::{MessageEnvelope, TraceContext}; - 20
- 21
/// Build a plain-text JSON envelope for testing. - 22
fn envelope(source: &str, event_type: &str, payload: serde_json::Value) -> MessageEnvelope { - 23
let data = serde_json::to_vec(&payload).unwrap(); - 24
MessageEnvelope::new( - 25
source, - 26
event_type, - 27
"test_ws", - 28
"test_agent", - 29
1, - 30
MessageEnvelope::GENESIS_HASH, - 31
data, - 32
Some(TraceContext::new_root()), - 33
) - 34
} - 35
- 36
// --------------------------------------------------------------------------- - 37
// InMemoryBus integration tests - 38
// --------------------------------------------------------------------------- - 39
- 40
#[tokio::test] - 41
async fn inmemory_roundtrip() { - 42
let bus = InMemoryBus::new(); - 43
let mut sub = EventSubscriber::subscribe(&bus, "vak.events.test.>") - 44
.await - 45
.unwrap(); - 46
- 47
EventPublisher::publish( - 48
&bus, - 49
"vak.events.test.session1.message", - 50
envelope("test", "message", serde_json::json!({ "msg": "hello" })), - 51
) - 52
.await - 53
.unwrap(); - 54
- 55
let received = tokio::time::timeout(Duration::from_secs(1), sub.recv()) - 56
.await - 57
.unwrap() - 58
.unwrap(); - 59
- 60
let payload: serde_json::Value = serde_json::from_slice(&received.payload).unwrap(); - 61
assert_eq!(payload["msg"], "hello"); - 62
assert_eq!(received.lineage.workspace_id, "test_ws"); - 63
assert_eq!(received.trace.traceparent.split('-').count(), 4); - 64
} - 65
- 66
#[tokio::test] - 67
async fn inmemory_work_queue_claim_and_ack() { - 68
let bus = InMemoryBus::new(); - 69
let task_id = WorkQueue::enqueue( - 70
&bus, - 71
"vak.work.test.coder", - 72
envelope("test", "task", serde_json::json!({ "job": "build" })), - 73
) - 74
.await - 75
.unwrap(); - 76
- 77
let claimed = WorkQueue::claim( - 78
&bus, - 79
"vak.work.test.coder", - 80
"worker-1", - 81
Duration::from_secs(1), - 82
) - 83
.await - 84
.unwrap() - 85
.expect("should have a pending task"); - 86
- 87
assert_eq!(claimed.task_id, task_id); - 88
assert_eq!(claimed.attempts, 1); - 89
claimed.ack_handle.ack().await.unwrap(); - 90
} - 91
- 92
#[tokio::test] - 93
async fn inmemory_work_queue_dlq_on_nack() { - 94
let bus = InMemoryBus::new(); - 95
let task_id = WorkQueue::enqueue( - 96
&bus, - 97
"vak.work.test.worker", - 98
envelope("test", "task", serde_json::json!({ "job": "poison" })), - 99
) - 100
.await - 101
.unwrap(); - 102
- 103
let claimed = WorkQueue::claim( - 104
&bus, - 105
"vak.work.test.worker", - 106
"worker-1", - 107
Duration::from_secs(1), - 108
) - 109
.await - 110
.unwrap() - 111
.expect("should have a pending task"); - 112
assert_eq!(claimed.task_id, task_id); - 113
- 114
// Nack with retry=false → routed to DLQ - 115
claimed.ack_handle.nack(false).await.unwrap(); - 116
- 117
let dlq = bus.dead_letters(); - 118
assert_eq!(dlq.len(), 1); - 119
assert_eq!(dlq[0].failure_reason, "retry threshold exhausted"); - 120
} - 121
- 122
#[tokio::test] - 123
async fn inmemory_metrics_tracks_publishes_and_subscribes() { - 124
let bus = InMemoryBus::new(); - 125
let metrics = bus.metrics(); - 126
- 127
EventPublisher::publish( - 128
&bus, - 129
"vak.events.test.metrics", - 130
envelope("test", "metrics", serde_json::json!({ "n": 1 })), - 131
) - 132
.await - 133
.unwrap(); - 134
- 135
let _sub = EventSubscriber::subscribe(&bus, "vak.events.test.metrics") - 136
.await - 137
.unwrap(); - 138
- 139
// Wait for subscriber to register - 140
tokio::time::sleep(Duration::from_millis(50)).await; - 141
- 142
let snap = metrics.snapshot(); - 143
assert!(snap.published_count >= 1); - 144
} - 145
- 146
#[tokio::test] - 147
async fn inmemory_subscription_pattern_filtering() { - 148
let bus = InMemoryBus::new(); - 149
let mut sub = EventSubscriber::subscribe(&bus, "vak.events.ws_alpha.>") - 150
.await - 151
.unwrap(); - 152
- 153
EventPublisher::publish( - 154
&bus, - 155
"vak.events.ws_alpha.sess_001.message", - 156
envelope("test", "msg", serde_json::json!({})), - 157
) - 158
.await - 159
.unwrap(); - 160
- 161
let received = tokio::time::timeout(Duration::from_secs(1), sub.recv()) - 162
.await - 163
.unwrap() - 164
.unwrap(); - 165
assert_eq!(received.lineage.workspace_id, "test_ws"); - 166
- 167
// A message on a different workspace should not be received. - 168
let other = EventPublisher::publish( - 169
&bus, - 170
"vak.events.ws_beta.sess_001.message", - 171
envelope("test", "msg", serde_json::json!({})), - 172
); - 173
// Should complete without the subscriber receiving it. - 174
other.await.unwrap(); - 175
tokio::time::sleep(Duration::from_millis(10)).await; - 176
} - 177
- 178
#[tokio::test] - 179
async fn inmemory_claims_remove_from_queue() { - 180
let bus = InMemoryBus::new(); - 181
for i in 0..3 { - 182
WorkQueue::enqueue( - 183
&bus, - 184
"vak.work.test.batch", - 185
envelope("test", "task", serde_json::json!({ "i": i })), - 186
) - 187
.await - 188
.unwrap(); - 189
} - 190
- 191
// Claim all three - 192
for _ in 0..3 { - 193
let task = WorkQueue::claim(&bus, "vak.work.test.batch", "w", Duration::from_secs(1)) - 194
.await - 195
.unwrap(); - 196
assert!(task.is_some()); - 197
} - 198
// Queue should now be empty - 199
let empty = WorkQueue::claim(&bus, "vak.work.test.batch", "w", Duration::from_millis(10)) - 200
.await - 201
.unwrap(); - 202
assert!(empty.is_none()); - 203
} - 204
- 205
// --------------------------------------------------------------------------- - 206
// NatsBus integration tests (require live NATS) - 207
// --------------------------------------------------------------------------- - 208
- 209
#[tokio::test] - 210
#[ignore = "requires a live NATS server on nats://127.0.0.1:4222"] - 211
async fn nats_connects_and_roundtrips() { - 212
let config = NatsConfig { - 213
url: "nats://127.0.0.1:4222".into(), - 214
..NatsConfig::default() - 215
}; - 216
let bus = NatsBus::connect(config).await.unwrap(); - 217
- 218
let mut sub = EventSubscriber::subscribe(&bus, "vak.events.test.nats.>") - 219
.await - 220
.unwrap(); - 221
- 222
EventPublisher::publish( - 223
&bus, - 224
"vak.events.test.nats.roundtrip", - 225
envelope( - 226
"test", - 227
"nats.msg", - 228
serde_json::json!({ "msg": "from-nats" }), - 229
), - 230
) - 231
.await - 232
.unwrap(); - 233
- 234
let received = tokio::time::timeout(Duration::from_secs(3), sub.recv()) - 235
.await - 236
.unwrap() - 237
.unwrap(); - 238
let payload: serde_json::Value = serde_json::from_slice(&received.payload).unwrap(); - 239
assert_eq!(payload["msg"], "from-nats"); - 240
} - 241
- 242
#[tokio::test] - 243
#[ignore = "requires a live NATS server on nats://127.0.0.1:4222"] - 244
async fn nats_rejects_bad_credentials() { - 245
let config = NatsConfig { - 246
url: "nats://127.0.0.1:4222".into(), - 247
credentials_jwt: Some("invalid.jwt.token".into()), - 248
nkey_seed: Some("SUAFISVULNERABLEFAKESEED".into()), - 249
connect_timeout: Duration::from_secs(3), - 250
}; - 251
let result = NatsBus::connect(config).await; - 252
assert!(result.is_err(), "bad credentials must fail"); - 253
} - 254
- 255
#[tokio::test] - 256
#[ignore = "requires a live NATS server on nats://127.0.0.1:4222"] - 257
async fn nats_work_queue_claim_and_ack() { - 258
let config = NatsConfig { - 259
url: "nats://127.0.0.1:4222".into(), - 260
..NatsConfig::default() - 261
}; - 262
let bus = NatsBus::connect(config).await.unwrap(); - 263
- 264
let task_id = WorkQueue::enqueue( - 265
&bus, - 266
"vak.work.test.nats_coder", - 267
envelope("test", "nats.task", serde_json::json!({ "job": "build" })), - 268
) - 269
.await - 270
.unwrap(); - 271
- 272
let claimed = WorkQueue::claim( - 273
&bus, - 274
"vak.work.test.nats_coder", - 275
"worker-1", - 276
Duration::from_secs(5), - 277
) - 278
.await - 279
.unwrap(); - 280
assert!(claimed.is_some(), "should claim the enqueued task"); - 281
let claimed = claimed.unwrap(); - 282
assert_eq!(claimed.task_id, task_id); - 283
claimed.ack_handle.ack().await.unwrap(); - 284
} - 285
- 286
/// Verify that `BusError` implements `std::error::Error` and `Display`. - 287
#[test] - 288
fn bus_error_is_display_and_error() { - 289
let err = BusError::Connection("nats down".into()); - 290
assert!(err.to_string().contains("nats down")); - 291
assert!(std::error::Error::source(&err).is_none()); - 292
} - 293
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.