Compare commits

..

2 Commits

Author SHA1 Message Date
3fbae1bd95 v1.1.0 — Gzip canvas compression + throttled sends
All checks were successful
CI / ci (push) Successful in 6m29s
- Canvas snapshots gzip-compressed before sending: 1MB raw -> ~1-5KB
  for a mostly-white canvas, ~100-500KB for a busy one
- Semaphore throttles concurrent canvas sends (default 20) to prevent
  memory spikes during connection bursts
- Frontend decompresses gzip via browser DecompressionStream API
- New config: MAX_CONCURRENT_CANVAS_SENDS (default 20)
- Added flate2 dependency for gzip encoding
2026-08-18 02:28:47 +02:00
31f0fce7a9 Switch default transport to WebSocket
- Dockerfile builds with --features websocket and sets VITE_TRANSPORT=websocket
- Server default feature changed from socketio to websocket
- Binary canvas transfer (1MB) instead of Socket.IO JSON (~3-4MB per connection)
- Significantly reduces memory pressure under concurrent connections
2026-08-18 02:09:09 +02:00
11 changed files with 118 additions and 23 deletions

62
Cargo.lock generated
View File

@@ -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",

View File

@@ -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"

View File

@@ -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

View File

@@ -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)?,

View File

@@ -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 }

View File

@@ -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(

View File

@@ -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()

View File

@@ -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)]

View File

@@ -4,7 +4,7 @@ version.workspace = true
edition.workspace = true
[features]
default = ["socketio"]
default = ["websocket"]
socketio = ["dep:socketio", "dep:socketioxide"]
websocket = ["dep:websocket"]

View File

@@ -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))
}

View File

@@ -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;
}