WebSocket Endpoint
สิ่งที่จะสร้าง
หัวข้อที่มีชื่อว่า “สิ่งที่จะสร้าง”realtime/handlers.rs: ws_upgrade handler ที่อยู่เบื้องหลัง GET /ws/boards/:id?token=<jwt> และ handle_socket ฟังก์ชันที่ ws_upgrade ส่งต่อ connection ที่รับเข้ามาแต่ละอันไปให้ทำงาน ตัว ws_upgrade รันชุดตรวจสอบ signature, expiry และ Redis allowlist ชุดเดียวกับที่ AuthUser รันบนทุก REST request บวกกับการตรวจสอบสมาชิกภาพของบอร์ด ตั้งแต่ก่อน upgrade connection ด้วยซ้ำ — WebSocket ไปยังบอร์ดที่คุณอ่านไม่ได้จึงโดนปฏิเสธเข้มงวดพอ ๆ กับ GET /boards/:id
มีสองชิ้นเล็ก ๆ ที่มาด้วยกัน ทั้งคู่จำเป็นแค่เพื่อให้ ws_upgrade คอมไพล์ผ่าน: realtime/hub.rs มี struct Hub ที่เก็บ tokio::sync::broadcast::Sender<String> หนึ่งตัวต่อบอร์ดที่มี subscriber ที่ยังเชื่อมต่ออยู่อย่างน้อยหนึ่งราย และฟิลด์ใหม่หนึ่งฟิลด์บน AppState คือ hub: Arc<realtime::hub::Hub> เคียงข้างกับ db, redis, และ config เมื่อจบบทนี้ สองแท็บ browser ที่เชื่อมต่อ WebSocket ของบอร์ดเดียวกันจะเห็น event ของกันและกันได้แล้ว — ตราบใดที่ทั้งคู่คุยกับ backend process เดียวกัน redis-backplane จะเอาข้อจำกัดสุดท้ายนี้ออกไป
AuthUser extractor แบบ FromRequestParts จาก middleware นำมาใช้ซ้ำที่นี่ไม่ได้ เพราะ extractor ตัวนั้นออกแบบมาให้ Axum รันก่อน handler ที่ป้องกันไว้ ทุกตัว ในฐานะ extractor ของ HTTP request ปกติ แต่การ upgrade WebSocket คือ request แบบ GET ที่ไม่มีช่องให้ใส่ bearer token แบบที่ REST call มี: browser ที่สร้าง object WebSocket แบบเนทีฟไม่มี API สำหรับแนบ header ใด ๆ เข้ากับ handshake ws_upgrade จึงฝังชุดตรวจสอบสามอย่างเดียวกับที่ AuthUser รันไว้เอง — signature และ expiry ผ่าน jwt::verify แล้วตามด้วย Redis allowlist — โดยอ่าน token จาก query parameter แทนที่จะเป็น header
Query parameter ไม่ใช่ header และทำไมนี่ถึงเป็น trade-off จริง ตัวสร้าง WebSocket ของ browser คือ new WebSocket(url, protocols?) — ไม่มี argument สำหรับ header เลยแม้แต่น้อย เหลือทางเลือกที่ใช้งานได้จริงอยู่สองทางสำหรับยืนยันตัวตน WebSocket handshake จาก browser: พึ่งพา cookie (ซึ่ง TaskFlow ไม่ได้ใช้ — jwt เป็น bearer token โดยการออกแบบ ไม่ใช่ cookie-based session) หรือใส่ token ไว้ในที่ที่ URL พาไปได้ query parameter เป็นทางเลือกที่ง่ายที่สุดในบรรดาทางเลือกที่พา token ผ่าน URL (อีกทางคือเข้ารหัสไว้ใน header Sec-WebSocket-Protocol ผ่าน argument protocols ซึ่งใช้ได้จริงแต่เป็นการใช้ที่แปลกสำหรับค่าที่ไม่ใช่ชื่อ protocol จริง ๆ) trade-off นี้มีจริงและควรพูดตรง ๆ: token ใน URL อาจไปโผล่ใน access log ของเซิร์ฟเวอร์, ประวัติ browser, และ request log ของ reverse proxy ใด ๆ — ที่ที่ header ไม่เคยไปโผล่โดยดีฟอลต์ TaskFlow ยอมรับสิ่งนี้เพราะ token ตัวนี้คือ JWT อายุสั้น (24 ชั่วโมง) ที่เพิกถอนได้เป็นรายตัวตัวเดียวกับที่ jwt สร้างไว้แล้ว ไม่ใช่ credential ใหม่ที่อ่อนไหวกว่าเดิม ระบบจริงที่จัดการข้อมูลอ่อนไหวกว่านี้น่าจะสร้าง “connect ticket” แยกต่างหาก ใช้ครั้งเดียว อายุนาทีเดียว ทำขึ้นมาเฉพาะสำหรับ WS handshake แทนที่จะใช้ bearer token อเนกประสงค์ซ้ำ และจะทำให้แน่ใจว่า token ถูกตัดออกจากบรรทัด access log ก่อนเขียนลงไฟล์
Hub มีอยู่เพราะ broadcast::Sender ต้องอยู่ที่ไหนสักที่ที่ใช้ร่วมกันได้ระหว่าง WebSocket connection ทุกตัวที่ดูบอร์ดเดียวกัน — handle_socket รันหนึ่งครั้งต่อ socket ที่เชื่อมต่อเข้ามา ดังนั้น state ต่อ socket จึงเป็นที่ที่รายการ subscriber ของ บอร์ด หนึ่งอยู่ไม่ได้ Hub คือ registry ที่ใช้ร่วมกันนี้ หนึ่งรายการต่อบอร์ด สร้างแบบ lazy ตอน subscriber รายแรก แล้วใช้ซ้ำโดย subscriber ทุกรายถัดไป ตอนนี้ Hub รู้จักแค่ socket ที่เชื่อมต่ออยู่กับ process นี้ เท่านั้น redis-backplane คือสิ่งที่ป้อน event ที่มาจาก backend instance อีกตัวป้อนเข้ามา
ข้อดีข้อเสีย
หัวข้อที่มีชื่อว่า “ข้อดีข้อเสีย”broadcast::Sender<String> หนึ่งตัวต่อบอร์ด (สิ่งที่เราใช้) เทียบกับ broadcast::Sender<BoardEvent> ตัวเดียวทั่วทั้งแอปที่ทุก socket กรองเองฝั่งไคลเอนต์
- ข้อดี: socket ที่ดูบอร์ด A จะไม่มีวันได้รับ traffic ของบอร์ด B เลยด้วยซ้ำ —
Hub::subscribe(board_id)ส่งคืน receiver ที่ผูกไว้กับ channel ของบอร์ดเดียวเท่านั้น ดังนั้นจึงไม่มี filter แบบif event.board_id == my_board_idต่อข้อความที่รันอยู่ใน hot loop ของทุก socket และบอร์ดที่คึกคักจะไม่ทำให้ socket ที่ดูบอร์ดอื่นที่เงียบสงบเกิดLaggederror (แต่ละบอร์ดมี buffer ของ channel เป็นของตัวเองอย่างอิสระ) - ข้อเสีย:
HashMap<Uuid, Sender>ของHubต้องใช้ lock (std::sync::Mutexที่ถือไว้สั้น ๆ แค่ตอน lookup/insert ของHashMapไม่เคยข้าม.await) ซึ่ง sender ตัวเดียวทั่วโลกไม่ต้องการเลย สำหรับสเกลของ TaskFlow — บอร์ดที่เปิดพร้อมกันหลักสิบถึงหลักพันต้น ๆ ไม่ใช่หลักล้าน —std::sync::Mutexที่ถือสั้น ๆ รอบ ๆ การ lookup entry ของHashMapไม่ใช่ประเด็นเรื่อง contention ระบบที่มี broadcast topic ที่แตกต่างกันจำนวนมากกว่านี้มากอาจหันไปใช้ sharded map แทน แต่นั่นคือการแก้ปัญหาที่สเกลของคอร์สนี้ไม่มี
ลงมือสร้าง
หัวข้อที่มีชื่อว่า “ลงมือสร้าง”1. realtime/hub.rs
หัวข้อที่มีชื่อว่า “1. realtime/hub.rs”สร้าง taskflow/backend/api/src/realtime/hub.rs:
use std::collections::HashMap;use std::sync::Mutex;
use tokio::sync::broadcast;use uuid::Uuid;
/// How many buffered events a lagging subscriber can fall behind by before/// `broadcast::Receiver::recv` starts returning `Lagged` instead of the/// oldest unread message.const CHANNEL_CAPACITY: usize = 128;
/// The in-process registry of live WebSocket subscribers, one/// `broadcast::Sender` per board that currently has at least one socket/// watching it.pub struct Hub { channels: Mutex<HashMap<Uuid, broadcast::Sender<String>>>,}
impl Hub { pub fn new() -> Self { Self { channels: Mutex::new(HashMap::new()), } }
/// Registers a new subscriber for `board_id`, creating that board's /// broadcast channel on first use. pub fn subscribe(&self, board_id: Uuid) -> broadcast::Receiver<String> { let mut channels = self.channels.lock().unwrap(); let sender = channels .entry(board_id) .or_insert_with(|| broadcast::channel(CHANNEL_CAPACITY).0); sender.subscribe() }}
impl Default for Hub { fn default() -> Self { Self::new() }}or_insert_with จะรัน closure — จัดสรร broadcast::channel ใหม่ — ก็ต่อเมื่อ board_id ยังไม่ใช่ key ที่มีอยู่แล้วเท่านั้น subscriber ทุกรายหลังจากรายแรกของบอร์ดหนึ่ง ๆ จะใช้ Sender ตัวเดิมซ้ำ แล้วเรียก .subscribe() เพื่อรับ Receiver ของตัวเองที่เป็นอิสระ broadcast::channel(128).0 ทิ้งครึ่ง Receiver ที่ตัวสร้างส่งคืนมาด้วย — Hub ต้องการแค่แจก receiver ผ่าน .subscribe() เท่านั้น ไม่เคยต้องอ่านจาก receiver เอง
2. realtime/mod.rs
หัวข้อที่มีชื่อว่า “2. realtime/mod.rs”สร้าง taskflow/backend/api/src/realtime/mod.rs:
pub mod handlers;pub mod hub;
use axum::{routing::get, Router};use serde::{Deserialize, Serialize};use uuid::Uuid;
use crate::state::AppState;
/// The single message shape published to a board's Redis channel and/// forwarded verbatim, as a JSON string, to every WebSocket subscribed to/// that board. See `protocol` for the full event-type catalog.#[derive(Debug, Serialize, Deserialize)]pub struct BoardEvent { pub r#type: String, #[serde(rename = "boardId")] pub board_id: Uuid, pub payload: serde_json::Value,}
pub fn routes() -> Router<AppState> { Router::new().route("/ws/boards/:id", get(handlers::ws_upgrade))}นี่คือรูปร่างของ BoardEvent แบบเป๊ะ ๆ จาก protocol ตอนนี้ลงไปอยู่ในโค้ดจริงแล้ว แต่ยังไม่มีอะไรสร้าง instance ขึ้นมาเลย — คนที่จะสร้างคือ realtime::publish ใน redis-backplane — ดังนั้นให้คาดหวัง warning dead_code บน BoardEvent หลังจากบทนี้ เหมือนกับที่ jwt::issue/jwt::verify เคยไม่ถูกใช้อยู่หนึ่งบทตอนกลับไปที่ jwt
3. realtime/handlers.rs
หัวข้อที่มีชื่อว่า “3. realtime/handlers.rs”สร้าง taskflow/backend/api/src/realtime/handlers.rs:
use axum::{ extract::{ ws::{Message, WebSocket, WebSocketUpgrade}, Path, Query, State, }, response::Response,};use deadpool_redis::redis::AsyncCommands;use serde::Deserialize;use tokio::sync::broadcast;use uuid::Uuid;
use crate::{ auth::jwt, boards::service as boards_service, error::{AppError, AppResult}, state::AppState,};
#[derive(Debug, Deserialize)]pub struct WsAuthQuery { pub token: String,}
/// `GET /ws/boards/:id?token=<jwt>` — upgrades to a WebSocket after the/// same three checks `AuthUser` runs on every REST request (signature,/// expiry, Redis allowlist), plus board membership.pub async fn ws_upgrade( State(state): State<AppState>, Path(board_id): Path<Uuid>, Query(query): Query<WsAuthQuery>, ws: WebSocketUpgrade,) -> AppResult<Response> { let claims = jwt::verify(&query.token, &state.config.jwt_secret).map_err(|_| AppError::Unauthorized)?;
let mut conn = state .redis .get() .await .map_err(|err| AppError::Internal(err.into()))?;
let session: Option<String> = conn .get(format!("auth:token:{}", claims.jti)) .await .map_err(|err| AppError::Internal(err.into()))?;
if session.is_none() { return Err(AppError::Unauthorized); }
let user_id = claims .sub .parse::<Uuid>() .map_err(|_| AppError::Unauthorized)?;
boards_service::assert_member(&state.db, user_id, board_id).await?;
Ok(ws.on_upgrade(move |socket| handle_socket(socket, state, board_id)))}
async fn handle_socket(mut socket: WebSocket, state: AppState, board_id: Uuid) { let mut rx = state.hub.subscribe(board_id);
loop { tokio::select! { event = rx.recv() => { match event { Ok(json) => { if socket.send(Message::Text(json)).await.is_err() { return; } } Err(broadcast::error::RecvError::Lagged(_)) => { continue; } Err(broadcast::error::RecvError::Closed) => { return; } } } incoming = socket.recv() => { match incoming { Some(Ok(Message::Close(_))) | None => { return; } Some(Ok(_)) => { // Text/Binary/Ping/Pong from the client — this // socket is read-only from the client's point of // view, so there's nothing to act on except noting // the connection is still alive. } Some(Err(_)) => { return; } } } } }}ไล่ดู ws_upgrade ตามลำดับ — รูปร่างเหมือน AuthUser::from_request_parts จาก middleware บวกอีกหนึ่งขั้นตอน:
- ตรวจสอบ JWT จาก query parameter
tokenเหมือนกับที่AuthUserตรวจสอบ bearer token ใน headerAuthorization - ตรวจสอบ Redis allowlist (
auth:token:{jti}) — token ที่ถูกเพิกถอนแล้ว (logout หรือการกระทำของ admin ในอนาคต) ก็โดนปฏิเสธที่นี่เช่นกัน ไม่ใช่แค่ตอน REST call - แปลง
subเป็นUuid— id ของผู้ใช้ที่ยืนยันตัวตนแล้ว - ตรวจสอบสมาชิกภาพของบอร์ด ด้วย
boards_service::assert_memberฟังก์ชันเดียวกับที่ REST handler ทุกตัวใน rest-api เรียกอยู่แล้ว token ที่ถูกต้องและยังไม่ถูกเพิกถอนของผู้ใช้ที่ไม่ได้อยู่ในบอร์ดนี้ยังคงถูกปฏิเสธ - Upgrade ส่ง
WebSocketที่สร้างเสร็จแล้วไปให้handle_socketพร้อมกับ clone ของstate(ต้นทุนถูก —AppStateเป็นArc-backed ทั้งหมด) และboard_id
tokio::select! ของ handle_socket แข่ง future สองตัวในทุกรอบของ loop:
rx.recv()— event ถัดไปที่Hub::subscribeจะส่งให้บอร์ดนี้ เมื่อรับสำเร็จก็ส่งต่อไปยังไคลเอนต์เป็นMessage::TextRecvError::Lagged(socket นี้ตามหลังเกินCHANNEL_CAPACITYข้อความ) ไม่ใช่ ข้อผิดพลาดร้ายแรง —continueจะดึงข้อความถัดไปที่มีต่อ; socket ที่ตามหลัง event ที่มาถี่ยิบ 128 ครั้งบนบอร์ดที่คึกคักมากจะพลาด state ระหว่างทางไปบ้าง แต่ยังทำงานต่อได้ และ reconnect จะอธิบายว่าทำไมนี่ถึงเป็นช่องว่างที่ยอมรับได้และเยียวยาตัวเองได้ ไม่ใช่บั๊กที่ต้องกำจัดRecvError::Closed(ทุกSenderของ channel นี้ถูก drop ไปแล้ว) จะจบ loop — สิ่งนี้เกิดขึ้นไม่ได้ในการออกแบบปัจจุบันของโมดูลนี้ เพราะHubเองถือSenderไว้ตลอด แต่ match ก็ต้องครอบคลุมทุกกรณีอยู่ดีsocket.recv()— อะไรก็ตามที่ไคลเอนต์ส่งมา frameCloseหรือ connection ที่ปิดไปแล้ว (None) จะจบ loop อย่างอื่น (Text,Binary,Ping,Pong) ตอนนี้ไม่ทำอะไรเลย: socket นี้ push-only จากฝั่งเซิร์ฟเวอร์ ไคลเอนต์จึงไม่มีอะไรต้องบอกกลับมา — แค่ได้รับ อะไรก็ได้ ก็เพียงพอที่จะรู้ว่า connection ยังมีชีวิตอยู่ ข้อผิดพลาดระดับ transport จะจบ loop
การ return จาก handle_socket — จาก branch ไหนก็ตาม — จะ drop rx ลดจำนวน receiver ของ Sender ของบอร์ดนั้นลงหนึ่ง ไม่มีขั้นตอนทำความสะอาดแยกต่างหากให้เขียนเพิ่ม: broadcast::Receiver ที่ถูก drop คือการยกเลิกการลงทะเบียนทั้งหมด
4. state.rs — เพิ่ม hub
หัวข้อที่มีชื่อว่า “4. state.rs — เพิ่ม hub”อัปเดต taskflow/backend/api/src/state.rs:
#[derive(Clone)]pub struct AppState { pub db: sqlx::PgPool, pub redis: deadpool_redis::Pool, pub config: std::sync::Arc<crate::config::Config>, pub hub: std::sync::Arc<crate::realtime::hub::Hub>,}5. เชื่อมเข้ากับ main.rs
หัวข้อที่มีชื่อว่า “5. เชื่อมเข้ากับ main.rs”mod auth;mod boards;mod cache;mod cards;mod columns;mod config;mod db;mod error;mod labels;mod middleware;mod realtime;mod state;
use std::net::SocketAddr;use std::sync::Arc;
use axum::{middleware::from_fn_with_state, routing::get, Json, Router};use config::Config;use middleware::rate_limit::rate_limit;use state::AppState;use tower_http::cors::CorsLayer;use tracing_subscriber::EnvFilter;
#[tokio::main]async fn main() -> anyhow::Result<()> { tracing_subscriber::fmt() .with_env_filter(EnvFilter::from_default_env()) .init();
let config = Config::from_env()?;
let db = db::create_pg_pool(&config.database_url).await?; let redis = db::create_redis_pool(&config.redis_url)?;
let state = AppState { db, redis, config: Arc::new(config), hub: Arc::new(realtime::hub::Hub::new()), };
let cors = CorsLayer::new().allow_origin( state .config .frontend_origin .parse::<axum::http::HeaderValue>() .expect("FRONTEND_ORIGIN must be a valid header value"), );
let app_port = state.config.app_port;
let app = Router::new() .route("/health", get(health)) .nest( "/auth", auth::routes().route_layer(from_fn_with_state(state.clone(), rate_limit)), ) .merge(boards::routes()) .merge(columns::routes()) .merge(cards::routes()) .merge(labels::routes()) .merge(realtime::routes()) .layer(cors) .with_state(state);
let listener = tokio::net::TcpListener::bind(format!("0.0.0.0:{app_port}")).await?; tracing::info!("listening on {}", listener.local_addr()?);
axum::serve( listener, app.into_make_service_with_connect_info::<SocketAddr>(), ) .await?;
Ok(())}
async fn health() -> Json<serde_json::Value> { Json(serde_json::json!({ "status": "ok" }))}สามอย่างที่เปลี่ยนจากเวอร์ชันของ rate-limit: mod realtime; เป็น module declaration ใหม่, literal ของ AppState ได้ฟิลด์ hub: Arc::new(realtime::hub::Hub::new()) เพิ่มเข้ามา และ .merge(realtime::routes()) เพิ่ม route /ws/boards/:id ใหม่เข้าไปใน Router เดียวกับที่ resource module อื่น ๆ ทุกตัว merge เข้าไป
ตรวจสอบผล
หัวข้อที่มีชื่อว่า “ตรวจสอบผล”cargo check -p apiเปิดสแตกขึ้นมา ขอ token และบอร์ดมา (ใช้รูปแบบเดิมจากโมดูลก่อนหน้า) แล้วเชื่อมต่อด้วย WebSocket client — websocat เป็น CLI ที่สะดวกตัวหนึ่ง:
cd taskflow/infra && docker compose up -d db rediscd ../backend && cargo run -p api &
TOKEN=$(curl -s -X POST http://localhost:8080/auth/register \ -H "Content-Type: application/json" \ -d '{"email":"ada@example.com","password":"correct horse battery staple","display_name":"Ada"}' \ | python3 -c 'import json,sys; print(json.load(sys.stdin)["token"])')
BOARD_ID=$(curl -s -X POST http://localhost:8080/boards \ -H "Authorization: Bearer $TOKEN" -H "Content-Type: application/json" \ -d '{"title":"Sprint 12"}' | python3 -c 'import json,sys; print(json.load(sys.stdin)["id"])')
websocat "ws://localhost:8080/ws/boards/$BOARD_ID?token=$TOKEN"connection ยังคงเปิดอยู่ — token ที่ถูกต้องสำหรับบอร์ดที่คุณเป็นสมาชิกจะ upgrade สำเร็จแล้วก็แค่รอเฉย ๆ เพราะยังไม่มีอะไร publish event เลย ตรวจสอบเส้นทางที่ถูกปฏิเสธด้วย ซึ่งแต่ละแบบจะปิด handshake แทนที่จะ upgrade:
# Bad tokenwebsocat "ws://localhost:8080/ws/boards/$BOARD_ID?token=not-a-real-token"
# Valid token, board you're not a member ofOTHER_BOARD=$(curl -s -X POST http://localhost:8080/boards \ -H "Authorization: Bearer $OTHER_TOKEN" -H "Content-Type: application/json" \ -d '{"title":"Not yours"}' | python3 -c 'import json,sys; print(json.load(sys.stdin)["id"])')
websocat "ws://localhost:8080/ws/boards/$OTHER_BOARD?token=$TOKEN"ทั้งสองแบบปิดทันทีโดยไม่มีการ upgrade สำเร็จ ด้วยสองเทอร์มินัลที่รันคำสั่ง websocat แรกกับ $BOARD_ID เดียวกันทั้งคู่ การปิดตัวใดตัวหนึ่ง (Ctrl-C) จะทำให้อีกตัวยังเชื่อมต่ออยู่และไม่ได้รับผลกระทบ — ยืนยันว่า Hub ติดตาม subscriber แยกกันเป็นรายตัว ไม่ใช่แบบรวมทั้งบอร์ด
คุณได้สร้าง registry Hub ใน realtime/hub.rs — broadcast::Sender<String> หนึ่งตัวต่อบอร์ด สร้างแบบ lazy — และ handler ws_upgrade ใน realtime/handlers.rs ซึ่งรันชุดตรวจสอบ signature/expiry/allowlist เดียวกับที่ AuthUser รันบน REST call ทุกอัน โดยอ่าน token จาก query parameter แทนที่จะเป็น header เพราะตัวสร้าง WebSocket ของ browser ตั้ง header ไม่ได้ บวกกับการตรวจสอบสมาชิกภาพของบอร์ดก่อนที่จะ upgrade เลยด้วยซ้ำ tokio::select! loop ของ handle_socket แข่ง event ที่เข้ามาทาง broadcast กับ frame ที่เข้ามาจากไคลเอนต์ ส่งต่อฝ่ายแรกและเฝ้าดูฝ่ายหลังแค่สัญญาณ disconnect เท่านั้น สอง socket บน backend process เดียวกัน ที่ดูบอร์ดเดียวกัน สามารถถึงกันได้ผ่าน Hub แล้ว — แต่ยังไม่มีอะไร publish event เข้าไปเลย และ backend process ตัวที่สองก็จะมี Hub เป็นของตัวเองแยกต่างหากโดยสิ้นเชิง ไม่มีทางได้ยิน traffic ของตัวแรก redis-backplane จะสร้างทั้งสองอย่าง: การเรียก publish ที่ mutation ทุกตัวเรียก และ background task ที่ป้อนข้อมูลจาก Redis ที่ทำให้ instance ของ Hub บนต่าง process เห็นตรงกัน