ข้ามไปยังเนื้อหา

Redis Backplane

สองส่วนที่เพิ่มเข้าไปใน realtime/hub.rsHub::route เมธอดส่วนตัวที่ป้อนข้อความเข้าไปยัง local channel ของบอร์ดที่ถูกต้อง และ run_subscriber ฟังก์ชันที่ spawn ครั้งเดียวตอน startup ซึ่งถือ Redis connection เฉพาะหนึ่งเส้น PSUBSCRIBE ไปยัง board:* แล้วเรียก route สำหรับทุกข้อความที่เข้ามา หนึ่งส่วนที่เพิ่มเข้าไปใน realtime/mod.rspublish ฟังก์ชันที่ mutation ทุกตัวเรียกเพื่อกระจายสิ่งที่เพิ่งทำไป และ dependency ใหม่หนึ่งตัว futures-util ที่จำเป็นสำหรับ Stream ที่ subscriber loop ในบทนี้ต้องใช้อ่าน

จากนั้นคือผลตอบแทน: realtime::publish ถูกเชื่อมเข้ากับ move_card — แทนที่คอมเมนต์ // Module 7 wires realtime broadcast of card.moved here ที่ move-reorder ทิ้งไว้เมื่อสองโมดูลก่อน — และเข้ากับ card/column mutation อื่น ๆ ทุกตัวที่มี event type ตรงกันในรายการของ protocol

ws-endpoint ทิ้งช่องว่างหนึ่งไว้โดยตั้งใจ: Hub รู้จักแค่ socket ที่เชื่อมต่อกับ process เดียวกัน ที่ Hub ตัวนั้นรันอยู่เท่านั้น ลอง cargo run -p api สองครั้งบนสองพอร์ตต่างกัน (หรือในโปรดักชัน รันสอง instance หลัง load balancer) แล้วคุณจะได้ Hub สองตัวที่เป็นอิสระต่อกันโดยสิ้นเชิง แต่ละตัวมี HashMap ของบอร์ดและ sender เป็นของตัวเอง ไม่แชร์อะไรกันเลย การย้ายการ์ดที่ instance A จัดการไม่มีทางไปถึง socket ที่บังเอิญเชื่อมต่อกับ instance B ได้ ทั้งที่ทั้งคู่กำลังดูบอร์ดเดียวกันเป๊ะ

นี่คือเหตุผลที่ backplane มีอยู่: บางอย่างที่ทั้งสอง instance แชร์กันอยู่แล้ว ที่ไม่ใช่หน่วยความจำใน process TaskFlow มีสิ่งที่ backend instance ทุกตัวเชื่อมต่ออยู่อย่างเป็นอิสระจากกันอยู่แล้วพอดีหนึ่งอย่าง — Redis — และ Redis pub/sub ถูกสร้างมาเพื่อปัญหารูปแบบนี้พอดี: PUBLISH จากไคลเอนต์ตัวไหนก็ตามจะไปถึงไคลเอนต์ที่ subscribe อยู่ทุกตัวในขณะนั้น โดยไม่รู้เลยว่าฝ่ายไหนเป็น process ไหน publish คือฝั่ง “instance ไหนก็ได้ request ไหนก็ได้” ของสมการนี้ run_subscriber คือฝั่ง “ทุก instance ตลอดเวลา” รวมกันแล้ว การ์ดที่ถูกย้ายโดย request ที่ instance A บังเอิญจัดการ จะไปถึง socket ที่เชื่อมต่อกับ instance B ในแบบเดียวกันเป๊ะกับที่ไปถึง socket ของ instance A เอง — run_subscriber ไม่ได้แยกกรณีพิเศษว่า “event นี้มาจาก process ของฉันเองหรือเปล่า” เพราะจากมุมมองของ Redis ไม่มีสิ่งที่เรียกว่า PUBLISH “ของตัวเอง” เลย — subscriber ทุกตัว รวมถึง instance ที่ publish เอง จะได้รับทุกข้อความบน pattern ที่สมัครไว้

