mocra-core 0.4.0

The mocra crawler framework runtime: errors, cache, utilities, domain models, downloader, data-plane queue, coordination, scheduler and engine.
Documentation
#![allow(unused)]
use deadpool_redis::redis;
use deadpool_redis::redis::AsyncCommands;
use serde_json::{Value, json};
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};

/// Read-only monitoring API over Redis-backed event data.
///
/// This helper is used by external observability tools to query counters,
/// progress snapshots, error summaries, and time-series metrics persisted by
/// event handlers.
pub struct RedisEventMonitor {
    redis_pool: Arc<deadpool_redis::Pool>,
    key_prefix: String,
}

impl RedisEventMonitor {
    async fn scan_keys(
        conn: &mut deadpool_redis::Connection,
        pattern: &str,
    ) -> Result<Vec<String>, Box<dyn std::error::Error + Send + Sync>> {
        let mut cursor: u64 = 0;
        let mut keys = Vec::new();

        loop {
            let (next_cursor, batch): (u64, Vec<String>) = redis::cmd("SCAN")
                .arg(cursor)
                .arg("MATCH")
                .arg(pattern)
                .arg("COUNT")
                .arg(500)
                .query_async(&mut *conn)
                .await?;
            keys.extend(batch);
            if next_cursor == 0 {
                break;
            }
            cursor = next_cursor;
        }

        Ok(keys)
    }

    pub fn new(redis_pool: Arc<deadpool_redis::Pool>, key_prefix: String) -> Self {
        Self {
            redis_pool,
            key_prefix,
        }
    }

    /// Returns aggregated real-time counters and derived success-rate metrics.
    pub async fn get_system_stats(
        &self,
    ) -> Result<Value, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;

        // Read high-level counters.
        let request_success: i64 = conn
            .get(format!("{}:stats:request_success_count", self.key_prefix))
            .await
            .unwrap_or(0);
        let request_fail: i64 = conn
            .get(format!("{}:stats:request_fail_count", self.key_prefix))
            .await
            .unwrap_or(0);
        let task_fail: i64 = conn
            .get(format!("{}:stats:task_fail_count", self.key_prefix))
            .await
            .unwrap_or(0);

        // Read per-event counters for selected core event keys.
        let mut event_counts = json!({});
        let event_types = [
            "engine.task_model.started",
            "engine.task_model.completed",
            "engine.task_model.failed",
            "engine.request_publish.started",
            "engine.request_publish.completed",
            "engine.request_publish.failed",
            "engine.download.started",
            "engine.download.completed",
            "engine.download.failed",
            "engine.parser_task_model.started",
            "engine.parser_task_model.completed",
            "engine.parser_task_model.failed",
        ];

        for event_type in &event_types {
            let count: i64 = conn
                .get(format!("{}:stats:count:{}", self.key_prefix, event_type))
                .await
                .unwrap_or(0);
            event_counts[event_type] = json!(count);
        }

