hara_native/vm/machine/
async_runtime.rs1use 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}