Compare commits
2 Commits
f652785acb
...
3fbae1bd95
| Author | SHA1 | Date | |
|---|---|---|---|
| 3fbae1bd95 | |||
| 31f0fce7a9 |
62
Cargo.lock
generated
62
Cargo.lock
generated
@@ -2,6 +2,12 @@
|
||||
# It is not intended for manual editing.
|
||||
version = 4
|
||||
|
||||
[[package]]
|
||||
name = "adler2"
|
||||
version = "2.0.1"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "320119579fcad9c21884f5c4861d16174d0e06250625266f50fe6898340abefa"
|
||||
|
||||
[[package]]
|
||||
name = "allocator-api2"
|
||||
version = "0.2.21"
|
||||
@@ -10,7 +16,7 @@ checksum = "683d7910e743518b0e34f1186f92494becacb047c7b6bf616c96772180fef923"
|
||||
|
||||
[[package]]
|
||||
name = "api-types"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"domain",
|
||||
"serde",
|
||||
@@ -18,7 +24,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "application"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"domain",
|
||||
"thiserror 2.0.20",
|
||||
@@ -160,7 +166,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "canvas-file"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"config",
|
||||
"domain",
|
||||
@@ -186,14 +192,14 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "config"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"thiserror 2.0.20",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "config-env"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"config",
|
||||
]
|
||||
@@ -222,6 +228,15 @@ dependencies = [
|
||||
"libc",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "crc32fast"
|
||||
version = "1.5.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "9481c1c90cbf2ac953f07c8d4a58aa3945c425b7185c9154d67a65e4230da511"
|
||||
dependencies = [
|
||||
"cfg-if",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "crossbeam-utils"
|
||||
version = "0.8.19"
|
||||
@@ -322,7 +337,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "domain"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"thiserror 2.0.20",
|
||||
]
|
||||
@@ -385,6 +400,16 @@ version = "1.0.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "877a4ace8713b0bcf2a4e7eec82529c029f1d0619886d18145fea96c3ffe5c0f"
|
||||
|
||||
[[package]]
|
||||
name = "flate2"
|
||||
version = "1.1.9"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "843fba2746e448b37e26a819579957415c8cef339bf08564fe8b7ddbd959573c"
|
||||
dependencies = [
|
||||
"crc32fast",
|
||||
"miniz_oxide",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "fnv"
|
||||
version = "1.0.7"
|
||||
@@ -634,7 +659,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "http-axum"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"axum",
|
||||
"config",
|
||||
@@ -841,6 +866,16 @@ dependencies = [
|
||||
"unicase",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "miniz_oxide"
|
||||
version = "0.8.9"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "1fa76a2c86f704bdb222d66965fb3d63269ce38518b83cb0575fca855ebb6316"
|
||||
dependencies = [
|
||||
"adler2",
|
||||
"simd-adler32",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "mio"
|
||||
version = "1.2.2"
|
||||
@@ -1208,7 +1243,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "server"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"application",
|
||||
"axum",
|
||||
@@ -1275,6 +1310,12 @@ dependencies = [
|
||||
"libc",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "simd-adler32"
|
||||
version = "0.3.10"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "3a219298ac11a56ea9a6d2120044824d6f01aeb034955e7af7bc16858527deea"
|
||||
|
||||
[[package]]
|
||||
name = "slab"
|
||||
version = "0.4.9"
|
||||
@@ -1305,7 +1346,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "socketio"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"api-types",
|
||||
"application",
|
||||
@@ -1823,11 +1864,12 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "websocket"
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
dependencies = [
|
||||
"application",
|
||||
"axum",
|
||||
"domain",
|
||||
"flate2",
|
||||
"futures",
|
||||
"serde",
|
||||
"serde_json",
|
||||
|
||||
@@ -24,7 +24,7 @@ default-members = [
|
||||
resolver = "2"
|
||||
|
||||
[workspace.package]
|
||||
version = "1.0.0"
|
||||
version = "1.1.0"
|
||||
edition = "2024"
|
||||
|
||||
[workspace.dependencies]
|
||||
@@ -37,6 +37,7 @@ config-env = { path = "crates/adapters/config-env" }
|
||||
http-axum = { path = "crates/adapters/http-axum" }
|
||||
socketio = { path = "crates/adapters/socketio" }
|
||||
websocket = { path = "crates/adapters/websocket" }
|
||||
flate2 = "1"
|
||||
futures = "0.3"
|
||||
|
||||
axum = "0.8.9"
|
||||
|
||||
@@ -3,6 +3,7 @@ WORKDIR /app/painter-js
|
||||
COPY painter-js/package.json painter-js/bun.lock ./
|
||||
RUN bun install --frozen-lockfile
|
||||
COPY painter-js/ .
|
||||
ENV VITE_TRANSPORT=websocket
|
||||
RUN bun run build
|
||||
|
||||
FROM rust:1-alpine AS builder
|
||||
@@ -13,10 +14,11 @@ COPY Cargo.toml Cargo.lock ./
|
||||
COPY crates/ ./crates/
|
||||
|
||||
RUN mkdir -p painter-js/dist && echo '<html></html>' > painter-js/dist/index.html
|
||||
RUN cargo build --release --bin server 2>&1 || true
|
||||
RUN cargo build --release --bin server --no-default-features --features websocket 2>&1 || true
|
||||
|
||||
COPY --from=frontend /app/painter-js/dist ./painter-js/dist
|
||||
RUN touch crates/adapters/http-axum/src/routes.rs && cargo build --release --bin server
|
||||
RUN touch crates/adapters/http-axum/src/routes.rs && \
|
||||
cargo build --release --bin server --no-default-features --features websocket
|
||||
|
||||
FROM scratch
|
||||
COPY --from=builder /app/target/release/server /painter
|
||||
|
||||
@@ -14,6 +14,7 @@ const DEFAULT_BROADCAST_CAPACITY: usize = 1024;
|
||||
const DEFAULT_SNAPSHOT_INTERVAL_SECS: u64 = 300;
|
||||
const DEFAULT_SNAPSHOT_MAX: usize = 5;
|
||||
const DEFAULT_SNAPSHOT_DIR: &str = "snapshots/";
|
||||
const DEFAULT_MAX_CONCURRENT_CANVAS_SENDS: usize = 20;
|
||||
|
||||
pub struct EnvConfigSource;
|
||||
|
||||
@@ -24,6 +25,10 @@ impl ConfigSource for EnvConfigSource {
|
||||
address: env_or("ADDRESS", DEFAULT_ADDRESS),
|
||||
port: parse_env("PORT", DEFAULT_PORT)?,
|
||||
enable_cors: parse_bool_env("ENABLE_CORS", true),
|
||||
max_concurrent_canvas_sends: parse_env(
|
||||
"MAX_CONCURRENT_CANVAS_SENDS",
|
||||
DEFAULT_MAX_CONCURRENT_CANVAS_SENDS,
|
||||
)?,
|
||||
},
|
||||
canvas: CanvasConfig {
|
||||
width: parse_env("CANVAS_WIDTH", DEFAULT_CANVAS_WIDTH)?,
|
||||
|
||||
@@ -7,6 +7,7 @@ edition.workspace = true
|
||||
application = { workspace = true }
|
||||
domain = { workspace = true }
|
||||
axum = { workspace = true, features = ["ws"] }
|
||||
flate2 = { workspace = true }
|
||||
futures = { workspace = true }
|
||||
serde = { workspace = true }
|
||||
serde_json = { workspace = true }
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
use std::io::Write;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::Ordering;
|
||||
|
||||
@@ -5,8 +6,10 @@ use axum::extract::State;
|
||||
use axum::extract::ws::{Message, WebSocket, WebSocketUpgrade};
|
||||
use axum::response::IntoResponse;
|
||||
use domain::{BroadcastEvent, BroadcastSubscription, Color, Position};
|
||||
use flate2::Compression;
|
||||
use flate2::write::GzEncoder;
|
||||
use futures::{SinkExt, StreamExt, stream::SplitSink};
|
||||
use tokio::sync::mpsc;
|
||||
use tokio::sync::{Semaphore, mpsc};
|
||||
use tracing::{error, info};
|
||||
|
||||
use crate::WsState;
|
||||
@@ -36,7 +39,7 @@ async fn handle_connection(socket: WebSocket, state: Arc<WsState>) {
|
||||
// Subscribe before snapshotting to avoid missing updates
|
||||
let subscription = state.app_state.broadcaster().subscribe();
|
||||
|
||||
if !send_canvas_snapshot(&mut sender, &state.app_state).await {
|
||||
if !send_canvas_snapshot(&mut sender, &state.app_state, &state.canvas_send_semaphore).await {
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -67,10 +70,39 @@ async fn handle_connection(socket: WebSocket, state: Arc<WsState>) {
|
||||
application::soldiers::disconnect::execute(&state.app_state, &connection_id);
|
||||
}
|
||||
|
||||
async fn send_canvas_snapshot(sender: &mut WsSender, state: &AppState) -> bool {
|
||||
async fn send_canvas_snapshot(
|
||||
sender: &mut WsSender,
|
||||
state: &AppState,
|
||||
semaphore: &Semaphore,
|
||||
) -> bool {
|
||||
let Ok(_permit) = semaphore.acquire().await else {
|
||||
return false;
|
||||
};
|
||||
|
||||
let pixels = application::canvas::get_state::execute(state);
|
||||
let bytes = Color::collect_as_bytes(&pixels);
|
||||
sender.send(Message::Binary(bytes.into())).await.is_ok()
|
||||
let raw_bytes = Color::collect_as_bytes(&pixels);
|
||||
|
||||
let Some(compressed) = gzip_compress(&raw_bytes) else {
|
||||
error!("Failed to compress canvas snapshot");
|
||||
return false;
|
||||
};
|
||||
|
||||
info!(
|
||||
"Sending canvas snapshot: {}KB raw -> {}KB gzip",
|
||||
raw_bytes.len() / 1024,
|
||||
compressed.len() / 1024
|
||||
);
|
||||
|
||||
sender
|
||||
.send(Message::Binary(compressed.into()))
|
||||
.await
|
||||
.is_ok()
|
||||
}
|
||||
|
||||
fn gzip_compress(data: &[u8]) -> Option<Vec<u8>> {
|
||||
let mut encoder = GzEncoder::new(Vec::new(), Compression::fast());
|
||||
encoder.write_all(data).ok()?;
|
||||
encoder.finish().ok()
|
||||
}
|
||||
|
||||
async fn run_send_loop(
|
||||
|
||||
@@ -6,16 +6,19 @@ use std::sync::atomic::AtomicU64;
|
||||
|
||||
use application::AppState;
|
||||
use axum::{Router, routing::get};
|
||||
use tokio::sync::Semaphore;
|
||||
|
||||
pub(crate) struct WsState {
|
||||
app_state: Arc<AppState>,
|
||||
connection_counter: AtomicU64,
|
||||
canvas_send_semaphore: Arc<Semaphore>,
|
||||
}
|
||||
|
||||
pub fn build_router(state: Arc<AppState>) -> Router {
|
||||
pub fn build_router(state: Arc<AppState>, max_concurrent_canvas_sends: usize) -> Router {
|
||||
let ws_state = Arc::new(WsState {
|
||||
app_state: state,
|
||||
connection_counter: AtomicU64::new(0),
|
||||
canvas_send_semaphore: Arc::new(Semaphore::new(max_concurrent_canvas_sends)),
|
||||
});
|
||||
|
||||
Router::new()
|
||||
|
||||
@@ -24,6 +24,7 @@ pub struct ServerConfig {
|
||||
pub address: String,
|
||||
pub port: u16,
|
||||
pub enable_cors: bool,
|
||||
pub max_concurrent_canvas_sends: usize,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
|
||||
@@ -4,7 +4,7 @@ version.workspace = true
|
||||
edition.workspace = true
|
||||
|
||||
[features]
|
||||
default = ["socketio"]
|
||||
default = ["websocket"]
|
||||
socketio = ["dep:socketio", "dep:socketioxide"]
|
||||
websocket = ["dep:websocket"]
|
||||
|
||||
|
||||
@@ -127,7 +127,7 @@ fn build_app(
|
||||
state: Arc<AppState>,
|
||||
config: &AppConfig,
|
||||
) -> Result<axum::Router, Box<dyn std::error::Error>> {
|
||||
let ws_router = websocket::build_router(state);
|
||||
let ws_router = websocket::build_router(state, config.server.max_concurrent_canvas_sends);
|
||||
let http_router = http_axum::build_router(config.server.enable_cors, &config.rate_limit)?;
|
||||
Ok(ws_router.merge(http_router))
|
||||
}
|
||||
|
||||
@@ -13,6 +13,13 @@ const createSocketIoTransport = () => {
|
||||
};
|
||||
};
|
||||
|
||||
const decompressGzip = async (buffer) => {
|
||||
const stream = new Blob([buffer])
|
||||
.stream()
|
||||
.pipeThrough(new DecompressionStream("gzip"));
|
||||
return new Response(stream).arrayBuffer();
|
||||
};
|
||||
|
||||
const createWebSocketTransport = () => {
|
||||
const handlers = {};
|
||||
const wsUrl = isDebug
|
||||
@@ -38,9 +45,10 @@ const createWebSocketTransport = () => {
|
||||
|
||||
ws.addEventListener("open", () => dispatch("connect"));
|
||||
|
||||
ws.addEventListener("message", (event) => {
|
||||
ws.addEventListener("message", async (event) => {
|
||||
if (event.data instanceof ArrayBuffer) {
|
||||
const pixels = new Uint32Array(event.data);
|
||||
const decompressed = await decompressGzip(event.data);
|
||||
const pixels = new Uint32Array(decompressed);
|
||||
dispatch("canvas_state", Array.from(pixels));
|
||||
return;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user