PSUBSCRIBE board:* ตัวเดียว ไม่ใช่ SUBSCRIBE board:{id} แยกต่อบอร์ดที่กำลังใช้งานอยู่ เป็นอีกทางเลือกที่จงใจเลือกไว้ตรงนี้ ทุกบอร์ดที่เคยมี socket เชื่อมต่อเข้ามายัง instance ไหนก็ตาม อยู่ในความครอบคลุมของ pattern subscription เดียวนี้ตั้งแต่ process เริ่มทำงาน โดยไม่ต้องทำอะไรเพิ่ม — ไม่มีบัญชี “subscriber รายแรกของบอร์ดหนึ่งกระตุ้น Redis SUBSCRIBE ใหม่” ให้ต้องทำให้ถูกต้อง และไม่มีความเสี่ยงที่ socket จะเชื่อมต่อเข้ามาก่อนที่ SUBSCRIBE call ที่ socket นั้นพึ่งพาจะเสร็จ (ซึ่งถ้าเกิดเช่นนั้น event ที่ publish ในช่วงนั้นจะหายไปเงียบ ๆ) ต้นทุนคือ run_subscriber จะได้รับ traffic ของทุกบอร์ดบนทุก instance ไม่ว่า instance นั้นจะมี socket ดูบอร์ดนั้นอยู่หรือไม่ — Hub::route คือจุดที่การกรองนี้เกิดขึ้นจริง อย่างประหยัด เป็นแค่ HashMap lookup ที่ไม่ทำอะไรเลยถ้าไม่มีใคร local กำลังฟังอยู่

background task ที่ PSUBSCRIBE board:* ตัวเดียว (สิ่งที่เราใช้) เทียบกับ SUBSCRIBE/UNSUBSCRIBE ต่อบอร์ดตอน socket เชื่อมต่อและตัดการเชื่อมต่อ

  • ข้อดี: Redis connection เดียว เปิดครั้งเดียวตอน startup ตลอดอายุของ process — ไม่มีการเปิดปิด connection ให้ปั่นป่วน ไม่มีบัญชี subscribe/unsubscribe ต่อบอร์ดที่ต้องคอยซิงก์กับ HashMap ของ Hub เอง และไม่มีช่วงเวลาที่ socket เชื่อมต่อเข้ามายังบอร์ดหนึ่งก่อนที่ SUBSCRIBE call ของบอร์ดนั้นจะเสร็จสมบูรณ์จริง ๆ กับ Redis (ซึ่งถ้าเป็นเช่นนั้น event ที่ publish ในช่วงเวลานั้นจะหายไปเงียบ ๆ)
  • ข้อเสีย: ทุก instance ได้รับ traffic ของทุกบอร์ดจาก Redis ไม่ว่าจะมีความต้องการ local หรือไม่ เป็นงานที่เสียเปล่าจริง ๆ สำหรับบอร์ดที่ไม่มีใครบน instance นั้นดูอยู่ — สำหรับสเกลของ TaskFlow (แอประดับคอร์ส ไม่ใช่ระบบที่กระจาย topic ความถี่สูงหลายล้านตัว) overhead นั้นคือ HashMap::get ต่อข้อความ ไม่ใช่ต้นทุนที่มีนัยสำคัญ ระบบที่มี topic จำนวนมหาศาลที่แทบไม่ overlap กันอาจหันไปใช้ SUBSCRIBE แบบ dynamic ต่อ topic แม้จะต้องเพิ่มบัญชีให้ยุ่งยากขึ้น เพื่อไม่ต้องจ่ายค่า traffic ที่ไม่มี local subscriber เลย บอร์ดของ TaskFlow ไม่ถึงสเกลที่ trade-off นี้จะพลิกกลับ

