use std::cell::RefCell;
use js_sys::{Object, Reflect};
use wasm_bindgen::prelude::*;
use web_sys::Worker;
const MP_MAX_PEERS: i32 = 8;
const MESH_FRESH_SECS: u64 = 40; const MESH_BEAT_TICKS: u32 = 3;
enum MpRole {
Host {
peers: Vec<(String, i32, crate::app::webrtc::Peer)>,
},
Joiner {
peer: crate::app::webrtc::Peer,
},
Mesh {
peers: Vec<Option<(String, crate::app::webrtc::Peer)>>,
connecting: Vec<bool>,
},
}
struct MpSession {
role: MpRole,
_gw: crate::wallet::GeneratedWallet, #[allow(dead_code)]
room: String,
}
thread_local! {
static MP_SESSION: RefCell<Option<MpSession>> = const { RefCell::new(None) };
}
fn joiner_id_from(addr: &[u8; 20]) -> String {
let mut s = String::with_capacity(8);
for b in &addr[0..4] {
s.push_str(&format!("{b:02x}"));
}
s
}
pub(crate) async fn mp_connect(worker: Worker, code: i32, is_host: bool) {
mp_teardown(); let room = format!("mp-{code}");
let gw = crate::wallet::generate();
let signer = gw.signer.clone();
if is_host {
MP_SESSION.with(|s| {
*s.borrow_mut() = Some(MpSession {
role: MpRole::Host { peers: Vec::new() },
_gw: gw,
room: room.clone(),
});
});
mp_post_status(&worker, 0, 0, 1);
wasm_bindgen_futures::spawn_local(mp_host_accept_loop(worker, room, signer, 1, false));
return;
}
let joiner_id = joiner_id_from(&crate::wallet::address(&signer));
let worker_for_msg = worker.clone();
let on_msg = move |bytes: Vec<u8>| mp_dispatch_peer_frame(&worker_for_msg, 0, &bytes);
match crate::app::webrtc::Peer::offer_to_host(&room, &joiner_id, &signer, on_msg).await {
Ok(peer) => {
for _ in 0..150 {
if peer.is_open() {
break;
}
crate::runtime::sleep_ms(100).await;
}
let connected = i32::from(peer.is_open());
let mut self_index = 1;
for _ in 0..4 {
let roster =
crate::registry::signal_get_joiners(&room).await.unwrap_or_default();
if let Some(pos) = roster.iter().position(|id| id == &joiner_id) {
self_index = pos as i32 + 1;
break;
}
crate::runtime::sleep_ms(300).await;
}
MP_SESSION.with(|s| {
*s.borrow_mut() = Some(MpSession {
role: MpRole::Joiner { peer },
_gw: gw,
room: room.clone(),
});
});
mp_post_status(&worker, connected, self_index, if connected == 1 { 2 } else { 1 });
}
Err(e) => {
web_sys::console::warn_1(&JsValue::from_str(&format!("mp join failed: {e:?}")));
mp_post_status(&worker, 0, -1, 0);
}
}
}
async fn mp_host_accept_loop(
worker: Worker,
room: String,
signer: k256::ecdsa::SigningKey,
idx_base: i32,
skip_first: bool,
) {
loop {
let is_host = MP_SESSION
.with(|s| matches!(s.borrow().as_ref().map(|x| &x.role), Some(MpRole::Host { .. })));
if !is_host {
return;
}
let joiners = crate::registry::signal_get_joiners(&room).await.unwrap_or_default();
for (roster_pos, jid) in joiners.iter().enumerate() {
if skip_first && roster_pos == 0 {
continue; }
let idx = roster_pos as i32 + idx_base;
let known = MP_SESSION.with(|s| {
if let Some(MpSession { role: MpRole::Host { peers }, .. }) = s.borrow().as_ref() {
peers.iter().any(|(id, _, _)| id == jid)
} else {
true
}
});
if known || idx >= MP_MAX_PEERS {
continue;
}
let worker_for_msg = worker.clone();
let on_msg = move |bytes: Vec<u8>| mp_dispatch_peer_frame(&worker_for_msg, idx, &bytes);
match crate::app::webrtc::Peer::answer_joiner(&room, jid, &signer, on_msg).await {
Ok(peer) => {
let count = MP_SESSION.with(|s| {
if let Some(MpSession { role: MpRole::Host { peers }, .. }) =
s.borrow_mut().as_mut()
{
peers.push((jid.clone(), idx, peer));
peers.len() as i32 + 1
} else {
1
}
});
mp_post_status(&worker, i32::from(count >= 2), 0, count);
}
Err(e) => {
web_sys::console::warn_1(&JsValue::from_str(&format!(
"mp answer_joiner failed: {e:?}"
)));
}
}
}
crate::runtime::sleep_ms(2000).await;
}
}
pub(crate) async fn mp_connect_mesh(worker: Worker, code: i32) {
mp_teardown();
let room = format!("mp-{code}");
let gw = crate::wallet::generate();
let signer = gw.signer.clone();
let addr_bytes = crate::wallet::address(&signer);
let my_id = joiner_id_from(&addr_bytes);
let my_addr = crate::encoding::bytes_to_hex_str(&addr_bytes);
let my_slot = match mesh_claim_slot(&room, &signer, &my_id, &my_addr).await {
Ok(s) => s,
Err(e) => {
web_sys::console::warn_1(&JsValue::from_str(&format!("mesh claim failed: {e}")));
mp_post_status(&worker, 0, -1, 0);
return;
}
};
MP_SESSION.with(|s| {
*s.borrow_mut() = Some(MpSession {
role: MpRole::Mesh {
peers: (0..MP_MAX_PEERS).map(|_| None).collect(),
connecting: (0..MP_MAX_PEERS).map(|_| false).collect(),
},
_gw: gw,
room: room.clone(),
});
});
mp_post_status(&worker, 0, my_slot, 1); wasm_bindgen_futures::spawn_local(mesh_loop(worker, room, signer, my_id, my_addr, my_slot));
}
async fn mesh_claim_slot(
room: &str,
signer: &k256::ecdsa::SigningKey,
my_id: &str,
my_addr: &str,
) -> Result<i32, String> {
for _ in 0..6 {
let ms = crate::registry::signal_get_slots(room).await?;
if let Some(pos) = ms.slots.iter().position(|e| {
e.as_ref().map(|x| x.addr.eq_ignore_ascii_case(my_addr)).unwrap_or(false)
}) {
return Ok(pos as i32); }
let free = ms.slots.iter().position(|e| match e {
None => true,
Some(x) => ms.now.saturating_sub(x.ts) > MESH_FRESH_SECS,
});
let idx = free.ok_or_else(|| "arena full (8 players)".to_string())?;
let mut next = ms.slots.clone();
next[idx] = Some(crate::registry::SlotEntry {
id: my_id.to_string(),
addr: my_addr.to_string(),
ts: ms.now,
});
let now = (js_sys::Date::now() / 1000.0) as u64;
match crate::registry::signal_put_slots(signer, now, room, &next, idx, ms.sha.as_deref()).await {
Ok(crate::registry::PutSlots::Written) => return Ok(idx as i32),
Ok(crate::registry::PutSlots::Conflict) => crate::runtime::sleep_ms(250).await,
Err(e) => return Err(e),
}
}
Err("slot claim contention".to_string())
}
async fn mesh_loop(
worker: Worker,
room: String,
signer: k256::ecdsa::SigningKey,
my_id: String,
my_addr: String,
my_slot: i32,
) {
let mut tick: u32 = 0;
loop {
let is_mesh = MP_SESSION
.with(|s| matches!(s.borrow().as_ref().map(|x| &x.role), Some(MpRole::Mesh { .. })));
if !is_mesh {
return;
}
let ms = match crate::registry::signal_get_slots(&room).await {
Ok(m) => m,
Err(_) => {
crate::runtime::sleep_ms(4000).await;
tick += 1;
continue;
}
};
let still_mine = ms
.slots
.get(my_slot as usize)
.and_then(|e| e.as_ref())
.map(|x| x.addr.eq_ignore_ascii_case(&my_addr))
.unwrap_or(false);
if !still_mine && tick > 0 {
mp_teardown();
mp_post_status(&worker, 0, -1, 0);
return;
}
if tick % MESH_BEAT_TICKS == 0 {
let mut next = ms.slots.clone();
next[my_slot as usize] = Some(crate::registry::SlotEntry {
id: my_id.clone(),
addr: my_addr.clone(),
ts: ms.now,
});
let now = (js_sys::Date::now() / 1000.0) as u64;
let _ = crate::registry::signal_put_slots(
&signer, now, &room, &next, my_slot as usize, ms.sha.as_deref(),
)
.await;
}
for q in 0..(MP_MAX_PEERS as usize) {
if q as i32 == my_slot {
continue;
}
let entry = ms.slots[q].as_ref();
let fresh = entry
.map(|x| ms.now.saturating_sub(x.ts) <= MESH_FRESH_SECS)
.unwrap_or(false);
if !fresh {
continue;
}
let their_id = entry.map(|x| x.id.clone()).unwrap_or_default();
let skip = MP_SESSION.with(|s| {
if let Some(MpSession { role: MpRole::Mesh { peers, connecting, .. }, .. }) =
s.borrow().as_ref()
{
connecting[q] || peers[q].as_ref().map(|(id, _)| id == &their_id).unwrap_or(false)
} else {
true
}
});
if skip {
continue;
}
MP_SESSION.with(|s| {
if let Some(MpSession { role: MpRole::Mesh { connecting, .. }, .. }) =
s.borrow_mut().as_mut()
{
connecting[q] = true;
}
});
wasm_bindgen_futures::spawn_local(mesh_connect_one(
worker.clone(), room.clone(), signer.clone(), my_slot, q as i32, their_id,
));
}
tick += 1;
crate::runtime::sleep_ms(4000).await;
}
}
async fn mesh_connect_one(
worker: Worker,
room: String,
signer: k256::ecdsa::SigningKey,
my_slot: i32,
q: i32,
their_id: String,
) {
let worker_for_msg = worker.clone();
let on_msg = move |bytes: Vec<u8>| mp_dispatch_peer_frame(&worker_for_msg, q, &bytes);
let result = if my_slot < q {
crate::app::webrtc::Peer::mesh_offer(&room, my_slot, q, &signer, on_msg).await
} else {
crate::app::webrtc::Peer::mesh_answer(&room, q, my_slot, &signer, on_msg).await
};
MP_SESSION.with(|s| {
if let Some(MpSession { role: MpRole::Mesh { peers, connecting, .. }, .. }) =
s.borrow_mut().as_mut()
{
connecting[q as usize] = false;
if let Ok(peer) = result {
peers[q as usize] = Some((their_id, peer));
}
}
});
for _ in 0..100 {
let open = MP_SESSION.with(|s| {
matches!(s.borrow().as_ref().map(|x| &x.role), Some(MpRole::Mesh { peers, .. })
if peers.iter().flatten().any(|(_, p)| p.is_open()))
});
if open {
break;
}
crate::runtime::sleep_ms(100).await;
}
let (connected, total) = MP_SESSION.with(|s| {
if let Some(MpSession { role: MpRole::Mesh { peers, .. }, .. }) = s.borrow().as_ref() {
let open = peers.iter().flatten().filter(|(_, p)| p.is_open()).count() as i32;
(i32::from(open > 0), open + 1)
} else {
(0, 0)
}
});
mp_post_status(&worker, connected, my_slot, total);
}
fn mp_post_status(worker: &Worker, connected: i32, self_index: i32, peer_count: i32) {
let m = Object::new();
let _ = Reflect::set(&m, &JsValue::from_str("type"), &JsValue::from_str("mp:status"));
let _ = Reflect::set(&m, &JsValue::from_str("connected"), &JsValue::from_f64(connected as f64));
let _ = Reflect::set(&m, &JsValue::from_str("selfIndex"), &JsValue::from_f64(self_index as f64));
let _ = Reflect::set(&m, &JsValue::from_str("peerCount"), &JsValue::from_f64(peer_count as f64));
let _ = worker.post_message(&m);
}
pub(crate) fn mp_read_int_array(data: &JsValue, field: &str) -> Vec<i32> {
Reflect::get(data, &JsValue::from_str(field))
.ok()
.map(|v| {
js_sys::Array::from(&v)
.iter()
.map(|x| x.as_f64().unwrap_or(0.0) as i32)
.collect()
})
.unwrap_or_default()
}
pub(crate) fn mp_send(deltas: Option<Vec<i32>>, events: Option<Vec<i32>>) {
let json = if let Some(d) = deltas {
format!("{{\"d\":{}}}", mp_ints_json(&d))
} else if let Some(ev) = events {
format!("{{\"e\":{}}}", mp_ints_json(&ev))
} else {
return;
};
MP_SESSION.with(|s| {
if let Some(sess) = s.borrow().as_ref() {
match &sess.role {
MpRole::Host { peers } => {
for (_, _, p) in peers {
if p.is_open() {
let _ = p.send_game(json.as_bytes());
}
}
}
MpRole::Joiner { peer } => {
if peer.is_open() {
let _ = peer.send_game(json.as_bytes());
}
}
MpRole::Mesh { peers, .. } => {
for slot in peers.iter().flatten() {
if slot.1.is_open() {
let _ = slot.1.send_game(json.as_bytes());
}
}
}
}
}
});
}
fn mp_ints_json(v: &[i32]) -> String {
let mut s = String::from("[");
for (i, n) in v.iter().enumerate() {
if i > 0 {
s.push(',');
}
s.push_str(&n.to_string());
}
s.push(']');
s
}
fn mp_dispatch_peer_frame(worker: &Worker, peer_index: i32, bytes: &[u8]) {
let text = match std::str::from_utf8(bytes) {
Ok(t) => t,
Err(_) => return,
};
let v: serde_json::Value = match serde_json::from_str(text) {
Ok(v) => v,
Err(_) => return,
};
let trust_p_tag = MP_SESSION
.with(|s| matches!(s.borrow().as_ref().map(|x| &x.role), Some(MpRole::Joiner { .. })));
let origin = if trust_p_tag {
v.get("p")
.and_then(|x| x.as_i64())
.map(|x| x as i32)
.unwrap_or(peer_index)
} else {
peer_index
};
let m = Object::new();
let _ = Reflect::set(&m, &JsValue::from_str("type"), &JsValue::from_str("mp:peer"));
let _ = Reflect::set(&m, &JsValue::from_str("peer"), &JsValue::from_f64(origin as f64));
if let Some(d) = v.get("d").and_then(|x| x.as_array()) {
let arr = js_sys::Array::new();
for n in d {
arr.push(&JsValue::from_f64(n.as_i64().unwrap_or(0) as f64));
}
let _ = Reflect::set(&m, &JsValue::from_str("deltas"), &arr);
}
if let Some(ev) = v.get("e").and_then(|x| x.as_array()) {
let arr = js_sys::Array::new();
for n in ev {
arr.push(&JsValue::from_f64(n.as_i64().unwrap_or(0) as f64));
}
let _ = Reflect::set(&m, &JsValue::from_str("events"), &arr);
}
let _ = worker.post_message(&m);
if peer_index >= 1 && v.get("p").is_none() && v.is_object() {
mp_relay_from_host(peer_index, &v);
}
}
fn mp_relay_from_host(origin_idx: i32, v: &serde_json::Value) {
MP_SESSION.with(|s| {
if let Some(MpSession { role: MpRole::Host { peers }, .. }) = s.borrow().as_ref() {
let mut tagged = v.clone();
tagged["p"] = serde_json::json!(origin_idx);
let bytes = tagged.to_string();
for (_, pidx, p) in peers.iter() {
if *pidx != origin_idx && p.is_open() {
let _ = p.send_game(bytes.as_bytes());
}
}
}
});
}
pub(crate) fn mp_teardown() {
MP_SESSION.with(|s| *s.borrow_mut() = None);
}