use crate::common::context::PipelineContext;
use crate::common::interface::middleware_manager::MiddlewareManager;
use crate::common::model::data::DataEvent;
use crate::common::model::message::{TaskErrorEvent, TaskEvent};
use crate::common::model::{ModuleConfig, Response};
use crate::common::status_tracker::ErrorDecision;
use crate::engine::events::{
DataMiddlewareEvent, DataStoreEvent, EventBus, EventEnvelope, EventPhase, EventType,
ModuleGenerateEvent, ParserEvent,
};
use crate::engine::processors::event_processor::{EventAwareTypedChain, EventProcessorTrait};
use crate::engine::task::TaskManager;
use crate::errors::{DataMiddlewareError, Error, Result};
use async_trait::async_trait;
use crate::cacheable::{CacheAble, CacheService};
use crate::common::model::login_info::LoginInfo;
use crate::common::processors::processor::{
ProcessorContext, ProcessorResult, ProcessorTrait, RetryPolicy,
};
use crate::common::processors::processor_chain::ErrorStrategy;
use crate::engine::chain::backpressure::{BackpressureSendState, send_with_backpressure};
use crate::engine::task::module::Module;
use crate::queue::{QueueManager, QueuedItem};
use dashmap::DashMap;
use futures::StreamExt;
use log::{debug, error, info, warn};
use metrics::counter;
use serde_json::json;
use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
pub struct ResponseModuleProcessor {
#[allow(dead_code)]
task_manager: Arc<TaskManager>,
cache_service: Arc<CacheService>,
state: Arc<PipelineContext>,
config_cache: Arc<DashMap<String, (Arc<ModuleConfig>, Instant)>>,
}
#[async_trait]
impl ProcessorTrait<Response, (Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>)>
for ResponseModuleProcessor
{
fn name(&self) -> &'static str {
"ResponseModuleProcessor"
}
async fn process(
&self,
input: Response,
context: ProcessorContext,
) -> ProcessorResult<(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>)> {
match self
.state
.status_tracker
.should_task_continue(&input.task_runtime_id())
.await
{
Ok(ErrorDecision::Continue) => {
}
Ok(ErrorDecision::Terminate(reason)) => {
error!(
"[ResponseModuleProcessor] task terminated before parsing: task_id={} reason={}",
input.task_runtime_id(),
reason
);
self.state
.status_tracker
.release_module_locker(&input.module_runtime_id())
.await;
return ProcessorResult::FatalFailure(
crate::errors::ModuleError::TaskMaxError(reason.into()).into(),
);
}
Err(e) => {
warn!(
"[ResponseModuleProcessor] task error check failed, continue anyway: task_id={} error={}",
input.task_runtime_id(),
e
);
}
_ => {}
}
match self
.state
.status_tracker
.should_module_continue(&input.module_runtime_id())
.await
{
Ok(ErrorDecision::Continue) => {
}
Ok(ErrorDecision::Terminate(reason)) => {
error!(
"[ResponseModuleProcessor] module terminated before parsing: module_id={} reason={}",
input.module_runtime_id(),
reason
);
self.state
.status_tracker
.release_module_locker(&input.module_runtime_id())
.await;
return ProcessorResult::FatalFailure(
crate::errors::ModuleError::ModuleMaxError(reason.into()).into(),
);
}
Err(e) => {
warn!(
"[ResponseModuleProcessor] module error check failed, continue anyway: module_id={} error={}",
input.module_runtime_id(),
e
);
}
_ => {}
}
let task: Result<(Arc<Module>, Option<LoginInfo>)> =
self.task_manager.load_module_with_response(&input).await;
match task {
Ok((module, login_info)) => {
let module_id = module.id();
let now = Instant::now();
let mut should_prune_cache = false;
let cached_config = if let Some(entry) = self.config_cache.get(&module_id) {
let (cfg, expires_at) = entry.value();
if now < *expires_at {
Some(cfg.clone())
} else {
should_prune_cache = true;
None
}
} else {
None
};
if should_prune_cache {
self.config_cache.remove(&module_id);
}
let config = if let Some(c) = cached_config {
c
} else {
match ModuleConfig::sync(&module_id, &self.cache_service).await {
Ok(Some(config)) => {
let config = Arc::new(config);
self.config_cache.insert(
module_id,
(config.clone(), Instant::now() + Duration::from_secs(10)),
);
config
}
_ => module.config.clone(),
}
};
ProcessorResult::Success((input, module, config, login_info))
}
Err(e) => {
warn!(
"[ResponseModuleProcessor] load_with_response failed, will retry: account={} platform={} request_id={} err={e}",
input.account, input.platform, input.id
);
ProcessorResult::RetryableFailure(
context
.retry_policy
.unwrap_or(RetryPolicy::default().with_reason(e.to_string())),
)
}
}
}
async fn pre_process(&self, input: &Response, _context: &ProcessorContext) -> Result<()> {
if self.state.config.read().await.download_config.enable_locker {
debug!(
"[ResponseModuleProcessor] lock module before parsing: module_id={}",
input.module_runtime_id()
);
self.state
.status_tracker
.lock_module(&input.module_runtime_id())
.await;
}
Ok(())
}
}
#[async_trait]
impl EventProcessorTrait<Response, (Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>)>
for ResponseModuleProcessor
{
fn pre_status(&self, input: &Response) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::ModuleGenerate,
EventPhase::Started,
ModuleGenerateEvent::from(input),
))
}
fn finish_status(
&self,
input: &Response,
_output: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::ModuleGenerate,
EventPhase::Completed,
ModuleGenerateEvent::from(input),
))
}
fn working_status(&self, input: &Response) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::ModuleGenerate,
EventPhase::Started,
ModuleGenerateEvent::from(input),
))
}
fn error_status(&self, input: &Response, err: &Error) -> Option<EventEnvelope> {
Some(EventEnvelope::engine_error(
EventType::ModuleGenerate,
EventPhase::Failed,
ModuleGenerateEvent::from(input),
err,
))
}
fn retry_status(&self, input: &Response, retry_policy: &RetryPolicy) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::ModuleGenerate,
EventPhase::Retry,
json!({
"data": ModuleGenerateEvent::from(input),
"retry_count": retry_policy.current_retry,
"reason": retry_policy.reason.clone().unwrap_or_default(),
}),
))
}
}
pub struct ResponseParserProcessor {
#[allow(dead_code)]
task_manager: Arc<TaskManager>,
queue_manager: Arc<QueueManager>,
state: Arc<PipelineContext>,
cache_service: Arc<CacheService>,
event_bus: Option<Arc<EventBus>>,
}
impl ResponseParserProcessor {
async fn emit_semantic_event(&self, event_type: EventType, payload: serde_json::Value) {
if let Some(event_bus) = &self.event_bus {
let _ = event_bus
.publish(EventEnvelope::engine(
event_type,
EventPhase::Completed,
payload,
))
.await;
}
}
async fn send_with_backpressure<T>(
&self,
tx: &tokio::sync::mpsc::Sender<QueuedItem<T>>,
item: QueuedItem<T>,
queue_kind: &'static str,
log_context: &str,
) -> bool {
match send_with_backpressure(tx, item).await {
Ok(BackpressureSendState::Direct) => true,
Ok(BackpressureSendState::RecoveredFromFull) => {
counter!("mocra_parser_chain_backpressure_total", "queue" => queue_kind, "reason" => "queue_full").increment(1);
warn!(
"[ResponseParserProcessor] queue full, fallback to awaited send: queue={} context={} remaining_capacity={}",
queue_kind,
log_context,
tx.capacity()
);
true
}
Err(err) => {
if err.after_full {
counter!("mocra_parser_chain_backpressure_total", "queue" => queue_kind, "reason" => "queue_full").increment(1);
warn!(
"[ResponseParserProcessor] queue full before close: queue={} context={} remaining_capacity={}",
queue_kind,
log_context,
tx.capacity()
);
}
counter!("mocra_parser_chain_backpressure_total", "queue" => queue_kind, "reason" => "queue_closed").increment(1);
error!(
"[ResponseParserProcessor] queue closed before send: queue={} context={}",
queue_kind, log_context
);
false
}
}
}
}
#[async_trait]
impl ProcessorTrait<(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>), Vec<DataEvent>>
for ResponseParserProcessor
{
fn name(&self) -> &'static str {
"ResponseParserProcessor"
}
async fn process(
&self,
input: (Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
context: ProcessorContext,
) -> ProcessorResult<Vec<DataEvent>> {
info!(
"[ResponseParserProcessor] start parse: account={} platform={} module={} request_id={} module_id={}",
input.0.account,
input.0.platform,
input.0.module,
input.0.id,
input.0.module_id()
);
let module = input.1.clone();
let config = input.2.clone();
let login_info = input.3.clone();
match serde_json::to_value(config.as_ref()) {
Ok(cfg_val) => {
context
.metadata
.write()
.await
.insert("config".to_string(), cfg_val);
}
Err(e) => {
warn!(
"[ResponseParserProcessor] failed to serialize module config into context: request_id={} module_id={} error={}",
input.0.id,
input.0.module_id(),
e
);
}
}
let task_id = input.0.task_runtime_id();
let module_id = input.0.module_runtime_id();
let request_id = input.0.id.to_string();
let account = input.0.account.clone();
let platform = input.0.platform.clone();
let module_name = input.0.module.clone();
let data = module.parser(input.0.clone(), Some(config)).await;
let mut data = match data {
Ok(d) => {
let has_next_task = !d.parser_task.is_empty();
let has_error = d.error_task.is_some();
info!(
"[ResponseParserProcessor] parser returned: request_id={} data_len={} has_next_task={} has_error={}",
request_id,
d.data.len(),
has_next_task,
has_error
);
self.state
.status_tracker
.record_parse_success(&request_id)
.await
.ok();
d
}
Err(e) => {
warn!(
"[ResponseParserProcessor] parser error: account={} platform={} module={} request_id={} error={e}",
account, platform, module_name, request_id
);
match self
.state
.status_tracker
.record_parse_error(&task_id, &module_id, &request_id, &e)
.await
{
Ok(ErrorDecision::Continue) | Ok(ErrorDecision::RetryAfter(_)) => {
debug!(
"[ResponseParserProcessor] will retry parsing: request_id={}",
request_id
);
return ProcessorResult::RetryableFailure(
context
.retry_policy
.unwrap_or(RetryPolicy::default().with_reason(e.to_string())),
);
}
Ok(ErrorDecision::Skip) => {
warn!(
"[ResponseParserProcessor] skip parse after max retries: request_id={}",
request_id
);
return ProcessorResult::Success(vec![]);
}
Ok(ErrorDecision::Terminate(reason)) => {
error!("[ResponseParserProcessor] terminate: {}", reason);
return ProcessorResult::FatalFailure(e);
}
Err(err) => {
error!("[ResponseParserProcessor] error tracker failed: {}", err);
return ProcessorResult::FatalFailure(e);
}
}
}
};
let parser_task_queue = self.queue_manager.get_parser_task_push_channel();
for task in data.parser_task.drain(..) {
let task_account = task.account_task.account.clone();
let task_platform = task.account_task.platform.clone();
let task_run_id = task.run_id;
let task_module_id = task.context.module_id.clone();
let task_step_idx = task.context.step_idx;
let task_prefix_request = task.prefix_request;
let target_module_name = task.account_task.module.as_ref().and_then(|v| v.first());
let is_same_module =
target_module_name.is_none_or(|name| name == module.module.name().as_str());
let is_same_context = task.account_task.account == module.account.name
&& task.account_task.platform == module.platform.name;
if is_same_module && is_same_context {
debug!(
"[ResponseParserProcessor] Optimizing: Generating requests locally for same module"
);
let mut module_clone = (*module).clone();
module_clone.pending_ctx = Some(task.context.clone());
module_clone.prefix_request = task.prefix_request;
module_clone.run_id = task.run_id;
let task_meta = task.metadata.clone();
match module_clone.generate(task_meta, login_info.clone()).await {
Ok(mut stream) => {
let request_queue = self.queue_manager.get_request_push_channel();
let mut generated_count = 0;
while let Some(req) = stream.next().await {
let req_id = req.id.to_string();
let req_module_id = req.module_id();
let context = format!(
"generated_request request_id={} module_id={}",
req_id, req_module_id
);
if self
.send_with_backpressure(
&request_queue,
QueuedItem::new(req),
"request",
&context,
)
.await
{
generated_count += 1;
}
}
info!(
"[ResponseParserProcessor] Locally generated {} requests",
generated_count
);
self.emit_semantic_event(
EventType::ParserTaskProduced,
json!({
"account": task_account,
"platform": task_platform,
"run_id": task_run_id,
"module_id": task_module_id,
"step_idx": task_step_idx,
"prefix_request": task_prefix_request,
"path": "local_generate"
}),
)
.await;
self.emit_semantic_event(
EventType::ModuleStepAdvanced,
json!({
"account": task_account,
"platform": task_platform,
"run_id": task_run_id,
"module_id": task_module_id,
"step_idx": task_step_idx,
"prefix_request": task_prefix_request,
"generated_count": generated_count,
"mode": "local_generate"
}),
)
.await;
}
Err(e) => {
error!(
"[ResponseParserProcessor] Failed to generate requests locally: {}, falling back to queue",
e
);
let context = format!(
"parser_task_fallback account={} platform={} run_id={} module_id={}",
task.account_task.account,
task.account_task.platform,
task.run_id,
task.context.module_id.as_deref().unwrap_or("unknown")
);
if !self
.send_with_backpressure(
&parser_task_queue,
QueuedItem::new(task),
"parser_task",
&context,
)
.await
{
error!(
"[ResponseParserProcessor] failed to send parser task fallback: {}",
context
);
} else {
self.emit_semantic_event(
EventType::ParserTaskProduced,
json!({
"account": task_account,
"platform": task_platform,
"run_id": task_run_id,
"module_id": task_module_id,
"step_idx": task_step_idx,
"prefix_request": task_prefix_request,
"path": "fallback_queue"
}),
)
.await;
self.emit_semantic_event(
EventType::ModuleStepFallback,
json!({
"account": task_account,
"platform": task_platform,
"run_id": task_run_id,
"module_id": task_module_id,
"step_idx": task_step_idx,
"prefix_request": task_prefix_request,
"reason": "local_generate_failed",
"error": e.to_string(),
"path": "parser_task_queue"
}),
)
.await;
}
}
}
} else {
let context = format!(
"parser_task account={} platform={} run_id={} module_id={}",
task.account_task.account,
task.account_task.platform,
task.run_id,
task.context.module_id.as_deref().unwrap_or("unknown")
);
if !self
.send_with_backpressure(
&parser_task_queue,
QueuedItem::new(task),
"parser_task",
&context,
)
.await
{
error!(
"[ResponseParserProcessor] failed to send parser task: {}",
context
);
} else {
self.emit_semantic_event(
EventType::ParserTaskProduced,
json!({
"account": task_account,
"platform": task_platform,
"run_id": task_run_id,
"module_id": task_module_id,
"step_idx": task_step_idx,
"prefix_request": task_prefix_request,
"path": "parser_task_queue"
}),
)
.await;
}
}
}
if let Some(mut msg) = data.error_task {
warn!(
"[ResponseParserProcessor] recorded response error for request_id={}, message={}",
input.0.id, msg.error_msg
);
msg.prefix_request = input.0.prefix_request;
msg.run_id = input.0.run_id;
let queue = self.queue_manager.get_error_push_channel();
let context = format!(
"error_task request_id={} module_id={}",
input.0.id,
input.0.module_id()
);
if !self
.send_with_backpressure(&queue, QueuedItem::new(msg), "error", &context)
.await
{
error!(
"[ResponseParserProcessor] failed to send parser error: {}",
context
);
} else {
self.emit_semantic_event(
EventType::ErrorTaskProduced,
json!({
"request_id": input.0.id,
"account": input.0.account,
"platform": input.0.platform,
"module_id": input.0.module_id(),
"step_idx": input.0.context.step_idx,
"prefix_request": input.0.prefix_request
}),
)
.await;
}
}
data.data.iter_mut().for_each(|x| {
x.account = input.0.account.clone();
x.platform = input.0.platform.clone();
if x.data_middleware.is_empty() {
x.data_middleware = input.0.data_middleware.clone();
}
x.request_id = input.0.id;
x.meta = input.0.metadata.clone();
});
ProcessorResult::Success(data.data)
}
async fn post_process(
&self,
input: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
_output: &Vec<DataEvent>,
_context: &ProcessorContext,
) -> Result<()> {
let config = self.state.config.read().await;
if config.download_config.enable_session {
if let Some(request_hash) = &input.0.request_hash {
input.0.send(request_hash, &self.cache_service).await.ok();
}
}
if config.download_config.enable_locker {
self.state
.status_tracker
.release_module_locker(&input.0.module_runtime_id())
.await;
debug!(
"[ResponseParserProcessor] released module lock after parsing: module_id={}",
input.0.module_runtime_id()
);
}
Ok(())
}
async fn handle_error(
&self,
input: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
error: Error,
_context: &ProcessorContext,
) -> ProcessorResult<Vec<DataEvent>> {
error!(
"[ResponseParserProcessor] fatal parser error: request_id={} module_id={} error={}",
input.0.id,
input.0.module_id(),
error
);
let timestamp = match SystemTime::now().duration_since(UNIX_EPOCH) {
Ok(duration) => duration.as_secs(),
Err(e) => {
warn!(
"[ResponseParserProcessor] system time before UNIX_EPOCH, fallback to 0: {}",
e
);
0
}
};
let error_metadata = input
.0
.metadata
.task
.as_object()
.cloned()
.unwrap_or_default();
let error_task = TaskErrorEvent {
id: input.0.id,
account_task: TaskEvent {
account: input.0.account.clone(),
platform: input.0.platform.clone(),
module: Some(vec![input.0.module.clone()]),
run_id: input.0.run_id,
priority: input.0.priority,
},
error_msg: error.to_string(),
timestamp,
metadata: error_metadata,
context: input.0.context.clone(),
run_id: input.0.run_id,
prefix_request: input.0.prefix_request,
};
let queue = self.queue_manager.get_error_push_channel();
let context = format!(
"handle_error request_id={} module_id={}",
input.0.id,
input.0.module_id()
);
if !self
.send_with_backpressure(&queue, QueuedItem::new(error_task), "error", &context)
.await
{
error!(
"[ResponseParserProcessor] failed to enqueue ErrorTaskModel: {}",
context
);
}
if self.state.config.read().await.download_config.enable_locker {
self.state
.status_tracker
.release_module_locker(&input.0.module_runtime_id())
.await;
}
ProcessorResult::FatalFailure(error)
}
}
impl
EventProcessorTrait<
(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
Vec<DataEvent>,
> for ResponseParserProcessor
{
fn pre_status(
&self,
input: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::Parser,
EventPhase::Started,
ParserEvent::from(&input.0),
))
}
fn finish_status(
&self,
input: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
_output: &Vec<DataEvent>,
) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::Parser,
EventPhase::Completed,
ParserEvent::from(&input.0),
))
}
fn working_status(
&self,
input: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::Parser,
EventPhase::Started,
ParserEvent::from(&input.0),
))
}
fn error_status(
&self,
input: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
err: &Error,
) -> Option<EventEnvelope> {
Some(EventEnvelope::engine_error(
EventType::Parser,
EventPhase::Failed,
ParserEvent::from(&input.0),
err,
))
}
fn retry_status(
&self,
input: &(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>),
retry_policy: &RetryPolicy,
) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::Parser,
EventPhase::Retry,
json!({
"data": ParserEvent::from(&input.0),
"retry_count": retry_policy.current_retry,
"reason": retry_policy.reason.clone().unwrap_or_default(),
}),
))
}
}
pub struct DataMiddlewareProcessor {
middleware_manager: Arc<MiddlewareManager>,
}
#[async_trait]
impl ProcessorTrait<DataEvent, DataEvent> for DataMiddlewareProcessor {
fn name(&self) -> &'static str {
"DataMiddlewareProcessor"
}
async fn process(
&self,
input: DataEvent,
context: ProcessorContext,
) -> ProcessorResult<DataEvent> {
let start = std::time::Instant::now();
let config = context
.metadata
.read()
.await
.get("config")
.map(|x| {
serde_json::from_value::<ModuleConfig>(x.clone()).unwrap_or_else(|e| {
warn!(
"[DataMiddlewareProcessor] config conversion failed, using default: request_id={} module={} error={}",
input.request_id,
input.module,
e
);
ModuleConfig::default()
})
});
let modified_data = self.middleware_manager.handle_data(input, &config).await;
let elapsed_ms = start.elapsed().as_millis();
if elapsed_ms > 10 {
info!(
"[DataMiddlewareProcessor] SLOW middleware execution: {} ms",
elapsed_ms
);
}
match modified_data {
Some(data) => ProcessorResult::Success(data),
None => ProcessorResult::FatalFailure(
DataMiddlewareError::EmptyData("data dropped by data middleware".into()).into(),
),
}
}
}
impl EventProcessorTrait<DataEvent, DataEvent> for DataMiddlewareProcessor {
fn pre_status(&self, input: &DataEvent) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::MiddlewareBefore,
EventPhase::Started,
DataMiddlewareEvent::from(input),
))
}
fn finish_status(&self, input: &DataEvent, output: &DataEvent) -> Option<EventEnvelope> {
let mut event: DataMiddlewareEvent = input.into();
event.after_size = output.size().into();
Some(EventEnvelope::engine(
EventType::MiddlewareBefore,
EventPhase::Completed,
event,
))
}
fn working_status(&self, input: &DataEvent) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::MiddlewareBefore,
EventPhase::Started,
DataMiddlewareEvent::from(input),
))
}
fn error_status(&self, input: &DataEvent, err: &Error) -> Option<EventEnvelope> {
Some(EventEnvelope::engine_error(
EventType::MiddlewareBefore,
EventPhase::Failed,
DataMiddlewareEvent::from(input),
err,
))
}
fn retry_status(&self, input: &DataEvent, retry_policy: &RetryPolicy) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::MiddlewareBefore,
EventPhase::Retry,
json!({
"data": DataMiddlewareEvent::from(input),
"retry_count": retry_policy.current_retry,
"reason": retry_policy.reason.clone().unwrap_or_default(),
}),
))
}
}
pub struct DataStoreProcessor {
middleware_manager: Arc<MiddlewareManager>,
outcomes: Arc<crate::engine::runner::StageCounter>,
}
#[async_trait]
impl ProcessorTrait<DataEvent, ()> for DataStoreProcessor {
fn name(&self) -> &'static str {
"DataStoreProcessor"
}
async fn process(&self, input: DataEvent, context: ProcessorContext) -> ProcessorResult<()> {
info!(
"[DataStoreProcessor] start store: request_id={} account={} platform={} module={} size={}",
input.request_id,
input.account,
input.platform,
input.module,
input.size()
);
let mut middleware = vec![];
if let Some(retry_policy) = &context.retry_policy
&& let Some(m_val) = retry_policy.meta.get("middleware")
&& let Some(m) = m_val.as_array()
{
middleware = m
.iter()
.filter_map(|x| x.as_str())
.map(|x| x.to_string())
.collect::<Vec<String>>();
}
let config = context
.metadata
.read()
.await
.get("config")
.map(|x| {
serde_json::from_value::<ModuleConfig>(x.clone()).unwrap_or_else(|e| {
warn!(
"[DataStoreProcessor] config conversion failed, using default: request_id={} module={} error={}",
input.request_id,
input.module,
e
);
ModuleConfig::default()
})
});
let request_id = input.request_id;
let res = if middleware.is_empty() {
self.middleware_manager
.handle_store_data(input, &config)
.await
} else {
self.middleware_manager
.handle_store_data_with_middleware(input, middleware, &config)
.await
};
if res.is_empty() {
self.outcomes.record_success();
info!(
"[DataStoreProcessor] store success, request_id={}",
request_id
);
ProcessorResult::Success(())
} else {
self.outcomes.record_failure();
let error_msg = res
.iter()
.map(|(m, e)| format!("Middleware: {m}, Error: {e:?}"))
.collect::<Vec<String>>()
.join("; ");
let mut retry_policy = context
.retry_policy
.unwrap_or_default()
.with_reason(error_msg);
let error_middleware = res.keys().map(|x| x.to_string()).collect::<Vec<String>>();
retry_policy.meta = serde_json::json!({ "middleware": error_middleware });
warn!(
"[DataStoreProcessor] request={}, store error, will retry: {}",
request_id,
retry_policy.reason.clone().unwrap_or_default()
);
ProcessorResult::RetryableFailure(retry_policy)
}
}
}
impl EventProcessorTrait<DataEvent, ()> for DataStoreProcessor {
fn pre_status(&self, input: &DataEvent) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::DataStore,
EventPhase::Started,
DataStoreEvent::from(input),
))
}
fn finish_status(&self, input: &DataEvent, _output: &()) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::DataStore,
EventPhase::Completed,
DataStoreEvent::from(input),
))
}
fn working_status(&self, input: &DataEvent) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::DataStore,
EventPhase::Started,
DataStoreEvent::from(input),
))
}
fn error_status(&self, input: &DataEvent, err: &Error) -> Option<EventEnvelope> {
Some(EventEnvelope::engine_error(
EventType::DataStore,
EventPhase::Failed,
DataStoreEvent::from(input),
err,
))
}
fn retry_status(&self, input: &DataEvent, retry_policy: &RetryPolicy) -> Option<EventEnvelope> {
Some(EventEnvelope::engine(
EventType::DataStore,
EventPhase::Retry,
json!({
"data": DataStoreEvent::from(input),
"retry_count": retry_policy.current_retry,
"reason": retry_policy.reason.clone().unwrap_or_default(),
}),
))
}
}
pub async fn create_parser_chain(
state: Arc<PipelineContext>,
#[allow(dead_code)] task_manager: Arc<TaskManager>,
middleware_manager: Arc<MiddlewareManager>,
queue_manager: Arc<QueueManager>,
event_bus: Option<Arc<EventBus>>,
cache_service: Arc<CacheService>,
store_outcomes: Arc<crate::engine::runner::StageCounter>,
) -> EventAwareTypedChain<Response, Vec<()>> {
let response_module_processor = ResponseModuleProcessor {
task_manager: task_manager.clone(),
cache_service: cache_service.clone(),
state: state.clone(),
config_cache: Arc::new(DashMap::new()),
};
let response_parser_processor = ResponseParserProcessor {
task_manager: task_manager.clone(),
queue_manager,
state,
cache_service: cache_service.clone(),
event_bus: event_bus.clone(),
};
let data_middleware_processor = DataMiddlewareProcessor {
middleware_manager: middleware_manager.clone(),
};
let data_store_processor = DataStoreProcessor {
middleware_manager,
outcomes: store_outcomes,
};
EventAwareTypedChain::<Response, Response>::new(event_bus)
.then::<(Response, Arc<Module>, Arc<ModuleConfig>, Option<LoginInfo>), _>(
response_module_processor,
)
.then::<Vec<DataEvent>, _>(response_parser_processor)
.then_map_vec_parallel_with_strategy_silent::<DataEvent, _>(
data_middleware_processor,
64,
ErrorStrategy::Skip,
)
.then_map_vec_parallel_with_strategy_silent::<(), _>(
data_store_processor,
64,
ErrorStrategy::Skip,
)
}