การ publish จาก mutation ทุกตัวผ่านฟังก์ชัน publish ที่อิง AppState ร่วมกัน (สิ่งที่เราใช้) เทียบกับ message queue (เช่น Redis Stream หรือ job queue เฉพาะทาง) ระหว่างการเขียนกับการกระจาย

  • ข้อดี: realtime::publish เป็นการเรียกฟังก์ชันแบบตรงไปตรงมา (ดูเหมือน synchronous จากมุมมองผู้เรียก) อยู่ติดกับการเรียก cache::invalidate ที่ invalidation เพิ่มเข้าไปในฟังก์ชันเดิมชุดเดียวกันแล้ว — เป็นแค่ .await? อีกหนึ่งตัวในเส้นทางจัดการ request ที่เป็น async อยู่แล้ว ไม่มี consumer process แยกต่างหาก ไม่มี queue infrastructure ให้ต้องรัน มอนิเตอร์ หรือคิดเรื่อง ordering guarantee การกระจาย event ของการย้ายการ์ดหนึ่งครั้งเกิดขึ้นภายใน request เดียวกับที่บันทึกการย้ายนั้นไว้ โดยใช้ Redis connection ที่ request มีเหตุผลจะถืออยู่แล้ว
  • ข้อเสีย: ถ้า PUBLISH call เองล้มเหลว (Redis สะดุดชั่วครู่) realtime::publish จะคืน Err และเพราะเรียกด้วย ? ทันทีหลัง cache::invalidate ทั้ง request จะล้มเหลวด้วย 500 ทั้งที่การเขียนลง Postgres ที่อยู่ข้างใต้ commit สำเร็จไปแล้วจริง ๆ — ไคลเอนต์อาจเห็น error สำหรับ mutation ที่ในความเป็นจริงมีผลไปแล้ว การออกแบบแบบ queue สามารถแยก “write สำเร็จไหม” ออกจาก “broadcast สำเร็จไหม” แล้ว retry การ broadcast อย่างอิสระได้ TaskFlow ยอมรับการผูกติดกันนี้: PUBLISH ที่ล้มเหลวหมายความแค่ว่า event ตัวนี้ตัวเดียว ไม่ถูกกระจายแบบ live — ไคลเอนต์ที่ reconnect หรือโหลดใหม่ยังคงเห็นการเปลี่ยนแปลงถูกต้องผ่าน GET /boards/:id ธรรมดา ซึ่งอ่านจาก cache-หรือ-Postgres ไม่ใช่จากประวัติ broadcast ใด ๆ เลเยอร์ realtime เป็นชั้นความสะดวกสำหรับอัปเดตแบบ live ที่วางทับ REST API ที่ถูกต้องอยู่แล้วอย่างเป็นอิสระ เลเยอร์ realtime ไม่เคยเป็น source of truth
Terminal window
cd taskflow/backend
cargo add futures-util -p api

คำสั่งนี้เพิ่ม futures-util = "0.3" ไว้ใต้ [dependencies] ใน api/Cargo.toml — จำเป็นสำหรับ trait StreamExt ซึ่งเปลี่ยน Redis pub/sub connection ดิบให้กลายเป็นสิ่งที่ loop ของ run_subscriber เรียก .next().await ได้

อัปเดต taskflow/backend/api/src/realtime/hub.rs:

use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use futures_util::StreamExt;
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()
}
/// Routes a message that arrived over the Redis backplane into the
/// matching board's local broadcast channel, if anyone's currently
/// listening for it.
fn route(&self, board_id: Uuid, payload: String) {
let channels = self.channels.lock().unwrap();
if let Some(sender) = channels.get(&board_id) {
// `send` returns `Err` when there are no receivers — that's not
// a failure, it just means nobody on *this* instance is
// watching this board right now.
let _ = sender.send(payload);
}
}
}
impl Default for Hub {
fn default() -> Self {
Self::new()
}
}
/// Spawned once at startup. Holds a single dedicated Redis connection,
/// `PSUBSCRIBE`s to every board's channel with one pattern, and routes each
/// incoming message into the matching board's `Hub` entry.
pub async fn run_subscriber(hub: Arc<Hub>, redis_url: String) -> anyhow::Result<()> {
let client = redis::Client::open(redis_url)?;
let mut pubsub = client.get_async_pubsub().await?;
pubsub.psubscribe("board:*").await?;
let mut messages = pubsub.on_message();
while let Some(msg) = messages.next().await {
let channel = msg.get_channel_name();
let Some(id) = channel.strip_prefix("board:") else {
continue;
};
let Ok(board_id) = id.parse::<Uuid>() else {
continue;
};
let Ok(payload) = msg.get_payload::<String>() else {
continue;
};
hub.route(board_id, payload);
}
Ok(())
}

