- 1
//! Comprehensive regression test suite for `vak-bus`. - 2
//! - 3
//! Validates zero-trust cryptography, Merkle causal chaining, W3C trace context, - 4
//! subject ACLs, pub/sub fan-out, competing consumer work queues, and Dead-Letter Queue (DLQ). - 5
- 6
use std::time::Duration; - 7
- 8
use vak_bus::bus::{EventPublisher, EventSubscriber, InMemoryBus, WorkQueue}; - 9
use vak_bus::crypto::{CryptoError, decrypt_envelope, encrypt_envelope}; - 10
use vak_bus::envelope::{MessageEnvelope, TraceContext}; - 11
use vak_bus::subjects::{AclPolicy, Subject}; - 12
- 13
#[tokio::test] - 14
async fn test_trace_context_propagation() { - 15
let root = TraceContext::new_root(); - 16
assert_eq!(root.trace_id().len(), 32); - 17
assert_eq!(root.span_id().len(), 16); - 18
- 19
let child = root.child_span(); - 20
assert_eq!(child.trace_id(), root.trace_id()); - 21
assert_ne!(child.span_id(), root.span_id()); - 22
} - 23
- 24
#[tokio::test] - 25
async fn test_envelope_merkle_chaining() { - 26
let env1 = MessageEnvelope::new( - 27
"vak://agent/coder-1", - 28
"vak.work.task.started", - 29
"ws_alpha", - 30
"agent_01", - 31
1, - 32
MessageEnvelope::GENESIS_HASH, - 33
b"step 1 data".to_vec(), - 34
None, - 35
); - 36
- 37
let env1_hash = env1.compute_hash(); - 38
- 39
// Valid successor in causal chain - 40
let env2 = MessageEnvelope::new( - 41
"vak://agent/coder-1", - 42
"vak.work.task.progress", - 43
"ws_alpha", - 44
"agent_01", - 45
2, - 46
&env1_hash, - 47
b"step 2 data".to_vec(), - 48
Some(env1.trace.child_span()), - 49
); - 50
- 51
assert!(MessageEnvelope::verify_merkle_link(&env1, &env2)); - 52
- 53
// Invalid sequence number breaks link - 54
let env2_bad_seq = MessageEnvelope::new( - 55
"vak://agent/coder-1", - 56
"vak.work.task.progress", - 57
"ws_alpha", - 58
"agent_01", - 59
3, // Should be 2! - 60
&env1_hash, - 61
b"step 2 data".to_vec(), - 62
None, - 63
); - 64
assert!(!MessageEnvelope::verify_merkle_link(&env1, &env2_bad_seq)); - 65
- 66
// Forged predecessor hash breaks link - 67
let env2_bad_hash = MessageEnvelope::new( - 68
"vak://agent/coder-1", - 69
"vak.work.task.progress", - 70
"ws_alpha", - 71
"agent_01", - 72
2, - 73
"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef", - 74
b"step 2 data".to_vec(), - 75
None, - 76
); - 77
assert!(!MessageEnvelope::verify_merkle_link(&env1, &env2_bad_hash)); - 78
} - 79
- 80
#[tokio::test] - 81
async fn test_crypto_envelope_round_trip_and_tampering() { - 82
let secret = b"super-strong-workspace-secret-key-32b!"; - 83
let original_payload = b"sensitive proprietary agent output".to_vec(); - 84
- 85
let mut env = MessageEnvelope::new( - 86
"vak://agent/secure", - 87
"vak.work.output", - 88
"ws_prod", - 89
"agent_sec", - 90
1, - 91
MessageEnvelope::GENESIS_HASH, - 92
original_payload.clone(), - 93
None, - 94
); - 95
- 96
// Encrypt - 97
encrypt_envelope(&mut env, secret, "v1").expect("encryption succeeds"); - 98
assert!(env.encrypted); - 99
assert_ne!(env.payload, original_payload); - 100
- 101
// Decrypt - 102
decrypt_envelope(&mut env, secret).expect("decryption succeeds"); - 103
assert!(!env.encrypted); - 104
assert_eq!(env.payload, original_payload); - 105
- 106
// Re-encrypt and tamper with payload - 107
encrypt_envelope(&mut env, secret, "v1").unwrap(); - 108
env.payload[0] ^= 0x42; // Corrupt ciphertext - 109
let err = decrypt_envelope(&mut env, secret).unwrap_err(); - 110
assert!(matches!(err, CryptoError::DecryptionFailed)); - 111
} - 112
- 113
#[tokio::test] - 114
async fn test_subjects_and_acl_enforcement() { - 115
let subj = Subject::EventsTokens { - 116
workspace_id: "ws1".into(), - 117
session_id: "sess1".into(), - 118
}; - 119
assert_eq!(subj.to_subject_string(), "vak.events.ws1.sess1.tokens"); - 120
- 121
let parsed = Subject::parse("vak.work.ws1.coder.task").unwrap(); - 122
assert_eq!( - 123
parsed, - 124
Subject::WorkTask { - 125
workspace_id: "ws1".into(), - 126
role: "coder".into(), - 127
} - 128
); - 129
- 130
let acl = AclPolicy::for_worker("ws1", "sess1", "coder_agent", "coder"); - 131
assert!(acl.can_publish("vak.events.ws1.sess1.tokens")); - 132
assert!(acl.can_publish("vak.events.ws1.sess1.terminal")); - 133
assert!(!acl.can_publish("vak.work.ws1.coder.task")); // Workers cannot forge work - 134
assert!(acl.can_subscribe("vak.work.ws1.coder.task")); - 135
assert!(acl.can_subscribe("vak.agent.ws1.coder_agent.inbox")); - 136
assert!(!acl.can_subscribe("vak.approvals.ws1.sess1.request")); // Cannot snoop approvals - 137
} - 138
- 139
#[tokio::test] - 140
async fn test_ephemeral_pub_sub_fanout() { - 141
let bus = InMemoryBus::new(); - 142
- 143
let mut rx1 = bus - 144
.subscribe("vak.events.ws_test.*.telemetry") - 145
.await - 146
.expect("sub 1"); - 147
let mut rx2 = bus.subscribe("vak.events.ws_test.>").await.expect("sub 2"); - 148
- 149
let env = MessageEnvelope::new( - 150
"vak://test", - 151
"vak.events.ws_test.sess_a.telemetry", - 152
"ws_test", - 153
"agent_a", - 154
1, - 155
MessageEnvelope::GENESIS_HASH, - 156
b"{\"rss_bytes\": 45000000}".to_vec(), - 157
None, - 158
); - 159
- 160
bus.publish("vak.events.ws_test.sess_a.telemetry", env.clone()) - 161
.await - 162
.expect("publish succeeds"); - 163
- 164
// Both subscribers should receive the envelope - 165
let received1 = tokio::time::timeout(Duration::from_millis(500), rx1.recv()) - 166
.await - 167
.expect("timeout") - 168
.expect("received"); - 169
assert_eq!(received1.id, env.id); - 170
- 171
let received2 = tokio::time::timeout(Duration::from_millis(500), rx2.recv()) - 172
.await - 173
.expect("timeout") - 174
.expect("received"); - 175
assert_eq!(received2.id, env.id); - 176
} - 177
- 178
#[tokio::test] - 179
async fn test_durable_work_queue_claim_and_ack() { - 180
let bus = InMemoryBus::new(); - 181
let queue = "code_review"; - 182
- 183
let env1 = MessageEnvelope::new( - 184
"vak://lead", - 185
"vak.work.task", - 186
"ws_test", - 187
"lead_agent", - 188
1, - 189
MessageEnvelope::GENESIS_HASH, - 190
b"task 1: review pull request".to_vec(), - 191
None, - 192
); - 193
let env2 = MessageEnvelope::new( - 194
"vak://lead", - 195
"vak.work.task", - 196
"ws_test", - 197
"lead_agent", - 198
2, - 199
env1.compute_hash(), - 200
b"task 2: run integration tests".to_vec(), - 201
None, - 202
); - 203
- 204
let id1 = bus.enqueue(queue, env1).await.expect("enqueue 1"); - 205
let id2 = bus.enqueue(queue, env2).await.expect("enqueue 2"); - 206
- 207
// Worker 1 claims first task - 208
let claim1 = bus - 209
.claim(queue, "worker_1", Duration::from_millis(500)) - 210
.await - 211
.expect("claim 1") - 212
.expect("task available"); - 213
assert_eq!(claim1.task_id, id1); - 214
claim1.ack_handle.ack().await.expect("ack 1 succeeds"); - 215
- 216
// Worker 2 claims second task - 217
let claim2 = bus - 218
.claim(queue, "worker_2", Duration::from_millis(500)) - 219
.await - 220
.expect("claim 2") - 221
.expect("task available"); - 222
assert_eq!(claim2.task_id, id2); - 223
claim2.ack_handle.ack().await.expect("ack 2 succeeds"); - 224
- 225
// Queue is now empty - 226
let claim3 = bus - 227
.claim(queue, "worker_1", Duration::from_millis(100)) - 228
.await - 229
.expect("claim empty"); - 230
assert!(claim3.is_none()); - 231
} - 232
- 233
#[tokio::test] - 234
async fn test_dead_letter_queue_on_max_retries() { - 235
let bus = InMemoryBus::new(); - 236
let queue = "fragile_work"; - 237
- 238
let env = MessageEnvelope::new( - 239
"vak://lead", - 240
"vak.work.task", - 241
"ws_test", - 242
"lead", - 243
1, - 244
MessageEnvelope::GENESIS_HASH, - 245
b"poison task".to_vec(), - 246
None, - 247
); - 248
- 249
let task_id = bus.enqueue(queue, env).await.expect("enqueue"); - 250
- 251
// Attempt 1: Nack with retry - 252
let claim1 = bus - 253
.claim(queue, "w1", Duration::from_millis(500)) - 254
.await - 255
.unwrap() - 256
.unwrap(); - 257
assert_eq!(claim1.attempts, 1); - 258
claim1.ack_handle.nack(true).await.expect("nack 1"); - 259
- 260
// Attempt 2: Nack with retry - 261
let claim2 = bus - 262
.claim(queue, "w2", Duration::from_millis(500)) - 263
.await - 264
.unwrap() - 265
.unwrap(); - 266
assert_eq!(claim2.attempts, 2); - 267
claim2.ack_handle.nack(true).await.expect("nack 2"); - 268
- 269
// Attempt 3: Exhausted! Reaches max_retries (3) and diverts to DLQ - 270
let claim3 = bus - 271
.claim(queue, "w3", Duration::from_millis(500)) - 272
.await - 273
.unwrap() - 274
.unwrap(); - 275
assert_eq!(claim3.attempts, 3); - 276
claim3.ack_handle.nack(true).await.expect("nack 3 (dlq)"); - 277
- 278
// Queue is now empty (poison pill removed) - 279
let empty = bus - 280
.claim(queue, "w1", Duration::from_millis(100)) - 281
.await - 282
.unwrap(); - 283
assert!(empty.is_none()); - 284
- 285
// Verified in DLQ - 286
let dlq = bus.dead_letters(); - 287
assert_eq!(dlq.len(), 1); - 288
assert_eq!(dlq[0].envelope_id, task_id); - 289
assert_eq!(dlq[0].attempts, 3); - 290
assert!(dlq[0].failure_reason.contains("exhausted")); - 291
- 292
// Metrics verified - 293
let metrics = bus.metrics().snapshot(); - 294
assert_eq!(metrics.dead_letter_count, 1); - 295
} - 296
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.