elf_utils 0.1.4

elf_rust utils
Documentation
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);
    // cargo run -- --payload '{"tasks":["abc","a"]}'
    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> {
    // report result
    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);
            //block_in_place fix: Cannot drop a runtime in a context where blocking is not allowed.
            // This happens when a runtime is dropped from within an asynchronous context.
            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)
    }
}