- 1
//! Global event hub (docs/design/23-memory.md): a single `tokio::broadcast` - 2
//! channel that every subsystem writes to and the admin console SSE endpoint - 3
//! reads from. Events are append-only and lossy by design — slow consumers - 4
//! get `Lagged` errors, not backpressure. - 5
- 6
#![allow(dead_code)] - 7
- 8
use std::sync::OnceLock; - 9
- 10
use serde::Serialize; - 11
use std::sync::Arc; - 12
use tokio::sync::broadcast; - 13
- 14
/// Capacity of the global broadcast channel. Slow consumers that fall - 15
/// more than this many events behind receive `Lagged` and must reconnect. - 16
const HUB_CAPACITY: usize = 4096; - 17
- 18
#[derive(Debug, Clone, Serialize)] - 19
#[serde(tag = "type", content = "data")] - 20
pub enum SystemEvent { - 21
// ---- agent lifecycle (forwarded from AgentEvent) ---- - 22
Agent(AgentEventPayload), - 23
- 24
// ---- session lifecycle ---- - 25
SessionCreated { - 26
session_id: String, - 27
project_hash: String, - 28
}, - 29
SessionEntryAppended { - 30
session_id: String, - 31
entry_id: String, - 32
kind: String, - 33
}, - 34
- 35
// ---- config / gateway ---- - 36
ConfigChanged { - 37
label: String, - 38
detail: String, - 39
}, - 40
GatewayInbound { - 41
surface: String, - 42
who: String, - 43
preview: String, - 44
}, - 45
- 46
// ---- approvals ---- - 47
ApprovalRequested { - 48
id: String, - 49
session_id: String, - 50
tool: String, - 51
reason: String, - 52
}, - 53
ApprovalGranted { - 54
id: String, - 55
tool: String, - 56
}, - 57
ApprovalDenied { - 58
id: String, - 59
tool: String, - 60
}, - 61
- 62
// ---- security ---- - 63
SecurityEvent { - 64
kind: String, - 65
label: String, - 66
}, - 67
- 68
// ---- provider health ---- - 69
ProviderError { - 70
provider: String, - 71
model: String, - 72
error: String, - 73
}, - 74
RateLimit { - 75
provider: String, - 76
retry_after_secs: Option<u64>, - 77
}, - 78
- 79
// ---- heartbeat ---- - 80
Heartbeat, - 81
} - 82
- 83
#[derive(Debug, Clone, Serialize)] - 84
pub struct AgentEventPayload { - 85
pub summary: String, - 86
#[serde(skip_serializing_if = "Option::is_none")] - 87
pub detail: Option<String>, - 88
} - 89
- 90
/// Global event hub singleton. Created once at server startup; every - 91
/// subsystem clones the `Sender` side. Consumers subscribe via `subscribe()`. - 92
/// - 93
/// When a `ServerBus` has been attached via `set_server_bus`, `emit` also - 94
/// fans the event out to the distributed fabric (vak-bus, docs/design/53). - 95
/// This is fire-and-forget: the local broadcast is authoritative and the bus - 96
/// is supplementary for distributed subscribers. - 97
#[derive(Debug, Clone)] - 98
pub struct EventHub { - 99
tx: broadcast::Sender<SystemEvent>, - 100
bus: Option<Arc<crate::bus::ServerBus>>, - 101
} - 102
- 103
impl EventHub { - 104
/// Create a new hub. Only call this once at server startup. - 105
pub fn new() -> Self { - 106
let (tx, _) = broadcast::channel(HUB_CAPACITY); - 107
EventHub { tx, bus: None } - 108
} - 109
- 110
/// Attach a distributed event bus. Events emitted after this call are - 111
/// also published to the bus fabric for remote subscribers. - 112
pub fn set_server_bus(&mut self, bus: Arc<crate::bus::ServerBus>) { - 113
self.bus = Some(bus); - 114
} - 115
- 116
/// Emit an event. Never blocks — dropped if no receivers are alive. - 117
/// - 118
/// When a `ServerBus` is attached, the event is also published to the - 119
/// distributed fabric via a spawned tokio task (fire-and-forget). - 120
pub fn emit(&self, event: SystemEvent) { - 121
let _ = self.tx.send(event.clone()); - 122
if let Some(bus) = &self.bus { - 123
let session_id = event_session_id(&event); - 124
let bus = bus.clone(); - 125
tokio::spawn(async move { - 126
if let Err(e) = bus.emit(&event, session_id.as_deref()).await { - 127
eprintln!("vak-server: ServerBus emit failed: {e}"); - 128
} - 129
}); - 130
} - 131
} - 132
- 133
/// Create a new receiver. Lagged receivers get `Lagged` errors on - 134
/// `recv()` and should reconnect. - 135
pub fn subscribe(&self) -> broadcast::Receiver<SystemEvent> { - 136
self.tx.subscribe() - 137
} - 138
- 139
/// Downgrade to a weak sender for non-owning references. - 140
pub fn sender(&self) -> broadcast::Sender<SystemEvent> { - 141
self.tx.clone() - 142
} - 143
} - 144
- 145
impl Default for EventHub { - 146
fn default() -> Self { - 147
Self::new() - 148
} - 149
} - 150
- 151
impl EventHub { - 152
/// Expose the attached `ServerBus` status, if any. - 153
/// Used by `/config/bus` and `/ops/center` to report distributed-fabric - 154
/// readiness to the admin console. - 155
pub fn bus_status(&self) -> Option<serde_json::Value> { - 156
self.bus.as_ref().map(|b| b.status()) - 157
} - 158
} - 159
- 160
/// Extract the session_id from a SystemEvent, if it carries one. - 161
/// Used for routing events to the correct vak-bus subject. - 162
fn event_session_id(event: &SystemEvent) -> Option<String> { - 163
match event { - 164
SystemEvent::SessionCreated { session_id, .. } - 165
| SystemEvent::SessionEntryAppended { session_id, .. } - 166
| SystemEvent::ApprovalRequested { session_id, .. } => Some(session_id.clone()), - 167
_ => None, - 168
} - 169
} - 170
- 171
// --------------------------------------------------------------------------- - 172
// Global singleton accessor - 173
// --------------------------------------------------------------------------- - 174
- 175
static GLOBAL_HUB: OnceLock<EventHub> = OnceLock::new(); - 176
- 177
/// Initialize the global event hub. Call once at server startup before - 178
/// any subsystem tries `global()`. - 179
pub fn init_global() -> EventHub { - 180
GLOBAL_HUB.get_or_init(EventHub::new).clone() - 181
} - 182
- 183
/// Access the global hub. Returns `None` if `init_global()` was never called. - 184
pub fn global() -> Option<EventHub> { - 185
GLOBAL_HUB.get().cloned() - 186
} - 187
- 188
// --------------------------------------------------------------------------- - 189
// Convenience helpers - 190
// --------------------------------------------------------------------------- - 191
- 192
impl EventHub { - 193
pub fn emit_session_created(&self, session_id: &str, project_hash: &str) { - 194
self.emit(SystemEvent::SessionCreated { - 195
session_id: session_id.to_string(), - 196
project_hash: project_hash.to_string(), - 197
}); - 198
} - 199
- 200
pub fn emit_entry_appended(&self, session_id: &str, entry_id: &str, kind: &str) { - 201
self.emit(SystemEvent::SessionEntryAppended { - 202
session_id: session_id.to_string(), - 203
entry_id: entry_id.to_string(), - 204
kind: kind.to_string(), - 205
}); - 206
} - 207
- 208
pub fn emit_config_changed(&self, label: &str, detail: &str) { - 209
self.emit(SystemEvent::ConfigChanged { - 210
label: label.to_string(), - 211
detail: detail.to_string(), - 212
}); - 213
} - 214
- 215
pub fn emit_gateway_inbound(&self, surface: &str, who: &str, preview: &str) { - 216
self.emit(SystemEvent::GatewayInbound { - 217
surface: surface.to_string(), - 218
who: who.to_string(), - 219
preview: preview.to_string(), - 220
}); - 221
} - 222
- 223
pub fn emit_agent_summary(&self, summary: &str, detail: Option<String>) { - 224
self.emit(SystemEvent::Agent(AgentEventPayload { - 225
summary: summary.to_string(), - 226
detail, - 227
})); - 228
} - 229
- 230
pub fn emit_security(&self, kind: &str, label: &str) { - 231
self.emit(SystemEvent::SecurityEvent { - 232
kind: kind.to_string(), - 233
label: label.to_string(), - 234
}); - 235
} - 236
} - 237
- 238
#[cfg(test)] - 239
mod tests { - 240
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 241
- 242
use super::*; - 243
- 244
#[test] - 245
fn emit_and_receive() { - 246
let hub = EventHub::new(); - 247
let mut rx = hub.subscribe(); - 248
hub.emit(SystemEvent::SessionCreated { - 249
session_id: "s1".into(), - 250
project_hash: "abc".into(), - 251
}); - 252
let event = rx.try_recv().unwrap(); - 253
match event { - 254
SystemEvent::SessionCreated { session_id, .. } => { - 255
assert_eq!(session_id, "s1"); - 256
} - 257
_ => panic!("wrong event type"), - 258
} - 259
} - 260
- 261
#[test] - 262
fn multiple_receivers_get_independent_copies() { - 263
let hub = EventHub::new(); - 264
let mut rx1 = hub.subscribe(); - 265
let mut rx2 = hub.subscribe(); - 266
hub.emit(SystemEvent::Heartbeat); - 267
assert!(rx1.try_recv().is_ok()); - 268
assert!(rx2.try_recv().is_ok()); - 269
} - 270
- 271
#[test] - 272
fn lagged_receiver_gets_error() { - 273
let hub = EventHub::new(); - 274
let rx = hub.subscribe(); - 275
drop(rx); // drop receiver - 276
// Fill the buffer beyond capacity. - 277
for _ in 0..5000 { - 278
hub.emit(SystemEvent::Heartbeat); - 279
} - 280
// Re-subscribe; old position is lost. - 281
let mut rx = hub.subscribe(); - 282
// The first recv may be Lagged or a valid event — both ok. - 283
let result = rx.try_recv(); - 284
// We may get Lagged or a valid event depending on timing — both ok. - 285
// But at minimum, the hub still works. - 286
drop(result); - 287
hub.emit(SystemEvent::Heartbeat); - 288
assert!(rx.try_recv().is_ok()); - 289
} - 290
- 291
#[test] - 292
fn convenience_helpers_produce_correct_events() { - 293
let hub = EventHub::new(); - 294
let mut rx = hub.subscribe(); - 295
- 296
hub.emit_config_changed("mode", "restricted"); - 297
if let SystemEvent::ConfigChanged { label, detail } = rx.try_recv().unwrap() { - 298
assert_eq!(label, "mode"); - 299
assert_eq!(detail, "restricted"); - 300
} else { - 301
panic!("expected ConfigChanged"); - 302
} - 303
- 304
hub.emit_security("AuthFailure", "bad token"); - 305
if let SystemEvent::SecurityEvent { kind, label } = rx.try_recv().unwrap() { - 306
assert_eq!(kind, "AuthFailure"); - 307
assert_eq!(label, "bad token"); - 308
} else { - 309
panic!("expected SecurityEvent"); - 310
} - 311
} - 312
} - 313
- 314
// ---- per-session replay bus (docs/design/48-web-client.md §4.4) ------------ - 315
- 316
/// One agent event with the sequence number a client resumes from. - 317
#[derive(Debug, Clone)] - 318
pub struct SeqEvent { - 319
pub seq: u64, - 320
pub event: vak_agent::AgentEvent, - 321
} - 322
- 323
/// A session's live event channel, plus a bounded replay ring. - 324
/// - 325
/// A bare `tokio::broadcast` gives a late or reconnecting subscriber - 326
/// nothing: the events it missed are gone, not delayed. On loopback that - 327
/// is a rare race the client papers over with a ten-second reconciliation - 328
/// poll. Over a WAN — a laptop lid closing, a phone changing cell, a proxy - 329
/// idling out a stream — it is the common case, and a run's reply can be - 330
/// durably logged while the UI still says "Working". - 331
/// - 332
/// So every event carries a monotonic `seq`, the last `CAPACITY` of them - 333
/// are retained, and a reconnect with `Last-Event-ID` gets the gap. When - 334
/// the gap is older than the ring, the client is told to resync rather - 335
/// than being handed a silently incomplete stream — a missing event it - 336
/// does not know is missing is worse than an explicit reload. - 337
#[derive(Clone)] - 338
pub struct EventBus { - 339
tx: broadcast::Sender<SeqEvent>, - 340
ring: std::sync::Arc<std::sync::Mutex<std::collections::VecDeque<SeqEvent>>>, - 341
next: std::sync::Arc<std::sync::atomic::AtomicU64>, - 342
/// Count of subscribers created via [`Self::subscribe_internal`] — the - 343
/// handle's own in-process projector, never a real client. Held for the - 344
/// handle's whole lifetime, so it only ever grows; [`Self:: - 345
/// external_subscribers`] subtracts it from `receiver_count()` so - 346
/// "is anyone actually watching" (the 2s attach wait, idle eviction) - 347
/// answers about real clients instead of being permanently pinned above - 348
/// zero by the projector that is always there. - 349
internal_subscribers: std::sync::Arc<std::sync::atomic::AtomicUsize>, - 350
} - 351
- 352
impl EventBus { - 353
const CAPACITY: usize = 1024; - 354
- 355
pub fn new() -> Self { - 356
let (tx, _) = broadcast::channel(Self::CAPACITY); - 357
EventBus { - 358
tx, - 359
ring: std::sync::Arc::new(std::sync::Mutex::new( - 360
std::collections::VecDeque::with_capacity(Self::CAPACITY), - 361
)), - 362
next: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(1)), - 363
internal_subscribers: std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0)), - 364
} - 365
} - 366
- 367
/// Number and retain an event, then broadcast it. - 368
/// - 369
/// Retention happens BEFORE the broadcast so a subscriber that reads - 370
/// the ring immediately after receiving a live event cannot observe a - 371
/// ring that is missing it. - 372
pub fn send(&self, event: vak_agent::AgentEvent) -> usize { - 373
let seq = self.next.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - 374
let framed = SeqEvent { seq, event }; - 375
{ - 376
let mut ring = self - 377
.ring - 378
.lock() - 379
.unwrap_or_else(std::sync::PoisonError::into_inner); - 380
if ring.len() == Self::CAPACITY { - 381
ring.pop_front(); - 382
} - 383
ring.push_back(framed.clone()); - 384
} - 385
self.tx.send(framed).unwrap_or(0) - 386
} - 387
- 388
pub fn subscribe(&self) -> broadcast::Receiver<SeqEvent> { - 389
self.tx.subscribe() - 390
} - 391
- 392
/// Subscribe as the handle's own in-process projector rather than an - 393
/// external client. The receiver behaves identically — this only marks - 394
/// it so [`Self::external_subscribers`] can tell the two apart; the - 395
/// subscription still counts toward `receiver_count()`, and toward - 396
/// keeping `send` from silently landing on zero receivers. - 397
pub fn subscribe_internal(&self) -> broadcast::Receiver<SeqEvent> { - 398
self.internal_subscribers - 399
.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - 400
self.tx.subscribe() - 401
} - 402
- 403
pub fn receiver_count(&self) -> usize { - 404
self.tx.receiver_count() - 405
} - 406
- 407
/// Subscribers that are real clients (SSE connections), excluding the - 408
/// handle's own always-on internal projector. - 409
/// - 410
/// `receiver_count()` alone can never observe zero once - 411
/// [`Self::subscribe_internal`] has been called once, which made two - 412
/// callers that actually mean "is anyone external watching" — the 2s - 413
/// attach wait in `run_prompt` and idle-session eviction — permanently - 414
/// unable to see "no external subscribers" once the projector attached - 415
/// at handle creation. - 416
pub fn external_subscribers(&self) -> usize { - 417
self.receiver_count().saturating_sub( - 418
self.internal_subscribers - 419
.load(std::sync::atomic::Ordering::SeqCst), - 420
) - 421
} - 422
- 423
/// Events strictly after `seq`, oldest first. - 424
/// - 425
/// `None` means the ring no longer covers that point and the caller - 426
/// must resync from the durable transcript instead. Note the - 427
/// distinction from `Some(vec![])`, which means "you are up to date" — - 428
/// conflating the two is how a client silently loses a turn. - 429
pub fn replay_after(&self, seq: u64) -> Option<Vec<SeqEvent>> { - 430
let ring = self - 431
.ring - 432
.lock() - 433
.unwrap_or_else(std::sync::PoisonError::into_inner); - 434
match ring.front() { - 435
// Nothing retained yet: only "no events at all" is consistent - 436
// with any resume point, and that is exactly an empty replay. - 437
None => Some(Vec::new()), - 438
// The oldest retained event is already past the client's - 439
// resume point, so whatever sits between them is unrecoverable. - 440
Some(oldest) if oldest.seq > seq.saturating_add(1) => None, - 441
_ => Some(ring.iter().filter(|e| e.seq > seq).cloned().collect()), - 442
} - 443
} - 444
} - 445
- 446
impl Default for EventBus { - 447
fn default() -> Self { - 448
Self::new() - 449
} - 450
} - 451
- 452
#[cfg(test)] - 453
#[allow(clippy::unwrap_used)] - 454
mod bus_tests { - 455
use super::*; - 456
use vak_agent::AgentEvent; - 457
- 458
fn note(n: u32) -> AgentEvent { - 459
AgentEvent::RetryScheduled { - 460
attempt: n, - 461
delay_ms: 0, - 462
reason: String::new(), - 463
} - 464
} - 465
- 466
#[test] - 467
fn a_reconnect_receives_exactly_the_gap() { - 468
let bus = EventBus::new(); - 469
bus.send(note(1)); - 470
bus.send(note(2)); - 471
bus.send(note(3)); - 472
// Resuming after the first event yields the two it missed, in order. - 473
let gap = bus.replay_after(1).unwrap(); - 474
assert_eq!(gap.len(), 2); - 475
assert_eq!(gap[0].seq, 2); - 476
assert_eq!(gap[1].seq, 3); - 477
} - 478
- 479
#[test] - 480
fn being_up_to_date_is_not_the_same_as_having_lost_events() { - 481
let bus = EventBus::new(); - 482
bus.send(note(1)); - 483
// Caught up: an empty replay, NOT a resync. - 484
assert_eq!(bus.replay_after(1).unwrap().len(), 0); - 485
} - 486
- 487
#[test] - 488
fn a_resume_point_older_than_the_ring_asks_for_a_resync() { - 489
let bus = EventBus::new(); - 490
for i in 0..(EventBus::CAPACITY as u32 + 10) { - 491
bus.send(note(i)); - 492
} - 493
// Seq 1 fell out of the ring long ago; the caller must not be told - 494
// "here is the gap" when the gap cannot be produced. - 495
assert!(bus.replay_after(1).is_none()); - 496
// The newest events are still resumable. - 497
let newest = EventBus::CAPACITY as u64 + 5; - 498
assert!(bus.replay_after(newest).is_some()); - 499
} - 500
- 501
#[test] - 502
fn an_empty_bus_replays_nothing_rather_than_demanding_a_resync() { - 503
let bus = EventBus::new(); - 504
assert_eq!(bus.replay_after(0).unwrap().len(), 0); - 505
} - 506
- 507
#[test] - 508
fn a_max_resume_cursor_cannot_overflow_the_ring_boundary_check() { - 509
let bus = EventBus::new(); - 510
bus.send(note(1)); - 511
assert_eq!(bus.replay_after(u64::MAX).unwrap().len(), 0); - 512
} - 513
- 514
#[test] - 515
fn internal_subscriber_is_excluded_from_external_subscribers() { - 516
let bus = EventBus::new(); - 517
assert_eq!(bus.external_subscribers(), 0); - 518
let _internal = bus.subscribe_internal(); - 519
assert_eq!(bus.receiver_count(), 1); - 520
assert_eq!( - 521
bus.external_subscribers(), - 522
0, - 523
"the always-on internal projector must never look like a real client" - 524
); - 525
let external = bus.subscribe(); - 526
assert_eq!(bus.receiver_count(), 2); - 527
assert_eq!(bus.external_subscribers(), 1); - 528
drop(external); - 529
// tokio::broadcast decrements receiver_count synchronously on drop. - 530
assert_eq!(bus.external_subscribers(), 0); - 531
} - 532
- 533
#[test] - 534
fn multiple_external_subscribers_all_count() { - 535
let bus = EventBus::new(); - 536
let _internal = bus.subscribe_internal(); - 537
let _a = bus.subscribe(); - 538
let _b = bus.subscribe(); - 539
assert_eq!(bus.external_subscribers(), 2); - 540
} - 541
} - 542
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.