ไล่ดู run_subscriber:

  1. redis::Client::open(redis_url)redis::Client ธรรมดา ไม่ใช่ deadpool_redis::Pool connection นี้อยู่เฉพาะตัวและมีอายุยืนตลอดชีวิตของ process ไม่เคยถูกคืนกลับเข้า pool ดังนั้นกลไก pooling จะเพิ่มแค่ overhead ที่นี่โดยไม่มีอะไรให้ pool เลย
  2. client.get_async_pubsub().await? — เปิด connection เฉพาะสำหรับใช้ pub/sub ส่งคืนค่า PubSub ที่ subscribe channel ได้และ stream ข้อความที่เข้ามาได้
  3. pubsub.psubscribe("board:*").await? — pattern subscription เดียว ทำครั้งเดียว ที่ครอบคลุม channel ของทุกบอร์ดโดยปริยายตลอดอายุที่เหลือของ task นี้
  4. pubsub.on_message() — เปลี่ยน connection ให้กลายเป็น Stream ของ Msg ที่เข้ามา; StreamExt::next() (จาก futures_util ที่เป็น dependency ใหม่ของบทนี้) คือสิ่งที่ทำให้ loop while let ดึงข้อมูลออกมาได้
  5. ต่อข้อความ: get_channel_name() คืนชื่อ channel เต็ม ๆ ที่ข้อความมาถึง (เช่น board:3fa8...); strip_prefix("board:") ดึงเฉพาะส่วน id ออกมา และ channel ที่ไม่ตรงกับ pattern ด้วยเหตุผลใดก็ตาม (ไม่ควรเกิดขึ้น เพราะมี PSUBSCRIBE board:* แต่ let else ป้องกันไว้เผื่อไว้) จะข้ามไปแทนที่จะ panic get_payload::<String>() อ่านตัวข้อความ — สตริง JSON ที่ publish สร้างขึ้น — และการ parse ที่ล้มเหลวก็ถูกข้ามไปเช่นกัน ไม่ทำให้ subscriber loop ทั้งหมดล่ม
  6. hub.route(board_id, payload) — ส่งข้อความให้ Hub ซึ่งจะส่งต่อไปยัง socket local ของ process นี้ที่ดูบอร์ดนั้นอยู่ หรือไม่ทำอะไรเลยเงียบ ๆ ถ้าไม่มี

route เป็น fn ไม่ใช่ pub fn เพราะมีแค่ run_subscriber ในไฟล์เดียวกันนี้ที่เรียก จึงไม่มีเหตุผลจะขยาย visibility ให้กว้างกว่าขอบเขตของ module นี้

อัปเดต taskflow/backend/api/src/realtime/mod.rs:

pub mod handlers;
pub mod hub;
use axum::{routing::get, Router};
use deadpool_redis::redis::AsyncCommands;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use crate::{
error::{AppError, AppResult},
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))
}
/// Builds a `BoardEvent`, serializes it to JSON, and `PUBLISH`es it to
/// `board:{board_id}` — the one function every mutation in this course
/// calls, right after persisting and invalidating the cache, to broadcast
/// what just changed.
pub async fn publish(
state: &AppState,
board_id: Uuid,
event_type: &str,
payload: serde_json::Value,
) -> AppResult<()> {
let event = BoardEvent {
r#type: event_type.to_string(),
board_id,
payload,
};
let json = serde_json::to_string(&event).map_err(|err| AppError::Internal(err.into()))?;
let mut conn = state
.redis
.get()
.await
.map_err(|err| AppError::Internal(err.into()))?;
conn.publish::<_, _, ()>(format!("board:{board_id}"), json)
.await
.map_err(|err| AppError::Internal(err.into()))?;
Ok(())
}

