moirai_iter/execution/
async_ctx.rs1use super::base::ExecutionBase;
4use super::hybrid::owned_chunks;
5
6#[derive(Clone)]
8pub struct AsyncContext {
9 pub(super) batch_size: usize,
10 pub(super) max_concurrent: usize,
11}
12
13impl Default for AsyncContext {
14 fn default() -> Self {
15 Self::new()
16 }
17}
18
19impl AsyncContext {
20 pub fn new() -> Self {
22 Self {
23 batch_size: 100,
24 max_concurrent: 1000,
25 }
26 }
27
28 pub fn with_batch_size(batch_size: usize) -> Self {
30 Self {
31 batch_size,
32 max_concurrent: 1000,
33 }
34 }
35
36 pub fn with_max_concurrent(mut self, max_concurrent: usize) -> Self {
38 self.max_concurrent = max_concurrent;
39 self
40 }
41}
42
43impl AsyncContext {
44 pub fn execute_iter<T, F, R>(
46 &self,
47 items: Vec<T>,
48 func: F,
49 ) -> Result<Vec<R>, Box<dyn std::error::Error + Send + Sync>>
50 where
51 T: Send + 'static,
52 F: Fn(T) -> R + Send + Sync + 'static,
53 R: Send + 'static,
54 {
55 let mut results = Vec::with_capacity(items.len());
56
57 for batch in owned_chunks(items, self.batch_size) {
58 for item in batch {
59 let result = func(item);
60 results.push(result);
61 }
62 }
63
64 Ok(results)
65 }
66
67 pub fn execute<F, R>(&self, func: F) -> Result<R, Box<dyn std::error::Error + Send + Sync>>
69 where
70 F: FnOnce() -> R + Send,
71 R: Send,
72 {
73 Ok(func())
76 }
77}
78
79impl ExecutionBase for AsyncContext {
80 fn context_type(&self) -> &'static str {
81 "Async"
82 }
83}