#![allow(dead_code)]
use std::any::{Any, TypeId};
use std::future::Future;
use dashmap::{DashMap, mapref::entry::Entry};
use tokio::{
runtime::Runtime,
sync::{
mpsc::{UnboundedSender, UnboundedReceiver, unbounded_channel},
oneshot,
},
task::{JoinError, JoinHandle},
};
pub trait Task: Sized + Send + 'static {
type Output: Send + 'static;
fn run(self) -> impl Future<Output = Self::Output> + Send;
}
pub trait BatchTask: Sized + Send + 'static {
fn batch_run(list: Vec<Self>) -> impl Future<Output = ()> + Send;
}
type BatchItem<BT> = (BT, Option<oneshot::Sender<()>>);
pub struct TaskExecutor {
rt: &'static Runtime,
batch_senders: DashMap<TypeId, Box<dyn Any + Send + Sync>>,
}
impl TaskExecutor {
pub fn new(rt: &'static Runtime) -> Self {
TaskExecutor {
rt,
batch_senders: DashMap::new(),
}
}
pub async fn execute_waiting<T: Task>(&self, task: T) -> Result<T::Output, JoinError> {
let handle = self.rt.spawn(task.run());
handle.await
}
pub fn execute_detached<T: Task>(&self, task: T) -> JoinHandle<T::Output> {
self.rt.spawn(task.run())
}
fn ensure_sender<BT: BatchTask>(&self) -> UnboundedSender<BatchItem<BT>> {
let key = TypeId::of::<BT>();
match self.batch_senders.entry(key) {
Entry::Occupied(o) => o.get()
.downcast_ref::<UnboundedSender<BatchItem<BT>>>()
.expect("Type mismatch in batch_senders")
.clone(),
Entry::Vacant(v) => {
let (tx, rx) = unbounded_channel::<BatchItem<BT>>();
Self::spawn_worker::<BT>(self.rt, rx);
v.insert(Box::new(tx.clone()));
tx
}
}
}
fn spawn_worker<BT: BatchTask>(rt: &'static Runtime, mut rx: UnboundedReceiver<BatchItem<BT>>) {
rt.spawn(async move {
loop {
let Some((first_task, first_sig)) = rx.recv().await else { break; };
let mut tasks: Vec<BT> = vec![first_task];
let mut sigs: Vec<Option<oneshot::Sender<()>>> = vec![first_sig];
while let Ok((t, s)) = rx.try_recv() {
tasks.push(t);
sigs.push(s);
}
let join = tokio::spawn(async move {
BT::batch_run(tasks).await;
});
match join.await {
Ok(()) => {
for s in sigs.into_iter().flatten() {
let _ = s.send(());
}
}
Err(_panic) => {
}
}
}
});
}
pub fn execute_batch_detached<BT: BatchTask>(&self, batch_task: BT) {
let key = TypeId::of::<BT>();
let mut tx = self.ensure_sender::<BT>();
match tx.send((batch_task, None)) {
Ok(()) => {}
Err(e) => {
self.batch_senders.remove(&key);
tx = self.ensure_sender::<BT>();
let _ = tx.send(e.0);
}
}
}
pub async fn execute_batch_waiting<BT: BatchTask>(
&self,
batch_task: BT,
) -> Result<(), oneshot::error::RecvError> {
let key = TypeId::of::<BT>();
let (tx_oneshot, rx_oneshot) = oneshot::channel::<()>();
let mut tx = self.ensure_sender::<BT>();
match tx.send((batch_task, Some(tx_oneshot))) {
Ok(()) => {}
Err(e) => {
self.batch_senders.remove(&key);
tx = self.ensure_sender::<BT>();
let _ = tx.send(e.0);
}
}
rx_oneshot.await
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use std::time::Duration;
use tokio::runtime::Builder;
use tokio::time::{sleep, timeout};
static TEST_RT: OnceLock<Runtime> = OnceLock::new();
fn get_test_runtime() -> &'static Runtime {
TEST_RT.get_or_init(|| {
Builder::new_multi_thread()
.worker_threads(2)
.enable_all()
.build()
.unwrap()
})
}
#[derive(Debug, Clone)]
struct SimpleTask {
id: u32,
value: String,
}
impl Task for SimpleTask {
type Output = String;
async fn run(self) -> Self::Output {
format!("Task {} completed with value: {}", self.id, self.value)
}
}
#[derive(Debug, Clone)]
struct DelayedTask {
id: u32,
delay_ms: u64,
}
impl Task for DelayedTask {
type Output = u32;
async fn run(self) -> Self::Output {
sleep(Duration::from_millis(self.delay_ms)).await;
self.id
}
}
#[derive(Debug, Clone)]
struct PanicTask;
impl Task for PanicTask {
type Output = String;
async fn run(self) -> Self::Output {
panic!("This task panics intentionally");
}
}
#[derive(Debug, Clone)]
struct CounterTask {
counter: Arc<AtomicUsize>,
increment: usize,
}
impl BatchTask for CounterTask {
async fn batch_run(list: Vec<Self>) {
if list.is_empty() {
return;
}
let counter = &list[0].counter;
let total_increment: usize = list.iter().map(|task| task.increment).sum();
counter.fetch_add(total_increment, Ordering::SeqCst);
}
}
#[derive(Debug, Clone)]
struct LogTask {
message: String,
log_storage: Arc<Mutex<Vec<String>>>,
}
impl BatchTask for LogTask {
async fn batch_run(list: Vec<Self>) {
if list.is_empty() {
return;
}
let log_storage = list[0].log_storage.clone();
let mut storage = log_storage.lock().unwrap();
for task in list {
storage.push(format!("BATCH: {}", task.message));
}
}
}
#[derive(Debug, Clone)]
struct DelayedBatchTask {
id: u32,
delay_ms: u64,
results: Arc<Mutex<Vec<u32>>>,
}
impl BatchTask for DelayedBatchTask {
async fn batch_run(list: Vec<Self>) {
if list.is_empty() {
return;
}
let results = list[0].results.clone();
let delay_ms = list[0].delay_ms;
sleep(Duration::from_millis(delay_ms)).await;
let mut results_guard = results.lock().unwrap();
for task in list {
results_guard.push(task.id);
}
}
}
#[derive(Debug, Clone)]
struct PanicBatchTask;
impl BatchTask for PanicBatchTask {
async fn batch_run(_list: Vec<Self>) {
panic!("This batch task panics intentionally");
}
}
#[test]
fn test_task_executor_new() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
assert_eq!(executor.batch_senders.len(), 0);
}
#[test]
fn test_execute_waiting_simple() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let result = rt.block_on(async {
let task = SimpleTask {
id: 1,
value: "test".to_string(),
};
executor.execute_waiting(task).await
});
assert!(result.is_ok());
assert_eq!(result.unwrap(), "Task 1 completed with value: test");
}
#[test]
fn test_execute_waiting_delayed() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let result = rt.block_on(async {
let task = DelayedTask {
id: 42,
delay_ms: 50,
};
executor.execute_waiting(task).await
});
assert!(result.is_ok());
assert_eq!(result.unwrap(), 42);
}
#[test]
fn test_execute_waiting_panic() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let result = rt.block_on(async {
let task = PanicTask;
executor.execute_waiting(task).await
});
assert!(result.is_err());
assert!(result.unwrap_err().is_panic());
}
#[test]
fn test_execute_detached() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let result = rt.block_on(async {
let task = SimpleTask {
id: 2,
value: "detached".to_string(),
};
let handle = executor.execute_detached(task);
handle.await
});
assert!(result.is_ok());
assert_eq!(result.unwrap(), "Task 2 completed with value: detached");
}
#[test]
fn test_execute_detached_panic() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let result = rt.block_on(async {
let task = PanicTask;
let handle = executor.execute_detached(task);
handle.await
});
assert!(result.is_err());
assert!(result.unwrap_err().is_panic());
}
#[test]
fn test_execute_batch_detached() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let counter = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: 1,
});
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: 2,
});
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: 3,
});
sleep(Duration::from_millis(100)).await;
});
assert_eq!(counter.load(Ordering::SeqCst), 6);
}
#[test]
fn test_execute_batch_waiting() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let counter = Arc::new(AtomicUsize::new(0));
let result = rt.block_on(async {
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: 5,
});
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: 10,
});
executor
.execute_batch_waiting(CounterTask {
counter: counter.clone(),
increment: 15,
})
.await
});
assert!(result.is_ok());
assert_eq!(counter.load(Ordering::SeqCst), 30);
}
#[test]
fn test_batch_processing_with_storage() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let log_storage = Arc::new(Mutex::new(Vec::new()));
rt.block_on(async {
executor.execute_batch_detached(LogTask {
message: "First log".to_string(),
log_storage: log_storage.clone(),
});
executor.execute_batch_detached(LogTask {
message: "Second log".to_string(),
log_storage: log_storage.clone(),
});
executor
.execute_batch_waiting(LogTask {
message: "Third log".to_string(),
log_storage: log_storage.clone(),
})
.await
.unwrap();
});
let logs = log_storage.lock().unwrap();
assert_eq!(logs.len(), 3);
assert!(logs.contains(&"BATCH: First log".to_string()));
assert!(logs.contains(&"BATCH: Second log".to_string()));
assert!(logs.contains(&"BATCH: Third log".to_string()));
}
#[test]
fn test_delayed_batch_processing() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let results = Arc::new(Mutex::new(Vec::new()));
let result = rt.block_on(async {
timeout(
Duration::from_millis(500),
executor.execute_batch_waiting(DelayedBatchTask {
id: 1,
delay_ms: 100,
results: results.clone(),
}),
)
.await
});
assert!(result.is_ok());
assert!(result.unwrap().is_ok());
let results_guard = results.lock().unwrap();
assert_eq!(results_guard.len(), 1);
assert_eq!(results_guard[0], 1);
}
#[test]
fn test_multiple_batch_types() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let counter = Arc::new(AtomicUsize::new(0));
let log_storage = Arc::new(Mutex::new(Vec::new()));
rt.block_on(async {
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: 100,
});
executor.execute_batch_detached(LogTask {
message: "Mixed batch test".to_string(),
log_storage: log_storage.clone(),
});
executor
.execute_batch_waiting(CounterTask {
counter: counter.clone(),
increment: 200,
})
.await
.unwrap();
executor
.execute_batch_waiting(LogTask {
message: "Mixed batch test 2".to_string(),
log_storage: log_storage.clone(),
})
.await
.unwrap();
});
assert_eq!(counter.load(Ordering::SeqCst), 300);
let logs = log_storage.lock().unwrap();
assert_eq!(logs.len(), 2);
}
#[test]
fn test_concurrent_execution() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let result = rt.block_on(async {
let mut handles = Vec::new();
for i in 0..10 {
let handle = executor.execute_detached(DelayedTask {
id: i,
delay_ms: 20,
});
handles.push(handle);
}
let mut results = Vec::new();
for handle in handles {
results.push(handle.await.unwrap());
}
results.sort();
results
});
assert_eq!(result.len(), 10);
assert_eq!(result, (0..10).collect::<Vec<_>>());
}
#[test]
fn test_batch_error_handling() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let result = rt.block_on(async {
timeout(
Duration::from_millis(500),
executor.execute_batch_waiting(PanicBatchTask),
)
.await
});
assert!(result.is_ok());
assert!(result.unwrap().is_err());
}
#[test]
fn batch_panics_but_worker_survives() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let r = rt.block_on(async {
executor.execute_batch_waiting(PanicBatchTask).await
});
assert!(r.is_err());
#[derive(Clone)]
struct OkTask(Arc<AtomicUsize>);
impl BatchTask for OkTask {
async fn batch_run(list: Vec<Self>) {
let c = &list[0].0;
c.fetch_add(list.len(), Ordering::SeqCst);
}
}
let c = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
executor.execute_batch_waiting(OkTask(c.clone())).await.unwrap();
});
assert_eq!(c.load(Ordering::SeqCst), 1);
}
#[test]
fn test_batch_processor_reuse() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let counter = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
executor
.execute_batch_waiting(CounterTask {
counter: counter.clone(),
increment: 1,
})
.await
.unwrap();
executor
.execute_batch_waiting(CounterTask {
counter: counter.clone(),
increment: 2,
})
.await
.unwrap();
});
assert_eq!(counter.load(Ordering::SeqCst), 3);
assert_eq!(executor.batch_senders.len(), 1);
}
#[test]
fn test_empty_batch_handling() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let counter = Arc::new(AtomicUsize::new(0));
#[derive(Debug, Clone)]
struct EmptyBatchTask {
counter: Arc<AtomicUsize>,
}
impl BatchTask for EmptyBatchTask {
async fn batch_run(list: Vec<Self>) {
if list.is_empty() {
return;
}
list[0].counter.fetch_add(1, Ordering::SeqCst);
}
}
let result = rt.block_on(async {
executor
.execute_batch_waiting(EmptyBatchTask {
counter: counter.clone(),
})
.await
});
assert!(result.is_ok());
assert_eq!(counter.load(Ordering::SeqCst), 1);
}
#[test]
fn test_high_volume_batch_processing() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let counter = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
for i in 0..100 {
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: i,
});
}
executor
.execute_batch_waiting(CounterTask {
counter: counter.clone(),
increment: 0,
})
.await
.unwrap();
});
assert_eq!(counter.load(Ordering::SeqCst), 4950);
}
#[test]
fn test_mixed_execution_patterns() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
let counter = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
let individual_handle = executor.execute_detached(DelayedTask {
id: 999,
delay_ms: 50,
});
executor.execute_batch_detached(CounterTask {
counter: counter.clone(),
increment: 10,
});
let individual_result = individual_handle.await.unwrap();
assert_eq!(individual_result, 999);
executor
.execute_batch_waiting(CounterTask {
counter: counter.clone(),
increment: 20,
})
.await
.unwrap();
let batch_result = executor
.execute_waiting(SimpleTask {
id: 777,
value: "mixed".to_string(),
})
.await
.unwrap();
assert_eq!(batch_result, "Task 777 completed with value: mixed");
});
assert_eq!(counter.load(Ordering::SeqCst), 30);
}
#[test]
fn same_type_panics_then_succeeds() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct SometimesPanicTask {
counter: Arc<AtomicUsize>,
should_panic: bool,
}
impl BatchTask for SometimesPanicTask {
async fn batch_run(list: Vec<Self>) {
if list.iter().any(|t| t.should_panic) {
panic!("boom");
}
let c = &list[0].counter;
c.fetch_add(list.len(), Ordering::SeqCst);
}
}
let c = Arc::new(AtomicUsize::new(0));
let r = rt.block_on(async {
tokio::time::timeout(
Duration::from_millis(500),
executor.execute_batch_waiting(SometimesPanicTask { counter: c.clone(), should_panic: true }),
).await
}).expect("timeout waiting for panic batch");
assert!(r.is_err());
rt.block_on(async {
executor.execute_batch_waiting(SometimesPanicTask { counter: c.clone(), should_panic: false })
.await
.unwrap();
});
assert_eq!(c.load(Ordering::SeqCst), 1);
}
#[test]
fn multiple_waiters_all_get_recverror_on_panic() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct AllPanic;
impl BatchTask for AllPanic {
async fn batch_run(_list: Vec<Self>) {
panic!("batch exploded");
}
}
let (r1, r2, r3) = rt.block_on(async {
tokio::join!(
executor.execute_batch_waiting(AllPanic),
executor.execute_batch_waiting(AllPanic),
executor.execute_batch_waiting(AllPanic),
)
});
assert!(r1.is_err());
assert!(r2.is_err());
assert!(r3.is_err());
}
#[test]
fn panic_storm_then_recover() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct SometimesPanic {
ok_counter: Arc<AtomicUsize>,
panic: bool,
}
impl BatchTask for SometimesPanic {
async fn batch_run(list: Vec<Self>) {
if list.iter().any(|t| t.panic) { panic!("storm"); }
let c = &list[0].ok_counter;
c.fetch_add(list.len(), Ordering::SeqCst);
}
}
for _ in 0..10 {
let r = rt.block_on(async { executor.execute_batch_waiting(SometimesPanic { ok_counter: Arc::new(AtomicUsize::new(0)), panic: true }).await });
assert!(r.is_err());
}
let ok = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
executor.execute_batch_waiting(SometimesPanic { ok_counter: ok.clone(), panic: false }).await.unwrap();
});
assert_eq!(ok.load(Ordering::SeqCst), 1);
}
#[test]
fn interleaved_ok_panic_ok_same_type() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct FlipTask {
sum: Arc<AtomicUsize>,
panic: bool,
}
impl BatchTask for FlipTask {
async fn batch_run(list: Vec<Self>) {
if list.iter().any(|t| t.panic) { panic!("flip"); }
list[0].sum.fetch_add(list.len(), Ordering::SeqCst);
}
}
let sum = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
executor.execute_batch_waiting(FlipTask { sum: sum.clone(), panic: false }).await.unwrap();
});
assert_eq!(sum.load(Ordering::SeqCst), 1);
rt.block_on(async {
tokio::time::timeout(
Duration::from_millis(500),
executor.execute_batch_waiting(FlipTask { sum: sum.clone(), panic: true }),
).await
}).expect("timeout").unwrap_err();
rt.block_on(async {
executor.execute_batch_waiting(FlipTask { sum: sum.clone(), panic: false }).await.unwrap();
});
assert_eq!(sum.load(Ordering::SeqCst), 2);
}
#[test]
fn detached_and_waiters_mixed_on_panic() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct MixPanic;
impl BatchTask for MixPanic {
async fn batch_run(_list: Vec<Self>) {
panic!("oops");
}
}
rt.block_on(async {
for _ in 0..5 {
executor.execute_batch_detached(MixPanic);
}
let (a, b) = tokio::join!(
executor.execute_batch_waiting(MixPanic),
executor.execute_batch_waiting(MixPanic),
);
assert!(a.is_err());
assert!(b.is_err());
});
}
#[test]
fn batch_coalesces_many_tasks_lower_bound() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct RecordLen {
last_len: Arc<AtomicUsize>,
}
impl BatchTask for RecordLen {
async fn batch_run(list: Vec<Self>) {
list[0].last_len.store(list.len(), Ordering::SeqCst);
}
}
let last = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
for _ in 0..200 {
executor.execute_batch_detached(RecordLen { last_len: last.clone() });
}
executor.execute_batch_waiting(RecordLen { last_len: last.clone() })
.await
.unwrap();
});
assert!(last.load(Ordering::SeqCst) >= 1);
}
#[test]
fn many_waiters_all_panic_no_hang() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct Boom;
impl BatchTask for Boom {
async fn batch_run(_list: Vec<Self>) { panic!("boom"); }
}
rt.block_on(async {
let mut results = Vec::new();
for _ in 0..50 {
let result = tokio::time::timeout(
Duration::from_secs(1),
executor.execute_batch_waiting(Boom)
).await;
results.push(result);
}
assert!(results.iter().all(|r| r.is_ok() && r.as_ref().unwrap().is_err()));
});
}
#[test]
fn dropped_waiter_does_not_break_worker() {
let rt = get_test_runtime();
#[derive(Clone)]
struct SlowOk(Arc<AtomicUsize>);
impl BatchTask for SlowOk {
async fn batch_run(list: Vec<Self>) {
tokio::time::sleep(Duration::from_millis(50)).await;
list[0].0.fetch_add(list.len(), Ordering::SeqCst);
}
}
let c1 = Arc::new(AtomicUsize::new(0));
let c2 = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
let rt_inner = get_test_runtime();
let executor1 = TaskExecutor::new(rt_inner);
let handle = tokio::spawn(async move {
executor1.execute_batch_waiting(SlowOk(c1)).await
});
tokio::time::sleep(Duration::from_millis(5)).await;
handle.abort();
let rt_inner2 = get_test_runtime();
let executor2 = TaskExecutor::new(rt_inner2);
executor2.execute_batch_waiting(SlowOk(c2.clone())).await.unwrap();
});
assert!(c2.load(Ordering::SeqCst) >= 1);
}
#[test]
fn multi_type_isolation_panic_does_not_affect_other_type() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct WillPanic;
impl BatchTask for WillPanic {
async fn batch_run(_list: Vec<Self>) { panic!("type A broke"); }
}
#[derive(Clone)]
struct TypeB(Arc<AtomicUsize>);
impl BatchTask for TypeB {
async fn batch_run(list: Vec<Self>) {
list[0].0.fetch_add(list.len(), Ordering::SeqCst);
}
}
let ra = rt.block_on(async { executor.execute_batch_waiting(WillPanic).await });
assert!(ra.is_err());
let bsum = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
executor.execute_batch_waiting(TypeB(bsum.clone())).await.unwrap();
});
assert_eq!(bsum.load(Ordering::SeqCst), 1);
}
#[test]
fn large_mixed_after_panics_recovers() {
let rt = get_test_runtime();
let executor = TaskExecutor::new(rt);
#[derive(Clone)]
struct MaybePanic {
ok: Arc<AtomicUsize>,
panic: bool,
}
impl BatchTask for MaybePanic {
async fn batch_run(list: Vec<Self>) {
if list.iter().any(|t| t.panic) { panic!("mixed"); }
list[0].ok.fetch_add(list.len(), Ordering::SeqCst);
}
}
for _ in 0..5 {
let r = rt.block_on(async { executor.execute_batch_waiting(MaybePanic { ok: Arc::new(AtomicUsize::new(0)), panic: true }).await });
assert!(r.is_err());
}
let ok = Arc::new(AtomicUsize::new(0));
rt.block_on(async {
for _ in 0..100 {
executor.execute_batch_detached(MaybePanic { ok: ok.clone(), panic: false });
}
executor.execute_batch_waiting(MaybePanic { ok: ok.clone(), panic: false }).await.unwrap();
});
assert!(ok.load(Ordering::SeqCst) >= 1);
}
}