use crate::downloader::DownloaderManager;
use crate::engine::events::EventBus;
use crate::engine::monitor::SystemMonitor;
use crate::engine::zombie;
use crate::queue::Identifiable;
use crate::utils::device_info::get_primary_local_ip;
use crate::common::policy::{DlqPolicy, PolicyResolver};
use crate::engine::chain::{
UnifiedTaskIngressChain, create_download_chain, create_parser_chain,
create_unified_task_ingress_chain,
};
use crate::engine::events::{EventEnvelope, EventPhase, EventType, HealthCheckEvent};
use metrics::counter;
use crate::common::state::State;
use crate::engine::chain::stream_chain::create_wss_download_chain;
use crate::queue::{QueueManager, QueuedItem};
use futures::{FutureExt, StreamExt};
use log::{error, info, warn};
use crate::common::interface::{
DataMiddlewareHandle, DataStoreMiddlewareHandle, DownloadMiddlewareHandle, MiddlewareManager,
ModuleTrait,
};
use crate::common::model::message::UnifiedTaskInput;
use crate::common::processors::processor::{ProcessorContext, RetryPolicy};
use crate::common::registry::NodeRegistry;
use crate::engine::runner::ProcessorRunner;
use crate::engine::scheduler::CronScheduler;
use crate::engine::task::TaskManager;
use crate::sync::LeaderElector;
use crate::utils::logger as app_logger;
use crate::utils::logger::{
LogOutputConfig as AppLogOutputConfig, LogSender as AppLogSender,
LoggerConfig as AppLoggerConfig, PrometheusConfig as AppPrometheusConfig,
};
use metrics_exporter_prometheus::{PrometheusBuilder, PrometheusHandle};
use mocra_proxy::ProxyManager;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::sync::{broadcast, watch};
use uuid::Uuid;
mod processors;
mod runtime;
#[cfg(test)]
mod tests;
pub struct Engine {
pub queue_manager: Arc<QueueManager>,
pub downloader_manager: Arc<DownloaderManager>,
pub task_manager: Arc<TaskManager>,
pub proxy_manager: Option<Arc<ProxyManager>>,
pub middleware_manager: Arc<MiddlewareManager>,
pub event_bus: Option<Arc<EventBus>>,
pub state: Arc<State>,
shutdown_tx: broadcast::Sender<()>,
pause_tx: watch::Sender<bool>,
pub prometheus_handle: Option<PrometheusHandle>,
pub node_registry: Arc<NodeRegistry>,
pub cron_scheduler: Arc<CronScheduler>,
pub inflight: crate::engine::runner::InflightCounters,
pub outcomes: crate::engine::runner::StageCounters,
}
impl Engine {
const NODE_HEARTBEAT_INTERVAL_SECS: u64 = 10;
const NODE_HEARTBEAT_TTL_SECS: u64 = Self::NODE_HEARTBEAT_INTERVAL_SECS * 3;
fn policy_event_label(event_type: &str) -> &'static str {
match event_type {
"task_model" => "task_model",
"download" => "download",
"parser_task_model" => "parser_task_model",
"system_error" => "system_error",
"parser" => "parser",
_ => "unknown",
}
}
fn policy_kind_label(kind: &crate::errors::ErrorKind) -> &'static str {
match kind {
crate::errors::ErrorKind::Request => "request",
crate::errors::ErrorKind::Response => "response",
crate::errors::ErrorKind::Command => "command",
crate::errors::ErrorKind::Service => "service",
crate::errors::ErrorKind::Proxy => "proxy",
crate::errors::ErrorKind::Download => "download",
crate::errors::ErrorKind::Queue => "queue",
crate::errors::ErrorKind::Orm => "orm",
crate::errors::ErrorKind::Task => "task",
crate::errors::ErrorKind::Module => "module",
crate::errors::ErrorKind::RateLimit => "rate_limit",
crate::errors::ErrorKind::ProcessorChain => "processor_chain",
crate::errors::ErrorKind::Parser => "parser",
crate::errors::ErrorKind::DataMiddleware => "data_middleware",
crate::errors::ErrorKind::DataStore => "data_store",
crate::errors::ErrorKind::DynamicLibrary => "dynamic_library",
crate::errors::ErrorKind::CacheService => "cache_service",
}
}
async fn handle_policy_failure<T>(
policy_resolver: &PolicyResolver,
queue_manager: &QueueManager,
topic: &str,
event_type: &str,
item: &T,
err: &crate::errors::Error,
ack_fn: &mut Option<crate::queue::AckFn>,
nack_fn: &mut Option<crate::queue::NackFn>,
) where
T: serde::Serialize + Identifiable + Send + Sync,
{
let decision =
policy_resolver.resolve_with_error("engine", Some(event_type), Some("failed"), err);
let action = if decision.policy.retryable {
"retry"
} else if decision.policy.dlq == DlqPolicy::Never {
"ack"
} else {
"dlq"
};
let event_label = Self::policy_event_label(event_type);
let kind_label = Self::policy_kind_label(err.kind());
counter!(
"mocra_policy_decisions_total",
"domain" => "engine",
"event_type" => event_label,
"phase" => "failed",
"kind" => kind_label,
"action" => action
)
.increment(1);
let reason = format!("{}: {}", decision.reason, err);
match action {
"retry" => {
if let Some(f) = nack_fn.take() {
let _ = f(reason).await;
}
}
"dlq" => {
let _ = queue_manager.send_to_dlq(topic, item, &reason).await;
if let Some(f) = ack_fn.take() {
let _ = f().await;
}
}
_ => {
if let Some(f) = ack_fn.take() {
let _ = f().await;
}
}
}
}
async fn handle_policy_retry<T>(
policy_resolver: &PolicyResolver,
queue_manager: &QueueManager,
topic: &str,
event_type: &str,
item: &T,
retry_policy: &RetryPolicy,
ack_fn: &mut Option<crate::queue::AckFn>,
nack_fn: &mut Option<crate::queue::NackFn>,
) where
T: serde::Serialize + Identifiable + Send + Sync,
{
let reason = retry_policy
.reason
.clone()
.unwrap_or_else(|| "retryable failure".to_string());
let err = crate::errors::Error::new(
crate::errors::ErrorKind::ProcessorChain,
Some(std::io::Error::other(reason.clone())),
);
let decision =
policy_resolver.resolve_with_error("engine", Some(event_type), Some("retry"), &err);
let action = if decision.policy.retryable {
"retry"
} else if decision.policy.dlq == DlqPolicy::Never {
"ack"
} else {
"dlq"
};
let event_label = Self::policy_event_label(event_type);
let kind_label = Self::policy_kind_label(err.kind());
counter!(
"mocra_policy_decisions_total",
"domain" => "engine",
"event_type" => event_label,
"phase" => "retry",
"kind" => kind_label,
"action" => action
)
.increment(1);
let reason = format!("{}: {}", decision.reason, reason);
match action {
"retry" => {
if let Some(f) = nack_fn.take() {
let _ = f(reason).await;
}
}
"dlq" => {
let _ = queue_manager.send_to_dlq(topic, item, &reason).await;
if let Some(f) = ack_fn.take() {
let _ = f().await;
}
}
_ => {
if let Some(f) = ack_fn.take() {
let _ = f().await;
}
}
}
}
fn init_queue_manager(cfg: &crate::common::model::config::Config) -> Arc<QueueManager> {
let log_topic = cfg.logger.as_ref().and_then(Self::first_mq_topic);
QueueManager::from_config_with_log_topic(cfg, log_topic.as_deref())
}
fn first_mq_topic(
logger: &crate::common::model::logger_config::LoggerConfig,
) -> Option<String> {
logger.outputs.iter().find_map(|output| match output {
crate::common::model::logger_config::LogOutputConfig::Mq { topic, .. } => {
Some(topic.clone())
}
_ => None,
})
}
fn build_app_logger_config(
logger: &crate::common::model::logger_config::LoggerConfig,
namespace: &str,
) -> AppLoggerConfig {
let mut config = AppLoggerConfig::for_app(namespace);
if let Some(enabled) = logger.enabled {
config.enabled = enabled;
}
if let Some(level) = &logger.level {
config.level = level.clone();
}
if let Some(format) = &logger.format {
if format.to_lowercase() != "text" {
eprintln!("logger.format only supports text for console/file, got {format}");
}
config.format = "text".to_string();
}
if let Some(include) = &logger.include {
config.include = include.clone();
}
if let Some(buffer) = logger.buffer {
config.buffer = buffer;
}
if let Some(interval) = logger.flush_interval_ms {
config.flush_interval_ms = interval;
}
config.outputs = logger
.outputs
.iter()
.map(|output| match output {
crate::common::model::logger_config::LogOutputConfig::Console => {
AppLogOutputConfig::Console
}
crate::common::model::logger_config::LogOutputConfig::File {
path,
rotation,
..
} => AppLogOutputConfig::File {
path: PathBuf::from(path),
rotation: rotation.clone(),
},
crate::common::model::logger_config::LogOutputConfig::Mq { format, .. } => {
if let Some(format) = format
&& format.to_lowercase() != "json"
{
eprintln!("logger.outputs.mq.format only supports json, got {format}");
}
AppLogOutputConfig::Mq
}
})
.collect();
if config.outputs.is_empty() {
config.outputs = AppLoggerConfig::default().outputs;
}
if let Some(prometheus) = &logger.prometheus
&& prometheus.enabled
{
config.prometheus = Some(AppPrometheusConfig { enabled: true });
}
config
}
fn base_level_from_filter(level: &str) -> Option<&str> {
level
.split([',', ';'])
.map(|value| value.trim())
.find(|value| !value.is_empty())
}
async fn setup_mq_log_sender(
logger: &crate::common::model::logger_config::LoggerConfig,
queue_manager: Arc<QueueManager>,
) -> Option<AppLogSender> {
let mq_output = logger.outputs.iter().find_map(|output| match output {
crate::common::model::logger_config::LogOutputConfig::Mq { buffer, .. } => {
Some(*buffer)
}
_ => None,
})?;
let buffer = mq_output.or(logger.buffer).unwrap_or(10000);
let level = logger
.level
.as_deref()
.and_then(Self::base_level_from_filter)
.unwrap_or("info")
.to_string();
let (sender, mut receiver) = tokio::sync::mpsc::channel(buffer);
let log_sender = AppLogSender::with_capacity(sender, level, buffer);
let queue_sender = queue_manager.get_log_push_channel();
tokio::spawn(async move {
while let Some(log) = receiver.recv().await {
let item = QueuedItem::new(log);
if let Err(e) = queue_sender.send(item).await {
eprintln!("Failed to forward log to queue: {e}");
}
}
});
Some(log_sender)
}
pub async fn new(
state: Arc<State>,
queue_manager: Option<Arc<QueueManager>>,
) -> crate::errors::Result<Self> {
let builder = PrometheusBuilder::new();
let prometheus_handle = builder.install_recorder().ok();
let event_bus = state
.config
.read()
.await
.event_bus
.as_ref()
.map(|conf| Arc::new(EventBus::new(conf.capacity, conf.concurrency)));
let (shutdown_tx, _shutdown_rx) = broadcast::channel(1);
let (pause_tx, _) = watch::channel(false);
let state_clone = Arc::clone(&state);
let pause_tx_clone = pause_tx.clone();
let mut shutdown_rx_poller = shutdown_tx.subscribe();
let pause_key = {
let ns = state_clone.cache_service.namespace();
if ns.is_empty() {
warn!(
"Cache namespace is empty; set config.name to avoid cross-app pause collisions"
);
"engine:pause".to_string()
} else {
format!("{ns}:engine:pause")
}
};
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(5));
loop {
tokio::select! {
_ = interval.tick() => {
let is_paused = matches!(state_clone.cache_service.get(&pause_key).await, Ok(Some(_)));
if *pause_tx_clone.borrow() != is_paused {
let _ = pause_tx_clone.send(is_paused);
if is_paused {
info!("Engine paused by global signal");
} else {
info!("Engine resumed by global signal");
}
}
}
_ = shutdown_rx_poller.recv() => {
info!("Engine pause poller shutting down");
break;
}
}
}
});
let task_manager = Arc::new(TaskManager::new(
&state.db,
Arc::clone(&state.cache_service),
state.cookie_service.clone(),
Arc::clone(&state.config),
));
let cfg = state.config.read().await.clone();
let _channel_config = cfg.channel_config.clone();
let namespace = cfg.name.clone();
crate::common::metrics::init_metrics(&namespace);
crate::common::metrics::set_node_up(true);
crate::common::metrics::set_component_health("engine", true);
let queue_manager = if let Some(qm) = queue_manager {
qm
} else {
Self::init_queue_manager(&cfg)
};
if let Some(logger_config) = &cfg.logger {
if logger_config.enabled.unwrap_or(true) {
let app_config = Self::build_app_logger_config(logger_config, &namespace);
let log_sender =
Self::setup_mq_log_sender(logger_config, queue_manager.clone()).await;
let _ = app_logger::init_logger(app_config).await;
if let Some(sender) = log_sender {
let _ = app_logger::set_log_sender(sender);
}
}
}
if let (Some(log_config), Some(event_bus)) = (&cfg.logger, event_bus.as_ref()) {
if log_config.enabled == Some(false) {
info!("Logger disabled; skipping EventBus log handlers");
} else {
use crate::common::model::logger_config::LogOutputConfig;
use crate::engine::events::handlers::{
console_handler::ConsoleLogHandler, queue_handler::QueueLogHandler,
};
for output in &log_config.outputs {
match output {
LogOutputConfig::Mq { .. } => {
let rx = event_bus.subscribe("*".to_string()).await;
QueueLogHandler::start(rx, queue_manager.clone(), "mq".to_string())
.await;
info!("Registered MQ Logger for EventBus");
}
LogOutputConfig::Console => {
let rx = event_bus.subscribe("*".to_string()).await;
let level = log_config
.level
.as_deref()
.and_then(Self::base_level_from_filter)
.unwrap_or("info")
.to_string();
ConsoleLogHandler::start(rx, level).await;
info!("Registered Console Logger for EventBus");
}
LogOutputConfig::File { .. } => {
info!(
"Registered File Logger for EventBus (Handled by Global Tracing)"
);
}
}
}
}
} else if cfg.logger.is_some() {
info!("EventBus disabled; skipping logger EventBus handlers");
}
let downloader_manager = DownloaderManager::new(
state.config.clone(),
state.limiter.clone(),
state.locker.clone(),
state.cache_service.clone(),
)
.await;
let proxy_manager = if let Some(proxy_config) = state.config.read().await.proxy.clone() {
Some(Arc::new(
ProxyManager::from_proxy_config(&proxy_config)
.await
.map_err(|e| {
crate::errors::Error::new(
crate::errors::ErrorKind::Service,
Some(format!("Failed to create ProxyManager: {}", e)),
)
})?,
))
} else {
None
};
let middleware_manager = MiddlewareManager::new();
let node_id = state
.config
.read()
.await
.crawler
.node_id
.clone()
.unwrap_or_else(|| Uuid::new_v4().to_string());
let node_registry = Arc::new(NodeRegistry::new(
state.cache_service.clone(),
node_id,
Duration::from_secs(Self::NODE_HEARTBEAT_TTL_SECS),
));
let leader_elector = if let Some(backend) = state.coordination.clone() {
let (elector, _) =
LeaderElector::new(Some(backend), format!("{}:leader:cron", namespace), 5000);
elector
} else {
let (elector, _) = LeaderElector::new(None, "".to_string(), 5000);
elector
};
let cron_scheduler = Arc::new(
CronScheduler::new(
task_manager.clone(),
state.clone(),
queue_manager.clone(),
shutdown_tx.subscribe(),
leader_elector,
)
.await,
);
Ok(Self {
queue_manager,
downloader_manager: Arc::new(downloader_manager),
task_manager,
proxy_manager,
middleware_manager: Arc::new(middleware_manager),
event_bus,
state,
shutdown_tx,
pause_tx,
prometheus_handle,
node_registry,
cron_scheduler,
inflight: crate::engine::runner::InflightCounters::default(),
outcomes: crate::engine::runner::StageCounters::default(),
})
}
pub async fn register_download_middleware(&self, middleware: DownloadMiddlewareHandle) {
self.middleware_manager
.register_download_middleware(middleware)
.await;
}
pub async fn register_data_middleware(&self, middleware: DataMiddlewareHandle) {
self.middleware_manager
.register_data_middleware(middleware)
.await;
}
pub async fn register_store_middleware(&self, middleware: DataStoreMiddlewareHandle) {
self.middleware_manager
.register_store_middleware(middleware)
.await;
}
pub async fn register_module(&self, module: Arc<dyn ModuleTrait>) {
self.task_manager.add_module(module).await;
}
}