use std::path::PathBuf;
use std::time::Instant;
use tokio::sync::broadcast;
use crate::types::{MacAddress, VmId};
use crate::vm::VmState;
#[derive(Clone, Debug)]
pub struct VmEvent {
pub timestamp: Instant,
pub vm_id: VmId,
pub kind: VmEventKind,
}
#[derive(Clone, Debug)]
pub enum VmEventKind {
StateChanged { from: VmState, to: VmState },
BootComplete { duration_ms: u64 },
ShutdownRequested,
Crashed { reason: String },
NetworkUp { mac: MacAddress },
NetworkDown { reason: String },
SnapshotCreated { path: PathBuf },
SnapshotRestored { path: PathBuf },
IpAssigned { ip: String },
ForceStop,
}
pub struct VmEventBus {
tx: broadcast::Sender<VmEvent>,
}
impl VmEventBus {
pub fn new(capacity: usize) -> Self {
let (tx, _) = broadcast::channel(capacity);
Self { tx }
}
pub fn emit(&self, event: VmEvent) {
let _ = self.tx.send(event);
}
pub fn state_changed(&self, vm_id: VmId, from: VmState, to: VmState) {
self.emit(VmEvent {
timestamp: Instant::now(),
vm_id,
kind: VmEventKind::StateChanged { from, to },
});
}
pub fn subscribe(&self) -> broadcast::Receiver<VmEvent> {
self.tx.subscribe()
}
}
impl Default for VmEventBus {
fn default() -> Self {
Self::new(256)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::types::VmId;
use crate::vm::VmState;
#[test]
fn event_bus_emits_to_subscriber() {
let bus = VmEventBus::default();
let mut rx = bus.subscribe();
bus.state_changed(VmId::from("test-vm"), VmState::Stopped, VmState::Running);
let event = rx.try_recv().unwrap();
assert_eq!(event.vm_id, VmId::from("test-vm"));
match event.kind {
VmEventKind::StateChanged { from, to } => {
assert_eq!(from, VmState::Stopped);
assert_eq!(to, VmState::Running);
}
other => panic!("expected StateChanged, got {other:?}"),
}
}
#[test]
fn event_bus_no_subscribers_does_not_panic() {
let bus = VmEventBus::default();
bus.state_changed(VmId::from("x"), VmState::Running, VmState::Stopped);
}
#[test]
fn event_bus_multiple_subscribers() {
let bus = VmEventBus::default();
let mut rx1 = bus.subscribe();
let mut rx2 = bus.subscribe();
bus.state_changed(VmId::from("vm"), VmState::Running, VmState::Paused);
assert!(rx1.try_recv().is_ok());
assert!(rx2.try_recv().is_ok());
}
}