publish เข้าถึง state.redisdeadpool_redis::Pool ตัวเดียวกับที่ cache.rs ใช้อยู่แล้ว — แทนที่จะใช้ connection เฉพาะที่ run_subscriber ถืออยู่ PUBLISH เป็นคำสั่งครั้งเดียวธรรมดาที่เข้ากับรูปแบบ pooled-connection ต่อ request ที่ Redis write อื่น ๆ ทุกตัวในคอร์สนี้ใช้อยู่แล้ว ไม่เหมือนกับ subscription อายุยืนที่ run_subscriber ต้องการ conn.publish::<_, _, ()>(...) ใช้ turbofish สำหรับทิ้ง return type แบบเดียวกับที่ set_ex/del ของ cache.rs เคยตั้งไว้แล้ว — PUBLISH ตอบกลับด้วยจำนวนไคลเอนต์ที่ได้รับข้อความ ซึ่งฟังก์ชันนี้ไม่มีประโยชน์ใช้เลย

อัปเดต taskflow/backend/api/src/cards/service.rs — ต้อง import crate::realtime ไว้ด้านบนของไฟล์ด้วย:

use sqlx::PgPool;
use uuid::Uuid;
use crate::{
boards::service as boards_service,
cache, columns,
error::{AppError, AppResult},
realtime,
state::AppState,
};
use super::model::Card;
use super::repo;
pub async fn create_card(
state: &AppState,
user_id: Uuid,
column_id: Uuid,
title: String,
description: Option<String>,
) -> AppResult<Card> {
let db = &state.db;
let column = columns::repo::find_column(db, column_id)
.await?
.ok_or(AppError::NotFound)?;
boards_service::assert_member(db, user_id, column.board_id).await?;
let position = repo::max_position(db, column_id).await?.unwrap_or(0.0) + 1.0;
let card = repo::insert_card(
db,
Uuid::new_v4(),
column_id,
&title,
description.as_deref(),
position,
)
.await?;
cache::invalidate(&state.redis, &format!("cache:board:{}", column.board_id)).await?;
realtime::publish(
state,
column.board_id,
"card.created",
serde_json::json!({ "columnId": column_id, "card": card }),
)
.await?;
Ok(card)
}
async fn card_board_id(db: &PgPool, card: &Card) -> AppResult<Uuid> {
let column = columns::repo::find_column(db, card.column_id)
.await?
.ok_or(AppError::NotFound)?;
Ok(column.board_id)
}
pub async fn get_card(db: &PgPool, user_id: Uuid, card_id: Uuid) -> AppResult<Card> {
let card = repo::find_card(db, card_id)
.await?
.ok_or(AppError::NotFound)?;
let board_id = card_board_id(db, &card).await?;
boards_service::assert_member(db, user_id, board_id).await?;
Ok(card)
}
pub async fn update_card(
state: &AppState,
user_id: Uuid,
card_id: Uuid,
title: Option<String>,
description: Option<String>,
) -> AppResult<Card> {
let db = &state.db;
let card = repo::find_card(db, card_id)
.await?
.ok_or(AppError::NotFound)?;
let board_id = card_board_id(db, &card).await?;
boards_service::assert_member(db, user_id, board_id).await?;
let updated = repo::update_card(db, card_id, title.as_deref(), description.as_deref())
.await?
.ok_or(AppError::NotFound)?;
cache::invalidate(&state.redis, &format!("cache:board:{board_id}")).await?;
realtime::publish(
state,
board_id,
"card.updated",
serde_json::json!({ "card": updated }),
)
.await?;
Ok(updated)
}
pub async fn delete_card(state: &AppState, user_id: Uuid, card_id: Uuid) -> AppResult<()> {
let db = &state.db;
let card = repo::find_card(db, card_id)
.await?
.ok_or(AppError::NotFound)?;
let board_id = card_board_id(db, &card).await?;
boards_service::assert_member(db, user_id, board_id).await?;
if repo::delete_card(db, card_id).await? {
cache::invalidate(&state.redis, &format!("cache:board:{board_id}")).await?;
realtime::publish(
state,
board_id,
"card.deleted",
serde_json::json!({ "cardId": card_id, "columnId": card.column_id }),
)
.await?;
Ok(())
} else {
Err(AppError::NotFound)
}
}
pub async fn move_card(
state: &AppState,
user_id: Uuid,
card_id: Uuid,
target_column_id: Uuid,
before_id: Option<Uuid>,
after_id: Option<Uuid>,
) -> AppResult<Card> {
let db = &state.db;
let card = repo::find_card(db, card_id)
.await?
.ok_or(AppError::NotFound)?;
let source_board_id = card_board_id(db, &card).await?;
boards_service::assert_member(db, user_id, source_board_id).await?;
let target_column = columns::repo::find_column(db, target_column_id)
.await?
.ok_or(AppError::NotFound)?;
if target_column.board_id != source_board_id {
return Err(AppError::Forbidden);
}
let before = match before_id {
Some(id) => Some(repo::find_card(db, id).await?.ok_or(AppError::NotFound)?),
None => None,
};
let after = match after_id {
Some(id) => Some(repo::find_card(db, id).await?.ok_or(AppError::NotFound)?),
None => None,
};
let position = match (&before, &after) {
(Some(before), Some(after)) => (before.position + after.position) / 2.0,
(Some(before), None) => before.position + 1.0,
(None, Some(after)) => after.position - 1.0,
(None, None) => 1.0,
};
let updated = repo::move_card(db, card_id, target_column_id, position)
.await?
.ok_or(AppError::NotFound)?;
cache::invalidate(&state.redis, &format!("cache:board:{source_board_id}")).await?;
realtime::publish(
state,
source_board_id,
"card.moved",
serde_json::json!({ "card": updated }),
)
.await?;
Ok(updated)
}

