use std::sync::Arc;
use crate::observability::models::{ActivityEvent, ActivitySnapshot, DeadLetterRecord, QueueStats};
use crate::runner::error::WorkerError;
use crate::storage::InspectionStorage;
use futures::stream::{Stream, StreamExt};
use uuid::Uuid;
#[derive(Clone)]
pub struct QueueInspector {
max_workers: Option<usize>,
backend: Arc<dyn InspectionStorage>,
}
impl QueueInspector {
pub fn new(backend: Arc<dyn InspectionStorage>) -> Self {
Self {
max_workers: None,
backend,
}
}
pub fn with_max_workers(mut self, max_workers: usize) -> Self {
self.max_workers = Some(max_workers);
self
}
pub fn event_stream(&self) -> impl Stream<Item = Result<ActivityEvent, WorkerError>> {
let backend_stream = self.backend.event_stream();
backend_stream.map(|r| {
r.map(Self::convert_backend_event)
.map_err(WorkerError::from)
})
}
pub async fn list_pending(
&self,
offset: usize,
limit: usize,
) -> Result<Vec<ActivitySnapshot>, WorkerError> {
let snapshots = self
.backend
.list_pending(offset, limit)
.await
.map_err(WorkerError::from)?;
Ok(snapshots
.into_iter()
.map(Self::convert_backend_snapshot)
.collect())
}
pub async fn list_processing(
&self,
offset: usize,
limit: usize,
) -> Result<Vec<ActivitySnapshot>, WorkerError> {
let snapshots = self
.backend
.list_processing(offset, limit)
.await
.map_err(WorkerError::from)?;
Ok(snapshots
.into_iter()
.map(Self::convert_backend_snapshot)
.collect())
}
pub async fn list_scheduled(
&self,
offset: usize,
limit: usize,
) -> Result<Vec<ActivitySnapshot>, WorkerError> {
let snapshots = self
.backend
.list_scheduled(offset, limit)
.await
.map_err(WorkerError::from)?;
Ok(snapshots
.into_iter()
.map(Self::convert_backend_snapshot)
.collect())
}
pub async fn list_dead_letter(
&self,
offset: usize,
limit: usize,
) -> Result<Vec<DeadLetterRecord>, WorkerError> {
let records = self
.backend
.list_dead_letter(offset, limit)
.await
.map_err(WorkerError::from)?;
Ok(records
.into_iter()
.map(Self::convert_backend_dead_letter)
.collect())
}
pub async fn list_completed(
&self,
offset: usize,
limit: usize,
) -> Result<Vec<ActivitySnapshot>, WorkerError> {
let snapshots = self
.backend
.list_completed(offset, limit)
.await
.map_err(WorkerError::from)?;
Ok(snapshots
.into_iter()
.map(Self::convert_backend_snapshot)
.collect())
}
pub async fn get_activity(
&self,
activity_id: Uuid,
) -> Result<Option<ActivitySnapshot>, WorkerError> {
self.backend
.get_activity(activity_id)
.await
.map(|opt| opt.map(Self::convert_backend_snapshot))
.map_err(WorkerError::from)
}
pub async fn get_result(
&self,
activity_id: Uuid,
) -> Result<Option<serde_json::Value>, WorkerError> {
self.backend
.get_result(activity_id)
.await
.map(|opt| opt.and_then(|r| r.data))
.map_err(WorkerError::from)
}
pub async fn recent_events(
&self,
activity_id: Uuid,
limit: usize,
) -> Result<Vec<ActivityEvent>, WorkerError> {
let events = self
.backend
.get_activity_events(activity_id, limit)
.await
.map_err(WorkerError::from)?;
Ok(events
.into_iter()
.map(Self::convert_backend_event)
.collect())
}
pub async fn stats(&self) -> Result<QueueStats, WorkerError> {
let backend_stats = self.backend.stats().await.map_err(WorkerError::from)?;
Ok(QueueStats {
pending_activities: backend_stats.pending,
processing_activities: backend_stats.processing,
critical_priority: backend_stats.by_priority.critical,
high_priority: backend_stats.by_priority.high,
normal_priority: backend_stats.by_priority.normal,
low_priority: backend_stats.by_priority.low,
scheduled_activities: backend_stats.scheduled,
dead_letter_activities: backend_stats.dead_letter,
max_workers: self.max_workers.or(backend_stats.max_workers),
})
}
fn convert_backend_snapshot(snapshot: crate::storage::ActivitySnapshot) -> ActivitySnapshot {
ActivitySnapshot {
id: snapshot.id,
activity_type: snapshot.activity_type,
payload: snapshot.payload,
priority: snapshot.priority,
status: snapshot.status,
created_at: snapshot.created_at,
scheduled_at: snapshot.scheduled_at,
started_at: snapshot.started_at,
completed_at: snapshot.completed_at,
current_worker_id: snapshot.current_worker_id,
last_worker_id: snapshot.last_worker_id,
retry_count: snapshot.retry_count,
max_retries: snapshot.max_retries,
timeout_seconds: snapshot.timeout_seconds,
retry_delay_seconds: snapshot.retry_delay_seconds,
metadata: snapshot.metadata,
last_error: snapshot.last_error,
last_error_at: snapshot.last_error_at,
status_updated_at: snapshot.status_updated_at,
score: None,
lease_deadline_ms: None,
processing_member: None,
idempotency_key: snapshot.idempotency_key,
}
}
fn convert_backend_dead_letter(record: DeadLetterRecord) -> DeadLetterRecord {
DeadLetterRecord {
activity: Self::convert_backend_snapshot(record.activity),
error: record.error,
failed_at: record.failed_at,
}
}
fn convert_backend_event(event: ActivityEvent) -> ActivityEvent {
event
}
}