2026-01-02 16:29:05 +01:00
|
|
|
use log::{error, warn};
|
2026-08-27 18:36:23 +05:00
|
|
|
use serde::{Deserialize, Serialize};
|
|
|
|
|
use std::{
|
|
|
|
|
collections::HashMap,
|
|
|
|
|
fs,
|
|
|
|
|
path::{Path, PathBuf},
|
|
|
|
|
time::{SystemTime, UNIX_EPOCH},
|
|
|
|
|
};
|
|
|
|
|
use tuliprox_repository::{build_input_storage_path, get_input_storage_path};
|
2026-01-02 16:29:05 +01:00
|
|
|
|
|
|
|
|
pub const STATUS_FILE: &str = "status.json";
|
|
|
|
|
|
|
|
|
|
#[derive(Serialize, Deserialize, Default, Debug, Clone, PartialEq)]
|
|
|
|
|
pub enum ClusterState {
|
|
|
|
|
#[default]
|
|
|
|
|
Ok,
|
|
|
|
|
Failed,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Serialize, Deserialize, Debug, Clone)]
|
|
|
|
|
pub struct ClusterStatus {
|
|
|
|
|
pub status: ClusterState,
|
|
|
|
|
pub timestamp: u64,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Serialize, Deserialize, Debug, Default, Clone)]
|
|
|
|
|
pub struct InputStatus {
|
|
|
|
|
#[serde(default)]
|
|
|
|
|
pub clusters: HashMap<String, ClusterStatus>,
|
|
|
|
|
}
|
|
|
|
|
|
2026-03-06 17:29:42 +01:00
|
|
|
pub async fn resolve_input_storage_path(storage_dir: &str, input_name: &str) -> PathBuf {
|
2026-08-27 18:36:23 +05:00
|
|
|
if let Ok(path) = get_input_storage_path(input_name, storage_dir).await {
|
|
|
|
|
path
|
|
|
|
|
} else {
|
2026-03-06 17:29:42 +01:00
|
|
|
build_input_storage_path(input_name, storage_dir)
|
2026-01-02 16:29:05 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn load_input_status(path: &Path) -> InputStatus {
|
|
|
|
|
let status_path = path.join(STATUS_FILE);
|
|
|
|
|
if status_path.exists() {
|
|
|
|
|
match fs::read_to_string(&status_path) {
|
|
|
|
|
Ok(content) => match serde_json::from_str(&content) {
|
|
|
|
|
Ok(status) => return status,
|
|
|
|
|
Err(e) => warn!("Failed to parse input status file {}: {e}", status_path.display()),
|
|
|
|
|
},
|
|
|
|
|
Err(e) => warn!("Failed to read input status file {}: {e}", status_path.display()),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
InputStatus::default()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn save_input_status(path: &Path, status: &InputStatus) {
|
|
|
|
|
let status_path = path.join(STATUS_FILE);
|
|
|
|
|
if let Some(parent) = status_path.parent() {
|
|
|
|
|
if let Err(e) = fs::create_dir_all(parent) {
|
2026-08-27 18:36:23 +05:00
|
|
|
error!("Failed to create input storage directory {}: {e}", parent.display());
|
|
|
|
|
return;
|
2026-01-02 16:29:05 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
match serde_json::to_string_pretty(status) {
|
|
|
|
|
Ok(content) => {
|
|
|
|
|
if let Err(e) = fs::write(&status_path, content) {
|
|
|
|
|
error!("Failed to write input status file {}: {e}", status_path.display());
|
|
|
|
|
}
|
2026-08-27 18:36:23 +05:00
|
|
|
}
|
2026-01-02 16:29:05 +01:00
|
|
|
Err(e) => error!("Failed to serialize input status: {e}"),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn is_cache_valid(status: &InputStatus, cluster: &str, cache_duration_seconds: u64) -> bool {
|
|
|
|
|
if cache_duration_seconds == 0 {
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
if let Some(cluster_status) = status.clusters.get(cluster) {
|
|
|
|
|
if cluster_status.status != ClusterState::Ok {
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
|
2026-08-21 17:50:12 +05:00
|
|
|
if now >= cluster_status.timestamp {
|
2026-08-27 18:36:23 +05:00
|
|
|
return now - cluster_status.timestamp < cache_duration_seconds;
|
2026-01-02 16:29:05 +01:00
|
|
|
}
|
|
|
|
|
// Timestamp in future? Invalid.
|
|
|
|
|
return false;
|
|
|
|
|
}
|
|
|
|
|
false
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn update_cluster_status(status: &mut InputStatus, cluster: &str, state: ClusterState) {
|
|
|
|
|
let now = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_secs();
|
2026-08-27 18:36:23 +05:00
|
|
|
status.clusters.insert(cluster.to_string(), ClusterStatus { status: state, timestamp: now });
|
2026-01-02 16:29:05 +01:00
|
|
|
}
|