การเรียก realtime::publish ใน move_card แทนที่คอมเมนต์ // Module 7 wires realtime broadcast of card.moved here ที่ move-reorder ทิ้งไว้ตรงจุดเดิมเป๊ะ — signature &AppState ที่บทนั้นวางไว้ล่วงหน้าสองโมดูล ในที่สุดก็ได้ใช้ตามเหตุผลที่เพิ่มเข้ามาตั้งแต่แรก

delete_card อ่าน card.column_id สำหรับ payload ของ event ก่อน branch if repo::delete_card(...)?card ดึงมาและตรวจสิทธิ์ไปแล้วตั้งแต่ต้นฟังก์ชัน ฟิลด์ต่าง ๆ จึงยังเป็นค่า Rust ที่ใช้งานได้ อธิบายแถวที่กำลังจะหายไป แม้ว่า DELETE จะรันไปแล้วจริง ๆ ก็ตาม

อัปเดต taskflow/backend/api/src/columns/service.rs — ต้อง import crate::realtime ด้วยเช่นกัน:

use uuid::Uuid;
use crate::{
boards::service as boards_service,
cache,
error::{AppError, AppResult},
realtime,
state::AppState,
};
use super::model::Column;
use super::repo;
pub async fn create_column(
state: &AppState,
user_id: Uuid,
board_id: Uuid,
title: String,
) -> AppResult<Column> {
let db = &state.db;
boards_service::assert_member(db, user_id, board_id).await?;
let position = repo::max_position(db, board_id).await?.unwrap_or(0.0) + 1.0;
let column = repo::insert_column(db, Uuid::new_v4(), board_id, &title, position).await?;
cache::invalidate(&state.redis, &format!("cache:board:{board_id}")).await?;
realtime::publish(
state,
board_id,
"column.created",
serde_json::json!({ "column": column }),
)
.await?;
Ok(column)
}
pub async fn update_column(
state: &AppState,
user_id: Uuid,
column_id: Uuid,
title: String,
) -> AppResult<Column> {
let db = &state.db;
let column = repo::find_column(db, column_id)
.await?
.ok_or(AppError::NotFound)?;
boards_service::assert_member(db, user_id, column.board_id).await?;
let updated = repo::update_title(db, column_id, &title)
.await?
.ok_or(AppError::NotFound)?;
cache::invalidate(&state.redis, &format!("cache:board:{}", column.board_id)).await?;
realtime::publish(
state,
column.board_id,
"column.updated",
serde_json::json!({ "column": updated }),
)
.await?;
Ok(updated)
}
pub async fn delete_column(state: &AppState, user_id: Uuid, column_id: Uuid) -> AppResult<()> {
let db = &state.db;
let column = repo::find_column(db, column_id)
.await?
.ok_or(AppError::NotFound)?;
boards_service::assert_member(db, user_id, column.board_id).await?;
if repo::delete_column(db, column_id).await? {
cache::invalidate(&state.redis, &format!("cache:board:{}", column.board_id)).await?;
realtime::publish(
state,
column.board_id,
"column.deleted",
serde_json::json!({ "columnId": column_id }),
)
.await?;
Ok(())
} else {
Err(AppError::NotFound)
}
}

