use std::time::Duration;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::app::models::FileInfo;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReservationStatus {
AlreadyExists,
Reserved,
ReservedByOther { worker_id: u32 },
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ReservationInfo {
pub worker_id: u32,
pub reserved_at: DateTime<Utc>,
pub status: ReservationState,
pub file_info: FileInfo,
pub retry_count: u32,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum ReservationState {
Reserved,
Downloading,
Completed,
Failed { error: String },
}
impl ReservationInfo {
pub fn new(worker_id: u32, file_info: FileInfo) -> Self {
Self {
worker_id,
reserved_at: Utc::now(),
status: ReservationState::Reserved,
file_info,
retry_count: 0,
}
}
pub fn is_timed_out(&self, timeout: Duration) -> bool {
let elapsed = Utc::now()
.signed_duration_since(self.reserved_at)
.to_std()
.unwrap_or(Duration::ZERO);
elapsed > timeout
}
pub fn mark_downloading(&mut self) {
self.status = ReservationState::Downloading;
}
pub fn mark_completed(&mut self) {
self.status = ReservationState::Completed;
}
pub fn mark_failed(&mut self, error: String) {
self.status = ReservationState::Failed { error };
self.retry_count += 1;
}
pub fn is_failed(&self) -> bool {
matches!(self.status, ReservationState::Failed { .. })
}
pub fn is_completed(&self) -> bool {
matches!(self.status, ReservationState::Completed)
}
pub fn is_downloading(&self) -> bool {
matches!(self.status, ReservationState::Downloading)
}
pub fn get_error(&self) -> Option<&str> {
match &self.status {
ReservationState::Failed { error } => Some(error),
_ => None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::app::hash::Md5Hash;
use crate::app::models::{DatasetFileInfo, QualityControlVersion};
use std::path::PathBuf;
fn create_test_file_info() -> FileInfo {
let hash = Md5Hash::from_hex("50c9d1c465f3cbff652be1509c2e2a4e").unwrap();
let dataset_info = DatasetFileInfo {
dataset_name: "test-dataset".to_string(),
version: "202407".to_string(),
county: Some("devon".to_string()),
station_id: Some("01381".to_string()),
station_name: Some("twist".to_string()),
quality_version: Some(QualityControlVersion::V1),
year: Some("1980".to_string()),
file_type: Some("data".to_string()),
};
FileInfo {
hash,
relative_path: "./data/test.csv".to_string(),
file_name: "test.csv".to_string(),
dataset_info,
manifest_version: Some(202507),
retry_count: 0,
last_attempt: None,
estimated_size: None,
destination_path: PathBuf::from("/tmp/test.csv"),
}
}
#[test]
fn test_reservation_creation() {
let file_info = create_test_file_info();
let reservation = ReservationInfo::new(42, file_info.clone());
assert_eq!(reservation.worker_id, 42);
assert_eq!(reservation.status, ReservationState::Reserved);
assert_eq!(reservation.file_info.hash, file_info.hash);
assert_eq!(reservation.retry_count, 0);
}
#[test]
fn test_reservation_state_transitions() {
let file_info = create_test_file_info();
let mut reservation = ReservationInfo::new(42, file_info);
assert_eq!(reservation.status, ReservationState::Reserved);
assert!(!reservation.is_downloading());
assert!(!reservation.is_completed());
assert!(!reservation.is_failed());
reservation.mark_downloading();
assert_eq!(reservation.status, ReservationState::Downloading);
assert!(reservation.is_downloading());
assert!(!reservation.is_completed());
assert!(!reservation.is_failed());
reservation.mark_completed();
assert_eq!(reservation.status, ReservationState::Completed);
assert!(!reservation.is_downloading());
assert!(reservation.is_completed());
assert!(!reservation.is_failed());
}
#[test]
fn test_reservation_failure() {
let file_info = create_test_file_info();
let mut reservation = ReservationInfo::new(42, file_info);
reservation.mark_failed("Network error".to_string());
assert!(reservation.is_failed());
assert_eq!(reservation.get_error(), Some("Network error"));
assert_eq!(reservation.retry_count, 1);
reservation.mark_failed("Another error".to_string());
assert_eq!(reservation.retry_count, 2);
assert_eq!(reservation.get_error(), Some("Another error"));
}
#[test]
fn test_reservation_timeout() {
let file_info = create_test_file_info();
let reservation = ReservationInfo::new(42, file_info);
assert!(!reservation.is_timed_out(Duration::from_secs(60)));
std::thread::sleep(Duration::from_millis(2));
assert!(reservation.is_timed_out(Duration::from_millis(1)));
}
#[test]
fn test_reservation_status_equality() {
assert_eq!(
ReservationStatus::AlreadyExists,
ReservationStatus::AlreadyExists
);
assert_eq!(ReservationStatus::Reserved, ReservationStatus::Reserved);
assert_eq!(
ReservationStatus::ReservedByOther { worker_id: 42 },
ReservationStatus::ReservedByOther { worker_id: 42 }
);
assert_ne!(
ReservationStatus::ReservedByOther { worker_id: 42 },
ReservationStatus::ReservedByOther { worker_id: 43 }
);
}
}