        Ok(json!({
            "summary": {
                "request_success_count": request_success,
                "request_fail_count": request_fail,
                "task_fail_count": task_fail,
                "request_success_rate": if request_success + request_fail > 0 {
                    (request_success as f64) / ((request_success + request_fail) as f64) * 100.0
                } else { 0.0 }
            },
            "event_counts": event_counts,
            "timestamp": SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs()
        }))
    }

    /// Returns progress information for a specific task id.
    pub async fn get_task_progress(
        &self,
        task_id: &str,
    ) -> Result<Option<Value>, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;
        let progress_key = format!("{}:progress:task:{}", self.key_prefix, task_id);

        let progress_data: Option<String> = conn.get(&progress_key).await?;
        match progress_data {
            Some(data) => Ok(Some(serde_json::from_str(&data)?)),
            None => Ok(None),
        }
    }

    /// Returns progress entries for all active tasks.
    pub async fn get_all_task_progress(
        &self,
    ) -> Result<Value, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;
        let pattern = format!("{}:progress:task:*", self.key_prefix);

        let keys = Self::scan_keys(&mut conn, &pattern).await?;

        let mut tasks = json!({});

        for key in keys {
            if let Ok(Some(data)) = conn.get::<&str, Option<String>>(&key).await
                && let Ok(task_data) = serde_json::from_str::<Value>(&data)
            {
                let task_id = key.split(':').next_back().unwrap_or("unknown");
                tasks[task_id] = task_data;
            }
        }

        Ok(json!({
            "tasks": tasks,
            "total_count": tasks.as_object().map(|o| o.len()).unwrap_or(0),
            "timestamp": SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs()
        }))
    }

    /// Returns recent download performance records and summary statistics.
    pub async fn get_download_performance(
        &self,
        limit: usize,
    ) -> Result<Value, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;
        let perf_key = format!("{}:stats:download_performance", self.key_prefix);

        let performance_data: Vec<String> = redis::cmd("LRANGE")
            .arg(&perf_key)
            .arg(0)
            .arg((limit as isize) - 1)
            .query_async(&mut *conn)
            .await?;

        let mut performances = Vec::new();
        let mut total_duration = 0u64;
        let mut total_size = 0u64;

        for data_str in performance_data {
            if let Ok(data) = serde_json::from_str::<Value>(&data_str) {
                if let (Some(duration), Some(size)) =
                    (data["duration_ms"].as_u64(), data["response_size"].as_u64())
                {
                    total_duration += duration;
                    total_size += size;
                }
                performances.push(data);
            }
        }

        let count = performances.len() as u64;
        Ok(json!({
            "recent_downloads": performances,
            "statistics": {
                "count": count,
                "avg_duration_ms": total_duration.checked_div(count).unwrap_or(0),
                "avg_response_size": total_size.checked_div(count).unwrap_or(0),
                "total_duration_ms": total_duration,
                "total_response_size": total_size
            },
            "timestamp": SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs()
        }))
    }

    /// Returns recent failure events ordered by descending timestamp.
    pub async fn get_error_details(
        &self,
        limit: usize,
    ) -> Result<Value, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;

        // Read recent failure-related event keys.
        let error_pattern = format!("{}:events:*failed", self.key_prefix);
        let keys = Self::scan_keys(&mut conn, &error_pattern).await?;

        let mut errors = Vec::new();
        for key in keys.into_iter().take(limit) {
            if let Ok(Some(data)) = conn.get::<&str, Option<String>>(&key).await
                && let Ok(error_data) = serde_json::from_str::<Value>(&data)
            {
                errors.push(error_data);
            }
        }

        // Sort by timestamp.
        errors.sort_by(|a, b| {
            let ts_a = a["timestamp"].as_u64().unwrap_or(0);
            let ts_b = b["timestamp"].as_u64().unwrap_or(0);
            ts_b.cmp(&ts_a)
        });

        Ok(json!({
            "recent_errors": errors,
            "count": errors.len(),
            "timestamp": SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs()
        }))
    }

    /// Returns time-series points for `event_type` in the latest `hours` window.
    pub async fn get_timeseries_data(
        &self,
        event_type: &str,
        hours: u64,
    ) -> Result<Value, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;
        let ts_key = format!("{}:timeseries:{}", self.key_prefix, event_type);

        let current_time = SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs();
        let start_time = current_time - (hours * 3600);

        let data: Vec<String> = redis::cmd("ZRANGEBYSCORE")
            .arg(&ts_key)
            .arg(start_time as f64)
            .arg(current_time as f64)
            .query_async(&mut *conn)
            .await?;

        let mut timeseries = Vec::new();
        for item in data {
            if let Ok(ts_data) = serde_json::from_str::<Value>(&item) {
                timeseries.push(ts_data);
            }
        }

        Ok(json!({
            "event_type": event_type,
            "time_range_hours": hours,
            "data": timeseries,
            "count": timeseries.len(),
            "timestamp": current_time
        }))
    }

    /// Returns a synthesized health report from recent system health events.
    pub async fn get_health_report(
        &self,
    ) -> Result<Value, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;

        // Read recent health-check events.
        let health_pattern = format!("{}:events:system.system_health.*", self.key_prefix);
        let keys = Self::scan_keys(&mut conn, &health_pattern).await?;

        let mut health_data = json!({});
        for key in keys {
            if let Ok(Some(data)) = conn.get::<&str, Option<String>>(&key).await
                && let Ok(health_info) = serde_json::from_str::<Value>(&data)
                && let Some(component) = health_info["data"]["payload"]["component"].as_str()
            {
                health_data[component] = health_info["data"].clone();
            }
        }

        // Compute overall status from component-level statuses.
        let mut healthy_components = 0;
        let mut total_components = 0;

        if let Some(components) = health_data.as_object() {
            total_components = components.len();
            for (_, component_data) in components {
                if component_data["status"].as_str() == Some("healthy") {
                    healthy_components += 1;
                }
            }
        }

        let overall_health = if total_components > 0 {
            if healthy_components == total_components {
                "healthy"
            } else if healthy_components > 0 {
                "degraded"
            } else {
                "unhealthy"
            }
        } else {
            "unknown"
        };

        Ok(json!({
            "overall_status": overall_health,
            "healthy_components": healthy_components,
            "total_components": total_components,
            "components": health_data,
            "timestamp": SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs()
        }))
    }

    /// Cleans up expired time-series points and returns deleted entry count.
    pub async fn cleanup_expired_data(
        &self,
    ) -> Result<u64, Box<dyn std::error::Error + Send + Sync>> {
        let mut conn = self
            .redis_pool
            .get()
            .await
            .map_err(|e| format!("Redis connection failed: {e}"))?;
        let current_time = SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs();
        let cutoff_time = current_time - (24 * 3600);

        let mut cleaned_count = 0u64;

        // Remove old points from each time-series key.
        let ts_pattern = format!("{}:timeseries:*", self.key_prefix);
        let ts_keys = Self::scan_keys(&mut conn, &ts_pattern).await?;

        for key in ts_keys {
            let removed: u64 = redis::cmd("ZREMRANGEBYSCORE")
                .arg(&key)
                .arg(0.0)
                .arg(cutoff_time as f64)
                .query_async(&mut *conn)
                .await?;
            cleaned_count += removed;
        }

        Ok(cleaned_count)
    }
}