* Copyright(c) 2024-2026 China Telecom Cloud Technologies Co., Ltd. All rights
* reserved. ctscat is licensed under Mulan PSL v2. You can use this software
* according to the terms and conditions of the Mulan PSL V2. You may obtain a
* copy of Mulan PSL v2 at: http://license.coscl.org.cn/MulanPSL2.
* THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY
* KIND, EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT,
* MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE. See the Mulan PSL v2 for
* more details.
*/
use anyhow::{Context, Result};
use std::collections::HashMap;
use tracing::info;
use super::parser::TaskConfig;
#[derive(Clone)]
pub struct TaskExecutor;
impl TaskExecutor {
pub fn new() -> Self {
Self
}
pub async fn execute_task(
&self,
task: &TaskConfig,
_variables: &HashMap<String, String>,
) -> Result<serde_yaml::Value> {
info!(" 执行任务类型: {}", task.task_type);
match task.task_type.as_str() {
"sync" => {
let tracking_id = task
.parameters
.get("tracking_id")
.and_then(|v| v.as_i64())
.context("sync 任务缺少 tracking_id 参数")?;
info!(" 执行 sync 任务: tracking_id = {}", tracking_id);
Ok(serde_yaml::Value::String(format!(
"Sync completed for tracking {}",
tracking_id
)))
}
"classify" => {
let limit = task
.parameters
.get("limit")
.and_then(|v| v.as_i64())
.unwrap_or(100);
info!(" 执行 classify 任务: limit = {}", limit);
Ok(serde_yaml::Value::String(format!(
"Classified {} items",
limit
)))
}
"compare" => {
let tracking_id = task
.parameters
.get("tracking_id")
.and_then(|v| v.as_i64())
.context("compare 任务缺少 tracking_id 参数")?;
info!(" 执行 compare 任务: tracking_id = {}", tracking_id);
Ok(serde_yaml::Value::String(format!(
"Comparison completed for tracking {}",
tracking_id
)))
}
"export" => {
let format = task
.parameters
.get("format")
.and_then(|v| v.as_str())
.unwrap_or("json");
info!(" 执行 export 任务: format = {}", format);
Ok(serde_yaml::Value::String(format!(
"Exported in {} format",
format
)))
}
"l0" => {
let package_id = task.parameters.get("package_id").and_then(|v| v.as_i64());
info!(" 执行 L0 任务: package_id = {:?}", package_id);
Ok(serde_yaml::Value::String(
"L0 polling completed".to_string(),
))
}
"snapshot" => {
let tracking_id = task
.parameters
.get("tracking_id")
.and_then(|v| v.as_i64())
.context("snapshot 任务缺少 tracking_id 参数")?;
info!(" 执行 snapshot 任务: tracking_id = {}", tracking_id);
Ok(serde_yaml::Value::String(format!(
"Snapshot operation completed for tracking {}",
tracking_id
)))
}
_ => Err(anyhow::anyhow!("未知的任务类型: {}", task.task_type)),
}
}
}
impl Default for TaskExecutor {
fn default() -> Self {
Self::new()
}
}