use crate::request::Request;
use std::{sync::Arc, time::Duration};
#[derive(Debug, Clone)]
pub struct HandledRequest {
pub request: Arc<Request>,
pub loaded_url: Option<url::Url>,
pub outcome: RequestFinalState,
pub response_status: Option<http::StatusCode>,
pub retry_count: u32,
pub duration: Duration,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum RequestFinalState {
Succeeded,
Failed,
Skipped,
}
#[derive(Debug, Clone, Default)]
#[non_exhaustive]
pub struct SystemSnapshot {
pub created_at: Option<time::OffsetDateTime>,
pub cpu_used_ratio: Option<f32>,
pub memory_used_bytes: Option<u64>,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub enum CrawlerEvent {
PersistState {
is_migrating: bool,
},
RequestFinished(HandledRequest),
RequestFailed {
request: Arc<Request>,
error: String,
},
SystemInfo(SystemSnapshot),
Aborting,
Exiting,
}
pub type EventStream = tokio::sync::broadcast::Receiver<CrawlerEvent>;
pub type ResultStream = tokio::sync::broadcast::Receiver<HandledRequest>;
#[derive(Debug, Clone)]
pub struct EventBus {
tx: tokio::sync::broadcast::Sender<CrawlerEvent>,
}
impl EventBus {
pub fn new(capacity: usize) -> Self {
assert!(capacity > 0, "event bus capacity must be greater than zero");
let (tx, _) = tokio::sync::broadcast::channel(capacity);
Self { tx }
}
pub fn subscribe(&self) -> EventStream {
self.tx.subscribe()
}
pub fn emit(&self, event: CrawlerEvent) {
let _ = self.tx.send(event);
}
pub fn subscriber_count(&self) -> usize {
self.tx.receiver_count()
}
}
impl Default for EventBus {
fn default() -> Self {
Self::new(1024)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn emit_without_subscribers_does_not_panic() {
EventBus::default().emit(CrawlerEvent::Exiting);
}
#[tokio::test]
async fn subscriber_receives_emitted_event() {
let bus = EventBus::default();
let mut subscriber = bus.subscribe();
bus.emit(CrawlerEvent::Aborting);
assert!(matches!(
subscriber.recv().await,
Ok(CrawlerEvent::Aborting)
));
}
#[tokio::test]
async fn all_subscribers_receive_the_same_event() {
let bus = EventBus::default();
let mut first = bus.subscribe();
let mut second = bus.subscribe();
bus.emit(CrawlerEvent::PersistState { is_migrating: true });
assert!(matches!(
first.recv().await,
Ok(CrawlerEvent::PersistState { is_migrating: true })
));
assert!(matches!(
second.recv().await,
Ok(CrawlerEvent::PersistState { is_migrating: true })
));
}
#[tokio::test]
async fn late_subscriber_does_not_receive_earlier_event() {
let bus = EventBus::default();
let mut existing = bus.subscribe();
bus.emit(CrawlerEvent::Exiting);
let mut late = bus.subscribe();
assert!(matches!(existing.recv().await, Ok(CrawlerEvent::Exiting)));
assert!(matches!(
late.try_recv(),
Err(tokio::sync::broadcast::error::TryRecvError::Empty)
));
}
}