async I/O กับ tokio::net — echo server ที่รับพันคอนเนกชันด้วย task ไม่ใช่ thread
สามบทแรกเป็นการปูพื้น: บท 1 พิสูจน์ว่า future ขี้เกียจและเขียน block_on เอง บท 2 เอา runtime ตัวจริงมาแทน บท 3 แตก task ออกไปด้วย spawn บทนี้คือบทที่ทั้งคอร์สเดินมาหา — เรากลับไปหา code ของ kaen-kvstore จาก #22 ตรงจุดที่มันมีเพดานต่ำที่สุด คือ accept loop ที่ส่งแต่ละคอนเนกชันเข้า ThreadPool ขนาดคงที่ 4 worker แล้วปล่อยให้คอนเนกชันนั้น ยึด worker ไว้ทั้งเส้น แล้วเปลี่ยนมันเป็น task
ข่าวดีคือ code แทบไม่เปลี่ยนรูปเลย tokio::net จงใจลอกรูปของ std::net มาแทบ 1 ต่อ 1 — TcpListener::bind ยังชื่อ bind accept() ยังคืน (TcpStream, SocketAddr) เหมือนเดิม สิ่งที่เพิ่มมาคือ .await ต่อท้ายทุกบรรทัดที่แตะ I/O และ thread::spawn กลายเป็น tokio::spawn ข่าวร้ายคือมีกำแพงสามอันซ่อนอยู่ในความคล้ายนั้น อันแรกดังมาก (compile ไม่ผ่านทันที) อันที่สองเงียบสนิท (compile ผ่าน แต่ protocol พังเงียบๆ) และอันที่สามคือ borrow checker ยื่นบิลมาเก็บตอนคุณแยก read กับ write ออกเป็น2 task
คอร์สนี้ ต่อยอด repo kaen-kvstore จาก #22 (code ตัวอย่างกำลังจัดทำ) — บทนี้คือจุดที่เรายก transport layer ทั้งชั้นขึ้นมาบน Tokio: accept loop, byte echo, และ wire format เดิมทั้งดุ้น คือ u32 little-endian นำหน้าความยาวแล้วตามด้วย payload ตัว logic ของ store (SET/GET, WAL, fsync) ยังไม่ถูกแตะในบทนี้เลย — เราเปลี่ยนแค่ วิธีรับ byte เข้ามา ซึ่งเป็นประเด็นทั้งหมดของ async
ทุก snippet pin ที่ rustc 1.97.1 (8bab26f4f 2026-07-14) · edition = “2024” · tokio 1.53.1 (ปล่อย 2026-07-20) และ [dependencies] ของบทนี้คือบรรทัดนี้เท่านั้น:
[dependencies]tokio = { version = "1.53.1", features = ["rt-multi-thread", "net", "io-util", "macros"] }สร้างเองได้ด้วย cargo new kvnet && cargo add tokio@1.53.1 --features rt-multi-thread,net,io-util,macros (เขียน tokio@1.53.1 ไม่ใช่ tokio@=1.53.1 — ตัวหลัง cargo add จะเขียนลง Cargo.toml เป็น version = "=1.53.1" ซึ่งไม่ตรงกับบรรทัดข้างบน · และ cargo new แจก src/main.rs มาให้ ลบทิ้งได้เลย เพราะ project นี้เป็น lib + หลาย binary ให้สร้าง src/lib.rs กับ folder src/bin/ ขึ้นมาแทน) ห้าม features = ["full"] เด็ดขาด สังเกตความต่างจากบท 1: ตอนนั้น scope ["rt", "macros"] resolve มาแค่ 7 crate และเราชี้ว่ายังไม่มี mio/socket2/libc พอเปิด "net" ปุ๊บ cargo tree ก็โผล่มาครบ — tokio, bytes, libc, mio, pin-project-lite, socket2, tokio-macros บวก proc-macro2/quote/syn/unicode-ident ที่เป็น build dep ของมาโคร รวม 11 crate บน Linux — ตัวเลข 11 นี้นับเฉพาะ dependency ที่ compile จริงบน target x86_64-unknown-linux-gnu (7 ตัวเดิมของบท 1 + bytes/libc/mio/socket2 อีกสี่) ถ้าไปนับรายการใน Cargo.lock จะได้ 15 รายการ ซึ่งตรงกับที่บท 2 พยากรณ์ไว้พอดี เพราะ lock นับ crate ของเราเองเข้าไปด้วย และยังบันทึกรายการฝั่ง wasi/windows ที่ไม่ได้ถูก compile บนเครื่องนี้เลย — สองตัวเลขนี้ถูกทั้งคู่ ต่างกันแค่ฐานนับ นี่คือหลักฐานที่จับต้องได้ว่า "net" = การดึง reactor ตัวจริงเข้ามา
project ของบทนี้มี6 file ที่เขียว ทุก file ผ่าน cargo build / cargo test / cargo clippy -- -D warnings แบบ zero-warnings บน target x86_64-unknown-linux-gnu และ รันจริง ทุกบรรทัด output ข้างล่างคือข้อความจริงที่พิมพ์ออกมา: src/lib.rs (read_frame / write_frame / read_frame_or_eof + test สองตัว) · src/bin/echo_server.rs · src/bin/echo_client.rs · src/bin/frame_server.rs · src/bin/frame_client.rs · src/bin/split_demo.rs บวกอีก2 file ❌ ที่มีไว้พังโดยเฉพาะแล้วลบทิ้ง คือ src/bin/bad_endian.rs กับ src/bin/bad_split.rs
เลขบรรทัดในข้อความ error ทุกก้อนของบทนี้ตรงกับ file ที่แสดงไว้เป๊ะ ฉะนั้น snippet ข้างล่างจึงไม่มีคอมเมนต์ชื่อ file คั่นหัว file — ชื่อ file อยู่ในย่อหน้าก่อนหน้าแทน
รูปเดิมจาก #22 ทุกบรรทัด แค่เติม .await
หัวข้อที่มีชื่อว่า “รูปเดิมจาก #22 ทุกบรรทัด แค่เติม .await”นี่คือ src/bin/echo_server.rs ทั้ง file วางข้างๆ code std::net ของ #22 แล้วอ่านทีละบรรทัดได้เลย:
use tokio::io::{AsyncReadExt, AsyncWriteExt};use tokio::net::TcpListener;
#[tokio::main]async fn main() -> std::io::Result<()> { let listener = TcpListener::bind("127.0.0.1:8080").await?; println!("echo server listening on {}", listener.local_addr()?); loop { let (mut socket, addr) = listener.accept().await?; println!("accepted {addr}"); tokio::spawn(async move { let mut buf = [0u8; 1024]; loop { match socket.read(&mut buf).await { Ok(0) => { println!("peer closed: {addr}"); return; } Ok(n) => { let _ = socket.write_all(&buf[..n]).await; } Err(_) => return, } } }); }}และ src/bin/echo_client.rs ที่จับคู่กัน:
use tokio::io::{AsyncReadExt, AsyncWriteExt};use tokio::net::TcpStream;
#[tokio::main]async fn main() -> std::io::Result<()> { let mut stream = TcpStream::connect("127.0.0.1:8080").await?; stream.write_all(b"hello").await?; let mut buf = [0u8; 1024]; let n = stream.read(&mut buf).await?; println!("echo_back={}", String::from_utf8_lossy(&buf[..n])); Ok(())}สั่ง cargo run --bin echo_server ค้างไว้ แล้วอีกเทอร์มินัล cargo run --bin echo_client ฝั่ง client พิมพ์:
echo_back=helloฝั่ง server พิมพ์ (เลข port ฝั่ง client เป็น ephemeral port จึงเปลี่ยนทุกครั้ง):
echo server listening on 127.0.0.1:8080accepted 127.0.0.1:47814peer closed: 127.0.0.1:47814ความต่างจาก #22 มีสามจุดเท่านั้น และจุดที่สามคือหัวใจ:
bindและacceptและreadและwrite_allทุกตัวมี.awaitต่อท้าย — เพราะทุกตัวคืน future ที่ไม่ทำอะไรจนกว่าจะถูก poll (บท 1)#[tokio::main]แทนfn main()เปล่าๆ — ต้องมี runtime มาหมุนให้ (บท 2)tokio::spawnแทนpool.execute(...)ที่ดึง OS thread — 1 connection = 1 task ไม่ใช่1 thread
ข้อ 3 คือเหตุผลทั้งหมดที่คุณอ่านคอร์สนี้ ThreadPool ของ #22 มี worker แค่ 4 ตัว และคอนเนกชันหนึ่งสายยึด worker ตัวหนึ่งไว้จนกว่าจะวางสาย — คอนเนกชันที่ 5 จึงต้องรอ ไม่ใช่เพราะ CPU ไม่พอ แต่เพราะ worker ทั้งสี่กำลัง block รอ byte อยู่เฉยๆ (จะขยาย pool ให้ใหญ่ขึ้นก็ได้ แต่แต่ละ worker คือ OS thread ที่ขอ stack จาก kernel เป็นเมกะ byte และให้ kernel scheduler สลับให้) ส่วน task ที่เอกสาร Tokio ระบุขนาดไว้ว่าเป็น “single allocation and 64 bytes” — คอนเนกชันที่กำลังรอ byte อยู่จึงมีต้นทุนแค่ struct ก้อนหนึ่งในหน่วยความจำ ไม่ใช่ OS thread ที่หลับอยู่ นี่คือที่มาของประโยค “รับพันคอนเนกชัน” ในหัวบท และเป็นเส้นที่ตรงกับ #22 พอดี: protocol เดิม โครงเดิม เปลี่ยนแค่หน่วยของการทำงานพร้อมกัน
tokio::net ให้รูปเดิมของ std::net มา แต่หน่วย concurrency เปลี่ยนจาก OS thread เป็น task — ที่เหลือทั้งบทคือราคาที่ต้องจ่ายให้การเปลี่ยนหน่วยนั้น
กำแพงที่หนึ่ง (ดัง): read ไม่ได้อยู่บน TcpStream
หัวข้อที่มีชื่อว่า “กำแพงที่หนึ่ง (ดัง): read ไม่ได้อยู่บน TcpStream”บรรทัดแรกสุดของ echo_server.rs คือ use tokio::io::{AsyncReadExt, AsyncWriteExt}; และมันไม่ใช่ของประดับ ลบบรรทัดนั้นทิ้ง (ทั้ง file จึงเลื่อนขึ้นหนึ่งบรรทัด) แล้ว cargo build ตอบกลับมาแบบนี้ (ตัดเหลือ error ตัวแรกและบรรทัดที่มีน้ำหนัก):
error[E0599]: no method named `read` found for struct `tokio::net::TcpStream` in the current scope --> src/bin/echo_server.rs:13:30 | 13 | match socket.read(&mut buf).await { | ^^^^ | ::: /home/nook/.cargo/registry/src/index.crates.io-1949cf8c6b5b557f/tokio-1.53.1/src/io/util/async_read_ext.rs:177:12 |177 | fn read<'a>(&'a mut self, buf: &'a mut [u8]) -> Read<'a, Self> | ---- the method is available for `tokio::net::TcpStream` here | = help: items from traits can only be used if the trait is in scopehelp: trait `AsyncReadExt` which provides `read` is implemented but not in scope; perhaps you want to import it | 1 + use tokio::io::AsyncReadExt; |อีก error หนึ่งตัวหน้าตาเหมือนกันเป๊ะสำหรับ write_all กับ AsyncWriteExt เติม use กลับไปแล้วเขียวทันที
กลไกเบื้องหลัง: เทรตฐานของ tokio คือ AsyncRead กับ AsyncWrite และสิ่งที่มันบังคับให้ implement มีแต่ method ทรง poll_* ล้วนๆ — AsyncRead มี method เดียวคือ poll_read ส่วน AsyncWrite มีสามตัวคือ poll_write / poll_flush / poll_shutdown ทุกตัวหน้าตาเหมือน poll ในบท 1 เป๊ะ (รับ Pin<&mut Self> กับ &mut Context คืน Poll) method ที่คุณอยากเรียกจริงๆ — read read_exact write_all flush read_u32_le — ไม่ได้อยู่บนเทรตฐาน แต่อยู่บน extension trait คือ AsyncReadExt และ AsyncWriteExt ซึ่งมี blanket impl ให้ทุกตัวที่ implement เทรตฐานอยู่แล้ว แปลว่า TcpStream มี method พวกนี้ครบตั้งแต่แรก — มันแค่ มองไม่เห็น จนกว่าเทรตจะเข้ามาอยู่ใน scope
C# ไม่มีอาการนี้เพราะ NetworkStream.ReadAsync เป็น method จริงบน class แต่คุณเคยเจอรูปแบบเดียวกันมาแล้ว: .Where() และ .Select() ไม่ compile จนกว่าจะมี using System.Linq; — extension method ต้องถูกดึงเข้า scope ก่อนถึงจะเห็น กรอบคิดที่ถูกคือ verb ของ async I/O ใน Rust เป็น extension method ทั้งหมด ต่างกันแค่ Rust บังคับเข้มกว่าและ IDE เดาให้น้อยกว่า
เมื่อไหร่ที่ task ถูกปลุก — reactor คือคนกดปุ่ม
หัวข้อที่มีชื่อว่า “เมื่อไหร่ที่ task ถูกปลุก — reactor คือคนกดปุ่ม”บท 1 ทิ้งคำถามไว้ข้อหนึ่ง: future ที่ตอบ Poll::Pending ต้องเก็บ Waker ไว้แล้วเรียก wake() ตอนไปต่อได้ — แต่ใน Countdown เราโกงด้วยการ wake_by_ref() ทันทีในรอบเดียวกัน คำถามคือ กับ socket จริง ใครเป็นคนเรียก wake()
คำตอบคือ reactor (บท 2 ประกอบมันเข้า runtime ไปแล้ว) ตอน socket.read(&mut buf).await แล้ว kernel ยังไม่มี byte ให้ สิ่งที่เกิดขึ้นไม่ใช่การรอเปล่าๆ แต่คือ: tokio เอา file descriptor ตัวนั้นไปลงทะเบียนกับ epoll ผ่าน crate mio (ตัวที่เพิ่งโผล่ใน cargo tree ตอนเปิด "net") พร้อมฝาก Waker ของ task ไว้ แล้วคืน Pending ให้ scheduler เอา worker thread ไปหมุน task อื่นต่อทันที เมื่อ kernel มี byte จริง epoll_wait ใน thread I/O driver ตื่นขึ้น มันค้นว่า fd นี้ผูกกับ Waker ตัวไหน แล้วเรียก wake() — task ถูกดันกลับเข้าคิว และถูก poll ซ้ำ คราวนี้ read คืน Ready
sequenceDiagram
participant T as task ของ connection
participant R as reactor คือ mio บน epoll
participant K as kernel
T->>R: poll แล้วยังไม่มี byte จึงฝาก Waker ไว้
R->>K: ลงทะเบียน fd นี้เข้า epoll
T-->>T: คืน Pending ตัว task หลับ worker thread ว่างไปหมุน task อื่น
K-->>R: epoll_wait ตื่น เพราะ fd พร้อมอ่านแล้ว
R->>T: เรียก wake บน Waker ที่ฝากไว้
T->>T: scheduler ดัน task กลับเข้าคิว แล้ว poll ซ้ำ คราวนี้ได้ Ready
คำบรรยายภาพ: เส้นทางการปลุก task ในโลก I/O จริง — task เรียก read แล้วยังไม่มีข้อมูล จึงฝาก Waker ไว้กับ reactor ซึ่งลงทะเบียน file descriptor นั้นกับ epoll ของ kernel task คืน Pending แล้วหลับ ปล่อยให้ worker thread ไปหมุนงานอื่น เมื่อ kernel บอกว่า fd พร้อม reactor เรียก wake ตัวเดิม scheduler จึงดัน task กลับเข้าคิวไป poll ซ้ำ ทั้งหมดนี้คือ handshake poll กับ wake อันเดิมของบท 1 เพียงแต่คนกดปุ่ม wake เปลี่ยนจาก thread unpark มาเป็น kernel
เทียบกับ .NET แล้วโครงนี้ เกือบ เหมือน — CLR ก็มี I/O completion port (หรือ epoll บน Linux) คอยรับสัญญาณจาก kernel เหมือนกัน ความต่างอยู่ที่ทิศทาง: ฝั่ง .NET เมื่อ I/O เสร็จ ระบบ push continuation ที่คุณลงทะเบียนไว้เข้าคิว ThreadPool ให้รันต่อ ส่วนฝั่ง Rust reactor แค่ เคาะประตูบอกว่า “กลับมาถามใหม่ได้แล้ว” ตัวงานจริงยังต้องรอ executor เดินมา poll ซ้ำอยู่ดี นี่คือ model pull ของบท 1 ที่ขยายมาถึงระดับ syscall
กำแพงที่สอง (เงียบ): ยก wire format ของ #22 มาทั้งดุ้น
หัวข้อที่มีชื่อว่า “กำแพงที่สอง (เงียบ): ยก wire format ของ #22 มาทั้งดุ้น”byte echo ข้างบนสวยแต่ไร้ประโยชน์ เพราะ kaen-kvstore ไม่ได้พูดภาษา byte เปล่า มันพูด length-prefixed-framinglength-prefixed-framingwire ของ #22: `u32` LE ความยาว ตามด้วย payload อ่านด้วย `read_exact` — u32 little-endian 4 byte บอกความยาว แล้วตามด้วย payload เท่านั้นพอดี วินัยนี้คือสิ่งที่ #22 ยืนยันมาแล้วว่า TCP ไม่รักษาขอบเขตข้อความให้ คุณต้องขีดเส้นเอง
async mirror ของมันสั้นกว่าที่คิด เพราะ AsyncReadExt แจก read_u32_le มาให้ตรงๆ นี่คือหัว file src/lib.rs ทั้งหมด:
use tokio::io::{AsyncReadExt, AsyncWriteExt};
pub async fn read_frame<R: AsyncReadExt + Unpin>(r: &mut R) -> std::io::Result<Vec<u8>> { let len = r.read_u32_le().await? as usize; let mut buf = vec![0u8; len]; r.read_exact(&mut buf).await?; Ok(buf)}
pub async fn write_frame<W: AsyncWriteExt + Unpin>(w: &mut W, bytes: &[u8]) -> std::io::Result<()> { w.write_u32_le(bytes.len() as u32).await?; w.write_all(bytes).await?; w.flush().await?; Ok(())}เทียบบรรทัดต่อบรรทัดกับ #22: read_u32_le().await? คือ read_exact(&mut [0u8; 4]) แล้ว u32::from_le_bytes รวมกันเป็นบรรทัดเดียว ส่วน read_exact(&mut buf).await? คือตัวเดิมเป๊ะ แค่ awaited วินัยไม่เปลี่ยนเลย protocol ไม่เปลี่ยนเลย เปลี่ยนแค่ว่ามันยอมปล่อย worker thread ระหว่างรอ
ข้อสังเกตเรื่อง bound: R: AsyncReadExt + Unpin ต้องมี Unpin เพราะ method พวกนี้ประกาศไว้ว่า where Self: Unpin — future ที่ read คืนมาต้องถือ &mut R ไว้ข้าม .await และจะทำแบบนั้นได้ต้องรู้ว่า R ย้ายที่ได้อย่างปลอดภัย TcpStream เป็น Unpin อยู่แล้ว จึงไม่มีอะไรต้องทำเพิ่ม
❌ version ดิบ: read_u32 แทน read_u32_le — compile ผ่าน แต่ protocol พัง
หัวข้อที่มีชื่อว่า “❌ version ดิบ: read_u32 แทน read_u32_le — compile ผ่าน แต่ protocol พัง”นี่คือกับดักที่อันตรายกว่า E0599 หลายเท่า เพราะ ไม่มี error ไม่มี warning ไม่มีอะไรเลย ลองเปลี่ยนฝั่ง server ให้อ่านด้วย read_u32() (ซึ่งเป็น big-endian ตามธรรมเนียม network byte order) ขณะที่ client ยังเขียนด้วย write_u32_le() ตาม wire ของ #22 — file ทดลองชื่อ src/bin/bad_endian.rs และมันคือ ❌ version ดิบ ของบทนี้:
use tokio::io::{AsyncReadExt, AsyncWriteExt};use tokio::net::{TcpListener, TcpStream};
#[tokio::main]async fn main() -> std::io::Result<()> { let listener = TcpListener::bind("127.0.0.1:8083").await?; let server = tokio::spawn(async move { let (mut socket, _) = listener.accept().await.unwrap(); // ผิด: read_u32 คือ big-endian ส่วน client เขียนด้วย write_u32_le let len = socket.read_u32().await.unwrap() as usize; println!("server thinks the frame is {len} bytes long"); let mut buf = vec![0u8; len]; let err = socket.read_exact(&mut buf).await.unwrap_err(); println!("read_exact failed: {:?}", err.kind()); });
let mut client = TcpStream::connect("127.0.0.1:8083").await?; client.write_u32_le(18).await?; client.write_all(b"SET greeting hello").await?; client.shutdown().await?; server.await.unwrap(); Ok(())}cargo run --bin bad_endian compile ผ่านสะอาด แล้วพิมพ์:
server thinks the frame is 301989888 bytes longread_exact failed: UnexpectedEof18 เขียนแบบ little-endian คือ byte 12 00 00 00 อ่านกลับแบบ big-endian ได้ 0x12000000 = 301,989,888 server จึงจอง Vec ขนาด ~288 MB ขึ้นมารอ payload ที่ไม่มีวันมา ใน file ทดลองนี้ client เรียก shutdown() ทันทีหลังส่ง 18 byte read_exact จึงชน EOF แล้วคืน UnexpectedEof ออกมาเลย โปรแกรมไม่ได้ค้างให้เห็น — ถ้า client เป็นของจริงที่ยังเปิดคอนเนกชันค้างไว้ server ตัวนี้จะรอต่อไปเรื่อยๆ พร้อมกับหน่วยความจำก้อนนั้นคาอยู่ในมือ นี่คือรูปที่โชคดี — server ตายเสียงดัง ในกรณีที่โชคร้ายกว่า (ตัวเลขความยาวที่อ่านผิดแล้วบังเอิญเล็ก) server จะอ่านไม่ครบ frame แล้ว byte ที่เหลือกลายเป็น “หัว frame ถัดไป” ทันที — stream desync แบบเงียบสนิทที่ debug ยากที่สุดแบบหนึ่ง
ทางกันคือกฎเดียว: ตัดสินใจเรื่อง endianness ครั้งเดียวแล้วเขียนมันลง test kaen-kvstore เลือก little-endian ตั้งแต่ #22 เพราะมันตรงกับ layout ของเครื่อง x86 และ ARM ที่รันอยู่จริง (และตรงกับ CLR ด้วย ซึ่งจะได้ใช้ในหัวข้อ C# ข้างล่าง) test คู่นี้ต่อท้าย src/lib.rs และเป็นตาข่ายที่จับกรณีข้างบนได้ทันที:
#[cfg(test)]mod tests { use super::*; use tokio::net::{TcpListener, TcpStream};
#[tokio::test] async fn frame_round_trips_over_tcp() -> std::io::Result<()> { let listener = TcpListener::bind("127.0.0.1:0").await?; let addr = listener.local_addr()?;
let server = tokio::spawn(async move { let (mut socket, _) = listener.accept().await.unwrap(); let frame = read_frame(&mut socket).await.unwrap(); write_frame(&mut socket, &frame).await.unwrap(); });
let mut client = TcpStream::connect(addr).await?; write_frame(&mut client, b"SET greeting hello").await?; let echoed = read_frame(&mut client).await?; assert_eq!(&echoed, b"SET greeting hello"); server.await.unwrap(); Ok(()) }
#[tokio::test] async fn truncated_frame_is_unexpected_eof() -> std::io::Result<()> { let listener = TcpListener::bind("127.0.0.1:0").await?; let addr = listener.local_addr()?;
let server = tokio::spawn(async move { let (mut socket, _) = listener.accept().await.unwrap(); read_frame(&mut socket).await.unwrap_err() });
let mut client = TcpStream::connect(addr).await?; client.write_u32_le(32).await?; client.write_all(b"only5").await?; client.shutdown().await?; drop(client);
let err = server.await.unwrap(); assert_eq!(err.kind(), std::io::ErrorKind::UnexpectedEof); Ok(()) }}cargo test ให้ผลนี้ (สังเกต bind("127.0.0.1:0") — ให้ kernel เลือก port ว่างให้ แล้วอ่านกลับด้วย local_addr() test จึงรันขนานกันได้โดยไม่ชนกันเอง · และเพราะมันรันขนานกัน ลำดับบรรทัด test … ok ที่คุณเห็นอาจสลับกับที่พิมพ์ไว้ข้างล่างนี้ เช่นเดียวกับเลข ephemeral port ที่เปลี่ยนทุกครั้ง สิ่งที่คงที่คือบรรทัดสรุป 2 passed; 0 failed):
running 2 teststest tests::truncated_frame_is_unexpected_eof ... oktest tests::frame_round_trips_over_tcp ... ok
test result: ok. 2 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out; finished in 0.01sกฎ correctness สองข้อ: Ok(0) กับ UnexpectedEof
หัวข้อที่มีชื่อว่า “กฎ correctness สองข้อ: Ok(0) กับ UnexpectedEof”test ตัวที่สองข้างบนตอกกฎข้อหนึ่งไว้แล้ว มาดูทั้งคู่พร้อมกัน เพราะมันคือคู่ที่คนสับสนกันมากที่สุด:
read()คืนOk(0)= peer ปิดคอนเนกชันอย่างสะอาด ไม่ใช่ error ไม่ใช่ timeout ไม่ใช่ “ยังไม่มีข้อมูล” (กรณีนั้นคือPendingซึ่งถูก.awaitกลืนไปแล้ว) เป็น EOF จริง นี่คือเหตุผลที่ echo server ข้างบน matchOk(0) => returnเป็น arm แรก ลองไล่ดูว่าถ้าลบ arm นั้นทิ้งจะเกิดอะไรขึ้น:Ok(n)จะรับn = 0ไปแทน แล้วเขียน&buf[..0]ซึ่งไม่ส่งอะไรเลย loop วนกลับไปreadใหม่ซึ่งคืนOk(0)อีกทันที (EOF จะคืนOk(0)ตลอดไป ไม่ใช่ครั้งเดียว) — loop จึงไม่มีจุดจบและไม่มีจุดไหนยอมPendingให้ scheduler ได้พัก นี่เป็นการไล่เหตุผลจากสัญญาของ API ไม่ใช่ตัวเลขที่บทนี้วัดมา — เราไม่ได้ ship file ❌ ตัวนี้ให้รันดู ต่างจากอีกสามกับดักข้างบนread_exactคืนErr(ErrorKind::UnexpectedEof)เมื่อ stream จบก่อนอ่านครบbuf.len()— คือกรณี “ประกาศว่าจะส่ง 32 byte แต่ส่งมา 5 แล้วปิด” ซึ่งเป็น ความผิดของ protocol ไม่ใช่การปิดปกติ
ปัญหาคือ read_frame ข้างบนแยกสองกรณีนี้ไม่ออก เพราะถ้า peer ปิดตรงขอบ frame พอดี (ซึ่งเป็นการปิดที่ถูกต้องสมบูรณ์) read_u32_le() ก็คืน UnexpectedEof เหมือนกัน server จะ log ว่า “protocol พัง” ทั้งที่ client แค่บอกลาอย่างสุภาพ ทางแก้คืออ่านหัว frame4 byte ด้วยมือ แล้วเช็ก Ok(0) เฉพาะตอนที่ยังไม่ได้อ่านอะไรเลย function นี้อยู่ใน src/lib.rs ต่อจาก write_frame:
pub async fn read_frame_or_eof<R: AsyncReadExt + Unpin>( r: &mut R,) -> std::io::Result<Option<Vec<u8>>> { let mut hdr = [0u8; 4]; let mut filled = 0; while filled < 4 { match r.read(&mut hdr[filled..]).await? { 0 if filled == 0 => return Ok(None), 0 => { return Err(std::io::Error::new( std::io::ErrorKind::UnexpectedEof, "header truncated", )); } n => filled += n, } } let len = u32::from_le_bytes(hdr) as usize; let mut buf = vec![0u8; len]; r.read_exact(&mut buf).await?; Ok(Some(buf))}Ok(None) = “จบอย่างสะอาด” · Err(UnexpectedEof) = “พังกลางทาง” 2 arm นี้แยกกันชัดแล้ว loop while filled < 4 จำเป็นเพราะ read() มีสิทธิ์คืนน้อยกว่าที่ขอ (short read) เสมอ — เหมือน NetworkStream.ReadAsync ใน C# เป๊ะ
server ที่ใช้มันคือ src/bin/frame_server.rs (ชื่อ crate ใน project นี้คือ kvnet ตามที่ cargo new kvnet ตั้งให้):
use kvnet::{read_frame_or_eof, write_frame};use tokio::net::TcpListener;
#[tokio::main]async fn main() -> std::io::Result<()> { let listener = TcpListener::bind("127.0.0.1:8081").await?; println!("frame server listening on {}", listener.local_addr()?); loop { let (mut socket, addr) = listener.accept().await?; tokio::spawn(async move { loop { match read_frame_or_eof(&mut socket).await { Ok(Some(frame)) => { println!("frame from {addr}: {} bytes", frame.len()); if write_frame(&mut socket, &frame).await.is_err() { return; } } Ok(None) => { println!("clean close at frame boundary: {addr}"); return; } Err(e) => { println!("broken frame from {addr}: {:?}", e.kind()); return; } } } }); }}และ client ที่จับคู่กันคือ src/bin/frame_client.rs:
use kvnet::{read_frame, write_frame};use tokio::net::TcpStream;
#[tokio::main]async fn main() -> std::io::Result<()> { let mut stream = TcpStream::connect("127.0.0.1:8081").await?; for cmd in [b"SET greeting hello".as_slice(), b"GET greeting".as_slice()] { write_frame(&mut stream, cmd).await?; let reply = read_frame(&mut stream).await?; println!("reply={}", String::from_utf8_lossy(&reply)); } Ok(())}รัน server ค้างไว้แล้วยิง client — ฝั่ง client:
reply=SET greeting helloreply=GET greetingฝั่ง server:
frame server listening on 127.0.0.1:8081frame from 127.0.0.1:46810: 18 bytesframe from 127.0.0.1:46810: 12 bytesclean close at frame boundary: 127.0.0.1:46810บรรทัดสุดท้ายคือรางวัลของหัวข้อนี้: การปิดที่สุภาพถูกจัดว่า สะอาด ไม่ใช่ error
สัญชาตญาณแรกของทุกคนคือ nc 127.0.0.1 8081 แล้วพิมพ์อะไรสักอย่าง — ใช้ไม่ได้ เพราะ nc ส่ง byte ดิบตามที่พิมพ์ ไม่มี u32 prefix นำหน้าให้ server จะเอา4 byte แรกของสิ่งที่คุณพิมพ์ไปตีความเป็น ความยาว ทันที ซึ่งกลายเป็นตัวเลขมั่วๆ ตัวหนึ่งเสมอ nc จึงใช้ได้เฉพาะกับ echo_server ที่พูด byte ดิบเท่านั้น สำหรับ protocol ที่มี framing ต้องมี client ที่พูดภาษาเดียวกัน ซึ่งคือเหตุผลที่บทนี้ ship frame_client.rs มาให้ครบคู่
กำแพงที่สาม: split() เข้า spawn ไม่ได้ ต้อง into_split()
หัวข้อที่มีชื่อว่า “กำแพงที่สาม: split() เข้า spawn ไม่ได้ ต้อง into_split()”จนถึงตรงนี้1 task อ่านแล้วเขียนสลับกันไป ซึ่งพอสำหรับ request/response แต่ทันทีที่คุณอยาก push ข้อมูลไปหา client โดยไม่รอ request (ซึ่งบท 8 จะต้องใช้) คุณต้องแยก read กับ write ออกเป็น2 task ที่เดินอิสระกัน และ borrow checker จะยื่นบิลมาเก็บทันที ลอง version ตรงไปตรงมาก่อน:
❌ version ดิบ ตัวนี้อยู่ใน file src/bin/bad_split.rs ทั้ง file ตามนี้ (เลขบรรทัดใน error ข้างล่างอ้าง file นี้ตรงๆ):
use tokio::io::{AsyncReadExt, AsyncWriteExt};use tokio::net::TcpListener;
#[tokio::main]async fn main() -> std::io::Result<()> { let listener = TcpListener::bind("127.0.0.1:8084").await?; let (mut socket, _) = listener.accept().await?; let (mut rd, mut wr) = socket.split();
tokio::spawn(async move { let mut buf = [0u8; 1024]; let _ = rd.read(&mut buf).await; }); tokio::spawn(async move { let _ = wr.write_all(b"hi").await; }); Ok(())}cargo build ตอบ (ตัดเหลือบรรทัดที่มีน้ำหนัก):
error[E0597]: `socket` does not live long enough --> src/bin/bad_split.rs:8:28 | 7 | let (mut socket, _) = listener.accept().await?; | ---------- binding `socket` declared here 8 | let (mut rd, mut wr) = socket.split(); | ^^^^^^ borrowed value does not live long enough 9 | 10 | / tokio::spawn(async move { 11 | | let mut buf = [0u8; 1024]; 12 | | let _ = rd.read(&mut buf).await; 13 | | }); | |______- argument requires that `socket` is borrowed for `'static`... 18 | } | - `socket` dropped here while still borrowed |note: requirement that the value outlives `'static` introduced here --> /home/nook/.cargo/registry/src/index.crates.io-1949cf8c6b5b557f/tokio-1.53.1/src/task/spawn.rs:176:28 |176 | F: Future + Send + 'static, | ^^^^^^^อ่าน note ท้ายสุดแล้วเรื่องจบในบรรทัดเดียว: tokio::spawn ต้องการ F: Future + Send + 'static (บท 3 ตอกไว้แล้ว) ส่วน split(&mut self) คืน ReadHalf<'a> กับ WriteHalf<'a> ที่ ยืม socket มา — มันมี lifetime ผูกกับ stack frame นี้ ไม่ใช่ 'static ตัว task ที่ spawn ไปอาจอยู่นานกว่า function ที่สร้างมัน compiler จึงปฏิเสธ
ทางแก้คือใช้ version owned:
let (mut rd, mut wr) = socket.into_split(); // (OwnedReadHalf, OwnedWriteHalf)into_split(self) กิน socket เข้าไปแล้วคืนสองครึ่งที่เป็นเจ้าของร่วมกัน (จ่าย heap allocation หนึ่งครั้งเป็นค่าผ่านทาง) ทั้งคู่จึงเป็น 'static และย้ายเข้า async move block ได้สบาย demo เต็มที่รันจริงคือ src/bin/split_demo.rs ซึ่งยัด server และ client ไว้ใน process เดียวเพื่อให้ output เรียงเหมือนเดิมทุกครั้ง:
use kvnet::{read_frame, write_frame};use tokio::io::{AsyncReadExt, AsyncWriteExt};use tokio::net::{TcpListener, TcpStream};
#[tokio::main]async fn main() -> std::io::Result<()> { let listener = TcpListener::bind("127.0.0.1:8082").await?; println!("split demo listening on {}", listener.local_addr()?);
let server = tokio::spawn(async move { let (socket, addr) = listener.accept().await.unwrap(); let (mut rd, mut wr) = socket.into_split();
let reader = tokio::spawn(async move { let mut buf = [0u8; 1024]; loop { match rd.read(&mut buf).await { Ok(0) => { println!("reader task: peer closed {addr}"); return; } Ok(n) => println!("reader task: got {n} bytes"), Err(_) => return, } } });
let writer = tokio::spawn(async move { write_frame(&mut wr, b"BANNER kaen-kvstore v4").await.unwrap(); println!("writer task: banner sent"); });
writer.await.unwrap(); reader.await.unwrap(); });
let mut client = TcpStream::connect("127.0.0.1:8082").await?; let banner = read_frame(&mut client).await?; println!("client: {}", String::from_utf8_lossy(&banner)); client.write_all(b"PING").await?; client.shutdown().await?; server.await.unwrap(); Ok(())}cargo run --bin split_demo พิมพ์:
split demo listening on 127.0.0.1:8082writer task: banner sentclient: BANNER kaen-kvstore v4reader task: got 4 bytesreader task: peer closed 127.0.0.1:46222สังเกตว่าเราเก็บ JoinHandle ของทั้ง reader และ writer ไว้แล้ว .await มันตามลำดับ (บท 3) — ไม่ใช่เพื่อความสวยงาม แต่เพราะถ้าไม่ทำ main จะจบก่อนที่ task จะได้พิมพ์ ทำให้ output ไม่แน่นอนทุกครั้งที่รัน
ตรงนี้คือจุดที่สัญชาตญาณจาก C# หลอกแรงที่สุดในบทนี้: ใน .NET คุณส่ง NetworkStream ตัวเดียวกันเข้าไปในสอง Task ได้ตรงๆ ไม่มีใครห้าม (จะ thread-safe หรือไม่เป็นเรื่องที่คุณต้องดูเอง) ส่วน Rust ห้ามที่ระดับ type ว่าจะมี &mut สองตัวชี้ไปที่เดียวกันไม่ได้ ทางเดียวคือแบ่งความเป็นเจ้าของออกจากกันจริงๆ ซึ่ง into_split() ทำให้ และผลพลอยได้คือคุณได้ความปลอดภัยฟรี — OwnedReadHalf ไม่มีทางเขียน OwnedWriteHalf ไม่มีทางอ่าน ความสับสนแบบ “2 task เขียนทับกันกลาง frame” หายไปตั้งแต่ตอน compile
ถ้าลด [dependencies] เหลือ features = ["rt", "net", "io-util", "macros"] แล้วยังเขียน #[tokio::main] เปล่าๆ ตาม code ทั้งบทนี้ cargo build ตอบทันที:
error: The default runtime flavor is `multi_thread`, but the `rt-multi-thread` feature is disabled. --> src/bin/echo_server.rs:4:1 |4 | #[tokio::main] | ^^^^^^^^^^^^^^ | = note: this error originates in the attribute macro `tokio::main` (in Nightly builds, run with -Z macro-backtrace for more info)ทางแก้มีสองทางและ ไม่เท่ากัน: เติม "rt-multi-thread" กลับเข้าไป (สิ่งที่บทนี้ทำ เพราะเราอยากได้ worker หลายตัวมารับหลายคอนเนกชัน) หรือเขียน #[tokio::main(flavor = "current_thread")] ซึ่งใช้ scheduler thread เดียวและต้องการแค่ "rt" — ตัวหลังยังรับพันคอนเนกชันได้ (เพราะ concurrency มาจาก task ไม่ใช่ thread) แต่ใช้ CPU ได้แค่คอร์เดียว 🔁 บท 2 เทียบ2 flavor นี้ไว้ครบแล้ว
แผนที่คำศัพท์จาก C#/.NET
หัวข้อที่มีชื่อว่า “แผนที่คำศัพท์จาก C#/.NET”| C# / .NET | Rust + tokio | จุดที่ต่างจริง |
|---|---|---|
TcpListener.Start() + await AcceptTcpClientAsync() | TcpListener::bind(addr).await + accept().await | ฝั่ง Rust bind เองก็ async และคืน io::Result ที่ต้องจัดการด้วย ? |
await new TcpClient().ConnectAsync(...) | TcpStream::connect(addr).await | ฝั่ง Rust ไม่มีชั้น TcpClient คร่อม NetworkStream — TcpStream คือ stream เลย |
NetworkStream.ReadAsync / WriteAsync | read / write_all บน AsyncReadExt / AsyncWriteExt | ต้อง use เทรตก่อน ไม่งั้น error[E0599] method หายไปเฉยๆ |
Stream.ReadExactlyAsync (มาใน .NET 7 · 2022-11-08) | read_exact | ความหมายตรงกัน รวมถึงการโยน/คืน error เมื่อ stream จบก่อนครบ |
BinaryPrimitives.ReadUInt32LittleEndian | read_u32_le / write_u32_le | CLR เป็น little-endian ดังนั้น BinaryWriter.Write(uint) แมตช์ write_u32_le เป๊ะ — wire ของ #22 ข้ามภาษาได้เลย |
ส่ง NetworkStream ตัวเดียวเข้าสอง Task ได้เลย | ต้อง into_split() ก่อน | borrow checker ห้าม &mut สองตัวชี้ที่เดียวกัน — split() ที่ยืมมาจะติด 'static ของ spawn |
| หนึ่งคอนเนกชัน = งานบน ThreadPool ที่มี thread หลักร้อย | หนึ่งคอนเนกชัน = task ขนาด “single allocation and 64 bytes” | Tokio มี worker เท่าจำนวนคอร์ แต่ task เป็นล้านได้ |
A — Future ขี้เกียจ: let f = socket.read(&mut buf); แล้วไม่ .await = ไม่มี syscall ไหนถูกยิงเลย socket ยังไม่ถูกแตะด้วยซ้ำ ต่างจาก ReadAsync ใน C# ที่เริ่มยิง I/O ทันทีตั้งแต่ก่อนคุณ await
B — std ไม่มี runtime: std::net::TcpStream มี read ที่ block thread ทั้งเส้น ส่วน tokio::net::TcpStream ต้องมี reactor คอยขับ และ reactor นั้นไม่ได้มาจาก std แต่มาจาก crate ที่เราเพิ่งดึงเข้ามา — หลักฐานอยู่ใน cargo tree ตรงๆ ว่าการเปิด feature "net" คือสิ่งที่พา mio / socket2 / libc เข้ามาใน project ทั้งที่บท 1 ด้วย scope ["rt", "macros"] ยังไม่มีสามตัวนี้เลย
C — function colouring: read_frame เป็น async fn function sync ตัวไหนก็เรียกมันไม่ได้ ผลคือ kaen-kvstore เดิมที่เป็น sync ทั้ง repo จะโดน async ไต่ขึ้นไปตาม call stack ทีละชั้นเมื่อเรายก transport มาไว้ตรงนี้ บท 7 ว่าด้วยเรื่องนี้ทั้งบท
D — ทุกอย่างในบทนี้รันจริง ไม่มีอันไหนเป็นไดอะแกรมอย่างเดียว (ยกเว้น snippet ที่ติดป้าย ❌ ซึ่งตั้งใจให้พัง): คำว่า “รับพันคอนเนกชัน” ในหัวบทเป็นการเล่าถึง model ไม่ใช่ตัวเลขที่เราวัด — บทนี้ไม่ได้ยิงโหลดพันสายใส่ server เพื่อพิสูจน์ สิ่งที่ยืนยันได้จริงมีสองอย่างคือขนาด task ที่เอกสาร Tokio ระบุไว้เอง (“single allocation and 64 bytes”) กับข้อเท็จจริงว่าคอนเนกชันที่รออยู่ไม่ได้ยึด OS thread ไว้อีกต่อไป ส่วน6 file ที่ระบุใน callout ด้านบนผ่าน build/run/test/clippy บน target native x86_64-unknown-linux-gnu (กลเม็ด musl + rust-lld ที่ #22/#23 ใช้ได้เพราะเป็น zero-crate ปลดระวางไปตั้งแต่บท 1 เพราะ tokio ลาก mio/socket2/libc เข้ามา) ส่วน bad_endian.rs compile ผ่านแล้ว รันจริง จนได้เลข 301,989,888 มา ส่วน bad_split.rs compile ไม่ผ่านโดยเจตนา เราจึงลงข้อความ error[E0597] ของมันไว้แทน output ทุก output ข้างบนคือข้อความจริงที่พิมพ์ออกมา และทุก error message คือข้อความจริงที่ rustc 1.97.1 ตอบกลับมา รวมถึง error[E0597] ของ split() — เราลอกมาตามที่ compiler พูด ไม่ได้ปรับให้ตรงกับที่เดาไว้ล่วงหน้า ข้อยกเว้นข้อเดียวที่เป็น การไล่เหตุผล ไม่ใช่ผลที่วัดมาคือกรณี “ลืม arm Ok(0) แล้ว loop ไม่จบ” ในหัวข้อกฎ correctness ซึ่งเราไม่ได้ ship file ให้รันดู ส่วนสิ่งที่ไม่ deterministic คือเลข ephemeral port ฝั่ง client กับลำดับบรรทัดของ cargo test ที่รันขนานกัน ซึ่งเปลี่ยนได้ทุกครั้งที่รัน
สรุปก่อนไปต่อ
หัวข้อที่มีชื่อว่า “สรุปก่อนไปต่อ”tokio::net ยกรูปของ std::net มาเกือบ 1 ต่อ 1 — bind / accept / connect ชื่อเดิม คืนค่าเดิม ต่างแค่มี .await ต่อท้ายและหนึ่งคอนเนกชันกลายเป็น1 task ไม่ใช่1 OS thread; verb ของ I/O ทั้งหมด (read write_all read_exact flush read_u32_le) อยู่บน extension trait จึงต้อง use tokio::io::{AsyncReadExt, AsyncWriteExt}; ไม่งั้นได้ error[E0599] ว่า method ไม่มีอยู่ ทั้งที่มันมีอยู่; reactor คือคนที่เอา fd ไปลงทะเบียนกับ epoll ผ่าน mio แล้วเรียก wake() ให้ task ตอน kernel บอกว่าพร้อม — คือ handshake เดิมของบท 1 ที่เปลี่ยนคนกดปุ่ม; wire ของ #22 ย้ายมาได้ทั้งดุ้นด้วย read_u32_le + read_exact และ endianness ผิดคือกับดักเงียบ ที่ compile ผ่านแล้วอ่านความยาว 18 เป็น 301,989,888; กฎ correctness สองข้อคือ read() คืน Ok(0) แปลว่าปิดสะอาด ส่วน read_exact ที่ขาดกลางคันคืน UnexpectedEof และถ้าอยากแยกสองอย่างนี้ออกจากกันจริงๆ ต้องอ่านหัว frame ด้วยมือแบบ read_frame_or_eof; และสุดท้าย read กับ write ที่อยู่คนละ task ต้องใช้ into_split() ที่คืน owned halves เพราะ split() คืนของที่ยืมมาซึ่งชน 'static ของ spawn ตรงๆ ด้วย error[E0597]
บท 5 เราจะเจอคำถามที่ .NET dev ถามเป็นข้อแรกเสมอ: accept loop ในบทนี้วนไม่รู้จบและไม่มีทางออก — แล้ว CancellationToken ของ Rust อยู่ไหน คำตอบคือไม่มี เพราะการยกเลิกใน async Rust คือการ drop future ทิ้ง บทหน้าเราจะใช้ select! แข่ง accept() กับสัญญาณหยุด เจอ timeout ที่คืน Result แทนการโยน exception และเจอกฎ cancellation safety ที่บอกว่า read_exact ของบทนี้ ไม่ปลอดภัย ที่จะวางไว้ใน select! ตรงๆ
บทนี้อิงต้นทางที่ลงวันที่กำกับ อ่านต่อได้โดยตรง:
- tokio 1.53.1 บน crates.io API (เข้าถึง 2026-07-27) — รุ่นที่ทั้งคอร์ส pin ไว้ ปล่อยเมื่อ 2026-07-20 และรายการ feature ที่มีให้เลือก รวมถึง
net,io-util,rt,rt-multi-thread,macros tokio::io::AsyncReadExt— docs.rs, tokio 1.53.1 (เข้าถึง 2026-07-27) — extension trait ที่แจกread/read_exact/read_u32_leและกติกาว่าreadคืนOk(0)เมื่อ peer ปิดคอนเนกชัน (patternOk(0) => returnของ echo server มาจากหน้านี้)AsyncReadExt::read_exact— docs.rs, tokio 1.53.1 (เข้าถึง 2026-07-27) — อ่านให้ครบbuf.len()หรือคืนErr(ErrorKind::UnexpectedEof)ถ้า stream จบก่อนTcpStream::into_split— docs.rs, tokio 1.53.1 (เข้าถึง 2026-07-27) — คืน owned halves ที่ย้ายข้าม task ได้ แลกกับ heap allocation หนึ่งครั้ง ต่างจากsplitที่คืนของยืมซึ่งย้ายเข้า task แยกไม่ได้tokio::task::spawn— docs.rs, tokio 1.53.1 (เข้าถึง 2026-07-27) — boundF: Future + Send + 'staticที่เป็นต้นเหตุของerror[E0597]ในหัวข้อsplit()#[tokio::main]— docs.rs, tokio 1.53.1 (เข้าถึง 2026-07-27) — default flavor คือmulti_threadจึงต้องเปิด featurert-multi-thread;flavor = "current_thread"ต้องการแค่rtกับmacros- Tokio Tutorial — Spawning (เข้าถึง 2026-07-27) — ที่มาของตัวเลข task ขนาด “single allocation and 64 bytes” และ pattern accept loop ที่ spawn 1 task ต่อหนึ่งคอนเนกชัน
Stream.ReadExactlyAsync— Microsoft Learn (เข้าถึง 2026-07-27) — คู่เทียบฝั่ง .NET ของread_exactเข้ามาใน .NET 7 (2022-11-08) ก่อนหน้านั้นต้องเขียน loop อ่านซ้ำเอง
เช็กความเข้าใจ — บทที่ 4
ข้อ 1 / 3คุณเขียน echo server ด้วย tokio::net::TcpStream แล้ว cargo build ตอบว่า error[E0599] no method named read found for struct tokio::net::TcpStream in the current scope สาเหตุคืออะไร?