use std::cell::{Cell, RefCell};
use std::collections::{BTreeMap, HashMap, HashSet};
use std::rc::Rc;
use std::sync::Arc;
use std::time::Duration;
use anyhow::Result;
use js_sys::{Function, Promise};
use wasm_bindgen::JsValue;
use wasm_bindgen_futures::{spawn_local, JsFuture};
use web_time::Instant;
use super::{
decode_wire, encode_wire, mesh_error, mesh_owner_error, AckKind, TicketInvite,
TicketMeshErrorCode, TicketMeshMember, TicketMeshOptions, TicketMeshRoster, Wire,
ADMISSION_LEASE, ADMISSION_REFRESH, GUEST_GRACE, HELLO_RETRY, INVITE_REFRESH,
INVITE_REFRESH_RETRY, MAX_RETRY, ROSTER_RESEND,
};
use crate::session_token::now_unix_ms;
use crate::Client;
const CLOSE_NOTICE_TIMEOUT_MS: u32 = 2_000;
const CONTROL_ACK_TIMEOUT_MS: u32 = 3_000;
const CONTROL_NOTICE_RETRY_MS: u32 = 500;
const CONNECT_ACTION_TIMEOUT_MS: u32 = 35_000;
enum Role {
Issuer {
members: BTreeMap<String, TicketMeshMember>,
unseen_since: HashMap<String, Instant>,
revision: u64,
last_roster_sent: Option<Instant>,
},
Guest {
roster: TicketMeshRoster,
local_ticket: String,
last_hello_sent: Option<Instant>,
local_ticket_issued_at: Instant,
last_ticket_refresh_attempt: Option<Instant>,
},
}
struct Retry {
next: Instant,
delay: Duration,
}
struct Inbound {
connection_id: String,
payload: Vec<u8>,
}
struct State {
role: Role,
inbound: Vec<Inbound>,
retries: HashMap<String, Retry>,
in_flight: HashSet<String>,
acks: HashSet<(String, AckKind)>,
invite_issued_at: Instant,
last_invite_refresh_attempt: Option<Instant>,
admission_revision: u64,
last_admission_update: Instant,
last_allowed_nodes: Vec<String>,
}
enum Action {
Connect { node_id: String, ticket: String },
Send { node_id: String, frame: String },
Invite { id: String, invitation: String },
}
impl Action {
fn value(self) -> serde_json::Value {
match self {
Self::Connect { node_id, ticket } => serde_json::json!({
"type": "connect", "nodeId": node_id, "ticket": ticket,
"timeoutMs": 30_000,
}),
Self::Send { node_id, frame } => serde_json::json!({
"type": "send", "nodeId": node_id, "frame": frame,
}),
Self::Invite { id, invitation } => serde_json::json!({
"type": "invite", "id": id, "invitation": invitation,
}),
}
}
}
pub(crate) struct WasmTicketMesh {
client: Arc<Client>,
invite: Rc<RefCell<TicketInvite>>,
closed: Rc<Cell<bool>>,
closing: Rc<Cell<bool>>,
state: Rc<RefCell<State>>,
handler: Function,
}
impl WasmTicketMesh {
pub async fn issue(
client: Arc<Client>,
id: &str,
options: TicketMeshOptions,
handler: Function,
) -> Result<Self> {
let invite = client.issue_ticket_invite(id, options).await?;
Ok(Self::start(
client,
invite,
Role::Issuer {
members: BTreeMap::new(),
unseen_since: HashMap::new(),
revision: 1,
last_roster_sent: None,
},
handler,
))
}
pub async fn join(client: Arc<Client>, invitation: &str, handler: Function) -> Result<Self> {
let invite = TicketInvite::parse(invitation)?;
if !invite.is_mesh() {
return Err(mesh_error(
TicketMeshErrorCode::InvalidInvitation,
"ticket invitation is not a mesh",
));
}
let local_node = client.current_node_id().await.ok_or_else(|| {
mesh_error(
TicketMeshErrorCode::Unavailable,
"ticket mesh requires an initialized endpoint",
)
})?;
let roster = TicketMeshRoster::new(&invite, &local_node)?;
let scope = format!("ticket:{}", invite.id());
client
.session_token_registry
.set_ticket_peer_limit(&scope, invite.max_peers())
.map_err(mesh_owner_error)?;
let local_ticket = match async {
client
.session_token_registry
.require_scope_peer_admission(&scope)
.map_err(anyhow::Error::msg)?;
client
.session_token_registry
.update_scope_peer_admission(
&scope,
1,
now_unix_ms() + ADMISSION_LEASE.as_millis() as u64,
roster.allowed_nodes(),
)
.map_err(anyhow::Error::msg)?;
let ticket = client.endpoint_ticket_with_token(&scope, 0).await?;
TicketInvite::new(
invite.id(),
&ticket,
TicketMeshOptions {
max_peers: invite.max_peers(),
},
)?;
Ok::<_, anyhow::Error>(ticket)
}
.await
{
Ok(ticket) => ticket,
Err(error) => {
client
.session_token_registry
.release_ticket_peer_admission(&scope);
client
.session_token_registry
.release_ticket_peer_limit(&scope);
return Err(error);
}
};
Ok(Self::start(
client,
invite,
Role::Guest {
roster,
local_ticket,
last_hello_sent: None,
local_ticket_issued_at: Instant::now(),
last_ticket_refresh_attempt: None,
},
handler,
))
}
fn start(client: Arc<Client>, invite: TicketInvite, role: Role, handler: Function) -> Self {
let issuer = invite.issuer_node().expect("validated ticket invitation");
let state = Rc::new(RefCell::new(State {
last_allowed_nodes: if matches!(role, Role::Guest { .. }) {
vec![issuer]
} else {
Vec::new()
},
role,
inbound: Vec::new(),
retries: HashMap::new(),
in_flight: HashSet::new(),
acks: HashSet::new(),
invite_issued_at: Instant::now(),
last_invite_refresh_attempt: None,
admission_revision: 1,
last_admission_update: Instant::now(),
}));
let closed = Rc::new(Cell::new(false));
let closing = Rc::new(Cell::new(false));
let current_invite = Rc::new(RefCell::new(invite));
let actor = Self {
client: client.clone(),
invite: current_invite.clone(),
closed: closed.clone(),
closing: closing.clone(),
state: state.clone(),
handler: handler.clone(),
};
spawn_local(async move {
while !closed.get() {
Self::tick(
&client,
¤t_invite,
&closed,
&closing,
&state,
&handler,
)
.await;
gloo_timers::future::sleep(Duration::from_secs(1)).await;
}
});
actor
}
pub fn invite(&self) -> String {
self.invite.borrow().encode()
}
pub fn issuer_node(&self) -> Result<String> {
self.invite.borrow().issuer_node()
}
pub fn receive(&self, connection_id: String, payload: Vec<u8>) {
if self.closed.get() || decode_wire(&payload).is_none() {
return;
}
let mut state = self.state.borrow_mut();
if state.inbound.len() < 64 {
state.inbound.push(Inbound {
connection_id,
payload,
});
}
}
pub async fn peers(&self) -> Vec<crate::client::PeerSessionSnapshot> {
if self.closed.get() {
return Vec::new();
}
let scope = format!("ticket:{}", self.invite.borrow().id());
self.client
.peer_sessions()
.await
.into_iter()
.filter(|peer| {
peer.settled_ready
&& !peer.logical_session_terminal
&& peer.scopes.iter().any(|candidate| candidate == &scope)
&& peer.node_id.is_some()
})
.collect()
}
pub async fn close(&self) {
if self.closed.get() {
return;
}
if self.closing.replace(true) {
while self.closing.get() {
gloo_timers::future::sleep(Duration::from_millis(25)).await;
}
return;
}
let (nodes, message, expected_ack) = {
let state = self.state.borrow();
match &state.role {
Role::Issuer { members, .. } => (
members.keys().cloned().collect::<Vec<_>>(),
Wire::Closed {
id: self.invite.borrow().id().to_string(),
},
AckKind::Close,
),
Role::Guest { .. } => (
self.issuer_node().into_iter().collect(),
Wire::Leave {
id: self.invite.borrow().id().to_string(),
},
AckKind::Leave,
),
}
};
let frame = String::from_utf8(encode_wire(&message)).expect("mesh wire is UTF-8");
let started = Instant::now();
while !self.closed.get()
&& started.elapsed() < Duration::from_millis(CONTROL_ACK_TIMEOUT_MS.into())
{
let pending = nodes
.iter()
.filter(|node| {
!self
.state
.borrow()
.acks
.contains(&((*node).clone(), expected_ack))
})
.cloned()
.collect::<Vec<_>>();
if pending.is_empty() {
break;
}
let notices = pending.into_iter().map(|node_id| {
run_action(
&self.handler,
Action::Send {
node_id,
frame: frame.clone(),
},
CONTROL_NOTICE_RETRY_MS,
)
});
futures::future::join_all(notices).await;
gloo_timers::future::sleep(Duration::from_millis(25)).await;
}
let scope = format!("ticket:{}", self.invite.borrow().id());
self.client.revoke_tokens_by_scope(&scope).await;
self.client
.session_token_registry
.release_ticket_peer_admission(&scope);
self.client
.session_token_registry
.release_ticket_peer_limit(&scope);
self.closed.set(true);
self.closing.set(false);
}
async fn tick(
client: &Arc<Client>,
shared_invite: &Rc<RefCell<TicketInvite>>,
closed: &Rc<Cell<bool>>,
closing: &Rc<Cell<bool>>,
state: &Rc<RefCell<State>>,
handler: &Function,
) {
let mut invite = shared_invite.borrow().clone();
let mut invite_changed = false;
let inbound = std::mem::take(&mut state.borrow_mut().inbound);
let scope = format!("ticket:{}", invite.id());
for event in inbound {
let Some(message) = decode_wire(&event.payload) else {
continue;
};
let Some(peer) = client.peer_session(&event.connection_id).await else {
continue;
};
if peer.active_connection_id.as_deref() != Some(event.connection_id.as_str())
|| !peer.settled_ready
|| peer.logical_session_terminal
|| !peer.scopes.iter().any(|candidate| candidate == &scope)
{
continue;
}
let Some(source) = peer.node_id else { continue };
let mut renewed = None;
let mut terminal = false;
let mut reply = None;
let mut observed_ack = None;
let mut state = state.borrow_mut();
match (&mut state.role, message) {
(
Role::Issuer {
members,
revision,
last_roster_sent,
..
},
Wire::Hello { id, ticket },
) if id == invite.id() => {
let valid = TicketInvite::new(
invite.id(),
&ticket,
TicketMeshOptions {
max_peers: invite.max_peers(),
},
)
.and_then(|grant| grant.issuer_node())
.is_ok_and(|node| node == source);
if !valid
|| (!members.contains_key(&source)
&& members.len() + 1 >= invite.max_peers() as usize)
{
continue;
}
if members.get(&source).is_none_or(|old| old.ticket != ticket) {
members.insert(
source.clone(),
TicketMeshMember {
node_id: source.clone(),
ticket,
},
);
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
}
(
Role::Issuer {
members,
unseen_since,
revision,
last_roster_sent,
},
Wire::Leave { id },
) if id == invite.id() => {
if members.remove(&source).is_some() {
unseen_since.remove(&source);
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
reply = Some(AckKind::Leave);
}
(
Role::Guest { roster, .. },
Wire::Roster {
id,
revision,
members,
issuer_invite,
},
) if id == invite.id() => {
let next = if issuer_invite == invite.encode() {
None
} else {
match invite.renewed(&issuer_invite) {
Ok(next) => Some(next),
Err(_) => continue,
}
};
if roster.accept(&source, revision, members).unwrap_or(false) {
renewed = next;
}
}
(Role::Guest { .. }, Wire::Closed { id })
if id == invite.id()
&& invite.issuer_node().ok().as_deref() == Some(source.as_str()) =>
{
terminal = true;
reply = Some(AckKind::Close);
}
(_, Wire::Ack { id, kind }) if id == invite.id() => {
observed_ack = Some((source.clone(), kind));
}
_ => {}
}
if let Some(ack) = observed_ack {
state.acks.insert(ack);
}
drop(state);
if let Some(kind) = reply {
let frame = String::from_utf8(encode_wire(&Wire::Ack {
id: invite.id().to_string(),
kind,
}))
.expect("mesh wire is UTF-8");
run_action(
handler,
Action::Send {
node_id: source,
frame,
},
CLOSE_NOTICE_TIMEOUT_MS,
)
.await;
}
if terminal {
closing.set(true);
client.revoke_tokens_by_scope(&scope).await;
client
.session_token_registry
.release_ticket_peer_admission(&scope);
client
.session_token_registry
.release_ticket_peer_limit(&scope);
closed.set(true);
closing.set(false);
return;
}
if let Some(next) = renewed {
*shared_invite.borrow_mut() = next.clone();
invite = next;
invite_changed = true;
}
}
if closing.get() {
return;
}
let peers: HashSet<String> = client
.peer_sessions()
.await
.into_iter()
.filter(|peer| {
peer.settled_ready
&& !peer.logical_session_terminal
&& peer.scopes.iter().any(|candidate| candidate == &scope)
})
.filter_map(|peer| peer.node_id)
.collect();
let now = Instant::now();
let refresh_issuer = {
let state = state.borrow();
matches!(state.role, Role::Issuer { .. })
&& now.duration_since(state.invite_issued_at) >= INVITE_REFRESH
&& state
.last_invite_refresh_attempt
.is_none_or(|attempt| now.duration_since(attempt) >= INVITE_REFRESH_RETRY)
};
if refresh_issuer {
state.borrow_mut().last_invite_refresh_attempt = Some(now);
if let Ok(ticket) = client.endpoint_ticket_with_token(&scope, 0).await {
if let Ok(next) = TicketInvite::new(
invite.id(),
&ticket,
TicketMeshOptions {
max_peers: invite.max_peers(),
},
)
.and_then(|next| invite.renewed(&next.encode()))
{
*shared_invite.borrow_mut() = next.clone();
invite = next;
invite_changed = true;
let mut state = state.borrow_mut();
state.invite_issued_at = now;
if let Role::Issuer {
revision,
last_roster_sent,
..
} = &mut state.role
{
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
}
}
}
let refresh_guest = {
let state = state.borrow();
matches!(&state.role,
Role::Guest { local_ticket_issued_at, last_ticket_refresh_attempt, .. }
if now.duration_since(*local_ticket_issued_at) >= INVITE_REFRESH
&& last_ticket_refresh_attempt.is_none_or(|attempt| now.duration_since(attempt) >= INVITE_REFRESH_RETRY)
)
};
if refresh_guest {
if let Role::Guest {
last_ticket_refresh_attempt,
..
} = &mut state.borrow_mut().role
{
*last_ticket_refresh_attempt = Some(now);
}
if let Ok(ticket) = client.endpoint_ticket_with_token(&scope, 0).await {
let local_node = client.current_node_id().await;
let valid = TicketInvite::new(
invite.id(),
&ticket,
TicketMeshOptions {
max_peers: invite.max_peers(),
},
)
.and_then(|grant| grant.issuer_node())
.is_ok_and(|node| local_node.as_deref() == Some(node.as_str()));
if valid {
if let Role::Guest {
local_ticket,
last_hello_sent,
local_ticket_issued_at,
..
} = &mut state.borrow_mut().role
{
*local_ticket = ticket;
*last_hello_sent = None;
*local_ticket_issued_at = now;
}
}
}
}
let mut actions = Vec::new();
if invite_changed {
actions.push(Action::Invite {
id: invite.id().to_string(),
invitation: invite.encode(),
});
}
{
let mut state = state.borrow_mut();
let mut wanted = Vec::new();
match &mut state.role {
Role::Issuer {
members,
unseen_since,
revision,
last_roster_sent,
} => {
let stale: Vec<_> = members
.keys()
.filter_map(|node| {
if peers.contains(node) {
unseen_since.remove(node);
None
} else if now
.duration_since(*unseen_since.entry(node.clone()).or_insert(now))
>= GUEST_GRACE
{
Some(node.clone())
} else {
None
}
})
.collect();
for node in stale {
members.remove(&node);
unseen_since.remove(&node);
*revision = revision.saturating_add(1);
*last_roster_sent = None;
}
if last_roster_sent.is_none_or(|sent| now.duration_since(sent) >= ROSTER_RESEND)
{
let frame = String::from_utf8(encode_wire(&Wire::Roster {
id: invite.id().into(),
revision: *revision,
members: members.values().cloned().collect(),
issuer_invite: invite.encode(),
}))
.expect("mesh wire is UTF-8");
for node in members.keys() {
actions.push(Action::Send {
node_id: node.clone(),
frame: frame.clone(),
});
}
*last_roster_sent = Some(now);
}
}
Role::Guest {
roster,
local_ticket,
last_hello_sent,
..
} => {
let issuer = invite.issuer_node().expect("validated ticket invitation");
wanted.push((issuer.clone(), invite.endpoint_ticket().to_string()));
wanted.extend(
roster
.desired()
.into_iter()
.map(|member| (member.node_id, member.ticket)),
);
if peers.contains(&issuer)
&& !roster.contains_local()
&& last_hello_sent
.is_none_or(|sent| now.duration_since(sent) >= HELLO_RETRY)
{
actions.push(Action::Send {
node_id: issuer,
frame: String::from_utf8(encode_wire(&Wire::Hello {
id: invite.id().into(),
ticket: local_ticket.clone(),
}))
.expect("mesh wire is UTF-8"),
});
*last_hello_sent = Some(now);
}
}
}
if let Role::Guest { roster, .. } = &state.role {
let allowed = roster.allowed_nodes();
if allowed != state.last_allowed_nodes
|| now.duration_since(state.last_admission_update) >= ADMISSION_REFRESH
{
state.admission_revision = state.admission_revision.saturating_add(1);
if client
.session_token_registry
.update_scope_peer_admission(
&scope,
state.admission_revision,
now_unix_ms() + ADMISSION_LEASE.as_millis() as u64,
allowed.clone(),
)
.is_err()
{
return;
}
state.last_allowed_nodes = allowed;
state.last_admission_update = now;
}
}
state
.retries
.retain(|node, _| wanted.iter().any(|(desired, _)| desired == node));
for (node_id, ticket) in wanted {
if peers.contains(&node_id) {
state.retries.remove(&node_id);
continue;
}
if state.in_flight.contains(&node_id) {
continue;
}
if state
.retries
.get(&node_id)
.is_some_and(|retry| now < retry.next)
{
continue;
}
let delay = state
.retries
.get(&node_id)
.map_or(Duration::from_millis(500), |retry| {
(retry.delay * 2).min(MAX_RETRY)
});
state.retries.insert(
node_id.clone(),
Retry {
next: now + delay,
delay,
},
);
state.in_flight.insert(node_id.clone());
actions.push(Action::Connect { node_id, ticket });
}
}
for action in actions {
let connecting = match &action {
Action::Connect { node_id, .. } => Some(node_id.clone()),
Action::Send { .. } | Action::Invite { .. } => None,
};
let value = match serde::Serialize::serialize(
&action.value(),
&serde_wasm_bindgen::Serializer::json_compatible(),
) {
Ok(value) => value,
Err(_) => {
if let Some(node_id) = connecting {
state.borrow_mut().in_flight.remove(&node_id);
}
continue;
}
};
let result = match handler.call1(&JsValue::UNDEFINED, &value) {
Ok(result) => result,
Err(_) => {
if let Some(node_id) = connecting {
state.borrow_mut().in_flight.remove(&node_id);
}
continue;
}
};
let timeout_ms = if connecting.is_some() {
CONNECT_ACTION_TIMEOUT_MS
} else {
CLOSE_NOTICE_TIMEOUT_MS
};
let state = state.clone();
spawn_local(async move {
let pending = JsFuture::from(Promise::resolve(&result));
let deadline = gloo_timers::future::TimeoutFuture::new(timeout_ms);
futures::pin_mut!(pending, deadline);
let _ = futures::future::select(pending, deadline).await;
if let Some(node_id) = connecting {
state.borrow_mut().in_flight.remove(&node_id);
}
});
}
}
}
impl Drop for WasmTicketMesh {
fn drop(&mut self) {
self.closing.set(false);
self.closed.set(true);
}
}
async fn run_action(handler: &Function, action: Action, timeout_ms: u32) {
let Ok(value) = serde::Serialize::serialize(
&action.value(),
&serde_wasm_bindgen::Serializer::json_compatible(),
) else {
return;
};
let Ok(result) = handler.call1(&JsValue::UNDEFINED, &value) else {
return;
};
let pending = JsFuture::from(Promise::resolve(&result));
let deadline = gloo_timers::future::TimeoutFuture::new(timeout_ms);
futures::pin_mut!(pending, deadline);
let _ = futures::future::select(pending, deadline).await;
}