use alloc::boxed::Box;
use alloc::collections::VecDeque;
use alloc::format;
use alloc::rc::Rc;
use alloc::string::String;
use alloc::vec::Vec;
use core::cell::RefCell;
use wasm_bindgen::JsCast;
use wasm_bindgen::closure::Closure;
use web_sys::{BinaryType, CloseEvent, Event, MessageEvent, WebSocket};
use super::super::core::{
CommandRefusal, DriverOutput, DriverPhase, ResponseExpectation, SocketCommand, SocketEvent,
SocketFailure, WebSocketFrameDriver,
};
use super::mirror::{BrowserMessageData, BrowserSocketAction, action_for_command};
pub type AdapterSignalSink = Box<dyn FnMut(AdapterSignal)>;
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum AdapterSignal {
Output(DriverOutput),
Fault(AdapterFault),
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum AdapterFault {
SeamReentered,
CommandExecutionFailed {
description: String,
},
}
#[derive(Clone, Debug, PartialEq, Eq, thiserror::Error)]
pub enum BrowserSocketError {
#[error("driver refused the command: {refusal:?}")]
CommandRefused {
refusal: CommandRefusal,
},
#[error("browser socket failure: {description}")]
Socket {
description: String,
},
#[error("single-threaded driver seam re-entered")]
SeamReentered,
#[error("driver contract violation: {description}")]
DriverContract {
description: String,
},
}
struct Seam {
driver: WebSocketFrameDriver,
last_event_detail: Option<String>,
}
struct SignalPort {
queue: RefCell<VecDeque<AdapterSignal>>,
sink: RefCell<AdapterSignalSink>,
}
impl SignalPort {
fn emit(&self, signal: AdapterSignal) {
match self.queue.try_borrow_mut() {
Ok(mut queue) => queue.push_back(signal),
Err(_) => return,
}
self.drain();
}
fn drain(&self) {
loop {
let Ok(mut sink) = self.sink.try_borrow_mut() else {
return;
};
let next = match self.queue.try_borrow_mut() {
Ok(mut queue) => queue.pop_front(),
Err(_) => return,
};
let Some(signal) = next else {
return;
};
sink(signal);
}
}
}
pub struct WebSysWebSocketSocket {
socket: WebSocket,
seam: Rc<RefCell<Seam>>,
port: Rc<SignalPort>,
_on_open: Closure<dyn FnMut(Event)>,
_on_message: Closure<dyn FnMut(MessageEvent)>,
_on_close: Closure<dyn FnMut(CloseEvent)>,
_on_error: Closure<dyn FnMut(Event)>,
}
impl WebSysWebSocketSocket {
pub fn open(url: &str, sink: AdapterSignalSink) -> Result<Self, BrowserSocketError> {
let mut driver = WebSocketFrameDriver::new();
let command = driver
.command_open()
.map_err(|refusal| BrowserSocketError::CommandRefused { refusal })?;
match command {
SocketCommand::Open => {}
SocketCommand::SendBinary(_) | SocketCommand::Close => {
return Err(BrowserSocketError::DriverContract {
description: format!("command_open emitted {command:?} instead of Open"),
});
}
}
let socket = WebSocket::new(url).map_err(|value| BrowserSocketError::Socket {
description: format!("browser refused websocket construction for {url}: {value:?}"),
})?;
socket.set_binary_type(BinaryType::Arraybuffer);
let seam = Rc::new(RefCell::new(Seam {
driver,
last_event_detail: None,
}));
let port = Rc::new(SignalPort {
queue: RefCell::new(VecDeque::new()),
sink: RefCell::new(sink),
});
let on_open = {
let seam = Rc::clone(&seam);
let port = Rc::clone(&port);
let socket = socket.clone();
Closure::wrap(Box::new(move |_event: Event| {
let negotiated = socket.extensions();
let detail = if negotiated.is_empty() {
None
} else {
Some(format!(
"browser negotiated websocket extensions in violation of the \
extension-free contract: {negotiated}"
))
};
let event = super::mirror::open_event(&negotiated);
dispatch(&seam, &socket, &port, event, detail);
}) as Box<dyn FnMut(Event)>)
};
let on_message = {
let seam = Rc::clone(&seam);
let port = Rc::clone(&port);
let socket = socket.clone();
Closure::wrap(Box::new(move |event: MessageEvent| {
let (data, detail) = classify_message_data(&event.data());
let event = super::mirror::message_event(data);
dispatch(&seam, &socket, &port, event, detail);
}) as Box<dyn FnMut(MessageEvent)>)
};
let on_close = {
let seam = Rc::clone(&seam);
let port = Rc::clone(&port);
let socket = socket.clone();
Closure::wrap(Box::new(move |event: CloseEvent| {
let detail = Some(format!(
"browser close event: code {}, wasClean {}, reason {:?}",
event.code(),
event.was_clean(),
event.reason(),
));
let event = super::mirror::close_event(event.was_clean());
dispatch(&seam, &socket, &port, event, detail);
}) as Box<dyn FnMut(CloseEvent)>)
};
let on_error = {
let seam = Rc::clone(&seam);
let port = Rc::clone(&port);
let socket = socket.clone();
Closure::wrap(Box::new(move |_event: Event| {
let detail = Some(String::from(
"browser websocket error event (opaque by specification)",
));
dispatch(&seam, &socket, &port, super::mirror::error_event(), detail);
}) as Box<dyn FnMut(Event)>)
};
socket.set_onopen(Some(on_open.as_ref().unchecked_ref()));
socket.set_onmessage(Some(on_message.as_ref().unchecked_ref()));
socket.set_onclose(Some(on_close.as_ref().unchecked_ref()));
socket.set_onerror(Some(on_error.as_ref().unchecked_ref()));
Ok(Self {
socket,
seam,
port,
_on_open: on_open,
_on_message: on_message,
_on_close: on_close,
_on_error: on_error,
})
}
pub fn send_binary(
&self,
bytes: Vec<u8>,
expectation: ResponseExpectation,
) -> Result<(), BrowserSocketError> {
let command = {
let mut seam = self
.seam
.try_borrow_mut()
.map_err(|_| BrowserSocketError::SeamReentered)?;
seam.driver
.command_send(bytes, expectation)
.map_err(|refusal| BrowserSocketError::CommandRefused { refusal })?
};
match action_for_command(command) {
Ok(BrowserSocketAction::SendBinary(bytes)) => {
if let Err(value) = self.socket.send_with_u8_array(&bytes) {
let description = format!("browser websocket send failed: {value:?}");
dispatch(
&self.seam,
&self.socket,
&self.port,
SocketEvent::Failed(SocketFailure::Transport),
Some(description.clone()),
);
return Err(BrowserSocketError::Socket { description });
}
Ok(())
}
Ok(action @ BrowserSocketAction::Close) => Err(BrowserSocketError::DriverContract {
description: format!("command_send mapped to a non-send action: {action:?}"),
}),
Err(refusal) => Err(BrowserSocketError::DriverContract {
description: format!("command_send emitted a non-send command: {refusal}"),
}),
}
}
pub fn close(&self) -> Result<(), BrowserSocketError> {
let command = {
let mut seam = self
.seam
.try_borrow_mut()
.map_err(|_| BrowserSocketError::SeamReentered)?;
seam.driver
.command_close()
.map_err(|refusal| BrowserSocketError::CommandRefused { refusal })?
};
match action_for_command(command) {
Ok(BrowserSocketAction::Close) => {
self.socket
.close()
.map_err(|value| BrowserSocketError::Socket {
description: format!("browser websocket close failed: {value:?}"),
})
}
Ok(action @ BrowserSocketAction::SendBinary(_)) => {
Err(BrowserSocketError::DriverContract {
description: format!("command_close mapped to a non-close action: {action:?}"),
})
}
Err(refusal) => Err(BrowserSocketError::DriverContract {
description: format!("command_close emitted a non-close command: {refusal}"),
}),
}
}
pub fn phase(&self) -> Result<DriverPhase, BrowserSocketError> {
self.seam
.try_borrow()
.map(|seam| seam.driver.phase())
.map_err(|_| BrowserSocketError::SeamReentered)
}
pub fn last_event_detail(&self) -> Result<Option<String>, BrowserSocketError> {
self.seam
.try_borrow()
.map(|seam| seam.last_event_detail.clone())
.map_err(|_| BrowserSocketError::SeamReentered)
}
}
impl Drop for WebSysWebSocketSocket {
fn drop(&mut self) {
self.socket.set_onopen(None);
self.socket.set_onmessage(None);
self.socket.set_onclose(None);
self.socket.set_onerror(None);
let live = !matches!(
self.seam.try_borrow().map(|seam| seam.driver.phase()),
Ok(DriverPhase::Terminated)
);
if live {
if let Err(value) = self.socket.close() {
self.port
.emit(AdapterSignal::Fault(AdapterFault::CommandExecutionFailed {
description: format!("browser websocket close on drop failed: {value:?}"),
}));
}
}
}
}
impl core::fmt::Debug for WebSysWebSocketSocket {
fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
formatter
.debug_struct("WebSysWebSocketSocket")
.finish_non_exhaustive()
}
}
fn classify_message_data(data: &wasm_bindgen::JsValue) -> (BrowserMessageData, Option<String>) {
if let Some(buffer) = data.dyn_ref::<js_sys::ArrayBuffer>() {
let bytes = js_sys::Uint8Array::new(buffer).to_vec();
return (BrowserMessageData::ArrayBuffer(bytes), None);
}
if data.as_string().is_some() {
let detail = String::from("peer sent a text message on the binary-only liminal route");
return (BrowserMessageData::Text, Some(detail));
}
if data.dyn_ref::<web_sys::Blob>().is_some() {
let detail =
String::from("browser delivered a Blob despite the pinned arraybuffer binaryType");
return (BrowserMessageData::Blob, Some(detail));
}
(
BrowserMessageData::Unrecognized,
Some(String::from("browser delivered unrecognized message data")),
)
}
fn dispatch(
seam: &Rc<RefCell<Seam>>,
socket: &WebSocket,
port: &Rc<SignalPort>,
event: SocketEvent,
detail: Option<String>,
) {
let step = if let Ok(mut seam) = seam.try_borrow_mut() {
if let Some(detail) = detail {
seam.last_event_detail = Some(detail);
}
seam.driver.handle_event(event)
} else {
port.emit(AdapterSignal::Fault(AdapterFault::SeamReentered));
return;
};
if let Some(command) = step.command {
match action_for_command(command) {
Ok(BrowserSocketAction::Close) => {
if let Err(value) = socket.close() {
port.emit(AdapterSignal::Fault(AdapterFault::CommandExecutionFailed {
description: format!("browser websocket close failed: {value:?}"),
}));
}
}
Ok(BrowserSocketAction::SendBinary(_)) => {
port.emit(AdapterSignal::Fault(AdapterFault::CommandExecutionFailed {
description: String::from(
"driver emitted SendBinary from handle_event, outside its contract",
),
}));
}
Err(refusal) => {
port.emit(AdapterSignal::Fault(AdapterFault::CommandExecutionFailed {
description: format!(
"driver emitted Open from handle_event, outside its contract: \
{refusal}"
),
}));
}
}
}
port.emit(AdapterSignal::Output(step.output));
}