diff --git a/internal/audio/audio.go b/internal/audio/audio.go index 5bd82033..ec69195f 100644 --- a/internal/audio/audio.go +++ b/internal/audio/audio.go @@ -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 安全)。 diff --git a/internal/mpsync/hub.go b/internal/mpsync/hub.go new file mode 100644 index 00000000..3ecb66ab --- /dev/null +++ b/internal/mpsync/hub.go @@ -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 +} diff --git a/internal/mpsync/hub_test.go b/internal/mpsync/hub_test.go new file mode 100644 index 00000000..1858b8d7 --- /dev/null +++ b/internal/mpsync/hub_test.go @@ -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("未收到方块变更广播") +} diff --git a/internal/netserver/netserver.go b/internal/netserver/netserver.go index 24763c55..035d1884 100644 --- a/internal/netserver/netserver.go +++ b/internal/netserver/netserver.go @@ -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() diff --git a/internal/world/world.go b/internal/world/world.go index ceb1aaaa..e0dbca93 100644 --- a/internal/world/world.go +++ b/internal/world/world.go @@ -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() } diff --git a/tools/zig.zip b/tools/zig.zip index c93e8339..155c45bc 100644 Binary files a/tools/zig.zip and b/tools/zig.zip differ