diff --git a/CourseCode/ch3/README.md b/CourseCode/ch3/README.md index f0d563bee12217c96b24d43caae26612f6e92bc9..3423e3d9fd57e819e32281bdef6d365f11de134a 100644 --- a/CourseCode/ch3/README.md +++ b/CourseCode/ch3/README.md @@ -25,7 +25,7 @@ ch3/ ├── exp3/game_sticky_packets/ # ③ 粘包灾难版(游戏演示) ├── exp3/step3_sticky_packets/ # ③ 粘包灾难版(纯文本) ├── exp3/step3_framing_demo/ # ③ 长度前缀+JSON 正确处理 - ├── exp3/TCP_reliable/ # ③ TCP 可靠性演示(server/player1/player2) + ├── exp3/TCP_reliable/ # ③ TCP 可靠性与断线重连状态冲突演示(server/player1/player2) ├── exp4/p2p_lockstep_host/ # ④ P2P 锁步 Host ├── exp4/p2p_lockstep_client/ # ④ P2P 锁步 Client ├── exp5/cs_blocking_server/ # ⑤ 阻塞服务器(对照组) @@ -157,7 +157,7 @@ go run ./cmd/exp3/step3_framing_demo/client.go **现象**:即使连续发送消息,服务端仍能按条稳定解包并逐条打印。 -#### 演示 4:TCP_reliable(三个终端,重连状态冲突演示) +#### 演示 4:TCP_reliable(三个终端,游戏化重连状态冲突演示) ```powershell # 终端1:服务器 @@ -171,8 +171,10 @@ go run ./cmd/exp3/TCP_reliable/player2.go ``` **观察点**: -1. 玩家1会收到“Player2 DEAD”状态广播。 -2. 玩家2断线后重连,服务端会检测到“新连接”,并打印状态冲突提示(用于讲解仅靠 TCP 连接可靠性不足以保证游戏状态一致性)。 +1. `server.go` 会以战场网格显示双方位置:玩家1 为 `1`,玩家2 为 `2`,服务器判死后显示为 `X`。 +2. `player1.go` 会显示攻击方视角,自动等待玩家2上线并发送攻击,随后收到 `STATE: Player2 DEAD`。 +3. `player2.go` 会先断开旧连接,再在按回车后使用新连接重返场景,本地界面仍显示 `100` 血。 +4. 实验最终展示的是:玩家1所在旧连接中的字节流可靠送达,但玩家2重连后若没有会话恢复与状态同步,仍会出现“服务器判死、客户端满血”的状态冲突。 --- diff --git a/CourseCode/ch3/cmd/README.md b/CourseCode/ch3/cmd/README.md index f719bdaa35fc149e460a74085d0bad2b2bf5fe00..e5867bc4bd837781b9440b20942570fadf65d5d3 100644 --- a/CourseCode/ch3/cmd/README.md +++ b/CourseCode/ch3/cmd/README.md @@ -49,6 +49,7 @@ cmd/ - `go run ./cmd/exp3/TCP_reliable/server.go` - `go run ./cmd/exp3/TCP_reliable/player1.go` - `go run ./cmd/exp3/TCP_reliable/player2.go` + - `TCP_reliable` 为三终端配合演示:服务器显示战场网格,玩家1 负责攻击,玩家2 负责断线与重连,用于展示“旧连接内传输可靠”与“新连接状态未恢复”是两件不同的事。 - exp4: - `go run ./cmd/exp4/p2p_lockstep_host` diff --git a/CourseCode/ch3/cmd/exp3/TCP_reliable/player1.go b/CourseCode/ch3/cmd/exp3/TCP_reliable/player1.go index 23462af8348d694621cc2060b2869e6d148f3f2a..61c062f2b1184fa0e1dcbbb394992f2802015b64 100644 --- a/CourseCode/ch3/cmd/exp3/TCP_reliable/player1.go +++ b/CourseCode/ch3/cmd/exp3/TCP_reliable/player1.go @@ -2,7 +2,24 @@ package main import ( "fmt" + "io" "net" + "strings" + "time" + "unicode/utf8" +) + +const ( + p1Reset = "\033[0m" + p1Bold = "\033[1m" + p1Dim = "\033[2m" + p1Red = "\033[91m" + p1Green = "\033[92m" + p1Yellow = "\033[93m" + p1Cyan = "\033[96m" + p1White = "\033[97m" + p1BgBlue = "\033[44m" + p1Cls = "\033[2J\033[H" ) func main() { @@ -12,24 +29,111 @@ func main() { } defer conn.Close() - hp := 100 - fmt.Printf("[玩家1] 成功上线,当前血量: %d HP\n", hp) - fmt.Println("[玩家1] 正在等待玩家2上线...") + renderPlayer1("等待目标", 100, 100, "等待玩家2上线", "未开火", "") + + signalBuf := make([]byte, len("P2_ONLINE")) + if _, err := io.ReadFull(conn, signalBuf); err != nil { + fmt.Println("读取上线信号失败:", err) + return + } - // 核心修改:阻塞等待服务器发送 P2_ONLINE 信号 - signalBuf := make([]byte, 9) - conn.Read(signalBuf) + renderPlayer1("锁定完成", 100, 100, "已发现玩家2", "瞄准中", "") + time.Sleep(700 * time.Millisecond) if string(signalBuf) == "P2_ONLINE" { - fmt.Println("[玩家1] 发现玩家2已上线!") - // 严格保证在玩家2上线后才执行这句 - fmt.Println("[玩家1] 对玩家2发起致命一击!") - conn.Write([]byte("ATTACK")) - } - - // 接收服务器广播的死亡状态 - buf := make([]byte, 23) - conn.Read(buf) - fmt.Printf("[玩家1] 收到状态: %s\n", string(buf)) - fmt.Println("[玩家1] 确认: 玩家2已死亡") -} \ No newline at end of file + if _, err := conn.Write([]byte("ATTACK")); err != nil { + fmt.Println("发送攻击失败:", err) + return + } + renderPlayer1("发起攻击", 100, 0, "攻击指令已发送", "已开火", "") + } + + msg := make([]byte, len("STATE: Player2 DEAD ")) + if _, err := io.ReadFull(conn, msg); err != nil { + fmt.Println("读取结算消息失败:", err) + return + } + + renderPlayer1("攻击生效", 100, 0, "服务器确认玩家2死亡", "已命中", strings.TrimSpace(string(msg))) +} + +func renderPlayer1(phase string, hp1, hp2 int, event, action, state string) { + var b strings.Builder + b.WriteString(p1Cls) + fmt.Fprintf(&b, "%s%s%s TCP Reliable Lab 玩家终端 %s\n\n", p1Bold, p1BgBlue, p1White, p1Reset) + fmt.Fprintf(&b, " 视角:%s玩家1%s\n", p1Cyan, p1Reset) + fmt.Fprintf(&b, " 阶段:%s%s%s\n\n", p1Yellow, phase, p1Reset) + + rows := []string{ + p1Line("角色", "玩家1", "定位", "攻击方"), + p1Line("网络", "在线", "重连", "无"), + p1Line("动作", action, "目标", targetState(hp2)), + p1Line("事件", event, "状态", "旧连接"), + } + p1Panel(&b, rows) + + fmt.Fprintf(&b, "\n 玩家1生命 [%s]\n", p1HP(hp1)) + fmt.Fprintf(&b, " 玩家2生命 [%s]\n", p1HP(hp2)) + if state != "" { + fmt.Fprintf(&b, "\n 状态包:%s%s%s\n", p1Green, state, p1Reset) + } + fmt.Fprintf(&b, " 现 象:%s攻击与死亡消息在这条连接中可靠送达%s\n", p1Green, p1Reset) + fmt.Print(b.String()) +} + +func p1Panel(b *strings.Builder, rows []string) { + b.WriteString(" +----------------------------------------+\n") + for _, row := range rows { + fmt.Fprintf(b, " | %s |\n", p1Pad(row, 38)) + } + b.WriteString(" +----------------------------------------+\n") +} + +func p1Line(k1, v1, k2, v2 string) string { + return fmt.Sprintf("%s:%s %s:%s", k1, v1, k2, v2) +} + +func p1Pad(s string, width int) string { + w := p1Width(s) + if w >= width { + return s + } + return s + strings.Repeat(" ", width-w) +} + +func p1Width(s string) int { + width := 0 + for len(s) > 0 { + r, size := utf8.DecodeRuneInString(s) + s = s[size:] + if r <= 127 { + width++ + } else { + width += 2 + } + } + return width +} + +func p1HP(hp int) string { + total := 16 + fill := hp * total / 100 + if fill < 0 { + fill = 0 + } + if fill > total { + fill = total + } + color := p1Green + if hp <= 30 { + color = p1Red + } + return color + strings.Repeat("|", fill) + p1Reset + p1Dim + strings.Repeat(".", total-fill) + p1Reset + fmt.Sprintf(" %3d", hp) +} + +func targetState(hp int) string { + if hp <= 0 { + return "已击倒" + } + return "存活" +} diff --git a/CourseCode/ch3/cmd/exp3/TCP_reliable/player2.go b/CourseCode/ch3/cmd/exp3/TCP_reliable/player2.go index 37f7f363e217d36bb8115c1e8c1fb04092b2fe86..da80651b902110803de7745d75c4a4cf42e9e76f 100644 --- a/CourseCode/ch3/cmd/exp3/TCP_reliable/player2.go +++ b/CourseCode/ch3/cmd/exp3/TCP_reliable/player2.go @@ -3,32 +3,48 @@ package main import ( "bufio" "fmt" + "io" "net" "os" + "strings" + "time" + "unicode/utf8" +) + +const ( + p2Reset = "\033[0m" + p2Bold = "\033[1m" + p2Dim = "\033[2m" + p2Red = "\033[91m" + p2Green = "\033[92m" + p2Yellow = "\033[93m" + p2Cyan = "\033[96m" + p2White = "\033[97m" + p2BgBlue = "\033[44m" + p2Cls = "\033[2J\033[H" ) func main() { - // 第一次连接 conn, err := net.Dial("tcp", "127.0.0.1:8888") if err != nil { panic(err) } - hp := 100 - fmt.Printf("[玩家2] 首次上线成功,当前血量: %d HP\n", hp) + hp1 := 100 + hp2 := 100 + renderPlayer2("首次上线", hp1, hp2, "在线", "未重连", "已进入战场", "") + time.Sleep(600 * time.Millisecond) - fmt.Println("[玩家2] 突然断网 (切断连接)!") - - // 发生真实断线:直接关闭底层 Socket + renderPlayer2("网络断开", hp1, hp2, "离线", "未重连", "旧连接已关闭", "") conn.Close() - // 阻塞终端,等待讲师手动输入回车 - fmt.Print("\n>>> 请在终端按下回车键,模拟玩家2唤醒重启客户端 <<<") + fmt.Print("\n>>> 按回车模拟玩家2重连 <<<") reader := bufio.NewReader(os.Stdin) reader.ReadBytes('\n') - // 唤醒后:发起全新的 TCP 连接 - fmt.Println("\n[玩家2] 客户端唤醒,正在重新连接服务器...") + renderPlayer2("客户端重启", hp1, hp2, "离线", "准备重连", "正在建立新连接", "") + time.Sleep(600 * time.Millisecond) + reconnectConn, err := net.Dial("tcp", "127.0.0.1:8888") if err != nil { fmt.Printf("[玩家2] 重连失败: %v\n", err) @@ -36,11 +52,84 @@ func main() { } defer reconnectConn.Close() - // 客户端本地状态重置 - fmt.Printf("[玩家2] 重连成功!当前血量: %d HP\n", hp) - - - // 保持进程存活,接收服务端的欢迎消息 - welcomeBuf := make([]byte, 12) - reconnectConn.Read(welcomeBuf) -} \ No newline at end of file + renderPlayer2("重新入场", hp1, hp2, "在线", "已重连", "本地仍显示100血", "") + + welcomeBuf := make([]byte, len("WELCOME_BACK")) + if _, err := io.ReadFull(reconnectConn, welcomeBuf); err == nil { + renderPlayer2("冲突出现", hp1, hp2, "在线", "已重连", "服务器接受新连接", strings.TrimSpace(string(welcomeBuf))) + } +} + +func renderPlayer2(phase string, hp1, hp2 int, netState, reconnectState, event, state string) { + var b strings.Builder + b.WriteString(p2Cls) + fmt.Fprintf(&b, "%s%s%s TCP Reliable Lab 玩家终端 %s\n\n", p2Bold, p2BgBlue, p2White, p2Reset) + fmt.Fprintf(&b, " 视角:%s玩家2%s\n", p2Cyan, p2Reset) + fmt.Fprintf(&b, " 阶段:%s%s%s\n\n", p2Yellow, phase, p2Reset) + + rows := []string{ + p2Line("角色", "玩家2", "定位", "观察方"), + p2Line("网络", netState, "重连", reconnectState), + p2Line("动作", "观察中", "目标", "玩家1"), + p2Line("事件", event, "状态", "本地满血"), + } + p2Panel(&b, rows) + + fmt.Fprintf(&b, "\n 玩家1生命 [%s]\n", p2HP(hp1)) + fmt.Fprintf(&b, " 玩家2生命 [%s]\n", p2HP(hp2)) + if state != "" { + fmt.Fprintf(&b, " 状态包:%s%s%s\n", p2Green, state, p2Reset) + } + fmt.Fprintf(&b, " 现 象:%s新连接建立后,本地角色仍按满血显示%s\n", p2Green, p2Reset) + fmt.Print(b.String()) +} + +func p2Panel(b *strings.Builder, rows []string) { + b.WriteString(" +----------------------------------------+\n") + for _, row := range rows { + fmt.Fprintf(b, " | %s |\n", p2Pad(row, 38)) + } + b.WriteString(" +----------------------------------------+\n") +} + +func p2Line(k1, v1, k2, v2 string) string { + return fmt.Sprintf("%s:%s %s:%s", k1, v1, k2, v2) +} + +func p2Pad(s string, width int) string { + w := p2Width(s) + if w >= width { + return s + } + return s + strings.Repeat(" ", width-w) +} + +func p2Width(s string) int { + width := 0 + for len(s) > 0 { + r, size := utf8.DecodeRuneInString(s) + s = s[size:] + if r <= 127 { + width++ + } else { + width += 2 + } + } + return width +} + +func p2HP(hp int) string { + total := 16 + fill := hp * total / 100 + if fill < 0 { + fill = 0 + } + if fill > total { + fill = total + } + color := p2Green + if hp <= 30 { + color = p2Red + } + return color + strings.Repeat("|", fill) + p2Reset + p2Dim + strings.Repeat(".", total-fill) + p2Reset + fmt.Sprintf(" %3d", hp) +} diff --git a/CourseCode/ch3/cmd/exp3/TCP_reliable/server.go b/CourseCode/ch3/cmd/exp3/TCP_reliable/server.go index d64189f2d7c81358c236589bdbdb7c6611ee4afe..93d0476b24ec8ae30267dd6a99dfd4cf1449eeaf 100644 --- a/CourseCode/ch3/cmd/exp3/TCP_reliable/server.go +++ b/CourseCode/ch3/cmd/exp3/TCP_reliable/server.go @@ -3,62 +3,247 @@ package main import ( "fmt" "net" + "strings" + "sync" "time" ) +const ( + sReset = "\033[0m" + sBold = "\033[1m" + sDim = "\033[2m" + sRed = "\033[91m" + sGreen = "\033[92m" + sYellow = "\033[93m" + sCyan = "\033[96m" + sWhite = "\033[97m" + sBgBlue = "\033[44m" + sCls = "\033[2J\033[H" +) + +type playerSnapshot struct { + HP int + Online bool + Alive bool +} + +type serverView struct { + mu sync.Mutex + phase string + player1 playerSnapshot + player2 playerSnapshot + ghost bool + conflict bool + logs []string +} + func main() { + view := &serverView{ + phase: "等待连接", + player1: playerSnapshot{HP: 100, Alive: true}, + player2: playerSnapshot{HP: 100, Alive: true}, + } + view.push("开始监听 127.0.0.1:8888") + + go view.renderLoop() + listener, err := net.Listen("tcp", "127.0.0.1:8888") if err != nil { panic(err) } defer listener.Close() - fmt.Println("[服务器] 启动成功,开始监听端口 8888...") + conn1, err := listener.Accept() + if err != nil { + panic(err) + } + view.mu.Lock() + view.phase = "玩家1已接入" + view.player1.Online = true + view.mu.Unlock() + view.push("玩家1 已连接") - // 1. 接收玩家1 - conn1, _ := listener.Accept() - fmt.Printf("[服务器] 玩家1连接: %s (初始HP: 100)\n", conn1.RemoteAddr().String()) + conn2, err := listener.Accept() + if err != nil { + panic(err) + } + view.mu.Lock() + view.phase = "玩家2已上线" + view.player2.Online = true + view.mu.Unlock() + view.push("玩家2 已连接") - // 2. 接收玩家2 - conn2, _ := listener.Accept() - fmt.Printf("[服务器] 玩家2连接: %s (初始HP: 100)\n", conn2.RemoteAddr().String()) + if _, err = conn1.Write([]byte("P2_ONLINE")); err == nil { + view.push("已通知玩家1开始攻击") + } - // 3. 核心修改:通知玩家1,玩家2已经上线,可以开始攻击了 - conn1.Write([]byte("P2_ONLINE")) + buf := make([]byte, 64) + if _, err = conn1.Read(buf); err != nil { + view.push("读取攻击指令失败") + return + } - // 4. 等待玩家1发起攻击的指令 - buf := make([]byte, 1024) - conn1.Read(buf) - fmt.Println("[服务器] 接收到玩家1攻击指令,判定玩家2 HP 归零...") + view.mu.Lock() + view.phase = "服务器判定死亡" + view.player2.HP = 0 + view.player2.Alive = false + view.mu.Unlock() + view.push("收到 ATTACK,玩家2 HP 归零") - // 稍微等待1秒,确保玩家2的终端已经真实断开底层连接 - time.Sleep(1 * time.Second) + time.Sleep(time.Second) msg := []byte("STATE: Player2 DEAD ") - - // 发送给玩家1 _, err1 := conn1.Write(msg) - fmt.Printf("[服务器] 发送给玩家1: 23 bytes, err=%v <- 成功\n", err1) - - // 发送给已经断线的玩家2 _, err2 := conn2.Write(msg) - fmt.Printf("[服务器] 发送给玩家2: 23 bytes, err=%v <- 也\"成功\"!\n", err2) - fmt.Println("[服务器] 广播完成。") - fmt.Println("--------------------------------------------------") - fmt.Println("[服务器] 进入持续监听状态,等待断线玩家唤醒并重连...") + view.mu.Lock() + view.phase = "广播死亡结果" + view.ghost = err2 == nil + view.mu.Unlock() + if err1 == nil { + view.push("玩家1 收到死亡广播") + } + if err2 == nil { + view.push("旧连接写入仍显示成功") + } else { + view.push("旧连接写入失败") + } - // 无限循环,保持一直 listen 的状态 for { conn3, err := listener.Accept() if err != nil { continue } - - fmt.Printf("\n[服务器] 收到新连接: %s (检测到玩家重连)\n", conn3.RemoteAddr().String()) - fmt.Println("[服务器] 严重状态冲突:服务器内存中该玩家已死,但新连接的客户端依然满血!") - - // 可以在这里回复重连成功的消息,保持连接存活 - conn3.Write([]byte("WELCOME_BACK")) - } -} \ No newline at end of file + + view.mu.Lock() + view.phase = "检测到重连" + view.player2.Online = true + view.conflict = true + view.mu.Unlock() + view.push("玩家2 以新连接重新入场") + + if _, err := conn3.Write([]byte("WELCOME_BACK")); err == nil { + view.push("已发送重连欢迎包") + } + } +} + +func (v *serverView) push(msg string) { + v.mu.Lock() + defer v.mu.Unlock() + v.logs = append(v.logs, msg) + if len(v.logs) > 5 { + v.logs = v.logs[len(v.logs)-5:] + } +} + +func (v *serverView) renderLoop() { + ticker := time.NewTicker(120 * time.Millisecond) + defer ticker.Stop() + + for range ticker.C { + v.mu.Lock() + frame := v.buildFrameLocked() + v.mu.Unlock() + fmt.Print(sCls + frame) + } +} + +func (v *serverView) buildFrameLocked() string { + var b strings.Builder + + b.WriteString(sCls) + fmt.Fprintf(&b, "%s%s%s TCP Reliable Lab 玩家终端 %s\n\n", sBold, sBgBlue, sWhite, sReset) + fmt.Fprintf(&b, " 视角:%s服务器%s\n", sCyan, sReset) + fmt.Fprintf(&b, " 阶段:%s%s%s\n\n", sYellow, v.phase, sReset) + + sBattlefield(&b, v) + + fmt.Fprintf(&b, "\n 玩家1生命 [%s]\n", sHP(v.player1.HP)) + fmt.Fprintf(&b, " 玩家2生命 [%s]\n", sHP(v.player2.HP)) + fmt.Fprintf(&b, "\n 现 象:%s%s%s\n", effectColor(v.ghost, v.conflict), effectText(v.ghost, v.conflict), sReset) + fmt.Fprintf(&b, "\n 记录:\n") + for _, log := range v.logs { + fmt.Fprintf(&b, " %s• %s%s\n", sDim, log, sReset) + } + + return b.String() +} + +func sBattlefield(b *strings.Builder, v *serverView) { + grid := [5][9]string{} + for y := range grid { + for x := range grid[y] { + grid[y][x] = "·" + } + } + + p1x, p1y := 1, 2 + p2x, p2y := 7, 2 + if v.player1.Online { + grid[p1y][p1x] = sCyan + "1" + sReset + } + + if v.player2.HP == 0 { + grid[p2y][p2x] = sRed + "X" + sReset + } else if v.player2.Online { + grid[p2y][p2x] = sYellow + "2" + sReset + } + + if v.player1.Online && v.player2.Online && v.player2.HP > 0 { + grid[p2y][3] = sYellow + ">" + sReset + grid[p2y][4] = sYellow + ">" + sReset + grid[p2y][5] = sYellow + ">" + sReset + } + + b.WriteString("\n 战场网格:\n") + b.WriteString(" +-------------------+\n") + for y := range grid { + var row strings.Builder + for x := range grid[y] { + row.WriteString(grid[y][x]) + if x != len(grid[y])-1 { + row.WriteString(" ") + } + } + fmt.Fprintf(b, " | %s |\n", row.String()) + } + b.WriteString(" +-------------------+\n") + fmt.Fprintf(b, " %s1%s=玩家1 %s2%s=玩家2 %sX%s=服务器判死\n", sCyan, sReset, sYellow, sReset, sRed, sReset) +} + +func sHP(hp int) string { + total := 16 + fill := hp * total / 100 + if fill < 0 { + fill = 0 + } + if fill > total { + fill = total + } + color := sGreen + if hp <= 30 { + color = sRed + } + return color + strings.Repeat("|", fill) + sReset + sDim + strings.Repeat(".", total-fill) + sReset + fmt.Sprintf(" %3d", hp) +} + +func effectText(ghost, conflict bool) string { + if conflict { + return "服务器记忆中玩家2已死,但新连接客户端仍以满血出现" + } + if ghost { + return "玩家2旧连接已断开,但写入时看起来仍然成功" + } + return "等待实验事件发生" +} + +func effectColor(ghost, conflict bool) string { + if conflict { + return sRed + } + if ghost { + return sYellow + } + return sGreen +} diff --git a/CourseCode/ch4/README.md b/CourseCode/ch4/README.md new file mode 100644 index 0000000000000000000000000000000000000000..23941bd5e17d8c94314ad7acb8bdb32a779b270d --- /dev/null +++ b/CourseCode/ch4/README.md @@ -0,0 +1,335 @@ +# 英雄集结演示实验 —— 代码与演示说明 + +> 本目录 (`ch4/`) 是一个**独立的 Go module**,围绕 Go 并发与后端架构设计,提供 6 个可运行的演示实验。 +> 其中实验 1~5 为单机并发教学演示;实验 6 使用 Docker 中的 Redis 与 Docker 中的 PostgreSQL,演示分层存储架构。 + +--- + +## 目录结构 + +```text +ch4/ +├── go.mod # 独立模块 ch4 +├── go.sum # 依赖校验文件 +├── README.md # ← 本文件 +├── 英雄集结演示实验.md # 实验要求与讲解要点 +└── cmd/ + ├── README.md # cmd 子目录说明 + ├── exp1/local_serial_loop_demo/ # ① 本地串行主循环 + ├── exp1/network_serial_server_demo/# ① 网络版:服务器串行收包 + ├── exp1/network_goroutine_server_demo/ # ① 网络版:服务器独立线程收包 + ├── exp1/network_event_driven_sync_demo/ # ① 网络版:事件驱动 + 增量同步 + ├── exp2/wrong/ # ② 无锁竞态错误版 + ├── exp2/right/ # ② Mutex 修复版 + ├── exp3/busy_wait/ # ③ 忙等错误版 + ├── exp3/cond_wait/ # ③ Cond 修复版 + ├── exp4/channel_semaphore/ # ④ Channel 信号量 + ├── exp4/channel_timeout_lock/ # ④ Channel 超时锁 + ├── exp4/rw_mutex/ # ④ RWMutex 读写分离 + ├── exp5/sync_pool_demo/ # ⑤ sync.Pool 性能优化 + └── exp6/storage_arch/ # ⑥ Redis + PostgreSQL 分层存储 +``` + +--- + +## 环境搭建 + +### 前提条件 + +| 依赖 | 版本 | 说明 | +|------|------|------| +| Go | ≥ 1.21 | 所有实验使用 Go 1.21+ | +| Docker Desktop | 已安装并启动 | 仅实验 6 需要,用于启动 Redis 与 PostgreSQL 容器 | + +### 首次构建 + +```powershell +# 进入 ch4 目录 +cd ch4 + +# 验证实验 1~6 +go build ./cmd/exp1/... +go build ./cmd/exp2/... +go build ./cmd/exp3/... +go build ./cmd/exp4/... +go build ./cmd/exp5/... +go build ./cmd/exp6/... +``` + + + +### 实验六运行前准备 + +实验 6 不是纯本地模拟,它依赖两个真实环境,并且这两个环境现在都通过 Docker Desktop 提供: + +- Redis:Docker 容器 `ch4-redis` +- PostgreSQL:Docker 容器 `ch4-postgres` + +#### 1. 启动 Docker Desktop + +先确认 Docker Desktop 已经启动;如果 Docker 没启动,下面的 `docker run` 和 `docker start` 都会失败。 + +检查命令: + +```powershell +docker ps +``` + +如果能正常返回容器列表,说明 Docker 已经可用。 + +#### 2. 启动 Redis 容器 + +第一次创建并启动: + +```powershell +docker run -d --name ch4-redis -p 6379:6379 redis:7 +``` + +如果容器已经创建过,直接启动: + +```powershell +docker start ch4-redis +``` + +检查 Redis 是否正常: + +```powershell +docker exec -it ch4-redis redis-cli ping +``` + +预期输出: + +```text +PONG +``` + +#### 3. 启动 PostgreSQL 容器 + +第一次创建并启动: + +```powershell +docker run -d --name ch4-postgres -e POSTGRES_USER=你的用户名 -e POSTGRES_PASSWORD=你的密码 -e POSTGRES_DB=postgres -p 5432:5432 postgres:16 +``` + +如果容器已经创建过,直接启动: + +```powershell +docker start ch4-postgres +``` + +这套配置对应的 PostgreSQL 连接串是: + +```text +postgres://你的用户名:你的密码@127.0.0.1:5432/postgres?sslmode=disable +``` + +#### 4. 设置实验六环境变量 + +在 PowerShell 中设置: + +```powershell +$env:REDIS_ADDR="127.0.0.1:6379" +$env:PG_DSN="postgres://你的用户名:你的密码@127.0.0.1:5432/postgres?sslmode=disable" +``` + +其中: + +- `REDIS_ADDR` 指向 Docker 中映射到本机 `6379` 的 Redis +- `PG_DSN` 指向 Docker 中映射到本机 `5432` 的 PostgreSQL +- 代码不会内置任何个人 PostgreSQL 账号密码,使用前请由每位使用者自行设置 `PG_DSN` + +#### 5. 运行实验六 + +```powershell +go run ./cmd/exp6/storage_arch +``` + +程序启动后会自动完成这些动作: + +- 连接 Redis +- 连接 PostgreSQL +- 自动创建 `players` 与 `game_configs` 两张表 +- 自动写入演示初始数据 +- 依次演示 `Write Through` 和 `Cache Aside` + +#### 6. 常见失败原因 + +- `docker` 命令报错:通常是 Docker Desktop 没有启动。 +- Redis 连接失败:通常是 `ch4-redis` 容器没启动,或 `6379` 端口未映射成功。 +- PostgreSQL 连接失败:通常是 `ch4-postgres` 容器没启动,或 `PG_DSN` 与容器账号密码不一致。 + +> 实验 6 依赖真实 Redis 与真实 PostgreSQL;实验 1~5 不需要额外环境。 + +--- + +## 各步骤演示操作 + +### Step 1 — 突破单线程瓶颈 + +**知识点**:串行阻塞、Goroutine 解耦、事件驱动、增量同步。 + +```powershell +go run ./cmd/exp1/local_serial_loop_demo +go run ./cmd/exp1/network_serial_server_demo server +go run ./cmd/exp1/network_serial_server_demo client +go run ./cmd/exp1/network_goroutine_server_demo server +go run ./cmd/exp1/network_goroutine_server_demo client +go run ./cmd/exp1/network_event_driven_sync_demo server +go run ./cmd/exp1/network_event_driven_sync_demo client +``` + +**观察点**: + +- `local_serial_loop_demo` 展示最原始的本地串行主循环,慢玩家会直接拖慢整帧。 +- `network_serial_server_demo` 用 2 个终端模拟真实客户端/服务器,但服务器仍按连接顺序串行收包。 +- `network_goroutine_server_demo` 保持真实网络分离,同时把每条连接的收包交给独立 goroutine,输出风格与串行网络版保持一致,方便直接对比。 +- `network_event_driven_sync_demo` 进一步展示真实网络下的事件驱动与增量同步:服务器按 tick 前进,只处理已经到达的输入,不等最慢玩家。 + +**课堂结论**: + +- 先看 `local_serial_loop_demo`,能理解“慢输入会拖帧”这个最基本现象。 +- 再看 `network_serial_server_demo`,能看清问题不只是“玩家慢”,而是服务器在串行等待某条慢连接。 +- 切到 `network_goroutine_server_demo` 后,慢连接仍然存在,但不会再让服务器主循环卡死在某一条连接上。 +- 最后看 `network_event_driven_sync_demo`,能补齐“只消费已到达事件、只同步变化状态”的思路。 + +--- + +### Step 2 — 临界区与数据竞争 + +**知识点**:Race Condition、临界区、`sync.Mutex`。 + +```powershell +go run ./cmd/exp2/wrong +go run ./cmd/exp2/right +``` + +**观察点**: + +- 无锁版会出现“唯一物品被多次领取”,而且重复领取人数不必每轮都一样。 +- 加锁版只允许一个玩家成功获得宝物。 +- 对照输出可以看到临界区保护前后的差异。 + +**课堂结论**: + +- 业务规则写得再正确,只要“检查是否可拿”和“标记已拿走”不在同一个临界区里,就仍然会出现竞态窗口。 +- `sync.Mutex` 修复的不是“谁先抢到”的业务逻辑,而是把检查与修改变成原子的一段。 + +--- + +### Step 3 — 告别忙等 + +**知识点**:忙等、`sync.Cond`、等待队列、虚假唤醒。 + +```powershell +go run ./cmd/exp3/busy_wait +go run ./cmd/exp3/cond_wait +``` + +**观察点**: + +- 忙等版会出现极高的空转次数。 +- `sync.Cond` 版会展示等待、唤醒和重新检查条件。 +- `for` 循环重检逻辑可用于讲解“为什么不能只用 `if`”。 + +**课堂结论**: + +- 忙等的问题不是功能错误,而是没有库存时线程仍在不停抢锁和检查,白白消耗 CPU。 +- `Cond.Wait()` 会先释放锁再休眠,被唤醒后重新抢锁,所以必须用 `for` 重检条件来防住虚假唤醒。 + +--- + +### Step 4 — 锁的进阶技巧与粒度优化 + + +**知识点**:Channel 信号量、超时锁、`sync.RWMutex`。 + +```powershell +go run ./cmd/exp4/channel_semaphore +go run ./cmd/exp4/channel_timeout_lock +go run ./cmd/exp4/rw_mutex +``` + +**观察点**: + +- `channel_semaphore` 展示如何用 Channel 控制最大并发数。 +- `channel_timeout_lock` 展示等待超时后如何快速降级,不把协程永远卡死。 +- `rw_mutex` 除了展示读多写少场景下的并发读取优势,也显式展示“写者排队后”和“写入进行中”到来的读者都会被阻塞。 +- 这一组实验用于对比“锁不仅只有 Mutex 一种用法”。 + +**课堂结论**: + +- Channel 不只适合传消息,也很适合表达“许可数”和“带超时的竞争”。 +- `RWMutex` 的重点不只是“多个读者可以并发”,还包括一旦写者开始排队,后续读者也要等待,避免写者长期饥饿。 + +--- + +### Step 5 — 高并发性能榨取 + +**知识点**:`sync.Pool`、对象复用、GC 压力、尾延迟。 + +```powershell +go run ./cmd/exp5/sync_pool_demo +``` + +**观察点**: + +- 对比 `new(bytes.Buffer)` 与 `sync.Pool` 的总耗时。 +- 观察分配次数、GC 次数或尾延迟差异。 +- 理解“池化不是为了炫技,而是为了减少重复分配”。 + +**课堂结论**: + +- `sync.Pool` 不保证每次都命中,但在高频临时对象场景里,通常能明显减少分配与 GC 压力。 +- 这里演示的是对象池,不是数据库连接池;两者都叫“池”,但适用资源类型和约束不同。 + +--- + +### Step 6 — 游戏数据分层存储架构 + +**知识点**:Redis、PostgreSQL、Write Through、Cache Aside。 + +**运行前环境要求**: + +- Docker Desktop 已启动 +- Redis 容器 `ch4-redis` 已运行并映射 `6379:6379` +- PostgreSQL 容器 `ch4-postgres` 已运行并映射 `5432:5432` +- PowerShell 中已设置 `PG_DSN`, $env:PG_DSN="postgres://账号:密码@127.0.0.1:5432/postgres?sslmode=disable" + +**运行命令**: + +```powershell +go run ./cmd/exp6/storage_arch +``` + +**观察点**: + +- `Write Through`:先写 PostgreSQL,再同步 Redis。 +- `Cache Aside`:先查 Redis,未命中时再查 PostgreSQL 并回填。 +- 终端会输出中文日志,便于课堂中逐步讲解数据流。 +- 程序会自动初始化表结构与演示数据,适合直接投屏讲解。 + +--- + +## 各步骤资源与依赖对照 + +| 步骤 | 额外依赖 | 说明 | +|------|----------|------| +| 1 | 无 | 单机并发演示 | +| 2 | 无 | 单机并发演示 | +| 3 | 无 | 单机并发演示 | +| 4 | 无 | 单机并发演示 | +| 5 | 无 | 单机性能演示 | +| 6 | Redis + PostgreSQL | Redis 与 PostgreSQL 都建议使用 Docker 容器启动 | + +--- + +## 核心知识点对照表 + +| 步骤 | 演示目标 | 关键结构/机制 | +|------|----------|----------------| +| 1 | 单线程阻塞 vs 并发解耦 | `goroutine`、事件驱动 | +| 2 | 竞态条件复现与修复 | `sync.Mutex` | +| 3 | 忙等替换为阻塞等待 | `sync.Cond` | +| 4 | 锁的高级技巧 | Channel、超时锁、`sync.RWMutex` | +| 5 | 高并发性能优化 | `sync.Pool`、GC、尾延迟 | +| 6 | 分层存储架构 | Redis、PostgreSQL、Write Through、Cache Aside | diff --git a/CourseCode/ch4/cmd/README.md b/CourseCode/ch4/cmd/README.md new file mode 100644 index 0000000000000000000000000000000000000000..f5c48afb0a2ef9d55d22d21d99a4afa53ad1f9e2 --- /dev/null +++ b/CourseCode/ch4/cmd/README.md @@ -0,0 +1,62 @@ +# ch4/cmd 目录索引(不改程序版) + +本目录仅做结构整理说明,不改动任何程序代码与运行路径。 + +## 目录总览 + +```text +cmd/ +├─ exp1/ +│ ├─ local_serial_loop_demo/ +│ ├─ network_serial_server_demo/ +│ ├─ network_goroutine_server_demo/ +│ └─ network_event_driven_sync_demo/ +├─ exp2/ +│ ├─ wrong/ +│ └─ right/ +├─ exp3/ +│ ├─ busy_wait/ +│ └─ cond_wait/ +├─ exp4/ +│ ├─ channel_semaphore/ +│ ├─ channel_timeout_lock/ +│ └─ rw_mutex/ +├─ exp5/ +│ └─ sync_pool_demo/ +└─ exp6/ + └─ storage_arch/ +``` + +## 运行入口 + +- exp1: + - `go run ./cmd/exp1/local_serial_loop_demo` + - `go run ./cmd/exp1/network_serial_server_demo server` + - `go run ./cmd/exp1/network_serial_server_demo client` + - `go run ./cmd/exp1/network_goroutine_server_demo server` + - `go run ./cmd/exp1/network_goroutine_server_demo client` + - `go run ./cmd/exp1/network_event_driven_sync_demo server` + - `go run ./cmd/exp1/network_event_driven_sync_demo client` +- exp2: + - `go run ./cmd/exp2/wrong` + - `go run ./cmd/exp2/right` +- exp3: + - `go run ./cmd/exp3/busy_wait` + - `go run ./cmd/exp3/cond_wait` +- exp4: + - `go run ./cmd/exp4/channel_semaphore` + - `go run ./cmd/exp4/channel_timeout_lock` + - `go run ./cmd/exp4/rw_mutex` +- exp5: + - `go run ./cmd/exp5/sync_pool_demo` +- exp6: + - `go run ./cmd/exp6/storage_arch` + +## 命名说明 + +- `exp1` 现在使用四个入口,对应“本地串行主循环版”“网络串行收包版”“网络独立线程收包版”“网络事件驱动 + 增量同步版”。 +- `exp2` 与 `exp3` 保持“错误版/修复版”的结构,便于课堂中前后对照。 +- `exp4` 将 Channel 技巧拆成“信号量”和“超时锁”两个独立入口,再配合 `RWMutex` 演示锁粒度。 +- `exp5` 聚焦 `sync.Pool` 的性能优化演示,不额外引入数据库连接池代码。 +- `exp6` 使用真实 Redis 与 PostgreSQL 演示分层存储架构。 +- 当前目录保持原路径不变,避免影响既有讲义、截图和运行命令。 diff --git a/CourseCode/ch4/cmd/exp1/local_serial_loop_demo/main.go b/CourseCode/ch4/cmd/exp1/local_serial_loop_demo/main.go new file mode 100644 index 0000000000000000000000000000000000000000..a07aff64361ebe42bbf407ca1c14f8d9e0beea7f --- /dev/null +++ b/CourseCode/ch4/cmd/exp1/local_serial_loop_demo/main.go @@ -0,0 +1,98 @@ +package main + +import ( + "fmt" + "time" +) + +const trackWidth = 20 + +type InputEvent struct { + PlayerID int + Action string + DeltaX int + Latency time.Duration +} + +func formatMS(d time.Duration) string { + return fmt.Sprintf("%.1fms", float64(d.Microseconds())/1000.0) +} + +func clamp(x, lo, hi int) int { + if x < lo { + return lo + } + if x > hi { + return hi + } + return x +} + +func renderPositions(pos map[int]int) string { + return fmt.Sprintf("P1(x=%d) P2(x=%d) P3(x=%d) P4(x=%d)", pos[1], pos[2], pos[3], pos[4]) +} + +func runFrameSingle( + frame int, + order []int, + events map[int]InputEvent, + positions map[int]int, + budget time.Duration, + dragHint bool, +) { + fmt.Printf("\n[Frame %d] 开始收集输入...\n", frame) + start := time.Now() + + for _, pid := range order { + ev := events[pid] + time.Sleep(ev.Latency) + waited := time.Since(start) + + note := "" + if pid == 4 { + note = " <- 卡顿源头" + } else if dragHint { + note = " <- 流畅玩家被拖累" + } + + fmt.Printf(" -> 收到玩家%d事件: %s(%+d) (累计等待%s)%s\n", + ev.PlayerID, ev.Action, ev.DeltaX, formatMS(waited), note) + + positions[ev.PlayerID] = clamp(positions[ev.PlayerID]+ev.DeltaX, 0, trackWidth) + } + + cost := time.Since(start) + fmt.Printf("[Frame %d] 位置快照: %s\n", frame, renderPositions(positions)) + fmt.Printf("[Frame %d] 结束, 耗时%s (目标<100ms)\n", frame, formatMS(cost)) + + if cost > budget { + fmt.Println(" 警告: 帧时间超标, 游戏体验卡顿") + fmt.Println(" 原因: 单线程串行收包, 慢事件阻塞了后续事件处理") + } +} + +func main() { + fmt.Println("=== 实验一:突破单线程瓶颈 / 本地串行主循环版 ===") + fmt.Println("场景: 4 名玩家同时上报输入,玩家4 处于“地铁断流”环境,延迟固定 500ms。") + fmt.Println("目标: 观察单线程主循环为什么会把慢玩家的延迟传染给整帧。") + + positions := map[int]int{1: 2, 2: 6, 3: 10, 4: 14} + + frame1 := map[int]InputEvent{ + 1: {PlayerID: 1, Action: "MOVE", DeltaX: +1, Latency: 11 * time.Millisecond}, + 2: {PlayerID: 2, Action: "MOVE", DeltaX: +1, Latency: 12 * time.Millisecond}, + 3: {PlayerID: 3, Action: "MOVE", DeltaX: -1, Latency: 13 * time.Millisecond}, + 4: {PlayerID: 4, Action: "MOVE", DeltaX: -1, Latency: 500 * time.Millisecond}, + } + runFrameSingle(1, []int{1, 2, 3, 4}, frame1, positions, 523*time.Millisecond, false) + + frame2 := map[int]InputEvent{ + 1: {PlayerID: 1, Action: "MOVE", DeltaX: +1, Latency: 11 * time.Millisecond}, + 2: {PlayerID: 2, Action: "MOVE", DeltaX: +1, Latency: 12 * time.Millisecond}, + 3: {PlayerID: 3, Action: "MOVE", DeltaX: -1, Latency: 13 * time.Millisecond}, + 4: {PlayerID: 4, Action: "MOVE", DeltaX: -1, Latency: 500 * time.Millisecond}, + } + runFrameSingle(2, []int{4, 1, 2, 3}, frame2, positions, 523*time.Millisecond, true) + + fmt.Println("\n[提示] 再运行 network_serial_server_demo、network_goroutine_server_demo 或 network_event_driven_sync_demo,对比不同层次的解耦方式。") +} diff --git a/CourseCode/ch4/cmd/exp1/network_event_driven_sync_demo/main.go b/CourseCode/ch4/cmd/exp1/network_event_driven_sync_demo/main.go new file mode 100644 index 0000000000000000000000000000000000000000..c7a53ec269770a5614071426948129d3ee38883d --- /dev/null +++ b/CourseCode/ch4/cmd/exp1/network_event_driven_sync_demo/main.go @@ -0,0 +1,478 @@ +package main + +import ( + "bufio" + "errors" + "flag" + "fmt" + "io" + "net" + "os" + "sort" + "strconv" + "strings" + "sync" + "time" +) + +const ( + trackWidth = 20 + defaultAddr = "127.0.0.1:9107" + expectedPlayers = 4 + tickInterval = 100 * time.Millisecond + totalTicks = 8 +) + +var logMu sync.Mutex + +type InputEvent struct { + PlayerID int + Seq int + Action string + DeltaX int + LocalAt time.Duration + Delay time.Duration +} + +type ClientSession struct { + ID int + Conn net.Conn + Reader *bufio.Reader + Writer *bufio.Writer +} + +func formatMS(d time.Duration) string { + return fmt.Sprintf("%.1fms", float64(d.Microseconds())/1000.0) +} + +func logf(format string, args ...any) { + logMu.Lock() + defer logMu.Unlock() + fmt.Printf(format, args...) +} + +func logln(args ...any) { + logMu.Lock() + defer logMu.Unlock() + fmt.Println(args...) +} + +func printDivider(title string) { + logMu.Lock() + defer logMu.Unlock() + fmt.Printf("\n========== %s ==========\n", title) +} + +func clamp(x, lo, hi int) int { + if x < lo { + return lo + } + if x > hi { + return hi + } + return x +} + +func renderPositions(pos map[int]int) string { + return fmt.Sprintf("P1(x=%d) P2(x=%d) P3(x=%d) P4(x=%d)", pos[1], pos[2], pos[3], pos[4]) +} + +func clientScript(playerID int) []InputEvent { + switch playerID { + case 1: + return []InputEvent{ + {PlayerID: 1, Seq: 1, Action: "MOVE", DeltaX: +1, LocalAt: 0 * time.Millisecond, Delay: 20 * time.Millisecond}, + {PlayerID: 1, Seq: 2, Action: "MOVE", DeltaX: +1, LocalAt: 200 * time.Millisecond, Delay: 20 * time.Millisecond}, + } + case 2: + return []InputEvent{ + {PlayerID: 2, Seq: 1, Action: "MOVE", DeltaX: +1, LocalAt: 0 * time.Millisecond, Delay: 30 * time.Millisecond}, + {PlayerID: 2, Seq: 2, Action: "MOVE", DeltaX: +1, LocalAt: 200 * time.Millisecond, Delay: 30 * time.Millisecond}, + } + case 3: + return []InputEvent{ + {PlayerID: 3, Seq: 1, Action: "MOVE", DeltaX: -1, LocalAt: 0 * time.Millisecond, Delay: 40 * time.Millisecond}, + {PlayerID: 3, Seq: 2, Action: "MOVE", DeltaX: -1, LocalAt: 200 * time.Millisecond, Delay: 40 * time.Millisecond}, + } + case 4: + return []InputEvent{ + {PlayerID: 4, Seq: 1, Action: "MOVE", DeltaX: -1, LocalAt: 0 * time.Millisecond, Delay: 500 * time.Millisecond}, + {PlayerID: 4, Seq: 2, Action: "MOVE", DeltaX: -1, LocalAt: 200 * time.Millisecond, Delay: 500 * time.Millisecond}, + } + default: + return nil + } +} + +type incomingEvent struct { + Event InputEvent + ReceivedAt time.Duration + Err error +} + +func runServer(addr string) error { + printDivider("实验一 / Network Event-Driven Server") + logln("监听地址:", addr) + logln("运行方式: 再打开 1 个终端执行 client 模式。") + logln("目标: 服务器按 tick 前进,不等最慢客户端,只消费已经到达的事件,并做增量同步。") + + ln, err := net.Listen("tcp", addr) + if err != nil { + return err + } + defer ln.Close() + + sessions := make(map[int]*ClientSession, expectedPlayers) + for len(sessions) < expectedPlayers { + conn, err := ln.Accept() + if err != nil { + return err + } + + session, err := acceptClient(conn) + if err != nil { + conn.Close() + return err + } + if _, exists := sessions[session.ID]; exists { + session.Conn.Close() + return fmt.Errorf("玩家%d 重复连接", session.ID) + } + sessions[session.ID] = session + logf("[连接] 玩家%-2d 已连接 (%d/%d)\n", session.ID, len(sessions), expectedPlayers) + } + defer closeSessions(sessions) + + for _, session := range sortedSessions(sessions) { + if err := writeLine(session.Writer, "BEGIN"); err != nil { + return err + } + } + + positions := map[int]int{1: 2, 2: 6, 3: 10, 4: 14} + eventsCh := make(chan incomingEvent, 32) + var recvWG sync.WaitGroup + start := time.Now() + + for _, session := range sortedSessions(sessions) { + recvWG.Add(1) + go func(session *ClientSession) { + defer recvWG.Done() + if err := receiveLoop(session, start, eventsCh); err != nil { + eventsCh <- incomingEvent{Err: err} + } + }(session) + } + + go func() { + recvWG.Wait() + close(eventsCh) + }() + + pending := make([]incomingEvent, 0, 8) + dirty := make(map[int]struct{}) + ticker := time.NewTicker(tickInterval) + defer ticker.Stop() + + nextTick := 1 + for nextTick <= totalTicks { + select { + case incoming, ok := <-eventsCh: + if !ok { + eventsCh = nil + continue + } + if incoming.Err != nil { + return incoming.Err + } + pending = append(pending, incoming) + + case tickTime := <-ticker.C: + printDivider(fmt.Sprintf("Tick %d / Network Event-Driven Server", nextTick)) + logf("[主循环] Tick %-2d 到达,当前时间=%s\n", nextTick, formatMS(tickTime.Sub(start))) + + if len(pending) == 0 { + logln("[处理] 本 tick 没有新事件,服务器继续前进,不等待慢客户端。") + } else { + sort.Slice(pending, func(i, j int) bool { + return pending[i].ReceivedAt < pending[j].ReceivedAt + }) + logln("[到达] 本 tick 前已到达的事件:") + for _, item := range pending { + ev := item.Event + logf(" - 玩家%-2d seq=%d %-4s %+d 输入=%-8s 延迟=%-8s 到达=%s\n", + ev.PlayerID, ev.Seq, ev.Action, ev.DeltaX, + formatMS(ev.LocalAt), formatMS(ev.Delay), formatMS(item.ReceivedAt)) + } + for _, item := range pending { + ev := item.Event + positions[ev.PlayerID] = clamp(positions[ev.PlayerID]+ev.DeltaX, 0, trackWidth) + dirty[ev.PlayerID] = struct{}{} + logf("[处理] 玩家%-2d seq=%d %-4s %+d -> 位置=%d\n", + ev.PlayerID, ev.Seq, ev.Action, ev.DeltaX, positions[ev.PlayerID]) + } + pending = pending[:0] + } + + if len(dirty) == 0 { + logln("[增量同步] 本 tick 无状态变化,不发送增量。") + } else { + playerIDs := make([]int, 0, len(dirty)) + for pid := range dirty { + playerIDs = append(playerIDs, pid) + } + sort.Ints(playerIDs) + logln("[增量同步] 本 tick 只同步发生变化的玩家:") + for _, pid := range playerIDs { + logf(" - 玩家%-2d x=%d\n", pid, positions[pid]) + } + clear(dirty) + } + + logf("[快照] Tick %-2d 当前世界状态: %s\n", nextTick, renderPositions(positions)) + nextTick++ + } + } + + for _, session := range sessions { + if err := writeLine(session.Writer, "DONE"); err != nil { + return err + } + } + + printDivider("Network Event-Driven Server 结束") + logln("提示: 这里展示的是“谁先到先处理”,慢客户端只会影响自己的更新时刻,不再阻塞下一 tick。") + return nil +} + +func receiveLoop(session *ClientSession, start time.Time, eventsCh chan<- incomingEvent) error { + for { + line, err := readLine(session.Reader) + if err != nil { + if errors.Is(err, io.EOF) { + return nil + } + return err + } + + ev, err := parseEventLine(line) + if err != nil { + return err + } + eventsCh <- incomingEvent{ + Event: ev, + ReceivedAt: time.Since(start), + } + } +} + +func acceptClient(conn net.Conn) (*ClientSession, error) { + reader := bufio.NewReader(conn) + writer := bufio.NewWriter(conn) + + line, err := readLine(reader) + if err != nil { + return nil, err + } + parts := strings.Fields(line) + if len(parts) != 2 || parts[0] != "HELLO" { + return nil, fmt.Errorf("非法握手消息: %q", line) + } + + playerID, err := strconv.Atoi(parts[1]) + if err != nil { + return nil, fmt.Errorf("非法玩家编号: %q", parts[1]) + } + return &ClientSession{ID: playerID, Conn: conn, Reader: reader, Writer: writer}, nil +} + +func runClient(addr string) error { + printDivider("实验一 / Network Event-Driven Client") + logln("连接服务器:", addr) + logf("当前进程会同时模拟 %d 名玩家,每名玩家保持一条独立连接。\n", expectedPlayers) + + var wg sync.WaitGroup + errCh := make(chan error, expectedPlayers) + for playerID := 1; playerID <= expectedPlayers; playerID++ { + wg.Add(1) + go func(playerID int) { + defer wg.Done() + if err := runOnePlayerClient(addr, playerID); err != nil { + errCh <- err + } + }(playerID) + } + wg.Wait() + close(errCh) + + for err := range errCh { + if err != nil { + return err + } + } + + printDivider("Network Event-Driven Client 结束") + logln("所有玩家脚本执行完毕。") + return nil +} + +func runOnePlayerClient(addr string, playerID int) error { + script := clientScript(playerID) + if len(script) == 0 { + return fmt.Errorf("玩家%d 没有预设脚本", playerID) + } + + conn, err := net.Dial("tcp", addr) + if err != nil { + return err + } + defer conn.Close() + + reader := bufio.NewReader(conn) + writer := bufio.NewWriter(conn) + + if err := writeLine(writer, fmt.Sprintf("HELLO %d", playerID)); err != nil { + return err + } + + for { + line, err := readLine(reader) + if err != nil { + if errors.Is(err, io.EOF) { + return nil + } + return err + } + + switch line { + case "BEGIN": + start := time.Now() + for _, ev := range script { + reportAt := ev.LocalAt + ev.Delay + wait := time.Until(start.Add(reportAt)) + if wait > 0 { + time.Sleep(wait) + } + msg := fmt.Sprintf("EVENT %d %d %s %d %d %d", + ev.PlayerID, ev.Seq, ev.Action, ev.DeltaX, ev.LocalAt.Milliseconds(), ev.Delay.Milliseconds()) + if err := writeLine(writer, msg); err != nil { + return err + } + logf("[上报] 玩家%-2d seq=%d 输入=%-8s 延迟=%-8s 实际上报=%s %-4s %+d\n", + playerID, ev.Seq, formatMS(ev.LocalAt), formatMS(ev.Delay), formatMS(reportAt), ev.Action, ev.DeltaX) + } + case "DONE": + return nil + default: + return fmt.Errorf("未知指令: %q", line) + } + } +} + +func parseEventLine(line string) (InputEvent, error) { + parts := strings.Fields(line) + if len(parts) != 7 || parts[0] != "EVENT" { + return InputEvent{}, fmt.Errorf("非法事件消息: %q", line) + } + + playerID, err := strconv.Atoi(parts[1]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法玩家编号: %q", parts[1]) + } + seq, err := strconv.Atoi(parts[2]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法序号: %q", parts[2]) + } + deltaX, err := strconv.Atoi(parts[4]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法位移: %q", parts[4]) + } + localAtMS, err := strconv.Atoi(parts[5]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法输入时刻: %q", parts[5]) + } + delayMS, err := strconv.Atoi(parts[6]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法延迟: %q", parts[6]) + } + + return InputEvent{ + PlayerID: playerID, + Seq: seq, + Action: parts[3], + DeltaX: deltaX, + LocalAt: time.Duration(localAtMS) * time.Millisecond, + Delay: time.Duration(delayMS) * time.Millisecond, + }, nil +} + +func sortedSessions(sessions map[int]*ClientSession) []*ClientSession { + ids := make([]int, 0, len(sessions)) + for id := range sessions { + ids = append(ids, id) + } + sort.Ints(ids) + ordered := make([]*ClientSession, 0, len(ids)) + for _, id := range ids { + ordered = append(ordered, sessions[id]) + } + return ordered +} + +func closeSessions(sessions map[int]*ClientSession) { + for _, session := range sessions { + session.Conn.Close() + } +} + +func readLine(reader *bufio.Reader) (string, error) { + line, err := reader.ReadString('\n') + if err != nil { + return "", err + } + return strings.TrimSpace(line), nil +} + +func writeLine(writer *bufio.Writer, line string) error { + if _, err := writer.WriteString(line + "\n"); err != nil { + return err + } + return writer.Flush() +} + +func usage() { + fmt.Println("用法:") + fmt.Println(" go run ./cmd/exp1/network_event_driven_sync_demo server [-addr 127.0.0.1:9107]") + fmt.Println(" go run ./cmd/exp1/network_event_driven_sync_demo client [-addr 127.0.0.1:9107]") + fmt.Println() + fmt.Println("建议打开 2 个终端: 1 个 server + 1 个 client。") +} + +func main() { + if len(os.Args) < 2 { + usage() + os.Exit(1) + } + + switch os.Args[1] { + case "server": + fs := flag.NewFlagSet("server", flag.ExitOnError) + addr := fs.String("addr", defaultAddr, "服务器监听地址") + fs.Parse(os.Args[2:]) + if err := runServer(*addr); err != nil { + fmt.Fprintf(os.Stderr, "server error: %v\n", err) + os.Exit(1) + } + case "client": + fs := flag.NewFlagSet("client", flag.ExitOnError) + addr := fs.String("addr", defaultAddr, "服务器地址") + fs.Parse(os.Args[2:]) + if err := runClient(*addr); err != nil { + fmt.Fprintf(os.Stderr, "client error: %v\n", err) + os.Exit(1) + } + default: + usage() + os.Exit(1) + } +} diff --git a/CourseCode/ch4/cmd/exp1/network_goroutine_server_demo/main.go b/CourseCode/ch4/cmd/exp1/network_goroutine_server_demo/main.go new file mode 100644 index 0000000000000000000000000000000000000000..32317e12a2343ad739cea3de7963e69096db9918 --- /dev/null +++ b/CourseCode/ch4/cmd/exp1/network_goroutine_server_demo/main.go @@ -0,0 +1,447 @@ +package main + +import ( + "bufio" + "errors" + "flag" + "fmt" + "io" + "net" + "os" + "sort" + "strconv" + "strings" + "sync" + "time" +) + +const ( + trackWidth = 20 + defaultAddr = "127.0.0.1:9105" + expectedPlayers = 4 + frameBudget = 100 * time.Millisecond +) + +var logMu sync.Mutex + +type InputEvent struct { + PlayerID int + Action string + DeltaX int + Latency time.Duration +} + +type ClientSession struct { + ID int + Conn net.Conn + Reader *bufio.Reader + Writer *bufio.Writer +} + +func formatMS(d time.Duration) string { + return fmt.Sprintf("%.1fms", float64(d.Microseconds())/1000.0) +} + +func logf(format string, args ...any) { + logMu.Lock() + defer logMu.Unlock() + fmt.Printf(format, args...) +} + +func logln(args ...any) { + logMu.Lock() + defer logMu.Unlock() + fmt.Println(args...) +} + +func printDivider(title string) { + logMu.Lock() + defer logMu.Unlock() + fmt.Printf("\n========== %s ==========\n", title) +} + +func clamp(x, lo, hi int) int { + if x < lo { + return lo + } + if x > hi { + return hi + } + return x +} + +func renderPositions(pos map[int]int) string { + return fmt.Sprintf("P1(x=%d) P2(x=%d) P3(x=%d) P4(x=%d)", pos[1], pos[2], pos[3], pos[4]) +} + +func defaultFrames(playerID int) map[int]InputEvent { + switch playerID { + case 1: + return map[int]InputEvent{ + 1: {PlayerID: 1, Action: "MOVE", DeltaX: +1, Latency: 11 * time.Millisecond}, + 2: {PlayerID: 1, Action: "MOVE", DeltaX: +1, Latency: 11 * time.Millisecond}, + } + case 2: + return map[int]InputEvent{ + 1: {PlayerID: 2, Action: "MOVE", DeltaX: +1, Latency: 12 * time.Millisecond}, + 2: {PlayerID: 2, Action: "MOVE", DeltaX: +1, Latency: 12 * time.Millisecond}, + } + case 3: + return map[int]InputEvent{ + 1: {PlayerID: 3, Action: "MOVE", DeltaX: -1, Latency: 13 * time.Millisecond}, + 2: {PlayerID: 3, Action: "MOVE", DeltaX: -1, Latency: 13 * time.Millisecond}, + } + case 4: + return map[int]InputEvent{ + 1: {PlayerID: 4, Action: "MOVE", DeltaX: -1, Latency: 500 * time.Millisecond}, + 2: {PlayerID: 4, Action: "MOVE", DeltaX: -1, Latency: 500 * time.Millisecond}, + } + default: + return map[int]InputEvent{} + } +} + +func runServer(addr string) error { + printDivider("实验一 / Network Goroutine Server") + logln("监听地址:", addr) + logln("运行方式: 再打开 1 个终端执行 client 模式。") + logln("目标: 让服务器把每条连接的收包工作交给独立 goroutine,再观察主循环耗时。") + + ln, err := net.Listen("tcp", addr) + if err != nil { + return err + } + defer ln.Close() + + sessions := make(map[int]*ClientSession, expectedPlayers) + for len(sessions) < expectedPlayers { + conn, err := ln.Accept() + if err != nil { + return err + } + + session, err := acceptClient(conn) + if err != nil { + conn.Close() + return err + } + if _, exists := sessions[session.ID]; exists { + session.Conn.Close() + return fmt.Errorf("玩家%d 重复连接", session.ID) + } + sessions[session.ID] = session + logf("[连接] 玩家%-2d 已连接 (%d/%d)\n", session.ID, len(sessions), expectedPlayers) + } + defer closeSessions(sessions) + + positions := map[int]int{1: 2, 2: 6, 3: 10, 4: 14} + for frame := 1; frame <= 2; frame++ { + if err := runConcurrentServerFrame(frame, positions, sessions); err != nil { + return err + } + } + + for _, session := range sessions { + if err := writeLine(session.Writer, "DONE"); err != nil { + return err + } + } + printDivider("Network Goroutine Server 结束") + logln("提示: 对比 network_serial_server_demo,可以直接看到慢客户端还在,但主循环已经不再排队等它。") + return nil +} + +func acceptClient(conn net.Conn) (*ClientSession, error) { + reader := bufio.NewReader(conn) + writer := bufio.NewWriter(conn) + + line, err := readLine(reader) + if err != nil { + return nil, err + } + parts := strings.Fields(line) + if len(parts) != 2 || parts[0] != "HELLO" { + return nil, fmt.Errorf("非法握手消息: %q", line) + } + + playerID, err := strconv.Atoi(parts[1]) + if err != nil { + return nil, fmt.Errorf("非法玩家编号: %q", parts[1]) + } + return &ClientSession{ID: playerID, Conn: conn, Reader: reader, Writer: writer}, nil +} + +type eventResult struct { + Event InputEvent + ReceivedAt time.Duration + Err error +} + +func runConcurrentServerFrame(frame int, positions map[int]int, sessions map[int]*ClientSession) error { + printDivider(fmt.Sprintf("Frame %d / Network Goroutine Server", frame)) + logf("[广播] START %d -> 所有客户端开始准备输入\n", frame) + for _, session := range sortedSessions(sessions) { + if err := writeLine(session.Writer, fmt.Sprintf("START %d", frame)); err != nil { + return err + } + } + + frameStart := time.Now() + results := make(chan eventResult, len(sessions)) + var wg sync.WaitGroup + + for _, session := range sortedSessions(sessions) { + wg.Add(1) + go func(session *ClientSession) { + defer wg.Done() + logf("[并发收包] Frame %-2d 玩家%-2d 的收包 goroutine 已启动\n", frame, session.ID) + line, err := readLine(session.Reader) + if err != nil { + results <- eventResult{Err: err} + return + } + + ev, err := parseEventLine(line) + if err != nil { + results <- eventResult{Err: err} + return + } + results <- eventResult{ + Event: ev, + ReceivedAt: time.Since(frameStart), + } + }(session) + } + dispatchCost := time.Since(frameStart) + logf("[主循环] Frame %-2d 收包任务分发完成: %s (目标<100ms)\n", frame, formatMS(dispatchCost)) + + go func() { + wg.Wait() + close(results) + }() + + received := make(map[int]eventResult, len(sessions)) + for result := range results { + if result.Err != nil { + return result.Err + } + received[result.Event.PlayerID] = result + logf("[收到] Frame %-2d 玩家%-2d %-4s %+d 客户端延迟=%-8s 到达时间=%s\n", + frame, result.Event.PlayerID, result.Event.Action, result.Event.DeltaX, + formatMS(result.Event.Latency), formatMS(result.ReceivedAt)) + } + + totalCollectCost := time.Since(frameStart) + logf("[摘要] Frame %-2d 所有消息收齐: %s\n", frame, formatMS(totalCollectCost)) + if dispatchCost > frameBudget { + logln("[现象] 这不应该发生;如果看到这里超时,说明主循环本身也被拖住了。") + } else { + logln("[现象] 主循环很快完成了收包任务分发,没有阻塞在某一条慢连接上。") + } + if totalCollectCost > frameBudget { + logln("[补充] 慢客户端仍然更晚到达,只是它被隔离在独立 goroutine 中,不再把主循环卡死。") + } + + applyOrder := []int{1, 2, 3, 4} + if frame == 2 { + applyOrder = []int{4, 1, 2, 3} + } + for _, pid := range applyOrder { + result := received[pid] + positions[pid] = clamp(positions[pid]+result.Event.DeltaX, 0, trackWidth) + } + logf("[摘要] Frame %-2d 位置快照: %s\n", frame, renderPositions(positions)) + return nil +} + +func runClient(addr string) error { + printDivider("实验一 / Network Goroutine Client") + logln("连接服务器:", addr) + logf("当前进程会同时模拟 %d 名玩家,每名玩家保持一条独立连接。\n", expectedPlayers) + + var wg sync.WaitGroup + errCh := make(chan error, expectedPlayers) + for playerID := 1; playerID <= expectedPlayers; playerID++ { + wg.Add(1) + go func(playerID int) { + defer wg.Done() + if err := runOnePlayerClient(addr, playerID); err != nil { + errCh <- err + } + }(playerID) + } + wg.Wait() + close(errCh) + + for err := range errCh { + if err != nil { + return err + } + } + printDivider("Network Goroutine Client 结束") + logln("所有玩家脚本执行完毕。") + return nil +} + +func runOnePlayerClient(addr string, playerID int) error { + frames := defaultFrames(playerID) + if len(frames) == 0 { + return fmt.Errorf("玩家%d 没有预设脚本,请使用 1~4 号玩家", playerID) + } + + conn, err := net.Dial("tcp", addr) + if err != nil { + return err + } + defer conn.Close() + + reader := bufio.NewReader(conn) + writer := bufio.NewWriter(conn) + + if err := writeLine(writer, fmt.Sprintf("HELLO %d", playerID)); err != nil { + return err + } + logf("[客户端%-2d] 已连接服务器并完成握手\n", playerID) + + for { + line, err := readLine(reader) + if err != nil { + if errors.Is(err, io.EOF) { + return nil + } + return err + } + + parts := strings.Fields(line) + if len(parts) == 0 { + continue + } + + switch parts[0] { + case "START": + if len(parts) != 2 { + return fmt.Errorf("非法 START 指令: %q", line) + } + frame, err := strconv.Atoi(parts[1]) + if err != nil { + return fmt.Errorf("非法帧编号: %q", parts[1]) + } + ev, ok := frames[frame] + if !ok { + return fmt.Errorf("玩家%d 缺少第%d帧脚本", playerID, frame) + } + + logf("[客户端%-2d] Frame %-2d 收到 START 模拟延迟=%s\n", playerID, frame, formatMS(ev.Latency)) + time.Sleep(ev.Latency) + + msg := fmt.Sprintf("EVENT %d %s %d %d", ev.PlayerID, ev.Action, ev.DeltaX, ev.Latency.Milliseconds()) + if err := writeLine(writer, msg); err != nil { + return err + } + logf("[客户端%-2d] Frame %-2d 已发送 %-4s %+d\n", playerID, frame, ev.Action, ev.DeltaX) + case "DONE": + logf("[客户端%-2d] 收到 DONE,退出\n", playerID) + return nil + default: + return fmt.Errorf("未知指令: %q", line) + } + } +} + +func parseEventLine(line string) (InputEvent, error) { + parts := strings.Fields(line) + if len(parts) != 5 || parts[0] != "EVENT" { + return InputEvent{}, fmt.Errorf("非法事件消息: %q", line) + } + + playerID, err := strconv.Atoi(parts[1]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法玩家编号: %q", parts[1]) + } + deltaX, err := strconv.Atoi(parts[3]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法位移: %q", parts[3]) + } + latencyMS, err := strconv.Atoi(parts[4]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法延迟: %q", parts[4]) + } + + return InputEvent{ + PlayerID: playerID, + Action: parts[2], + DeltaX: deltaX, + Latency: time.Duration(latencyMS) * time.Millisecond, + }, nil +} + +func sortedSessions(sessions map[int]*ClientSession) []*ClientSession { + ids := make([]int, 0, len(sessions)) + for id := range sessions { + ids = append(ids, id) + } + sort.Ints(ids) + ordered := make([]*ClientSession, 0, len(ids)) + for _, id := range ids { + ordered = append(ordered, sessions[id]) + } + return ordered +} + +func closeSessions(sessions map[int]*ClientSession) { + for _, session := range sessions { + session.Conn.Close() + } +} + +func readLine(reader *bufio.Reader) (string, error) { + line, err := reader.ReadString('\n') + if err != nil { + return "", err + } + return strings.TrimSpace(line), nil +} + +func writeLine(writer *bufio.Writer, line string) error { + if _, err := writer.WriteString(line + "\n"); err != nil { + return err + } + return writer.Flush() +} + +func usage() { + fmt.Println("用法:") + fmt.Println(" go run ./cmd/exp1/network_goroutine_server_demo server [-addr 127.0.0.1:9105]") + fmt.Println(" go run ./cmd/exp1/network_goroutine_server_demo client [-addr 127.0.0.1:9105]") + fmt.Println() + fmt.Println("建议打开 2 个终端: 1 个 server + 1 个 client。") +} + +func main() { + if len(os.Args) < 2 { + usage() + os.Exit(1) + } + + switch os.Args[1] { + case "server": + fs := flag.NewFlagSet("server", flag.ExitOnError) + addr := fs.String("addr", defaultAddr, "服务器监听地址") + fs.Parse(os.Args[2:]) + if err := runServer(*addr); err != nil { + fmt.Fprintf(os.Stderr, "server error: %v\n", err) + os.Exit(1) + } + case "client": + fs := flag.NewFlagSet("client", flag.ExitOnError) + addr := fs.String("addr", defaultAddr, "服务器地址") + fs.Parse(os.Args[2:]) + if err := runClient(*addr); err != nil { + fmt.Fprintf(os.Stderr, "client error: %v\n", err) + os.Exit(1) + } + default: + usage() + os.Exit(1) + } +} diff --git a/CourseCode/ch4/cmd/exp1/network_serial_server_demo/main.go b/CourseCode/ch4/cmd/exp1/network_serial_server_demo/main.go new file mode 100644 index 0000000000000000000000000000000000000000..a204edd5bfdcf0a48f1a5b2f87418cde9953caa5 --- /dev/null +++ b/CourseCode/ch4/cmd/exp1/network_serial_server_demo/main.go @@ -0,0 +1,427 @@ +package main + +import ( + "bufio" + "errors" + "flag" + "fmt" + "io" + "net" + "os" + "sort" + "strconv" + "strings" + "sync" + "time" +) + +const ( + trackWidth = 20 + defaultAddr = "127.0.0.1:9101" + expectedPlayers = 4 + frameBudget = 100 * time.Millisecond +) + +var logMu sync.Mutex + +type InputEvent struct { + PlayerID int + Action string + DeltaX int + Latency time.Duration +} + +type ClientSession struct { + ID int + Conn net.Conn + Reader *bufio.Reader + Writer *bufio.Writer +} + +func formatMS(d time.Duration) string { + return fmt.Sprintf("%.1fms", float64(d.Microseconds())/1000.0) +} + +func logf(format string, args ...any) { + logMu.Lock() + defer logMu.Unlock() + fmt.Printf(format, args...) +} + +func logln(args ...any) { + logMu.Lock() + defer logMu.Unlock() + fmt.Println(args...) +} + +func printDivider(title string) { + logMu.Lock() + defer logMu.Unlock() + fmt.Printf("\n========== %s ==========\n", title) +} + +func clamp(x, lo, hi int) int { + if x < lo { + return lo + } + if x > hi { + return hi + } + return x +} + +func renderPositions(pos map[int]int) string { + return fmt.Sprintf("P1(x=%d) P2(x=%d) P3(x=%d) P4(x=%d)", pos[1], pos[2], pos[3], pos[4]) +} + +func defaultFrames(playerID int) map[int]InputEvent { + switch playerID { + case 1: + return map[int]InputEvent{ + 1: {PlayerID: 1, Action: "MOVE", DeltaX: +1, Latency: 11 * time.Millisecond}, + 2: {PlayerID: 1, Action: "MOVE", DeltaX: +1, Latency: 11 * time.Millisecond}, + } + case 2: + return map[int]InputEvent{ + 1: {PlayerID: 2, Action: "MOVE", DeltaX: +1, Latency: 12 * time.Millisecond}, + 2: {PlayerID: 2, Action: "MOVE", DeltaX: +1, Latency: 12 * time.Millisecond}, + } + case 3: + return map[int]InputEvent{ + 1: {PlayerID: 3, Action: "MOVE", DeltaX: -1, Latency: 13 * time.Millisecond}, + 2: {PlayerID: 3, Action: "MOVE", DeltaX: -1, Latency: 13 * time.Millisecond}, + } + case 4: + return map[int]InputEvent{ + 1: {PlayerID: 4, Action: "MOVE", DeltaX: -1, Latency: 500 * time.Millisecond}, + 2: {PlayerID: 4, Action: "MOVE", DeltaX: -1, Latency: 500 * time.Millisecond}, + } + default: + return map[int]InputEvent{} + } +} + +func runServer(addr string) error { + printDivider("实验一 / Network Serial Server") + logln("监听地址:", addr) + logln("运行方式: 再打开 1 个终端执行 client 模式。") + logln("目标: 让“慢客户端 -> 服务器串行收包 -> 主循环被拖慢”的链路真实发生。") + + ln, err := net.Listen("tcp", addr) + if err != nil { + return err + } + defer ln.Close() + + sessions := make(map[int]*ClientSession, expectedPlayers) + for len(sessions) < expectedPlayers { + conn, err := ln.Accept() + if err != nil { + return err + } + + session, err := acceptClient(conn) + if err != nil { + conn.Close() + return err + } + if _, exists := sessions[session.ID]; exists { + session.Conn.Close() + return fmt.Errorf("玩家%d 重复连接", session.ID) + } + sessions[session.ID] = session + logf("[连接] 玩家%-2d 已连接 (%d/%d)\n", session.ID, len(sessions), expectedPlayers) + } + defer closeSessions(sessions) + + positions := map[int]int{1: 2, 2: 6, 3: 10, 4: 14} + frameOrders := map[int][]int{ + 1: {1, 2, 3, 4}, + 2: {4, 1, 2, 3}, + } + + for frame := 1; frame <= len(frameOrders); frame++ { + if err := runServerFrame(frame, frameOrders[frame], positions, sessions); err != nil { + return err + } + } + + for _, session := range sessions { + if err := writeLine(session.Writer, "DONE"); err != nil { + return err + } + } + printDivider("Server 结束") + logln("提示: 现在可以再运行 network_goroutine_server_demo,对比服务器把收包交给独立 goroutine 后,慢客户端是否还会让主循环排队。") + return nil +} + +func acceptClient(conn net.Conn) (*ClientSession, error) { + reader := bufio.NewReader(conn) + writer := bufio.NewWriter(conn) + + line, err := readLine(reader) + if err != nil { + return nil, err + } + parts := strings.Fields(line) + if len(parts) != 2 || parts[0] != "HELLO" { + return nil, fmt.Errorf("非法握手消息: %q", line) + } + + playerID, err := strconv.Atoi(parts[1]) + if err != nil { + return nil, fmt.Errorf("非法玩家编号: %q", parts[1]) + } + return &ClientSession{ + ID: playerID, + Conn: conn, + Reader: reader, + Writer: writer, + }, nil +} + +func runServerFrame( + frame int, + order []int, + positions map[int]int, + sessions map[int]*ClientSession, +) error { + printDivider(fmt.Sprintf("Frame %d / Server", frame)) + logf("[广播] START %d -> 所有客户端开始准备输入\n", frame) + for _, session := range sortedSessions(sessions) { + if err := writeLine(session.Writer, fmt.Sprintf("START %d", frame)); err != nil { + return err + } + } + + frameStart := time.Now() + for _, pid := range order { + session := sessions[pid] + logf("[等待] Frame %-2d 按顺序读取 玩家%-2d\n", frame, pid) + line, err := readLine(session.Reader) + if err != nil { + return err + } + + ev, err := parseEventLine(line) + if err != nil { + return err + } + if ev.PlayerID != pid { + return fmt.Errorf("收到的玩家编号与等待顺序不一致: want=%d got=%d", pid, ev.PlayerID) + } + + waited := time.Since(frameStart) + logf("[收到] Frame %-2d 玩家%-2d %-4s %+d 客户端延迟=%-8s 累计等待=%s\n", + frame, ev.PlayerID, ev.Action, ev.DeltaX, formatMS(ev.Latency), formatMS(waited)) + positions[ev.PlayerID] = clamp(positions[ev.PlayerID]+ev.DeltaX, 0, trackWidth) + } + + cost := time.Since(frameStart) + logf("[摘要] Frame %-2d 位置快照: %s\n", frame, renderPositions(positions)) + logf("[摘要] Frame %-2d 主循环耗时: %s (目标<100ms)\n", frame, formatMS(cost)) + if cost > frameBudget { + logln("[现象] 慢客户端尚未发来输入时,服务器主循环会一直阻塞在对应连接上。") + } + return nil +} + +func runClient(addr string) error { + printDivider("实验一 / Network Serial Client") + logln("连接服务器:", addr) + logf("当前进程会同时模拟 %d 名玩家,每名玩家保持一条独立连接。\n", expectedPlayers) + + var wg sync.WaitGroup + errCh := make(chan error, expectedPlayers) + + for playerID := 1; playerID <= expectedPlayers; playerID++ { + wg.Add(1) + go func(playerID int) { + defer wg.Done() + if err := runOnePlayerClient(addr, playerID); err != nil { + errCh <- err + } + }(playerID) + } + + wg.Wait() + close(errCh) + + for err := range errCh { + if err != nil { + return err + } + } + printDivider("Client 结束") + logln("所有玩家脚本执行完毕。") + return nil +} + +func runOnePlayerClient(addr string, playerID int) error { + frames := defaultFrames(playerID) + if len(frames) == 0 { + return fmt.Errorf("玩家%d 没有预设脚本,请使用 1~4 号玩家", playerID) + } + + conn, err := net.Dial("tcp", addr) + if err != nil { + return err + } + defer conn.Close() + + reader := bufio.NewReader(conn) + writer := bufio.NewWriter(conn) + + if err := writeLine(writer, fmt.Sprintf("HELLO %d", playerID)); err != nil { + return err + } + logf("[客户端%-2d] 已连接服务器并完成握手\n", playerID) + + for { + line, err := readLine(reader) + if err != nil { + if errors.Is(err, io.EOF) { + return nil + } + return err + } + + parts := strings.Fields(line) + if len(parts) == 0 { + continue + } + + switch parts[0] { + case "START": + if len(parts) != 2 { + return fmt.Errorf("非法 START 指令: %q", line) + } + frame, err := strconv.Atoi(parts[1]) + if err != nil { + return fmt.Errorf("非法帧编号: %q", parts[1]) + } + + ev, ok := frames[frame] + if !ok { + return fmt.Errorf("玩家%d 缺少第%d帧脚本", playerID, frame) + } + + logf("[客户端%-2d] Frame %-2d 收到 START 模拟延迟=%s\n", + playerID, frame, formatMS(ev.Latency)) + time.Sleep(ev.Latency) + + msg := fmt.Sprintf("EVENT %d %s %d %d", ev.PlayerID, ev.Action, ev.DeltaX, ev.Latency.Milliseconds()) + if err := writeLine(writer, msg); err != nil { + return err + } + logf("[客户端%-2d] Frame %-2d 已发送 %-4s %+d\n", playerID, frame, ev.Action, ev.DeltaX) + case "DONE": + logf("[客户端%-2d] 收到 DONE,退出\n", playerID) + return nil + default: + return fmt.Errorf("未知指令: %q", line) + } + } +} + +func parseEventLine(line string) (InputEvent, error) { + parts := strings.Fields(line) + if len(parts) != 5 || parts[0] != "EVENT" { + return InputEvent{}, fmt.Errorf("非法事件消息: %q", line) + } + + playerID, err := strconv.Atoi(parts[1]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法玩家编号: %q", parts[1]) + } + deltaX, err := strconv.Atoi(parts[3]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法位移: %q", parts[3]) + } + latencyMS, err := strconv.Atoi(parts[4]) + if err != nil { + return InputEvent{}, fmt.Errorf("非法延迟: %q", parts[4]) + } + + return InputEvent{ + PlayerID: playerID, + Action: parts[2], + DeltaX: deltaX, + Latency: time.Duration(latencyMS) * time.Millisecond, + }, nil +} + +func sortedSessions(sessions map[int]*ClientSession) []*ClientSession { + ids := make([]int, 0, len(sessions)) + for id := range sessions { + ids = append(ids, id) + } + sort.Ints(ids) + + ordered := make([]*ClientSession, 0, len(ids)) + for _, id := range ids { + ordered = append(ordered, sessions[id]) + } + return ordered +} + +func closeSessions(sessions map[int]*ClientSession) { + for _, session := range sessions { + session.Conn.Close() + } +} + +func readLine(reader *bufio.Reader) (string, error) { + line, err := reader.ReadString('\n') + if err != nil { + return "", err + } + return strings.TrimSpace(line), nil +} + +func writeLine(writer *bufio.Writer, line string) error { + if _, err := writer.WriteString(line + "\n"); err != nil { + return err + } + return writer.Flush() +} + +func usage() { + fmt.Println("用法:") + fmt.Println(" go run ./cmd/exp1/network_serial_server_demo server [-addr 127.0.0.1:9101]") + fmt.Println(" go run ./cmd/exp1/network_serial_server_demo client [-addr 127.0.0.1:9101]") + fmt.Println() + fmt.Println("建议打开 2 个终端: 1 个 server + 1 个 client。") +} + +func main() { + if len(os.Args) < 2 { + usage() + os.Exit(1) + } + + switch os.Args[1] { + case "server": + fs := flag.NewFlagSet("server", flag.ExitOnError) + addr := fs.String("addr", defaultAddr, "服务器监听地址") + fs.Parse(os.Args[2:]) + + if err := runServer(*addr); err != nil { + fmt.Fprintf(os.Stderr, "server error: %v\n", err) + os.Exit(1) + } + case "client": + fs := flag.NewFlagSet("client", flag.ExitOnError) + addr := fs.String("addr", defaultAddr, "服务器地址") + fs.Parse(os.Args[2:]) + + if err := runClient(*addr); err != nil { + fmt.Fprintf(os.Stderr, "client error: %v\n", err) + os.Exit(1) + } + default: + usage() + os.Exit(1) + } +} diff --git a/CourseCode/ch4/cmd/exp2/right/main.go b/CourseCode/ch4/cmd/exp2/right/main.go new file mode 100644 index 0000000000000000000000000000000000000000..21d0354fcc83603c885bad447ec408827fbe2d8f --- /dev/null +++ b/CourseCode/ch4/cmd/exp2/right/main.go @@ -0,0 +1,62 @@ +package main + +import ( + "fmt" + "sync" + "time" +) + +type NPC struct { + ID int + Active bool + mu sync.Mutex +} + +type Result struct { + PlayerID int + Success bool +} + +func main() { + fmt.Println("=== 实验二:临界区与数据竞争 / Mutex 修复版 ===") + fmt.Println("场景: 仍然是 5 名玩家同时抢同一个 NPC 的唯一掉落。") + fmt.Println("目标: 用互斥锁保护“检查 + 修改”,确保只有一个赢家。") + + npc := &NPC{ID: 999, Active: false} + var wg sync.WaitGroup + results := make(chan Result, 5) + + for i := 1; i <= 5; i++ { + wg.Add(1) + go func(playerID int) { + defer wg.Done() + + npc.mu.Lock() + defer npc.mu.Unlock() + + if npc.Active { + fmt.Printf("[失败] 玩家 %d 晚来一步,NPC %d 已消失。\n", playerID, npc.ID) + results <- Result{PlayerID: playerID, Success: false} + return + } + + time.Sleep(10 * time.Millisecond) + npc.Active = true + fmt.Printf("[掉落] 玩家 %d 触发成功!NPC %d 送出宝物\n", playerID, npc.ID) + results <- Result{PlayerID: playerID, Success: true} + }(i) + } + + wg.Wait() + close(results) + + successCount := 0 + for result := range results { + if result.Success { + successCount++ + } + } + + fmt.Printf("[统计] 成功领取人数 = %d\n", successCount) + fmt.Println("[提示] 对照无锁版的多轮统计,观察这里是否还能出现重复掉落。") +} diff --git a/CourseCode/ch4/cmd/exp2/wrong/main.go b/CourseCode/ch4/cmd/exp2/wrong/main.go new file mode 100644 index 0000000000000000000000000000000000000000..907cd2af8fa1b957efd1958b1e654906b6024b62 --- /dev/null +++ b/CourseCode/ch4/cmd/exp2/wrong/main.go @@ -0,0 +1,94 @@ +package main + +import ( + "fmt" + "math/rand" + "sync" + "time" +) + +type NPC struct { + ID int + Active bool +} + +type Result struct { + PlayerID int + Success bool + ArriveLag time.Duration + PickupCost time.Duration +} + +func runOneRound(round int, rng *rand.Rand) int { + npc := &NPC{ID: 999, Active: false} + var wg sync.WaitGroup + results := make(chan Result, 5) + startGate := make(chan struct{}) + + for i := 1; i <= 5; i++ { + arriveLag := time.Duration(rng.Intn(18)) * time.Millisecond + pickupCost := time.Duration(6+rng.Intn(10)) * time.Millisecond + wg.Add(1) + go func(playerID int, arriveLag, pickupCost time.Duration) { + defer wg.Done() + + <-startGate + time.Sleep(arriveLag) + + if !npc.Active { + fmt.Printf("[第%d轮] 玩家 %d 看到宝物仍在,准备拾取 (到达偏移=%s, 拾取耗时=%s)\n", + round, playerID, arriveLag, pickupCost) + time.Sleep(pickupCost) + npc.Active = true + fmt.Printf("[第%d轮][掉落] 玩家 %d 成功拿到唯一宝物\n", round, playerID) + results <- Result{ + PlayerID: playerID, + Success: true, + ArriveLag: arriveLag, + PickupCost: pickupCost, + } + return + } + fmt.Printf("[第%d轮][失败] 玩家 %d 到达时宝物已经被别人改写\n", round, playerID) + results <- Result{ + PlayerID: playerID, + Success: false, + ArriveLag: arriveLag, + PickupCost: pickupCost, + } + }(i, arriveLag, pickupCost) + } + + close(startGate) + wg.Wait() + close(results) + + successCount := 0 + for result := range results { + if result.Success { + successCount++ + } + } + return successCount +} + +func main() { + fmt.Println("=== 实验二:临界区与数据竞争 / 无锁错误版 ===") + fmt.Println("场景: 5 名玩家同时抢同一个 NPC 的唯一掉落。") + fmt.Println("目标: 观察“检查宝物是否存在”和“标记宝物已被拿走”分离后,会出现重复掉落。") + + rng := rand.New(rand.NewSource(time.Now().UnixNano())) + totalRounds := 5 + for round := 1; round <= totalRounds; round++ { + successCount := runOneRound(round, rng) + fmt.Printf("[第%d轮统计] 成功领取人数 = %d\n", round, successCount) + if successCount > 1 { + fmt.Println(" 现象: 本轮发生了重复掉落。") + } else { + fmt.Println(" 现象: 本轮看起来正常,但这并不代表代码没有竞态窗口。") + } + fmt.Println() + } + + fmt.Println("[提示] 再运行 right 版本,对比把“检查 + 修改”放进同一个互斥区后,统计结果会怎样变化。") +} diff --git a/CourseCode/ch4/cmd/exp3/busy_wait/main.go b/CourseCode/ch4/cmd/exp3/busy_wait/main.go new file mode 100644 index 0000000000000000000000000000000000000000..eb32e6cd246fadb23743bb0aa35d25f6eb5386fe --- /dev/null +++ b/CourseCode/ch4/cmd/exp3/busy_wait/main.go @@ -0,0 +1,56 @@ +package main + +import ( + "fmt" + "sync" + "time" +) + +type NPC struct { + ID int + Items int + mu sync.Mutex +} + +func playerTask(playerID int, npc *NPC, wg *sync.WaitGroup) { + defer wg.Done() + emptyRuns := 0 + + for { + npc.mu.Lock() + if npc.Items > 0 { + npc.Items-- + fmt.Printf("\n[成功] 玩家 %d 抢到宝物!结束交互。(总计白跑了 %d 次)\n", playerID, emptyRuns) + npc.mu.Unlock() + return + } + npc.mu.Unlock() + emptyRuns++ + } +} + +func main() { + fmt.Println("=== 实验三:告别忙等 / 仅用互斥锁的错误版 ===") + fmt.Println("场景: 宝物暂时为空,玩家不断重复 Lock -> Check -> Unlock。") + fmt.Println("目标: 观察 CPU 被“忙等”白白消耗,虽然业务最终能成功。") + + npc := &NPC{ID: 999, Items: 0} + var wg sync.WaitGroup + start := time.Now() + + for i := 1; i <= 3; i++ { + wg.Add(1) + go playerTask(i, npc, &wg) + } + + time.Sleep(100 * time.Millisecond) + + fmt.Printf("\n[系统] 宝物补货前,玩家已经忙等了约 %v。\n", time.Since(start)) + fmt.Println("[系统] NPC 开始投放 3 个宝物...") + npc.mu.Lock() + npc.Items += 3 + npc.mu.Unlock() + + wg.Wait() + fmt.Println("[提示] 再运行 cond_wait 版本,对比等待中的玩家是否还会继续空转。") +} diff --git a/CourseCode/ch4/cmd/exp3/cond_wait/main.go b/CourseCode/ch4/cmd/exp3/cond_wait/main.go new file mode 100644 index 0000000000000000000000000000000000000000..32843c17a5138afcc5194b375917d83880c35f78 --- /dev/null +++ b/CourseCode/ch4/cmd/exp3/cond_wait/main.go @@ -0,0 +1,61 @@ +package main + +import ( + "fmt" + "sync" + "time" +) + +type NPC struct { + ID int + Items int + mu sync.Mutex + cond *sync.Cond +} + +func playerGetTreasure(playerID int, npc *NPC, wg *sync.WaitGroup) { + defer wg.Done() + + npc.mu.Lock() + for npc.Items == 0 { + fmt.Printf("[等待] 宝物为空,玩家 %d 陷入沉睡,释放锁...\n", playerID) + npc.cond.Wait() + fmt.Printf("[唤醒] 玩家 %d 被叫醒,重新检查宝物库存...\n", playerID) + } + + npc.Items-- + fmt.Printf("[成功] 玩家 %d 取走宝物!剩余宝物数: %d\n", playerID, npc.Items) + npc.mu.Unlock() +} + +func npcRestock(npc *NPC, amount int) { + npc.mu.Lock() + npc.Items += amount + fmt.Printf("\n[系统] NPC 补货了 %d 个宝物!\n", amount) + npc.mu.Unlock() + + for i := 0; i < amount; i++ { + npc.cond.Signal() + } +} + +func main() { + fmt.Println("=== 实验三:告别忙等 / Cond 修复版 ===") + fmt.Println("场景: 玩家在库存为空时不再轮询,而是进入等待队列。") + fmt.Println("目标: 观察 Wait/Signal 的睡眠唤醒流程,并强调 for 重检条件。") + + npc := &NPC{ID: 999, Items: 0} + npc.cond = sync.NewCond(&npc.mu) + + var wg sync.WaitGroup + for i := 1; i <= 3; i++ { + wg.Add(1) + go playerGetTreasure(i, npc, &wg) + } + + time.Sleep(50 * time.Millisecond) + npcRestock(npc, 3) + + wg.Wait() + fmt.Println("[提示] 注意 Wait 前后的日志顺序,以及玩家被叫醒后还会再次检查库存。") +} diff --git a/CourseCode/ch4/cmd/exp4/channel_semaphore/main.go b/CourseCode/ch4/cmd/exp4/channel_semaphore/main.go new file mode 100644 index 0000000000000000000000000000000000000000..b0d636936e68c7b0c49ff4df6432f258ab5720d6 --- /dev/null +++ b/CourseCode/ch4/cmd/exp4/channel_semaphore/main.go @@ -0,0 +1,32 @@ +package main + +import ( + "fmt" + "sync" + "time" +) + +func main() { + fmt.Println("=== 实验四:锁的进阶技巧与粒度优化 / Channel 信号量 ===") + fmt.Println("目标: 使用容量为 3 的 Channel 限制同时工作的协程数量。") + + semaphore := make(chan struct{}, 3) + var wg sync.WaitGroup + + for i := 1; i <= 10; i++ { + wg.Add(1) + go func(id int) { + defer wg.Done() + + fmt.Printf("[Worker %d] 尝试获取许可...\n", id) + semaphore <- struct{}{} + fmt.Printf("[Worker %d] 获取许可,开始执行\n", id) + time.Sleep(500 * time.Millisecond) + fmt.Printf("[Worker %d] 工作完成,释放许可\n", id) + <-semaphore + }(i) + } + + wg.Wait() + fmt.Println("[提示] 观察同一时刻真正进入执行区的 Worker 数量是否被限制在 3 个以内。") +} diff --git a/CourseCode/ch4/cmd/exp4/channel_timeout_lock/main.go b/CourseCode/ch4/cmd/exp4/channel_timeout_lock/main.go new file mode 100644 index 0000000000000000000000000000000000000000..737e0c4f330fa93ece6742b5fe2c3ead97049e4f --- /dev/null +++ b/CourseCode/ch4/cmd/exp4/channel_timeout_lock/main.go @@ -0,0 +1,34 @@ +package main + +import ( + "fmt" + "time" +) + +func main() { + fmt.Println("=== 实验四:锁的进阶技巧与粒度优化 / Channel 超时锁 ===") + fmt.Println("目标: 用 select + time.After 模拟带超时的资源竞争与降级逻辑。") + + lock := make(chan struct{}, 1) + + go func() { + lock <- struct{}{} + fmt.Println("[玩家A] 抢到了资源,故意占用 5 秒") + time.Sleep(5 * time.Second) + fmt.Println("[玩家A] 处理完成,释放资源") + <-lock + }() + + time.Sleep(100 * time.Millisecond) + fmt.Println("[玩家B] 尝试获取资源,超时时间设置为 2 秒") + + select { + case lock <- struct{}{}: + fmt.Println("[玩家B] 成功获取资源") + <-lock + case <-time.After(2 * time.Second): + fmt.Println("[玩家B] 等待超时,转入降级逻辑,避免一直卡住") + } + + fmt.Println("[提示] 如果去掉超时分支,玩家B 会一直堵在这里。") +} diff --git a/CourseCode/ch4/cmd/exp4/rw_mutex/main.go b/CourseCode/ch4/cmd/exp4/rw_mutex/main.go new file mode 100644 index 0000000000000000000000000000000000000000..c002c5eccf7655c0ce181d9cd6a0d8fe5e09545d --- /dev/null +++ b/CourseCode/ch4/cmd/exp4/rw_mutex/main.go @@ -0,0 +1,71 @@ +package main + +import ( + "fmt" + "sync" + "time" +) + +var ( + data = map[string]int{"key": 100} + rw sync.RWMutex +) + +func formatMS(d time.Duration) string { + return fmt.Sprintf("%.1fms", float64(d.Microseconds())/1000.0) +} + +func Reader(name string, hold time.Duration, base time.Time, wg *sync.WaitGroup) { + defer wg.Done() + + arrive := time.Now() + fmt.Printf("[%s][读者 %s] 尝试获取读锁\n", formatMS(time.Since(base)), name) + rw.RLock() + fmt.Printf("[%s][读者 %s] 获取读锁成功,等待了 %s,读取数据: %d\n", + formatMS(time.Since(base)), name, formatMS(time.Since(arrive)), data["key"]) + time.Sleep(hold) + fmt.Printf("[%s][读者 %s] 释放读锁\n", formatMS(time.Since(base)), name) + rw.RUnlock() +} + +func Writer(name string, val int, hold time.Duration, base time.Time, wg *sync.WaitGroup) { + defer wg.Done() + + arrive := time.Now() + fmt.Printf("[%s][写者 %s] 尝试获取写锁\n", formatMS(time.Since(base)), name) + rw.Lock() + fmt.Printf("[%s][写者 %s] 获取写锁成功,等待了 %s,开始写入数据: %d\n", + formatMS(time.Since(base)), name, formatMS(time.Since(arrive)), val) + time.Sleep(hold) + data["key"] = val + fmt.Printf("[%s][写者 %s] 写入完成,准备释放写锁\n", formatMS(time.Since(base)), name) + rw.Unlock() +} + +func main() { + fmt.Println("=== 实验四:锁的进阶技巧与粒度优化 / RWMutex ===") + fmt.Println("目标: 观察读者并发、写者独占,以及“写者排队后/写入期间到来的读者”都会被阻塞。") + var wg sync.WaitGroup + base := time.Now() + + wg.Add(1) + go Reader("R1", 1200*time.Millisecond, base, &wg) + + wg.Add(1) + go Reader("R2", 1200*time.Millisecond, base, &wg) + + time.Sleep(100 * time.Millisecond) + wg.Add(1) + go Writer("W1", 999, 1*time.Second, base, &wg) + + time.Sleep(150 * time.Millisecond) + wg.Add(1) + go Reader("R3(写者排队后到达)", 400*time.Millisecond, base, &wg) + + time.Sleep(1500 * time.Millisecond) + wg.Add(1) + go Reader("R4(写入进行中到达)", 400*time.Millisecond, base, &wg) + + wg.Wait() + fmt.Println("[提示] 重点看 R3 和 R4 的“尝试获取读锁”与“真正拿到读锁”之间的等待时间。") +} diff --git a/CourseCode/ch4/cmd/exp5/sync_pool_demo/main.go b/CourseCode/ch4/cmd/exp5/sync_pool_demo/main.go new file mode 100644 index 0000000000000000000000000000000000000000..414a52031548661c41a7be29e8f06bad38f830f9 --- /dev/null +++ b/CourseCode/ch4/cmd/exp5/sync_pool_demo/main.go @@ -0,0 +1,124 @@ +package main + +import ( + "bytes" + "fmt" + "runtime" + "sort" + "sync" + "time" +) + +var mockData = bytes.Repeat([]byte("A"), 10*1024) + +var bufferPool = sync.Pool{ + New: func() interface{} { + return new(bytes.Buffer) + }, +} + +var dummy int + +//go:noinline +func doSomeWork() { + for j := 0; j < 500; j++ { + dummy += j + } +} + +func processWithoutPool(id int, latencies []time.Duration, wg *sync.WaitGroup) { + defer wg.Done() + start := time.Now() + + buf := new(bytes.Buffer) + buf.Grow(10 * 1024) + buf.Write(mockData) + _ = buf.Bytes() + + doSomeWork() + latencies[id] = time.Since(start) +} + +func processWithPool(id int, latencies []time.Duration, wg *sync.WaitGroup) { + defer wg.Done() + start := time.Now() + + buf := bufferPool.Get().(*bytes.Buffer) + buf.Reset() + buf.Write(mockData) + _ = buf.Bytes() + + doSomeWork() + bufferPool.Put(buf) + latencies[id] = time.Since(start) +} + +func printPercentiles(name string, latencies []time.Duration, startMem runtime.MemStats, startTime time.Time) { + var endMem runtime.MemStats + runtime.ReadMemStats(&endMem) + + sort.Slice(latencies, func(i, j int) bool { + return latencies[i] < latencies[j] + }) + + totalDuration := time.Since(startTime) + allocs := endMem.Mallocs - startMem.Mallocs + gcCount := endMem.NumGC - startMem.NumGC + + length := len(latencies) + p50 := latencies[length*50/100] + p90 := latencies[length*90/100] + p99 := latencies[length*99/100] + p999 := latencies[length*999/1000] + + fmt.Printf("\n[%s]\n", name) + fmt.Printf("总耗时 (涵盖调度): %v\n", totalDuration) + fmt.Printf("堆内存分配 (Mallocs): %d 次\n", allocs) + fmt.Printf("触发 GC 次数: %d 次\n", gcCount) + fmt.Println("------------- 延迟分布 (微秒/毫秒) -------------") + fmt.Printf("P50 (中位数): %.3f 微秒\n", float64(p50.Nanoseconds())/1000.0) + fmt.Printf("P90 延迟 : %.3f 微秒\n", float64(p90.Nanoseconds())/1000.0) + fmt.Printf("P99 延迟 : %.3f 微秒\n", float64(p99.Nanoseconds())/1000.0) + fmt.Printf("P99.9延迟 : %.3f 微秒 <-- 观察对比项\n", float64(p999.Nanoseconds())/1000.0) + fmt.Println("------------------------------------------------") +} + +func main() { + const numRequests = 20000 + var wg sync.WaitGroup + + fmt.Println("=== 实验五:高并发性能榨取(sync.Pool) ===") + fmt.Println("场景: 高频请求里频繁申请 10KB 临时缓冲,对比“每次 new”与“对象复用”。") + fmt.Println("说明: 本实验只做本地对象池演示,数据库连接池保留为概念扩展。") + + latenciesNoPool := make([]time.Duration, numRequests) + runtime.GC() + var m1 runtime.MemStats + runtime.ReadMemStats(&m1) + t1 := time.Now() + + for i := 0; i < numRequests; i++ { + wg.Add(1) + go processWithoutPool(i, latenciesNoPool, &wg) + } + wg.Wait() + printPercentiles("基准:无优化 (频繁 10KB 分配)", latenciesNoPool, m1, t1) + + time.Sleep(1 * time.Second) + + latenciesWithPool := make([]time.Duration, numRequests) + runtime.GC() + var m2 runtime.MemStats + runtime.ReadMemStats(&m2) + t2 := time.Now() + + for i := 0; i < numRequests; i++ { + wg.Add(1) + go processWithPool(i, latenciesWithPool, &wg) + } + wg.Wait() + printPercentiles("优化:+sync.Pool (复用 10KB 缓冲)", latenciesWithPool, m2, t2) + + fmt.Println("\n[提示] 重点对比两组结果里的 Mallocs、GC 次数和 P99.9 延迟。") + fmt.Println("[扩展] 真正的数据库连接池属于另一类资源池,本章只做对象池演示。") +} diff --git a/CourseCode/ch4/cmd/exp6/storage_arch/main.go b/CourseCode/ch4/cmd/exp6/storage_arch/main.go new file mode 100644 index 0000000000000000000000000000000000000000..803a4d779d3f6bdef9eb3efa12883f7fd13e6eb5 --- /dev/null +++ b/CourseCode/ch4/cmd/exp6/storage_arch/main.go @@ -0,0 +1,291 @@ +package main + +import ( + "context" + "database/sql" + "errors" + "fmt" + "os" + "time" + + _ "github.com/lib/pq" + "github.com/redis/go-redis/v9" +) + +var ctx = context.Background() + +type noopRedisLogger struct{} + +func (noopRedisLogger) Printf(context.Context, string, ...interface{}) {} + +type StorageDemo struct { + db *sql.DB + redis *redis.Client +} + +func defaultRedisAddr() string { + if addr := os.Getenv("REDIS_ADDR"); addr != "" { + return addr + } + return "127.0.0.1:6379" +} + +func defaultPGDSN() string { + if dsn := os.Getenv("PG_DSN"); dsn != "" { + return dsn + } + return "" +} + +func newStorageDemo(redisAddr, pgDSN string) (*StorageDemo, error) { + redis.SetLogger(noopRedisLogger{}) + + rdb := redis.NewClient(&redis.Options{Addr: redisAddr}) + if err := rdb.Ping(ctx).Err(); err != nil { + return nil, fmt.Errorf("连接 Redis %s 失败: %w", redisAddr, err) + } + + db, err := sql.Open("postgres", pgDSN) + if err != nil { + return nil, fmt.Errorf("打开 PostgreSQL 失败: %w", err) + } + + db.SetMaxOpenConns(5) + db.SetMaxIdleConns(5) + db.SetConnMaxLifetime(5 * time.Minute) + + pingCtx, cancel := context.WithTimeout(ctx, 3*time.Second) + defer cancel() + if err := db.PingContext(pingCtx); err != nil { + _ = db.Close() + return nil, fmt.Errorf("连接 PostgreSQL 失败: %w", err) + } + + demo := &StorageDemo{db: db, redis: rdb} + if err := demo.initSchema(); err != nil { + _ = demo.Close() + return nil, err + } + if err := demo.seedData(); err != nil { + _ = demo.Close() + return nil, err + } + + return demo, nil +} + +func (s *StorageDemo) initSchema() error { + statements := []string{ + `CREATE TABLE IF NOT EXISTS players ( + user_id TEXT PRIMARY KEY, + gold INTEGER NOT NULL + )`, + `CREATE TABLE IF NOT EXISTS game_configs ( + config_key TEXT PRIMARY KEY, + config_value TEXT NOT NULL + )`, + } + + for _, stmt := range statements { + if _, err := s.db.Exec(stmt); err != nil { + return fmt.Errorf("初始化 PostgreSQL 表结构失败: %w", err) + } + } + return nil +} + +func (s *StorageDemo) seedData() error { + if _, err := s.db.Exec(` + INSERT INTO players (user_id, gold) + VALUES ('player_1', 100) + ON CONFLICT (user_id) DO UPDATE SET gold = EXCLUDED.gold`); err != nil { + return fmt.Errorf("写入玩家初始数据失败: %w", err) + } + + if _, err := s.db.Exec(` + INSERT INTO game_configs (config_key, config_value) + VALUES ('drop_rate', '1.5') + ON CONFLICT (config_key) DO UPDATE SET config_value = EXCLUDED.config_value`); err != nil { + return fmt.Errorf("写入配置初始数据失败: %w", err) + } + + if err := s.redis.Set(ctx, "gold:player_1", "100", 0).Err(); err != nil { + return fmt.Errorf("写入 Redis 初始金币失败: %w", err) + } + if err := s.redis.Del(ctx, "cfg:drop_rate").Err(); err != nil { + return fmt.Errorf("清理 Redis 配置缓存失败: %w", err) + } + + return nil +} + +func (s *StorageDemo) Close() error { + var firstErr error + if err := s.redis.Close(); err != nil { + firstErr = err + } + if err := s.db.Close(); err != nil && firstErr == nil { + firstErr = err + } + return firstErr +} + +func (s *StorageDemo) deductGold(userID string, deductAmount int) error { + start := time.Now() + fmt.Printf("[Write Through] 开始扣除 %s 金币 %d...\n", userID, deductAmount) + + tx, err := s.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("开启 PostgreSQL 事务失败: %w", err) + } + + var currentGold int + err = tx.QueryRowContext(ctx, ` + UPDATE players + SET gold = gold - $1 + WHERE user_id = $2 + RETURNING gold`, deductAmount, userID).Scan(¤tGold) + if err != nil { + _ = tx.Rollback() + if errors.Is(err, sql.ErrNoRows) { + return fmt.Errorf("玩家 %s 不存在", userID) + } + return fmt.Errorf("更新 PostgreSQL 金币失败: %w", err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("提交 PostgreSQL 事务失败: %w", err) + } + + if err := s.redis.Set(ctx, "gold:"+userID, fmt.Sprintf("%d", currentGold), 0).Err(); err != nil { + return fmt.Errorf("同步 Redis 失败: %w", err) + } + + fmt.Printf("[Write Through] 扣除成功,当前金币=%d,耗时=%v\n\n", currentGold, time.Since(start)) + return nil +} + +func (s *StorageDemo) showGoldConsistency(userID string) { + var dbGold int + if err := s.db.QueryRowContext(ctx, `SELECT gold FROM players WHERE user_id = $1`, userID).Scan(&dbGold); err != nil { + fmt.Printf("[一致性检查] PostgreSQL 查询失败: %v\n", err) + return + } + + cacheGold, err := s.redis.Get(ctx, "gold:"+userID).Result() + if err != nil { + fmt.Printf("[一致性检查] Redis 查询失败: %v\n", err) + return + } + + fmt.Printf("[一致性检查] PostgreSQL.gold=%d, Redis.gold=%s\n\n", dbGold, cacheGold) +} + +func (s *StorageDemo) getGameConfig(key string) string { + start := time.Now() + cacheKey := "cfg:" + key + + val, err := s.redis.Get(ctx, cacheKey).Result() + if err == nil { + fmt.Printf("[Cache Aside 读] 缓存命中 %s=%s,耗时=%v\n", key, val, time.Since(start)) + return val + } + + fmt.Println("[Cache Aside 读] 缓存未命中,开始查询 PostgreSQL...") + + var dbVal string + if err := s.db.QueryRowContext(ctx, ` + SELECT config_value FROM game_configs WHERE config_key = $1`, key).Scan(&dbVal); err != nil { + fmt.Printf("[Cache Aside 读] PostgreSQL 查询失败: %v\n", err) + return "" + } + + if err := s.redis.Set(ctx, cacheKey, dbVal, 0).Err(); err != nil { + fmt.Printf("[Cache Aside 读] Redis 回填失败: %v\n", err) + } + + fmt.Printf("[Cache Aside 读] 已从 PostgreSQL 读取 %s=%s,耗时=%v\n\n", key, dbVal, time.Since(start)) + return dbVal +} + +func (s *StorageDemo) updateGameConfig(key, newVal string) { + start := time.Now() + fmt.Printf("[Cache Aside 写] 开始更新 %s=%s ...\n", key, newVal) + + if _, err := s.db.ExecContext(ctx, ` + UPDATE game_configs + SET config_value = $1 + WHERE config_key = $2`, newVal, key); err != nil { + fmt.Printf("[Cache Aside 写] PostgreSQL 更新失败: %v\n", err) + return + } + + if err := s.redis.Del(ctx, "cfg:"+key).Err(); err != nil { + fmt.Printf("[Cache Aside 写] 删除 Redis 缓存失败: %v\n", err) + return + } + + fmt.Printf("[Cache Aside 写] 更新成功,缓存已失效,耗时=%v\n\n", time.Since(start)) +} + +func printRunHints(redisAddr, pgDSN string) { + fmt.Println("运行前置条件:") + fmt.Printf("- Redis: %s(建议通过 Docker Desktop 启动)\n", redisAddr) + if pgDSN == "" { + fmt.Println("- PostgreSQL: 请先在 PowerShell 中设置 PG_DSN 环境变量") + } else { + fmt.Printf("- PostgreSQL: %s\n", pgDSN) + } + fmt.Println("- 程序会自动创建 players 和 game_configs 两张表,并写入演示初始数据。") + fmt.Println() +} + +func printInfraHelp() { + fmt.Println("Redis 启动示例(先启动 Docker Desktop 再执行):") + fmt.Println(`docker run -d --name ch4-redis -p 6379:6379 redis:7`) + fmt.Println() + fmt.Println("PostgreSQL 启动示例(第一次创建):") + fmt.Println(`docker run -d --name ch4-postgres -e POSTGRES_USER=你的用户名 -e POSTGRES_PASSWORD=你的密码 -e POSTGRES_DB=postgres -p 5432:5432 postgres:16`) + fmt.Println("如果容器已经创建过,可直接执行:docker start ch4-postgres") + fmt.Println() + fmt.Println("PostgreSQL 连接串示例:") + fmt.Println(`$env:PG_DSN="postgres://你的用户名:你的密码@127.0.0.1:5432/postgres?sslmode=disable"`) + fmt.Println("请把示例中的用户名、密码、主机、端口替换成你自己的 PostgreSQL 配置。") +} + +func main() { + fmt.Println("=== 实验六:分层存储架构(Redis + PostgreSQL) ===") + fmt.Println("目标:使用真实 Redis 与 PostgreSQL 演示 Write Through 和 Cache Aside。") + + redisAddr := defaultRedisAddr() + pgDSN := defaultPGDSN() + printRunHints(redisAddr, pgDSN) + if pgDSN == "" { + fmt.Println("[错误] 未设置 PG_DSN。为了避免在代码里写死个人账号密码,实验六要求每位使用者自行配置 PostgreSQL 连接串。") + printInfraHelp() + return + } + + demo, err := newStorageDemo(redisAddr, pgDSN) + if err != nil { + fmt.Printf("[错误] 初始化基础设施失败:%v\n", err) + printInfraHelp() + return + } + defer demo.Close() + + if err := demo.deductGold("player_1", 20); err != nil { + fmt.Printf("[错误] Write Through 执行失败:%v\n", err) + return + } + demo.showGoldConsistency("player_1") + + fmt.Println("--- 模拟配置读取与更新 ---") + demo.getGameConfig("drop_rate") + demo.getGameConfig("drop_rate") + demo.updateGameConfig("drop_rate", "2.0") + demo.getGameConfig("drop_rate") + demo.getGameConfig("drop_rate") + + fmt.Println("[结论] 核心资产适合 Write Through,读多写少的配置数据适合 Cache Aside。") +} diff --git a/CourseCode/ch4/go.mod b/CourseCode/ch4/go.mod new file mode 100644 index 0000000000000000000000000000000000000000..74268ed3bfaad8de7176757441330ba8d18424ea --- /dev/null +++ b/CourseCode/ch4/go.mod @@ -0,0 +1,14 @@ +module ch4 + +go 1.21 + +require ( + github.com/lib/pq v1.12.0 + github.com/redis/go-redis/v9 v9.18.0 +) + +require ( + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect + go.uber.org/atomic v1.11.0 // indirect +) diff --git "a/CourseCode/ch4/\350\213\261\351\233\204\351\233\206\347\273\223\346\274\224\347\244\272\345\256\236\351\252\214.md" "b/CourseCode/ch4/\350\213\261\351\233\204\351\233\206\347\273\223\346\274\224\347\244\272\345\256\236\351\252\214.md" index 9bf084651d13e837492a36e75f12e1e51fa92171..fe1e9bca1956632701372bcf19698d52c694b27b 100644 --- "a/CourseCode/ch4/\350\213\261\351\233\204\351\233\206\347\273\223\346\274\224\347\244\272\345\256\236\351\252\214.md" +++ "b/CourseCode/ch4/\350\213\261\351\233\204\351\233\206\347\273\223\346\274\224\347\244\272\345\256\236\351\252\214.md" @@ -44,21 +44,21 @@ 3. **粒度优化:** 对比“全局大锁”与读多写少场景下的 `sync.RWMutex` 。 - **成功标准:** 获取锁超时的玩家能正确执行降级逻辑(`doFallback()`)不被永远卡死 ;使用读写锁后,多个 `RLock()` 读者可以同时进入临界区执行读取 。 -### 实验五:高并发性能榨取(连接池与对象池 sync.Pool) +### 实验五:高并发性能榨取(对象池 sync.Pool) - **对应页码:** 第 54, 56, 58 页 - **测试目标:** 降低高并发下的 GC(垃圾回收)压力,复用临时对象以消除 STW(Stop The World)延迟抖动 。 - **操作步骤:** - 1. 在每秒过万次请求的网络收发中,每次都 `new(bytes.Buffer)` 创建新的数据包 。 - 2. 引入 `sync.Pool` 维护对象缓存,通过 `bufferPool.Get()` 获取并重置缓冲,处理完后 `Put()` 放回 。 - 3. 配置数据库长连接池 `SetMaxOpenConns(100)` 。 -- **成功标准:** 结合性能监控工具(pprof),观察到引入 `sync.Pool` 后,GC 耗时从 50ms 降至 10ms,P99 延迟显著下降 ;引入连接池后吞吐量(QPS)从 300 暴涨至 1200 。 + 1. 在高频请求场景中,每次都 `new(bytes.Buffer)` 创建新的 10KB 临时缓冲。 + 2. 引入 `sync.Pool` 维护对象缓存,通过 `bufferPool.Get()` 获取并重置缓冲,处理完后 `Put()` 放回。 + 3. 对比两版的总耗时、内存分配次数、GC 次数和尾延迟。 +- **成功标准:** 引入 `sync.Pool` 后,分配次数与 GC 次数显著下降,P99 或 P99.9 延迟明显优于基线版;本实验只讲对象池,数据库连接池保留为扩展概念。 -### 实验六:游戏数据分层存储架构(Redis + DB) +### 实验六:游戏数据分层存储架构(Redis + PostgreSQL) - **对应页码:** 第 61-63 页 - **测试目标:** 针对高频热数据(如实时排行)和低频冷数据(如配置、历史成就) ,实现正确的数据一致性同步策略。 - **操作步骤:** 1. **强一致场景(写穿 Write Through):** 编写金币扣除逻辑,必须先写 DB 成功,再同步更新 Redis,两者皆成功才返回 。 2. **最终一致场景(旁路缓存 Cache Aside):** 编写游戏配置表读取逻辑,先查缓存,miss 则查 DB 并回填 。写操作先写 DB,**然后直接删除缓存** 。 -- **成功标准:** 成功实现分层架构,热数据读取延迟控制在 <1ms 。 \ No newline at end of file +- **成功标准:** 成功实现分层架构,能够稳定观察到 Write Through 的双写一致性,以及 Cache Aside 的“首次 miss、回填、再次命中”流程。