use std::time::Duration;
use elf_rust::*;
use log::{error, info};
use serde_json::{Map, Value};
use crate::params::{parse_params, parse_params_from_str};
use crate::{rpt, vault};
use crate::fs::report_card;
use crate::rpt::Reporter;
pub async fn default_handle_script<F, Fut>(name: &str, script: F)
where
F: Fn(Context, Value) -> Fut + 'static + Send + Sync,
Fut: std::future::Future<Output=Result<String, String>> + 'static + Send,
{
default_handle_scripts(vec![MyScript::new(name, script)]).await;
}
pub async fn default_handle_scripts(scripts: Vec<impl IntoScriptConfigs + 'static>) {
let mut app = App::new();
app.add_scripts(scripts);
let mut executor = Executor::new();
executor.set_task(app.into());
executor.set_initializer(|ctx| {
let runtime = ctx.get(KEY_RUNTIME).unwrap();
info!("runtime {}", runtime);
if runtime != Runtime::Local.to_string() {
let err = vault::init();
if err.is_err() {
error!("vault init failed: {:?}", err.err().unwrap());
}
info!("vault init success");
}
let uk_name = std::env::var("UK_FUNCTION_NAME").unwrap_or("".to_string());
ctx.set("uk_name", uk_name.to_string());
info!("uk_name: {}", uk_name);
});
executor.set_before_action(before_action);
executor.set_after_action(after_action);
executor.run().await.expect("executor failed to run");
}
pub fn before_action(ctx: &Context, payload: Value) -> Value {
let params_res = parse_params(payload.clone());
if params_res.is_err() {
error!("parse params failed {:?}", params_res.err().unwrap());
return payload.clone();
}
let mut params = params_res.unwrap();
let uk_name = ctx.get("uk_name");
if params.valid() && uk_name.is_some() {
let rpt_host = std::env::var("REPORT_HOST").unwrap_or("".to_string());
tokio::task::block_in_place(|| {
let rpt = rpt::HttpReporter::new(rpt_host.to_string(), uk_name.unwrap());
let result = rpt.start(params.gen_uk().as_str());
match result {
Ok(lid) => {
params.set_lid(lid);
ctx.set("kp_params", params.to_string());
info!("start report success, lid: {}", lid);
}
Err(e) => {
error!("start report failed: {:?}", e);
}
}
});
}
payload
}
pub fn after_action(ctx: &Context, payload: Value, result: Result<String, String>) -> Result<String, String> {
if let Some(v) = ctx.get("kp_params") {
let params_res = parse_params_from_str(v.as_str());
if params_res.is_ok() {
let mut err_msg = "".to_string();
if result.is_err() {
err_msg = result.clone().err().unwrap().to_string();
}
let params = params_res.unwrap();
let uk_name = ctx.get("uk_name");
let rpt_host = std::env::var("REPORT_HOST").unwrap_or("".to_string());
tokio::task::block_in_place(|| {
let rpt = rpt::HttpReporter::new(rpt_host.to_string(), uk_name.unwrap());
let result = rpt.finish(params.gen_uk().as_str(), &err_msg);
match result {
Ok(_) => info!("finish report end success"),
Err(e) => error!("finish report end failed: {:?}", e)
}
});
} else {
error!("parse params data {:?} failed {:?}",v, params_res);
}
};
let uk_name = ctx.get("uk_name").unwrap_or("Elf Script ".to_string());
let user = std::env::var("TASK_MASTER").unwrap_or("".to_string());
match result.clone() {
Ok(v) => {
info!("task executed with payload: {}, result: {:?}", payload.to_string(), v);
let value = serde_json::from_str(result.as_ref().unwrap().as_str()).unwrap();
if let Value::Object(map) = value {
let empty_map = Map::new();
let failed_map = map.get("failed").and_then(|v| v.as_object()).unwrap_or(&empty_map);
let success_map = map.get("success").and_then(|v| v.as_object()).unwrap_or(&empty_map);
let cost_time_map = map.get("cost_time").and_then(|v| v.as_object()).unwrap_or(&empty_map);
let mut msg = "".to_string();
for (task, result) in success_map {
let cost_time_ms = cost_time_map.get(task).and_then(|v| v.as_i64()).unwrap_or(0);
let cost_time = format_duration(Duration::from_millis(cost_time_ms as u64));
msg.push_str(&format!("{} executed successfully\n{}\ncost {}\n", task, result, cost_time));
}
for (task, error) in failed_map {
msg.push_str(&format!("{} executed failed\nerr:{}\n", task, error.as_str().unwrap_or("unknown error")));
}
tokio::task::block_in_place(|| {
report_card(uk_name.as_str(), msg.as_str(), user.as_str())
});
}
}
Err(e) => {
error!("task executed with payload: {}, result: {:?}", payload.to_string(), e);
tokio::task::block_in_place(|| {
report_card(uk_name.as_str(), &format!("task executed failed\nerr:{}", e), user.as_str());
})
}
}
result.clone()
}
fn format_duration(duration: Duration) -> String {
let mins = duration.as_secs() / 60;
let secs = duration.as_secs() % 60;
let millis = duration.subsec_millis();
if mins > 0 {
format!("{} m {} s {} ms", mins, secs, millis)
} else if secs > 0 {
format!("{} s {} ms", secs, millis)
} else {
format!("{} ms", millis)
}
}