update_board และ attach_label/detach_label จงใจ ไม่ ถูกเชื่อมเข้ากับ realtime::publish และควรพูดถึงตรงนี้อย่างชัดเจน เหมือนกับที่ invalidation เคยอธิบายว่าทำไม create_board/delete_board/create_label/delete_label ถูกยกเว้นจาก cache invalidation รายการ event ของ protocol ไม่มีรายการ board.updated หรือ label.* เลย — ทุก event type ที่มีอยู่ตรงกับการเปลี่ยนแปลงของการ์ดหรือคอลัมน์ เพราะนั่นคือสองสิ่งที่มุมมอง live ของบอร์ด Kanban จริง ๆ วาดใหม่แบบ in-place การคิดค้น event type ใหม่ขึ้นมาครอบคลุมการเปลี่ยนชื่อบอร์ดหรือการเปลี่ยนแปลง label เป็นก้าวต่อไปที่สมเหตุสมผลจริง ๆ สำหรับบทเรียนในอนาคต (broadcast board.updated เพื่อให้ชื่อบอร์ดที่เปลี่ยนอัปเดตแบบ live ในทุกแท็บที่เปิดอยู่ สะท้อนเหตุผลการรวมไว้เชิงป้องกันแบบเดียวกับที่ invalidation ใช้กับการ invalidate cache ของ attach_label/detach_label) — แต่การเพิ่มเข้ามาตรงนี้ โดยไม่มีรายการที่ตรงกันในรายการอย่างเป็นทางการของบทนี้เอง จะหมายความว่า event handler ฝั่งไคลเอนต์จะได้รับ type ที่ไม่เคยถูกบอกให้คาดหวังไว้ โมดูลนี้ส่งมอบ event type ทั้งเจ็ดรายการที่ protocol ตั้งชื่อไว้เท่านั้น ไม่มากไปกว่านั้น

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 realtime::hub::Hub;
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 hub = Arc::new(Hub::new());
let state = AppState {
db,
redis,
config: Arc::new(config),
hub: hub.clone(),
};
tokio::spawn({
let hub = hub.clone();
let redis_url = state.config.redis_url.clone();
async move {
if let Err(err) = realtime::hub::run_subscriber(hub, redis_url).await {
tracing::error!(error = %err, "realtime redis subscriber exited");
}
}
});
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" }))
}

สองอย่างที่เปลี่ยนจากเวอร์ชันของ ws-endpoint: hub ถูกดึงออกมาเป็น binding let hub = Arc::new(Hub::new()); ของตัวเอง — แทนที่จะสร้างแบบ inline อยู่ใน literal ของ AppState — เพื่อให้ทั้ง state.hub และ task ที่ spawn ด้านล่างแชร์ Arc ตัวเดียวกันผ่าน .clone() (เป็นการ clone Arc ไม่ใช่ clone Hub — ต้นทุนถูก และ handle ทั้งสองชี้ไปที่ registry เดียวกันเป๊ะ) และบล็อก tokio::spawn รัน run_subscriber ตลอดอายุของ process โดย log (ผ่าน tracing::error!) แทนที่จะทำให้เซิร์ฟเวอร์ทั้งตัวล่มถ้า Redis subscriber จบด้วย error สักวัน — realtime ที่ล่มจะลดระดับเหลือแค่ “อัปเดตแบบ live หยุดมาถึง” ไม่ใช่ “API ล่ม” สอดคล้องกับการอภิปรายใน Pros & cons ด้านบนที่ว่า publish ล้มเหลวไม่ควรบล็อกการเขียนที่อยู่ข้างใต้

