- NodeAddr abstraction replacing SocketAddr across distribution crate (~30 files) - Core wasm32 compatibility: web-time, cfg-gated threads, wire encoding extraction - WebSocket transport (browser) with GatewayControl protocol - WebSocket gateway (native) with session routing and control protocol - BrowserRuntime API: JS actor support, connect/spawn/send/tick - Dashboard: extracted HTML to static files, gossip protocol panel with membership event log, visual fixes, click-to-filter, worker-colored rows Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
174 lines
5.4 KiB
Rust
174 lines
5.4 KiB
Rust
//! WebSocket gateway example.
|
|
//!
|
|
//! Starts a native node with some demo actors and a WebSocket gateway.
|
|
//! Browser clients connect via WebSocket to send/receive messages.
|
|
//!
|
|
//! ```bash
|
|
//! cargo run --example ws_gateway --features transport
|
|
//! ```
|
|
//!
|
|
//! Then open `crates/wasm/www/index.html` in a browser.
|
|
|
|
use std::sync::Arc;
|
|
use std::time::Duration;
|
|
|
|
use swactor::actor::{ActorAddress, ActorInterface};
|
|
use swactor::runtime::{Ctx, Runtime, RuntimeConfig};
|
|
use swactor::transport::{
|
|
Codec, CodecRegistry, NetworkMessage, TransportRouter,
|
|
};
|
|
use swactor::Error;
|
|
|
|
use swactor_gateway::WsGateway;
|
|
|
|
// ─── Messages ────────────────────────────────────────────────────────────────
|
|
|
|
/// A simple echo request: the payload + who to reply to.
|
|
#[derive(Clone, Debug)]
|
|
struct EchoRequest {
|
|
payload: String,
|
|
reply_to: ActorAddress,
|
|
}
|
|
|
|
impl NetworkMessage for EchoRequest {
|
|
fn type_tag() -> &'static str {
|
|
"example::EchoRequest"
|
|
}
|
|
}
|
|
|
|
/// Echo response — just the echoed payload.
|
|
#[derive(Clone, Debug)]
|
|
struct EchoResponse {
|
|
payload: String,
|
|
}
|
|
|
|
impl NetworkMessage for EchoResponse {
|
|
fn type_tag() -> &'static str {
|
|
"example::EchoResponse"
|
|
}
|
|
}
|
|
|
|
// ─── JSON codec for the demo messages ────────────────────────────────────────
|
|
|
|
struct JsonCodec;
|
|
|
|
impl Codec<EchoRequest> for JsonCodec {
|
|
fn encode(&self, msg: &EchoRequest) -> Result<Vec<u8>, Error> {
|
|
let mut buf = Vec::new();
|
|
buf.extend_from_slice(&msg.reply_to.0);
|
|
buf.extend_from_slice(msg.payload.as_bytes());
|
|
Ok(buf)
|
|
}
|
|
fn decode(&self, bytes: &[u8]) -> Result<EchoRequest, Error> {
|
|
if bytes.len() < 32 {
|
|
return Err(Error::from("EchoRequest too short"));
|
|
}
|
|
let mut addr = [0u8; 32];
|
|
addr.copy_from_slice(&bytes[..32]);
|
|
let payload = String::from_utf8_lossy(&bytes[32..]).to_string();
|
|
Ok(EchoRequest {
|
|
payload,
|
|
reply_to: ActorAddress(addr),
|
|
})
|
|
}
|
|
}
|
|
|
|
impl Codec<EchoResponse> for JsonCodec {
|
|
fn encode(&self, msg: &EchoResponse) -> Result<Vec<u8>, Error> {
|
|
Ok(msg.payload.as_bytes().to_vec())
|
|
}
|
|
fn decode(&self, bytes: &[u8]) -> Result<EchoResponse, Error> {
|
|
Ok(EchoResponse {
|
|
payload: String::from_utf8_lossy(bytes).to_string(),
|
|
})
|
|
}
|
|
}
|
|
|
|
// Also register the JsMessage type that the browser uses
|
|
#[derive(Clone, Debug)]
|
|
struct JsMessage {
|
|
type_tag: String,
|
|
payload: String,
|
|
}
|
|
|
|
impl NetworkMessage for JsMessage {
|
|
fn type_tag() -> &'static str {
|
|
"swactor::JsMessage"
|
|
}
|
|
}
|
|
|
|
struct JsMessageCodec;
|
|
|
|
impl Codec<JsMessage> for JsMessageCodec {
|
|
fn encode(&self, msg: &JsMessage) -> Result<Vec<u8>, Error> {
|
|
Ok(msg.payload.as_bytes().to_vec())
|
|
}
|
|
fn decode(&self, bytes: &[u8]) -> Result<JsMessage, Error> {
|
|
Ok(JsMessage {
|
|
type_tag: "js".to_string(),
|
|
payload: String::from_utf8_lossy(bytes).to_string(),
|
|
})
|
|
}
|
|
}
|
|
|
|
// ─── Echo actor ──────────────────────────────────────────────────────────────
|
|
|
|
struct EchoActor;
|
|
|
|
impl ActorInterface for EchoActor {
|
|
type Incoming = JsMessage;
|
|
type Response = ();
|
|
|
|
fn handle(&mut self, ctx: &Ctx, msg: JsMessage) {
|
|
println!("[echo] Received: type_tag={}, payload={}", msg.type_tag, msg.payload);
|
|
// Echo back to sender — in a real app, the payload would contain
|
|
// the reply_to address. For the demo, we just log it.
|
|
let _ = msg;
|
|
let _ = ctx;
|
|
}
|
|
}
|
|
|
|
// ─── Main ────────────────────────────────────────────────────────────────────
|
|
|
|
const WS_ADDR: &str = "127.0.0.1:9000";
|
|
|
|
fn main() {
|
|
println!("[gateway] Starting WebSocket gateway on ws://{WS_ADDR}");
|
|
println!("[gateway] Open crates/wasm/www/index.html in a browser to connect.\n");
|
|
|
|
// Build codec registry
|
|
let mut codecs = CodecRegistry::new();
|
|
codecs.register::<EchoRequest, _>(JsonCodec);
|
|
codecs.register::<EchoResponse, _>(JsonCodec);
|
|
codecs.register::<JsMessage, _>(JsMessageCodec);
|
|
let codecs = Arc::new(codecs);
|
|
|
|
// Build runtime with an echo actor
|
|
let transport_router = Arc::new(TransportRouter::new());
|
|
let mut rt = Runtime::new(RuntimeConfig {
|
|
num_threads: 1,
|
|
..RuntimeConfig::default()
|
|
});
|
|
rt.set_codec_registry(codecs.clone());
|
|
rt.set_transport_router(transport_router.clone());
|
|
|
|
let echo_addr = rt.spawn(EchoActor).unwrap();
|
|
rt.tick(); // drain spawn queue
|
|
println!("[gateway] EchoActor spawned at {echo_addr}");
|
|
|
|
// Start WebSocket gateway
|
|
let mut gateway = WsGateway::bind(WS_ADDR.parse().unwrap())
|
|
.expect("failed to bind WebSocket gateway");
|
|
println!("[gateway] Listening on ws://{}", gateway.local_addr());
|
|
println!();
|
|
|
|
// Event loop
|
|
loop {
|
|
let n = gateway.tick(&rt, &codecs, &transport_router);
|
|
if n > 0 {
|
|
println!("[gateway] Processed {n} envelope(s)");
|
|
}
|
|
rt.tick();
|
|
std::thread::sleep(Duration::from_millis(16)); // ~60 ticks/sec
|
|
}
|
|
}
|