use super::*;
const WORKER_AGENT_PRELUDE: &str = r#"
var $262 = {
global: this,
agent: {
report: function (m) { $262_agent_report(m); },
receiveBroadcast: function (f) { $262_agent_receiveBroadcast(f); },
leaving: function () {},
sleep: function (ms) { $262_agent_sleep(ms); },
monotonicNow: function () { return $262_agent_monotonicNow(); }
}
};
"#;
impl<'a> Interp<'a> {
pub(crate) fn agent_start(&mut self, src: NanBox) -> Result<NanBox, ExecError> {
let source = self.coerce_to_string(src)?;
let idx = self.create_realm_env();
let full = alloc::format!("{WORKER_AGENT_PRELUDE}{source}");
let saved = self.agent.current_agent_realm;
self.agent.current_agent_realm = Some(idx);
let result = self.eval_source_in_realm(idx, &full);
self.agent.current_agent_realm = saved;
match result {
Ok(_) | Err(ExecError::Throw(_)) => Ok(NanBox::undefined()),
Err(other) => Err(other),
}
}
pub(crate) fn agent_broadcast(&mut self, sab: NanBox) -> Result<NanBox, ExecError> {
let cbs = core::mem::take(&mut self.agent.broadcasts);
for (idx, cb) in cbs {
if idx < self.created_realms.len() {
let saved_intrinsics = self.realm.intrinsics_snapshot();
let saved_gt = self.global_this;
self.realm
.restore_intrinsics(self.created_realms[idx].intrinsics);
self.global_this = self.created_realms[idx].global_this;
let child_intl = core::mem::take(&mut self.created_realms[idx].intl_protos);
let saved_intl = self.realm.replace_intl_protos(child_intl);
let _ = self.call(cb, &[sab]);
self.created_realms[idx].intl_protos = self.realm.replace_intl_protos(saved_intl);
self.realm.restore_intrinsics(saved_intrinsics);
self.global_this = saved_gt;
} else {
let _ = self.call(cb, &[sab]);
}
}
Ok(NanBox::undefined())
}
pub(crate) fn agent_report(&mut self, msg: NanBox) -> Result<NanBox, ExecError> {
let s = self.coerce_to_string(msg)?;
self.agent.reports.push_back(s);
Ok(NanBox::undefined())
}
pub(crate) fn agent_get_report(&mut self) -> NanBox {
match self.agent.reports.pop_front() {
Some(s) => self.new_str(&s),
None => NanBox::null(),
}
}
pub(crate) fn agent_get_report_async(&mut self) -> Result<NanBox, ExecError> {
let v = self.agent_get_report();
let p = self.fresh_promise();
self.resolve_with(p, v);
Ok(NanBox::handle(p.to_raw()))
}
pub(crate) fn agent_receive_broadcast(&mut self, cb: NanBox) -> Result<NanBox, ExecError> {
let idx = self.agent.current_agent_realm.unwrap_or(usize::MAX);
self.agent.broadcasts.push((idx, cb));
Ok(NanBox::undefined())
}
pub(crate) fn agent_monotonic_now(&mut self) -> f64 {
self.virtual_now
}
fn atomics_wait_key(&self, ta: Handle, idx: usize, kind: u8) -> (u64, usize) {
let buffer = self.realm.typed_array_object(ta).map_or(0, |h| h.to_raw());
let size = if crate::nbexec::is_bigint_kind(kind) {
8
} else {
4
};
let offset = self.realm.typed_byte_offset(ta).unwrap_or(0);
(buffer, offset + idx * size)
}
pub(crate) fn atomics_notify(&mut self, ta: Handle, idx: usize, kind: u8, count: f64) -> usize {
let (buffer, byte_index) = self.atomics_wait_key(ta, idx, kind);
let mut to_resolve = Vec::new();
let mut remaining = count;
let mut i = 0;
while i < self.agent.waiters.len() {
if remaining <= 0.0 {
break;
}
let w = &self.agent.waiters[i];
if w.buffer == buffer && w.byte_index == byte_index {
to_resolve.push(self.agent.waiters.remove(i).promise);
remaining -= 1.0;
} else {
i += 1;
}
}
let woken = to_resolve.len();
for p in to_resolve {
let ok = self.new_str("ok");
self.settle(p, ok, true);
}
woken
}
pub(crate) fn atomics_wait_async(
&mut self,
ta: Handle,
idx: usize,
kind: u8,
equal: bool,
timeout: f64,
) -> NanBox {
let result = self.realm.new_object();
let resolved = |slf: &mut Self, s: &str| -> NanBox {
let v = slf.new_str(s);
slf.realm
.set_property(result, "async", NanBox::boolean(false));
slf.realm.set_property(result, "value", v);
NanBox::handle(result.to_raw())
};
if !equal {
return resolved(self, "not-equal");
}
let t = if timeout.is_nan() {
f64::INFINITY
} else {
timeout.max(0.0)
};
if t <= 0.0 {
return resolved(self, "timed-out");
}
let (buffer, byte_index) = self.atomics_wait_key(ta, idx, kind);
let promise = self.fresh_promise();
self.agent.waiters.push(AtomicsWaiter {
buffer,
byte_index,
promise,
});
if t.is_finite() {
let cb = self
.realm
.new_bound_native(N_ATOMICS_ASYNC_TIMEOUT, promise);
self.schedule_timer(t, NanBox::handle(cb.to_raw()), Vec::new());
}
self.realm
.set_property(result, "async", NanBox::boolean(true));
self.realm
.set_property(result, "value", NanBox::handle(promise.to_raw()));
NanBox::handle(result.to_raw())
}
pub(crate) fn atomics_wait_async_timeout(&mut self, promise: Handle) {
if let Some(pos) = self.agent.waiters.iter().position(|w| w.promise == promise) {
self.agent.waiters.remove(pos);
let v = self.new_str("timed-out");
self.settle(promise, v, true);
}
}
}