use std::collections::HashMap;
use std::sync::{Arc, Mutex, OnceLock};
use ani_rs::objects::{AniFnObject, GlobalRefCallback};
use ani_rs::AniEnv;
use request_client::RequestClient;
use request_core::info::{Faults, Progress, Response, WaitingReason};
use crate::api10::bridge::{self, Task};
#[ani_rs::native]
pub fn on_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("on_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
let callback = callback.into_global_callback(env).unwrap();
let coll = match event.as_str() {
"completed" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_complete.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![callback]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![]),
})
}
}
"pause" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_pause.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![callback]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![]),
})
}
}
"failed" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_fail.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![callback]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![]),
})
}
}
"remove" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_remove.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![callback]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![]),
})
}
}
"progress" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_progress.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![callback]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![]),
})
}
}
"resume" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_resume.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![callback]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![]),
})
}
}
_ => unimplemented!(),
};
RequestClient::get_instance().register_callback(task_id, coll.clone());
callback_mgr.tasks.lock().unwrap().insert(task_id, coll);
Ok(())
}
#[ani_rs::native]
pub fn on_response_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("on_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
let callback = callback.into_global_callback(env).unwrap();
let coll = match event.as_str() {
"response" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_response.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![callback]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![]),
})
}
}
_ => unimplemented!(),
};
RequestClient::get_instance().register_callback(task_id, coll.clone());
callback_mgr.tasks.lock().unwrap().insert(task_id, coll);
Ok(())
}
#[ani_rs::native]
pub fn on_fault_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("on_fault_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
let callback = callback.into_global_callback(env).unwrap();
let coll = match event.as_str() {
"faultOccur" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_fault.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![callback]),
on_wait: Mutex::new(vec![]),
})
}
}
_ => unimplemented!(),
};
RequestClient::get_instance().register_callback(task_id, coll.clone());
callback_mgr.tasks.lock().unwrap().insert(task_id, coll);
Ok(())
}
#[ani_rs::native]
pub fn on_wait_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("on_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
let callback = callback.into_global_callback(env).unwrap();
let coll = match event.as_str() {
"wait" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_wait.lock().unwrap().push(callback);
return Ok(());
} else {
Arc::new(CallbackColl {
on_progress: Mutex::new(vec![]),
on_complete: Mutex::new(vec![]),
on_pause: Mutex::new(vec![]),
on_resume: Mutex::new(vec![]),
on_remove: Mutex::new(vec![]),
on_fail: Mutex::new(vec![]),
on_response: Mutex::new(vec![]),
on_fault: Mutex::new(vec![]),
on_wait: Mutex::new(vec![callback]),
})
}
}
_ => unimplemented!(),
};
RequestClient::get_instance().register_callback(task_id, coll.clone());
callback_mgr.tasks.lock().unwrap().insert(task_id, coll);
Ok(())
}
#[ani_rs::native]
pub fn off_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("off_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
let callback: GlobalRefCallback<(bridge::Progress,)> =
callback.into_global_callback(env).unwrap();
match event.as_str() {
"completed" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_complete.lock().unwrap().retain(|x| *x != callback);
}
}
"pause" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_pause.lock().unwrap().retain(|x| *x != callback);
}
}
"failed" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_fail.lock().unwrap().retain(|x| *x != callback);
}
}
"remove" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_remove.lock().unwrap().retain(|x| *x != callback);
}
}
"progress" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_progress.lock().unwrap().retain(|x| *x != callback);
}
}
"resume" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_resume.lock().unwrap().retain(|x| *x != callback);
}
}
_ => unimplemented!(),
};
Ok(())
}
#[ani_rs::native]
pub fn off_response_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("off_response_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
let callback = callback.into_global_callback(env).unwrap();
match event.as_str() {
"response" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_response.lock().unwrap().retain(|x| *x != callback);
}
}
_ => unimplemented!(),
};
Ok(())
}
#[ani_rs::native]
pub fn off_fault_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("off_fault_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
let callback = callback.into_global_callback(env).unwrap();
match event.as_str() {
"faultOccur" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_fault.lock().unwrap().retain(|x| *x != callback);
}
}
_ => unimplemented!(),
};
Ok(())
}
#[ani_rs::native]
pub fn off_wait_event(
env: &AniEnv,
this: Task,
event: String,
callback: AniFnObject,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("off_wait_event called");
let callback_mgr = CallbackManager::get_instance();
let callback = callback.into_global_callback(env).unwrap();
match event.as_str() {
"wait" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_wait.lock().unwrap().retain(|x| *x != callback);
}
}
_ => unimplemented!(),
};
Ok(())
}
#[ani_rs::native]
pub fn off_events(
env: &AniEnv,
this: Task,
event: String,
) -> Result<(), ani_rs::business_error::BusinessError> {
let task_id = this.tid.parse().unwrap();
info!("off_fault_event called with event: {}", event);
let callback_mgr = CallbackManager::get_instance();
match event.as_str() {
"completed" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_complete.lock().unwrap().clear();
}
}
"pause" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_pause.lock().unwrap().clear();
}
}
"failed" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_fail.lock().unwrap().clear();
}
}
"remove" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_remove.lock().unwrap().clear();
}
}
"progress" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_progress.lock().unwrap().clear();
}
}
"resume" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_resume.lock().unwrap().clear();
}
}
"faultOccur" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_fault.lock().unwrap().clear();
}
}
"response" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_response.lock().unwrap().clear();
}
}
"wait" => {
if let Some(coll) = callback_mgr.tasks.lock().unwrap().get(&task_id) {
coll.on_wait.lock().unwrap().clear();
}
}
_ => unimplemented!(),
};
Ok(())
}
pub struct CallbackColl {
on_progress: Mutex<Vec<GlobalRefCallback<(bridge::Progress,)>>>,
on_complete: Mutex<Vec<GlobalRefCallback<(bridge::Progress,)>>>,
on_pause: Mutex<Vec<GlobalRefCallback<(bridge::Progress,)>>>,
on_resume: Mutex<Vec<GlobalRefCallback<(bridge::Progress,)>>>,
on_remove: Mutex<Vec<GlobalRefCallback<(bridge::Progress,)>>>,
on_fail: Mutex<Vec<GlobalRefCallback<(bridge::Progress,)>>>,
on_response: Mutex<Vec<GlobalRefCallback<(bridge::HttpResponse,)>>>,
on_fault: Mutex<Vec<GlobalRefCallback<(bridge::Faults,)>>>,
on_wait: Mutex<Vec<GlobalRefCallback<(bridge::WaitingReason,)>>>,
}
impl request_client::Callback for CallbackColl {
fn on_progress(&self, progress: &Progress) {
let callbacks = self.on_progress.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((progress.into(),));
}
}
fn on_completed(&self, progress: &Progress) {
let callbacks = self.on_complete.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((progress.into(),));
}
}
fn on_pause(&self, progress: &Progress) {
let callbacks = self.on_pause.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((progress.into(),));
}
}
fn on_resume(&self, progress: &Progress) {
let callbacks = self.on_resume.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((progress.into(),));
}
}
fn on_remove(&self, progress: &Progress) {
let callbacks = self.on_remove.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((progress.into(),));
}
}
fn on_response(&self, response: &Response) {
let callbacks = self.on_response.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((response.into(),));
}
}
fn on_failed(&self, progress: &Progress, _error_code: i32) {
let callbacks = self.on_fail.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((progress.into(),));
}
}
fn on_fault(&self, faults: Faults) {
let callbacks = self.on_fault.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((faults.into(),));
}
}
fn on_wait(&self, waiting_reason: WaitingReason) {
let callbacks = self.on_wait.lock().unwrap();
for callback in callbacks.iter() {
callback.execute((waiting_reason.into(),));
}
}
}
pub struct CallbackManager {
tasks: Mutex<HashMap<i64, Arc<CallbackColl>>>,
}
impl CallbackManager {
pub fn get_instance() -> &'static Self {
static INSTANCE: OnceLock<CallbackManager> = OnceLock::new();
INSTANCE.get_or_init(|| CallbackManager {
tasks: Mutex::new(HashMap::new()),
})
}
pub fn remove_task(&self, task_id: i64) {
self.tasks.lock().unwrap().remove(&task_id);
}
}