Terminal window
cargo check -p api

เปิดสแตกขึ้นมา เชื่อมต่อ websocat client เข้ากับบอร์ด (ใช้ $TOKEN/$BOARD_ID เดิมจาก ws-endpoint) แล้ว — จากอีกเทอร์มินัลหนึ่ง — ทำ mutation:

Terminal window
cd taskflow/infra && docker compose up -d db redis
cd ../backend && cargo run -p api &
websocat "ws://localhost:8080/ws/boards/$BOARD_ID?token=$TOKEN" &
sleep 1
curl -s -X POST http://localhost:8080/boards/$BOARD_ID/columns \
-H "Authorization: Bearer $TOKEN" -H "Content-Type: application/json" \
-d '{"title":"To Do"}' > /dev/null

เทอร์มินัลของ websocat จะพิมพ์ event ทันทีที่ request curl เสร็จสมบูรณ์:

{"type":"column.created","boardId":"...","payload":{"column":{"id":"...","boardId":"...","title":"To Do","position":1.0}}}

เพื่อดูว่า backplane เองทำงานจริง ๆ — ไม่ใช่แค่ Hub ภายใน process เดียวที่ ws-endpoint พิสูจน์ไปแล้ว — จำลอง backend instance ตัวที่สองด้วยการ publish ตรงไปยัง Redis เลย โดยข้าม API process ไปทั้งหมด:

Terminal window
docker compose exec redis redis-cli PUBLISH "board:$BOARD_ID" \
'{"type":"card.updated","boardId":"'"$BOARD_ID"'","payload":{"card":{"id":"simulated"}}}'

connection websocat ตัวเดิมพิมพ์ event นั้นออกมาด้วยเช่นกัน — PSUBSCRIBE board:* ของ run_subscriber ไม่แยกแยะระหว่างข้อความที่ publish โดย process api ของคอร์สนี้เอง กับข้อความที่ publish โดยไคลเอนต์อื่นใดก็ตามที่พูด protocol เดียวกัน คุณสมบัตินี้ทำให้ subscriber ทำงานเหมือนกันเป๊ะ ไม่ว่าผู้ publish อีกฝ่ายจะเป็น cargo run -p api instance ตัวที่สอง หรือ — ในที่นี้ — redis-cli ที่สวมบทบาทแทน

คุณได้ขยาย hub.rs ด้วย run_subscriber — Redis connection เฉพาะหนึ่งเส้น PSUBSCRIBE board:* หนึ่งครั้ง ส่งทุกข้อความที่เข้ามาไปยัง local board channel ที่ตรงกัน — และเพิ่ม realtime::publish ใน mod.rs ฟังก์ชันที่สร้าง BoardEvent แล้ว serialize ก่อน PUBLISH ไปที่ board:{board_id} คุณเชื่อม publish เข้ากับ move_card (ปิด TODO ที่ค้างมาสองโมดูลจาก move-reorder) และเข้ากับ card/column mutation อีกหกตัว ตรงกับรายการของ protocol เป๊ะ — และอธิบายว่าทำไม update_board กับ label mutation ถึงจงใจถูกละไว้ก่อน คุณได้เห็นว่าทำไม PSUBSCRIBE board:* แบบเดียวถึงดีกว่าบัญชี SUBSCRIBE/UNSUBSCRIBE ต่อบอร์ด และยืนยันว่า backplane ปิดช่องว่างข้าม instance ได้จริง โดยการ publish ตรงไปยัง Redis แล้วดูผลลัพธ์เดียวกันกับที่ api process ตัวที่สองจะสร้างขึ้น นั่นคือการอัปเดตแบบ realtime ที่ทำงานครบวงจร บน backend instance กี่ตัวก็ได้ ต่อไป reconnect จะทำให้ socket loop เองแข็งแรงขึ้น — ตรวจจับ connection ที่ตายไปโดยไม่มีการปิดแบบสะอาด และกำหนดว่าไคลเอนต์ต้องทำอะไรทันทีที่ reconnect หลังจากพลาด event ไปช่วงหนึ่งโดยสิ้นเชิง