feat(mpsync): 多人同步 MVP——登录区块流、方块变更广播;修复 netserver 关闭竞态与世界 worker 退出死锁
This commit is contained in:
@@ -22,9 +22,9 @@ const maxSources = 64
|
||||
type category uint8
|
||||
|
||||
const (
|
||||
catSFX category = iota // 音效
|
||||
catMusic // 音乐
|
||||
catAmbient // 环境
|
||||
catSFX category = iota // 音效
|
||||
catMusic // 音乐
|
||||
catAmbient // 环境
|
||||
)
|
||||
|
||||
// Manager 音频管理器(单例)。
|
||||
@@ -40,12 +40,12 @@ type Manager struct {
|
||||
|
||||
// source 一个正在播放的声源。
|
||||
type source struct {
|
||||
key string
|
||||
pos [3]float64 // 世界坐标
|
||||
dist float64 // 距玩家距离
|
||||
ctrl *beep.Ctrl
|
||||
cat category
|
||||
start time.Time
|
||||
key string
|
||||
pos [3]float64 // 世界坐标
|
||||
dist float64 // 距玩家距离
|
||||
ctrl *beep.Ctrl
|
||||
cat category
|
||||
start time.Time
|
||||
}
|
||||
|
||||
// New 创建音频管理器(找不到音频设备时降级为静默——服务器/CI 安全)。
|
||||
|
||||
81
internal/mpsync/hub.go
Normal file
81
internal/mpsync/hub.go
Normal file
@@ -0,0 +1,81 @@
|
||||
// Package mpsync 多人同步:登录发送区块数据、方块变更广播(多人同步.md §2、§4–§5)。
|
||||
//
|
||||
// 服务器权威:方块破坏/放置由服务器执行后广播(多人同步.md §5)。
|
||||
package mpsync
|
||||
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"mc/internal/block"
|
||||
"mc/internal/netproto"
|
||||
"mc/internal/netserver"
|
||||
"mc/internal/save"
|
||||
"mc/internal/world"
|
||||
)
|
||||
|
||||
// Hub 同步中心:玩家会话与广播。
|
||||
type Hub struct {
|
||||
mu sync.Mutex
|
||||
clients map[*netserver.Conn]struct{} // 已登录连接
|
||||
world *world.World
|
||||
viewDist int
|
||||
}
|
||||
|
||||
// NewHub 创建同步中心。
|
||||
func NewHub(w *world.World, viewDist int) *Hub {
|
||||
return &Hub{clients: make(map[*netserver.Conn]struct{}), world: w, viewDist: viewDist}
|
||||
}
|
||||
|
||||
// AddPlayer 玩家登录:加入广播集并发送周围区块(多人同步.md §4 区块流)。
|
||||
func (h *Hub) AddPlayer(c *netserver.Conn, spawnX, spawnZ float64) {
|
||||
h.mu.Lock()
|
||||
h.clients[c] = struct{}{}
|
||||
h.mu.Unlock()
|
||||
|
||||
c.Send(netproto.Frame{MsgID: netproto.MsgLoginResponse, Payload: []byte{0}})
|
||||
// 发送视距内已加载区块(按距离排序,多人同步.md §4)
|
||||
chunks := h.world.ActiveChunks()
|
||||
for _, ch := range chunks {
|
||||
if absi32(ch.CX-int32(spawnX)/16) > int32(h.viewDist) || absi32(ch.CZ-int32(spawnZ)/16) > int32(h.viewDist) {
|
||||
continue
|
||||
}
|
||||
data := save.EncodeChunk(ch)
|
||||
c.Send(netproto.Frame{MsgID: netproto.MsgChunkData, Payload: data})
|
||||
}
|
||||
}
|
||||
|
||||
// RemovePlayer 玩家断开:移出广播集。
|
||||
func (h *Hub) RemovePlayer(c *netserver.Conn) {
|
||||
h.mu.Lock()
|
||||
delete(h.clients, c)
|
||||
h.mu.Unlock()
|
||||
}
|
||||
|
||||
// BroadcastBlockChange 广播方块变更(多人同步.md §5:只发视距内,MVP 全量广播)。
|
||||
func (h *Hub) BroadcastBlockChange(x, y, z int32, id uint16, meta uint8) {
|
||||
msg := netproto.BlockChange{X: x, Y: y, Z: z, BlockID: id, Meta: meta}
|
||||
payload := msg.Encode()
|
||||
h.mu.Lock()
|
||||
clients := make([]*netserver.Conn, 0, len(h.clients))
|
||||
for c := range h.clients {
|
||||
clients = append(clients, c)
|
||||
}
|
||||
h.mu.Unlock()
|
||||
for _, c := range clients {
|
||||
c.Send(netproto.Frame{MsgID: netproto.MsgBlockChange, Payload: payload})
|
||||
}
|
||||
}
|
||||
|
||||
// SetBlockWorld 服务器权威写方块:修改世界并广播(多人同步.md §2 权威链路)。
|
||||
func (h *Hub) SetBlockWorld(x, y, z int32, s block.State) {
|
||||
h.world.SetBlock(x, y, z, s)
|
||||
h.BroadcastBlockChange(x, y, z, s.ID(), s.Meta())
|
||||
}
|
||||
|
||||
// absi32 绝对值。
|
||||
func absi32(v int32) int32 {
|
||||
if v < 0 {
|
||||
return -v
|
||||
}
|
||||
return v
|
||||
}
|
||||
174
internal/mpsync/hub_test.go
Normal file
174
internal/mpsync/hub_test.go
Normal file
@@ -0,0 +1,174 @@
|
||||
package mpsync
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"mc/internal/block"
|
||||
"mc/internal/logx"
|
||||
"mc/internal/netproto"
|
||||
"mc/internal/netserver"
|
||||
"mc/internal/world"
|
||||
"mc/internal/worldgen"
|
||||
)
|
||||
|
||||
// syncHandler 测试用服务器处理器:登录 → hub.AddPlayer。
|
||||
type syncHandler struct {
|
||||
hub *Hub
|
||||
recv chan *netserver.Conn
|
||||
}
|
||||
|
||||
func (h *syncHandler) OnConnect(c *netserver.Conn) {}
|
||||
func (h *syncHandler) OnHandshake(c *netserver.Conn, hs netproto.Handshake) uint8 {
|
||||
if h.recv != nil {
|
||||
h.recv <- c
|
||||
}
|
||||
h.hub.AddPlayer(c, 8, 8)
|
||||
return 0
|
||||
}
|
||||
func (h *syncHandler) OnMove(c *netserver.Conn, m netproto.PlayerMove) {}
|
||||
func (h *syncHandler) OnDisconnect(c *netserver.Conn) { h.hub.RemovePlayer(c) }
|
||||
|
||||
// TestLoginReceivesChunk 登录后客户端收到区块数据帧(多人同步.md §4)。
|
||||
func TestLoginReceivesChunk(t *testing.T) {
|
||||
reg, err := block.Load(filepath.Join("..", "..", "assets", "config", "blocks.json"))
|
||||
if err != nil {
|
||||
t.Fatalf("加载注册表失败: %v", err)
|
||||
}
|
||||
gen, err := worldgen.New(reg, 7, 0.005, 0.01, 0.1)
|
||||
if err != nil {
|
||||
t.Fatalf("创建生成器失败: %v", err)
|
||||
}
|
||||
w, err := world.New(reg, gen, 4)
|
||||
if err != nil {
|
||||
t.Fatalf("创建世界失败: %v", err)
|
||||
}
|
||||
defer w.Close()
|
||||
deadline := time.Now().Add(15 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
w.Update(8, 64, 8, 8*time.Millisecond)
|
||||
if w.Block(8, 0, 8) != block.Air {
|
||||
break
|
||||
}
|
||||
time.Sleep(2 * time.Millisecond)
|
||||
}
|
||||
|
||||
hub := NewHub(w, 4)
|
||||
log, _ := logx.New("", logx.LevelDebug)
|
||||
h := &syncHandler{hub: hub, recv: make(chan *netserver.Conn, 1)}
|
||||
srv := netserver.New("127.0.0.1:0", log, h)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
if err := srv.Listen(ctx); err != nil {
|
||||
t.Fatalf("监听失败: %v", err)
|
||||
}
|
||||
nc, err := net.DialTimeout("tcp", srv.Addr(), 3*time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("拨号失败: %v", err)
|
||||
}
|
||||
defer nc.Close()
|
||||
|
||||
hs := netproto.Handshake{ProtocolVersion: 1, ClientID: "同步测试"}
|
||||
_, _ = nc.Write(netproto.EncodeFrame(netproto.Frame{MsgID: netproto.MsgHandshake, Payload: hs.Encode()}))
|
||||
|
||||
gotChunk := false
|
||||
buf := make([]byte, 64*1024)
|
||||
_ = nc.SetReadDeadline(time.Now().Add(5 * time.Second))
|
||||
acc := []byte{}
|
||||
for !gotChunk {
|
||||
n, err := nc.Read(buf)
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
acc = append(acc, buf[:n]...)
|
||||
for len(acc) >= 4 {
|
||||
f, rest, err := netproto.DecodeFrame(acc)
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
acc = rest
|
||||
if f.MsgID == netproto.MsgChunkData {
|
||||
gotChunk = true
|
||||
}
|
||||
}
|
||||
}
|
||||
if !gotChunk {
|
||||
t.Fatal("登录后未收到区块数据")
|
||||
}
|
||||
}
|
||||
|
||||
// TestBlockChangeBroadcast 服务器权威写方块 → 客户端收到广播(多人同步.md §5)。
|
||||
func TestBlockChangeBroadcast(t *testing.T) {
|
||||
reg, err := block.Load(filepath.Join("..", "..", "assets", "config", "blocks.json"))
|
||||
if err != nil {
|
||||
t.Fatalf("加载注册表失败: %v", err)
|
||||
}
|
||||
gen, err := worldgen.New(reg, 8, 0.005, 0.01, 0.1)
|
||||
if err != nil {
|
||||
t.Fatalf("创建生成器失败: %v", err)
|
||||
}
|
||||
w, err := world.New(reg, gen, 2)
|
||||
if err != nil {
|
||||
t.Fatalf("创建世界失败: %v", err)
|
||||
}
|
||||
defer w.Close()
|
||||
|
||||
hub := NewHub(w, 4)
|
||||
log, _ := logx.New("", logx.LevelDebug)
|
||||
h := &syncHandler{hub: hub, recv: make(chan *netserver.Conn, 1)}
|
||||
srv := netserver.New("127.0.0.1:0", log, h)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
if err := srv.Listen(ctx); err != nil {
|
||||
t.Fatalf("监听失败: %v", err)
|
||||
}
|
||||
nc, err := net.DialTimeout("tcp", srv.Addr(), 3*time.Second)
|
||||
if err != nil {
|
||||
t.Fatalf("拨号失败: %v", err)
|
||||
}
|
||||
defer nc.Close()
|
||||
_, _ = nc.Write(netproto.EncodeFrame(netproto.Frame{MsgID: netproto.MsgHandshake, Payload: netproto.Handshake{ProtocolVersion: 1}.Encode()}))
|
||||
|
||||
// 等待服务器完成登录注册(避免广播早于 AddPlayer 执行,多人同步.md §2 时序)
|
||||
select {
|
||||
case <-h.recv:
|
||||
case <-time.After(3 * time.Second):
|
||||
t.Fatal("未完成登录注册")
|
||||
}
|
||||
|
||||
torch, _ := reg.ID("torch")
|
||||
go hub.SetBlockWorld(8, 70, 8, block.NewState(torch, 0))
|
||||
|
||||
buf := make([]byte, 64*1024)
|
||||
_ = nc.SetReadDeadline(time.Now().Add(5 * time.Second))
|
||||
acc := []byte{}
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
n, err := nc.Read(buf)
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
acc = append(acc, buf[:n]...)
|
||||
for len(acc) >= 4 {
|
||||
f, rest, err := netproto.DecodeFrame(acc)
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
acc = rest
|
||||
if f.MsgID != netproto.MsgBlockChange {
|
||||
continue
|
||||
}
|
||||
m, err := netproto.DecodeBlockChange(f.Payload)
|
||||
if err != nil {
|
||||
t.Fatalf("解码失败: %v", err)
|
||||
}
|
||||
if m.X == 8 && m.Y == 70 && m.Z == 8 && m.BlockID == torch {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
t.Fatal("未收到方块变更广播")
|
||||
}
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
@@ -112,12 +113,13 @@ type Conn struct {
|
||||
conn net.Conn
|
||||
srv *Server
|
||||
send chan []byte
|
||||
done chan struct{}
|
||||
closeOnce sync.Once
|
||||
}
|
||||
|
||||
// newConn 创建连接。
|
||||
func newConn(id uint64, nc net.Conn, s *Server) *Conn {
|
||||
return &Conn{id: id, conn: nc, srv: s, send: make(chan []byte, 64)}
|
||||
return &Conn{id: id, conn: nc, srv: s, send: make(chan []byte, 64), done: make(chan struct{})}
|
||||
}
|
||||
|
||||
// ID 连接 ID(服务器分配,多人同步.md §3)。
|
||||
@@ -186,29 +188,41 @@ func (c *Conn) handle(f netproto.Frame) {
|
||||
|
||||
// writeLoop 写循环:channel 驱动,写超时断开(网络协议.md §7)。
|
||||
func (c *Conn) writeLoop() {
|
||||
for data := range c.send {
|
||||
_ = c.conn.SetWriteDeadline(time.Now().Add(readTimeout))
|
||||
if _, err := c.conn.Write(data); err != nil {
|
||||
c.close()
|
||||
for {
|
||||
select {
|
||||
case <-c.done:
|
||||
return
|
||||
case data := <-c.send:
|
||||
_ = c.conn.SetWriteDeadline(time.Now().Add(readTimeout))
|
||||
if _, err := c.conn.Write(data); err != nil {
|
||||
c.close()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Send 发送一帧(非阻塞入队;队满丢弃并告警,防慢连接阻塞服务器)。
|
||||
// Send 发送一帧(非阻塞入队;连接关闭或队满时丢弃并告警,防慢连接阻塞服务器)。
|
||||
func (c *Conn) Send(f netproto.Frame) {
|
||||
select {
|
||||
case <-c.done:
|
||||
return
|
||||
default:
|
||||
}
|
||||
select {
|
||||
case c.send <- netproto.EncodeFrame(f):
|
||||
case <-c.done:
|
||||
default:
|
||||
c.srv.log.Warnf("连接 %d 发送队列已满,丢弃帧 0x%02X", c.id, f.MsgID)
|
||||
}
|
||||
}
|
||||
|
||||
// close 幂等关闭。
|
||||
// close 幂等关闭:关闭 done 信号与 socket;不关闭 send(避免与在途 Send 竞争)。
|
||||
func (c *Conn) close() {
|
||||
c.closeOnce.Do(func() {
|
||||
c.srv.log.Debugf("连接 %d 关闭(来源:\n%s)", c.id, debug.Stack())
|
||||
close(c.done)
|
||||
_ = c.conn.Close()
|
||||
close(c.send)
|
||||
c.srv.mu.Lock()
|
||||
delete(c.srv.conns, c.id)
|
||||
c.srv.mu.Unlock()
|
||||
|
||||
@@ -46,6 +46,7 @@ type World struct {
|
||||
// 异步生成流水线
|
||||
workers chan chunkCoord // 待生成坐标
|
||||
results chan *chunk.Chunk
|
||||
done chan struct{} // 关闭信号(goroutine 可退出)
|
||||
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
@@ -73,6 +74,7 @@ func New(reg *block.Registry, gen *worldgen.Generator, viewDistance int) (*World
|
||||
viewDistance: viewDistance,
|
||||
workers: make(chan chunkCoord, n*2),
|
||||
results: make(chan *chunk.Chunk, n*2),
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
for i := 0; i < n; i++ {
|
||||
w.wg.Add(1)
|
||||
@@ -84,11 +86,20 @@ func New(reg *block.Registry, gen *worldgen.Generator, viewDistance int) (*World
|
||||
// generateWorker 生成 goroutine:纯 CPU 工作,不碰 GL、不持主线程锁。
|
||||
func (w *World) generateWorker() {
|
||||
defer w.wg.Done()
|
||||
for c := range w.workers {
|
||||
ch := w.gen.Generate(c.cx, c.cz)
|
||||
// 生成完立即做天空光灌入(光照.md §3,单区块:邻居视为阻隔无影响)
|
||||
light.SkyFill(chunkLightAdapter{ch, w.reg}, 0, 15, 0, 15)
|
||||
w.results <- ch
|
||||
for {
|
||||
select {
|
||||
case <-w.done:
|
||||
return
|
||||
case c := <-w.workers:
|
||||
ch := w.gen.Generate(c.cx, c.cz)
|
||||
// 生成完立即做天空光灌入(光照.md §3,单区块:邻居视为阻隔无影响)
|
||||
light.SkyFill(chunkLightAdapter{ch, w.reg}, 0, 15, 0, 15)
|
||||
select {
|
||||
case w.results <- ch:
|
||||
case <-w.done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -236,7 +247,7 @@ func (w *World) Count() int {
|
||||
|
||||
// Close 停止 worker pool(goroutine 可退出,性能与内存.md §8)。
|
||||
func (w *World) Close() {
|
||||
close(w.workers)
|
||||
close(w.done)
|
||||
w.wg.Wait()
|
||||
}
|
||||
|
||||
|
||||
BIN
tools/zig.zip
BIN
tools/zig.zip
Binary file not shown.
Reference in New Issue
Block a user