Skip to main content

hara_native/vm/machine/
async_runtime.rs

1//! Cooperative promise scheduling for bytecode machines.
2
3use super::*;
4
5impl Machine {
6    fn retain_async_child(
7        scheduler: &Rc<RefCell<AsyncScheduler>>,
8        machine: Machine,
9        result: Promise,
10        pending: Promise,
11    ) {
12        let id = {
13            let mut state = scheduler.borrow_mut();
14            let id = state.next_id;
15            state.next_id = id.wrapping_add(1);
16            state.children.insert(
17                id,
18                AsyncChild {
19                    machine,
20                    result: result.downgrade(),
21                    pending: pending.clone(),
22                },
23            );
24            id
25        };
26        let weak = Rc::downgrade(scheduler);
27        pending.on_settle(Rc::new(move |state| {
28            if let Some(scheduler) = weak.upgrade() {
29                scheduler.borrow_mut().ready.push_back((id, state));
30            }
31        }));
32    }
33
34    fn finish_async(
35        scheduler: &Rc<RefCell<AsyncScheduler>>,
36        mut machine: Machine,
37        result: Promise,
38        outcome: VmOutcome,
39    ) {
40        match outcome {
41            VmOutcome::Returned(value) => {
42                #[cfg(feature = "tracing-jit")]
43                store_program_jit(&machine.program.clone(), std::mem::take(&mut machine.jit));
44                settle_result(&result, Ok(value));
45            }
46            VmOutcome::Failed(error) => {
47                #[cfg(feature = "tracing-jit")]
48                store_program_jit(&machine.program.clone(), std::mem::take(&mut machine.jit));
49                result.reject(error.message);
50            }
51            VmOutcome::Suspended(pending) => {
52                Self::retain_async_child(scheduler, machine, result, pending)
53            }
54            VmOutcome::Yielded(_) => {
55                #[cfg(feature = "tracing-jit")]
56                store_program_jit(&machine.program.clone(), std::mem::take(&mut machine.jit));
57                result.reject("coroutine/yield used outside of a coroutine");
58            }
59        }
60    }
61
62    fn poll_scheduler(scheduler: &Rc<RefCell<AsyncScheduler>>) -> usize {
63        {
64            let mut state = scheduler.borrow_mut();
65            if state.polling {
66                return 0;
67            }
68            state.polling = true;
69        }
70        let pending = scheduler
71            .borrow()
72            .children
73            .values()
74            .map(|child| child.pending.clone())
75            .collect::<Vec<_>>();
76        for promise in pending {
77            promise.state();
78        }
79        let mut count = 0;
80        loop {
81            let Some((id, state)) = scheduler.borrow_mut().ready.pop_front() else {
82                break;
83            };
84            let Some(mut child) = scheduler.borrow_mut().children.remove(&id) else {
85                continue;
86            };
87            let Some(result) = child.result.upgrade() else {
88                child.pending.cancel();
89                continue;
90            };
91            count += 1;
92            let outcome = child.machine.resume(state);
93            Self::finish_async(scheduler, child.machine, result, outcome);
94        }
95        scheduler.borrow_mut().polling = false;
96        count
97    }
98
99    fn cancel_async_result(scheduler: &Rc<RefCell<AsyncScheduler>>, identity: usize) {
100        let ids = scheduler
101            .borrow()
102            .children
103            .iter()
104            .filter_map(|(id, child)| {
105                child
106                    .result
107                    .upgrade()
108                    .is_some_and(|candidate| candidate.identity_address() == identity)
109                    .then_some(*id)
110            })
111            .collect::<Vec<_>>();
112        for id in ids {
113            let child = { scheduler.borrow_mut().children.remove(&id) };
114            if let Some(child) = child {
115                child.pending.cancel();
116            }
117        }
118    }
119
120    pub(super) fn spawn_async(&self, mut machine: Machine) -> Promise {
121        let scheduler = self
122            .scheduler
123            .upgrade()
124            .expect("root VM owns its async scheduler");
125        machine.scheduler = Rc::downgrade(&scheduler);
126        machine.scheduler_owner = None;
127        let result = Promise::new();
128        let context = crate::core::NativeCallbackContext::capture();
129        let poll = scheduler.clone();
130        let poll_context = context.clone();
131        result.set_poller(Rc::new(move || {
132            poll_context.with(|| Self::poll_scheduler(&poll));
133        }));
134        let wait = scheduler.clone();
135        let wait_context = context.clone();
136        result.set_waiter(Rc::new(move || {
137            wait_context.with(|| Self::poll_scheduler(&wait));
138        }));
139        let cancel_scheduler = scheduler.clone();
140        let cancel_identity = result.identity_address();
141        result.set_cancel_hook(Rc::new(move || {
142            Self::cancel_async_result(&cancel_scheduler, cancel_identity);
143        }));
144        let outcome = machine.run();
145        Self::finish_async(&scheduler, machine, result.clone(), outcome);
146        result
147    }
148
149    pub fn poll_async(&mut self) -> usize {
150        self.scheduler
151            .upgrade()
152            .map(|scheduler| Self::poll_scheduler(&scheduler))
153            .unwrap_or(0)
154    }
155}
156
157pub(super) fn async_result(mut machine: Machine) -> Promise {
158    let outcome = machine.run();
159    async_result_from_outcome(machine, outcome)
160}
161
162pub(super) fn async_result_from_outcome(mut machine: Machine, outcome: VmOutcome) -> Promise {
163    let scheduler = machine
164        .scheduler
165        .upgrade()
166        .or_else(|| machine.scheduler_owner.clone())
167        .expect("VM owns its async scheduler");
168    machine.scheduler = Rc::downgrade(&scheduler);
169    machine.scheduler_owner = None;
170    let result = Promise::new();
171    let context = crate::core::NativeCallbackContext::capture();
172    let poll = scheduler.clone();
173    let poll_context = context.clone();
174    result.set_poller(Rc::new(move || {
175        poll_context.with(|| Machine::poll_scheduler(&poll));
176    }));
177    let wait = scheduler.clone();
178    let wait_context = context.clone();
179    result.set_waiter(Rc::new(move || {
180        wait_context.with(|| Machine::poll_scheduler(&wait));
181    }));
182    let cancel_scheduler = scheduler.clone();
183    let cancel_identity = result.identity_address();
184    result.set_cancel_hook(Rc::new(move || {
185        Machine::cancel_async_result(&cancel_scheduler, cancel_identity);
186    }));
187    Machine::finish_async(&scheduler, machine, result.clone(), outcome);
188    result
189}