use std::collections::HashMap;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use uuid::Uuid;
use crate::activity::activity::ActivityStatus;
use crate::ActivityPriority;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ActivitySnapshot {
pub id: Uuid,
pub activity_type: String,
pub payload: Value,
pub priority: ActivityPriority,
pub status: ActivityStatus,
pub created_at: DateTime<Utc>,
pub scheduled_at: Option<DateTime<Utc>>,
pub started_at: Option<DateTime<Utc>>,
pub completed_at: Option<DateTime<Utc>>,
pub current_worker_id: Option<String>,
pub last_worker_id: Option<String>,
pub retry_count: u32,
pub max_retries: u32,
pub timeout_seconds: u64,
pub retry_delay_seconds: u64,
pub metadata: HashMap<String, String>,
pub last_error: Option<String>,
pub last_error_at: Option<DateTime<Utc>>,
pub status_updated_at: DateTime<Utc>,
#[serde(skip_serializing_if = "Option::is_none")]
pub score: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub lease_deadline_ms: Option<i64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub processing_member: Option<String>,
pub idempotency_key: Option<String>,
}
impl ActivitySnapshot {
pub fn update_status(&mut self, status: ActivityStatus, timestamp: DateTime<Utc>) {
self.status = status;
self.status_updated_at = timestamp;
}
pub fn mark_started(&mut self, started_at: DateTime<Utc>) {
self.started_at = Some(started_at);
}
pub fn mark_completed(&mut self, completed_at: DateTime<Utc>) {
self.completed_at = Some(completed_at);
}
pub fn set_current_worker_id(&mut self, worker_id: Option<String>) {
self.current_worker_id = worker_id;
}
pub fn set_last_worker_id(&mut self, worker_id: Option<String>) {
self.last_worker_id = worker_id;
}
pub fn set_last_error(&mut self, error: Option<String>, at: Option<DateTime<Utc>>) {
self.last_error = error;
self.last_error_at = at;
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub enum ActivityEventType {
Enqueued,
Scheduled,
Dequeued,
Started,
Completed,
Failed,
Retrying,
DeadLetter,
Requeued,
LeaseExtended,
ResultStored,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ActivityEvent {
pub activity_id: Uuid,
pub timestamp: DateTime<Utc>,
pub event_type: ActivityEventType,
pub worker_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub detail: Option<Value>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct DeadLetterRecord {
pub activity: ActivitySnapshot,
pub error: String,
pub failed_at: DateTime<Utc>,
}
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct QueueStats {
pub pending_activities: u64,
pub processing_activities: u64,
pub critical_priority: u64,
pub high_priority: u64,
pub normal_priority: u64,
pub low_priority: u64,
pub scheduled_activities: u64,
pub dead_letter_activities: u64,
pub max_workers: Option<usize>,
}