use proc_macro::*;
use quote::{format_ident, quote, ToTokens};
use syn::{
self, parse_macro_input, punctuated::Punctuated, ItemFn, ItemImpl, ItemStruct, Meta,
};
#[proc_macro_attribute]
pub fn cron(att: TokenStream, code: TokenStream) -> TokenStream {
let args = parse_macro_input!(att with Punctuated::<Meta, syn::Token![,]>::parse_terminated);
let args = args.into_iter().map(|x| {
x.require_name_value()
.map(|x| {
let arg_name = x.path.to_token_stream().to_string();
let arg_val = x.value.to_token_stream().to_string();
(arg_name, arg_val.replace("\"", ""))
})
.unwrap()
});
let (arg_1_name, cron_expr) = args.clone().peekable().nth(0).unwrap();
let (arg_2_name, timeout) = args.peekable().nth(1).unwrap();
if arg_1_name == "expr" && arg_2_name == "timeout" {
let parsed = syn::parse::<ItemFn>(code.clone());
if parsed.is_ok() {
let origin_function = parsed.clone().unwrap().to_token_stream();
let ident = parsed.clone().unwrap().sig.ident;
let job_name = ident.to_string();
let new_code = quote! {
#origin_function
inventory::submit! {
JobBuilder::global_job(#job_name, #ident, #cron_expr, #timeout)
}
};
return new_code.into();
} else if let Some(error) = parsed.err() {
println!("parse Error: {}", error);
}
}
code
}
#[proc_macro_attribute]
pub fn cron_obj(_att: TokenStream, code: TokenStream) -> TokenStream {
let item_struct = syn::parse::<ItemStruct>(code.clone()).unwrap();
let r#struct = item_struct.to_token_stream();
let method_jobs = format_ident!("CRONFRAME_METHOD_JOBS_{}", item_struct.ident);
let function_jobs = format_ident!("CRONFRAME_FUNCTION_JOBS_{}", item_struct.ident);
let cf_fn_jobs_flag = format_ident!("CF_FN_JOBS_FLAG_{}", item_struct.ident);
let cf_fn_jobs_channels = format_ident!("CF_FN_JOBS_CHANNELS_{}", item_struct.ident);
let mut tmp = r#struct.to_string();
let struct_name = item_struct.ident;
if tmp.contains("{") {
tmp.insert_str(
tmp.chars().count() - 1,
"tx: Option<crossbeam_channel::Sender<String>>",
);
} else {
tmp.insert_str(
tmp.chars().count() - 1,
"{tx: Option<crossbeam_channel::Sender<String>>}",
);
tmp = (&tmp[0..tmp.len() - 1].to_string()).clone();
}
let struct_edited: proc_macro2::TokenStream = tmp.parse().unwrap();
let type_name = struct_name.clone().into_token_stream().to_string();
let mut cron_obj_fn = String::from("fn new_cron_obj(");
if !item_struct.fields.is_empty() {
let mut tmp = item_struct.fields.iter().map(|x| {
let field_name = x.ident.to_token_stream().to_string();
let field_type = x.ty.to_token_stream().to_string();
format!("{field_name} : {field_type},")
});
for _ in 0..item_struct.fields.len() {
cron_obj_fn.push_str(&tmp.next().unwrap());
}
}
cron_obj_fn.push_str(") -> ");
cron_obj_fn.push_str(type_name.as_str());
cron_obj_fn.push_str("{");
cron_obj_fn.push_str(type_name.as_str());
cron_obj_fn.push_str("{");
if !item_struct.fields.is_empty() {
let mut tmp = item_struct.fields.iter().map(|x| {
let field_name = x.ident.to_token_stream().to_string();
format!("{field_name},")
});
for _ in 0..item_struct.fields.len() {
cron_obj_fn.push_str(&tmp.next().unwrap());
}
}
cron_obj_fn.push_str("tx: None");
cron_obj_fn.push_str("}");
cron_obj_fn.push_str("}");
let cron_job_fn_tokens: proc_macro2::TokenStream = cron_obj_fn.parse().unwrap();
let new_code = quote! {
#struct_edited
static #cf_fn_jobs_flag: Mutex<bool> = Mutex::new(false);
static #cf_fn_jobs_channels: once_cell::sync::Lazy<(crossbeam_channel::Sender<String>, crossbeam_channel::Receiver<String>)> = once_cell::sync::Lazy::new(|| crossbeam_channel::bounded(1));
impl Drop for #struct_name {
fn drop(&mut self) {
if self.tx.is_some(){
let _= self.tx.as_ref().unwrap().send("JOB_DROP".to_string());
}
}
}
impl #struct_name {
#cron_job_fn_tokens
fn cf_drop(&self) {
if *#cf_fn_jobs_flag.lock().unwrap(){
for func in #function_jobs{
let _= #cf_fn_jobs_channels.0.send("JOB_DROP".to_string());
}
*#cf_fn_jobs_flag.lock().unwrap() = false;
}
}
}
#[distributed_slice]
static #method_jobs: [fn(Arc<Box<dyn Any + Send + Sync>>) -> JobBuilder<'static>];
#[distributed_slice]
static #function_jobs: [fn() -> JobBuilder<'static>];
};
new_code.into()
}
#[proc_macro_attribute]
pub fn cron_impl(_att: TokenStream, code: TokenStream) -> TokenStream {
let item_impl = syn::parse::<ItemImpl>(code.clone()).unwrap();
let r#impl = item_impl.to_token_stream();
let impl_items = item_impl.items.clone();
let impl_type = item_impl.self_ty.to_token_stream();
let method_jobs = format_ident!("CRONFRAME_METHOD_JOBS_{impl_type}");
let function_jobs = format_ident!("CRONFRAME_FUNCTION_JOBS_{impl_type}");
let mut new_code = quote! {
#r#impl
};
let mut count = 0;
for item in impl_items {
let item_token = item.to_token_stream();
let item_fn_parsed = syn::parse::<ItemFn>(item_token.into());
let item_fn_id = item_fn_parsed.clone().unwrap().sig.ident;
let helper = format_ident!("cron_helper_{}", item_fn_id);
let linkme_deserialize = format_ident!("LINKME_{}_{count}", item_fn_id);
let new_code_tmp = if check_self(&item_fn_parsed) {
// method job
quote! {
#[distributed_slice(#method_jobs)]
static #linkme_deserialize: fn(_self: Arc<Box<dyn Any + Send + Sync>>)-> JobBuilder<'static> = #impl_type::#helper;
}
} else {
// function job
quote! {
#[distributed_slice(#function_jobs)]
static #linkme_deserialize: fn()-> JobBuilder<'static> = #impl_type::#helper;
}
};
new_code.extend(new_code_tmp.into_iter());
count += 1;
}
let type_name = impl_type.to_string();
let cf_fn_jobs_flag = format_ident!("CF_FN_JOBS_FLAG_{}", type_name);
let cf_fn_jobs_channels = format_ident!("CF_FN_JOBS_CHANNELS_{}", type_name);
let gather_fn = quote! {
impl #impl_type{
pub fn cf_gather_mt(&mut self, frame: Arc<CronFrame>){
info!("Collecting Method Jobs from {}", #type_name);
if !#method_jobs.is_empty(){
let life_channels = crossbeam_channel::bounded(1);
self.tx = Some(life_channels.0.clone());
for method_job in #method_jobs {
let job_builder = (method_job)(Arc::new(Box::new(self.clone())));
let mut cron_job = job_builder.build();
cron_job.life_channels = Some(life_channels.clone());
info!("Found Method Job \"{}\" from {}.", cron_job.name, #type_name);
frame.cron_jobs.lock().unwrap().push(cron_job);
}
info!("Method Jobs from {} Collected.", #type_name);
} else {
info!("Not Method Jobs from {} has been found.", #type_name);
}
}
pub fn cf_gather_fn(frame: Arc<CronFrame>){
info!("Collecting Function Jobs from {}", #type_name);
if !#function_jobs.is_empty(){
let fn_flag = *#cf_fn_jobs_flag.lock().unwrap();
if !fn_flag {
for function_job in #function_jobs {
let job_builder = (function_job)();
let mut cron_job = job_builder.build();
cron_job.life_channels = Some(#cf_fn_jobs_channels.clone());
info!("Found Function Job \"{}\" from {}.", cron_job.name, #type_name);
frame.cron_jobs.lock().unwrap().push(cron_job);
}
info!("Function Jobs from {} Collected.", #type_name);
*#cf_fn_jobs_flag.lock().unwrap() = true;
}
} else {
info!("Not Function Jobs from {} has been found.", #type_name);
}
}
pub fn cf_gather(&mut self, frame: Arc<CronFrame>){
self.cf_gather_mt(frame.clone());
Self::cf_gather_fn(frame.clone());
}
}
};
new_code.extend(gather_fn.into_iter());
new_code.into()
}
#[proc_macro_attribute]
pub fn fn_job(att: TokenStream, code: TokenStream) -> TokenStream {
let parsed = syn::parse::<ItemFn>(code.clone());
if check_self(&parsed) {
}
let args = parse_macro_input!(att with Punctuated::<Meta, syn::Token![,]>::parse_terminated);
let args = args.into_iter().map(|x| {
x.require_name_value()
.map(|x| {
let arg_name = x.path.to_token_stream().to_string();
let arg_val = x.value.to_token_stream().to_string();
(arg_name, arg_val.replace("\"", ""))
})
.unwrap()
});
let (arg_1_name, cron_expr) = args.clone().peekable().nth(0).unwrap();
let (arg_2_name, timeout) = args.peekable().nth(1).unwrap();
if arg_1_name != "expr" && arg_2_name != "timeout" {
return code;
}
let origin_function = parsed.clone().unwrap().to_token_stream();
let ident = parsed.clone().unwrap().sig.ident;
let job_name = ident.to_string();
let helper = format_ident!("cron_helper_{}", ident);
let new_code = quote! {
#origin_function
fn #helper() -> JobBuilder<'static> {
JobBuilder::function_job(#job_name, Self::#ident, #cron_expr, #timeout)
}
};
new_code.into()
}
#[proc_macro_attribute]
pub fn mt_job(att: TokenStream, code: TokenStream) -> TokenStream {
let parsed = syn::parse::<ItemFn>(code.clone());
if !check_self(&parsed) {
}
let args = parse_macro_input!(att with Punctuated::<Meta, syn::Token![,]>::parse_terminated);
let args = args.into_iter().map(|x| {
x.require_name_value()
.map(|x| {
let arg_name = x.path.to_token_stream().to_string();
let arg_val = x.value.to_token_stream().to_string();
(arg_name, arg_val.replace("\"", ""))
})
.unwrap()
});
let (arg_1_name, cron_expr) = args.clone().peekable().nth(0).unwrap();
if arg_1_name != "expr" {
}
let origin_method = parsed.clone().unwrap().to_token_stream();
let ident = parsed.clone().unwrap().sig.ident;
let job_name = ident.to_string();
let block = parsed.clone().unwrap().block;
let cronframe_method = format_ident!("cron_method_{}", ident);
let helper = format_ident!("cron_helper_{}", ident);
let expr = format_ident!("expr");
let tout = format_ident!("tout");
let block_string = block.clone().into_token_stream().to_string();
let mut block_string_edited = block_string.replace("self.", "cronframe_self.");
block_string_edited.insert_str(
1,
"let cron_frame_instance = arg.clone();
let cronframe_self = (*cron_frame_instance).downcast_ref::<Self>().unwrap();",
);
let block_edited: proc_macro2::TokenStream = block_string_edited.parse().unwrap();
let mut new_code = quote! {
#origin_method
fn #cronframe_method(arg: Arc<Box<dyn Any + Send + Sync>>) #block_edited
};
let helper_code = quote! {
fn #helper(arg: Arc<Box<dyn Any + Send + Sync>>) -> JobBuilder<'static> {
let instance = arg.clone();
let this_obj = (*instance).downcast_ref::<Self>().unwrap();
let #expr = this_obj.cron_expr.expr();
let #tout = format!("{}", this_obj.cron_expr.timeout());
let instance = arg.clone();
JobBuilder::method_job(#job_name, Self::#cronframe_method, #expr.clone(), #tout, instance)
}
};
let helper_code_edited = helper_code
.clone()
.into_token_stream()
.to_string()
.replace("cron_expr", &cron_expr);
let block_edited: proc_macro2::TokenStream = helper_code_edited.parse().unwrap();
new_code.extend(block_edited.into_iter());
new_code.into()
}
fn check_self(parsed: &Result<ItemFn, syn::Error>) -> bool {
if !parsed.clone().unwrap().sig.inputs.is_empty()
&& parsed
.clone()
.unwrap()
.sig
.inputs
.first()
.unwrap()
.to_token_stream()
.to_string()
== "self"
{
true
} else {
false
}
}