From 616d30bc089aa49c3341250fd93a379a7be79729 Mon Sep 17 00:00:00 2001 From: kevin Date: Wed, 8 Jul 2026 19:14:48 +0800 Subject: [PATCH] add debug logging for VPN packet flow diagnosis Add comprehensive console logging at key data path points: - [TUN-RX]: per-packet TUN read (size, proto, src/dst IP, slow write detection) - [TUN-TX]: per-packet TUN write (size, success/error) - [WS-RX]: per-packet WebSocket read from client (size, proto, src/dst IP) - [WS-TX]: per-packet WebSocket write to client (size, lockWaitMs, writeMs) - [ROUTE]: routing decisions (TUN->client, client->tun, client->relay, anti-spoof) - [CONN]: connection lifecycle (init, ready, stats every 10s, disconnect with final stats) - [STATS]: global traffic stats every 10s (clients, rx/tx bytes and rates) - [PING]: WebSocket ping failures - [WS]: WebSocket upgrade with client IP and auth method - [TUN]: MTU set confirmation, VPN start/stop environment info --- internal/vpn/handler.go | 9 +++-- internal/vpn/service.go | 55 ++++++++++++++++++++++++++---- internal/vpn/switch.go | 39 ++++++++++++++++++++-- internal/vpn/tun_darwin.go | 9 ++++- internal/vpn/tun_linux.go | 9 ++++- internal/vpn/tunnel.go | 68 +++++++++++++++++++++++++++++++++----- 6 files changed, 167 insertions(+), 22 deletions(-) diff --git a/internal/vpn/handler.go b/internal/vpn/handler.go index 56fbb86..05aad8e 100644 --- a/internal/vpn/handler.go +++ b/internal/vpn/handler.go @@ -31,14 +31,16 @@ var upgrader = websocket.Upgrader{ func HandleWS(c *gin.Context) { tokenStr := c.Query("token") + clientIP := c.ClientIP() conn, err := upgrader.Upgrade(c.Writer, c.Request, nil) if err != nil { - log.Printf("WebSocket 升级失败: %v", err) + log.Printf("[WS] upgrade failed clientIP=%s err=%v", clientIP, err) return } if tokenStr != "" { + log.Printf("[WS] upgrade clientIP=%s auth=jwt", clientIP) claims, err := middleware.ParseToken(tokenStr) if err != nil { sendJSON(conn, authResponse{Type: "auth_err", Message: "令牌无效或已过期"}) @@ -55,9 +57,10 @@ func HandleWS(c *gin.Context) { return } - user, err := authenticate(conn, db.DB, c.ClientIP()) + log.Printf("[WS] upgrade clientIP=%s auth=password", clientIP) + user, err := authenticate(conn, db.DB, clientIP) if err != nil { - log.Printf("认证读取失败: %v", err) + log.Printf("[WS] auth read failed clientIP=%s err=%v", clientIP, err) conn.Close() return } diff --git a/internal/vpn/service.go b/internal/vpn/service.go index 3f80e84..279c5c1 100644 --- a/internal/vpn/service.go +++ b/internal/vpn/service.go @@ -155,10 +155,13 @@ func (s *VpnService) ApplySettings(settings model.VpnSetting, reservations4, res s.mu.Unlock() go s.serveTUN() - log.Printf("VPN 服务已启动: tun=%s subnet=%s server=%s", tun.Name(), ipNet.String(), serverIP.String()) + go s.startStatsLogger() + log.Printf("[TUN] created name=%s mtu=%d", tun.Name(), settings.MTU) + log.Printf("[TUN] ipv4 addr=%s/%d subnet=%s server=%s", serverIP.String(), prefix, ipNet.String(), serverIP.String()) if ipNet6 != nil { - log.Printf(" IPv6: subnet=%s server=%s", ipNet6.String(), serverIP6.String()) + log.Printf("[TUN] ipv6 addr=%s/%d subnet=%s server=%s", serverIP6.String(), prefix6, ipNet6.String(), serverIP6.String()) } + log.Printf("[VPN] started c2c=%v bufSize=%d", settings.AllowClientToClient, settings.MTU+64) return nil } @@ -174,17 +177,49 @@ func (s *VpnService) serveTUN() { for { n, err := tun.Iface.Read(packet) if err != nil { - log.Printf("TUN 读取结束: %v", err) + log.Printf("[TUN-RX] read ended: %v", err) close(done) return } if n < 1 { continue } - targets := switchx.RouteFromTUN(packet[:n]) - for _, t := range targets { - _ = t.WritePacket(packet[:n]) + pkt := packet[:n] + srcIP, destIP, ok := parseIPAddrs(pkt) + if ok { + log.Printf("[TUN-RX] size=%d proto=%s src=%s dst=%s", n, protoName(pkt), srcIP, destIP) + } else { + log.Printf("[TUN-RX] size=%d parse=fail", n) } + iterStart := time.Now() + targets := switchx.RouteFromTUN(pkt) + for _, t := range targets { + _ = t.WritePacket(pkt) + } + if iterDur := time.Since(iterStart); iterDur > 10*time.Millisecond { + log.Printf("[TUN-RX] slow write iterMs=%d targets=%d size=%d", iterDur.Milliseconds(), len(targets), n) + } + } +} + +func (s *VpnService) startStatsLogger() { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + var lastRx, lastTx int64 + for range ticker.C { + if !s.Running() { + return + } + rx, tx := s.TotalLiveTraffic() + rxRate := rx - lastRx + txRate := tx - lastTx + s.mu.RLock() + n := len(s.clients) + s.mu.RUnlock() + log.Printf("[STATS] clients=%d totalRx=%d totalTx=%d intervalRx=%d(%dB/s) intervalTx=%d(%dB/s)", + n, rx, tx, rxRate, rxRate/10, txRate, txRate/10) + lastRx = rx + lastTx = tx } } @@ -210,7 +245,7 @@ func (s *VpnService) Stop() error { <-done } } - log.Printf("VPN 服务已停止") + log.Printf("[VPN] stopped") return nil } @@ -242,9 +277,15 @@ func (s *VpnService) WriteToTUN(packet []byte) error { tun := s.tun s.mu.RUnlock() if tun == nil { + log.Printf("[TUN-TX] size=%d ok=false err=TUN not ready", len(packet)) return errors.New("TUN 未就绪") } _, err := tun.Iface.Write(packet) + if err != nil { + log.Printf("[TUN-TX] size=%d ok=false err=%v", len(packet), err) + } else { + log.Printf("[TUN-TX] size=%d ok=true", len(packet)) + } return err } diff --git a/internal/vpn/switch.go b/internal/vpn/switch.go index 6159190..7a4a5e8 100644 --- a/internal/vpn/switch.go +++ b/internal/vpn/switch.go @@ -1,6 +1,7 @@ package vpn import ( + "log" "net" "sync" @@ -106,6 +107,16 @@ func parseIPAddrs(packet []byte) (src, dest net.IP, ok bool) { return nil, nil, false } +func protoName(packet []byte) string { + if waterutil.IsIPv4(packet) { + return "IPv4" + } + if waterutil.IsIPv6(packet) { + return "IPv6" + } + return "unknown" +} + func (s *PacketSwitch) allowC2C() bool { s.mu.RLock() v := s.allowClientToClient @@ -116,43 +127,67 @@ func (s *PacketSwitch) allowC2C() bool { func (s *PacketSwitch) RouteFromClient(src SwitchConn, packet []byte) []SwitchConn { srcIP, dest, ok := parseIPAddrs(packet) if !ok { + log.Printf("[ROUTE] client->? parse failed size=%d", len(packet)) return nil } + proto := protoName(packet) // anti-spoof: enforce assigned source IP by version if srcIP != nil { if srcIP.To4() != nil { if !srcIP.Equal(src.AssignedIP()) { + log.Printf("[ROUTE] anti-spoof drop src=%s expected=%s dst=%s proto=%s size=%d", + srcIP, src.AssignedIP(), dest, proto, len(packet)) return nil } } else { assigned6 := src.AssignedIP6() if assigned6 == nil || !srcIP.Equal(assigned6) { + log.Printf("[ROUTE] anti-spoof drop src=%s expected=%s dst=%s proto=%s size=%d", + srcIP, assigned6, dest, proto, len(packet)) return nil } } } if dest.IsGlobalUnicast() { if c := s.findByIP(dest); c != nil && s.allowC2C() { + log.Printf("[ROUTE] client->relay src=%s dst=%s proto=%s size=%d", + src.AssignedIP(), dest, proto, len(packet)) return []SwitchConn{c} } + log.Printf("[ROUTE] client->tun src=%s dst=%s proto=%s size=%d", + src.AssignedIP(), dest, proto, len(packet)) return nil } if s.allowC2C() { - return s.allExcept(src) + targets := s.allExcept(src) + log.Printf("[ROUTE] client->broadcast src=%s dst=%s proto=%s size=%d targets=%d", + src.AssignedIP(), dest, proto, len(packet), len(targets)) + return targets } + log.Printf("[ROUTE] client->drop(broadcast,c2c off) src=%s dst=%s proto=%s size=%d", + src.AssignedIP(), dest, proto, len(packet)) return nil } func (s *PacketSwitch) RouteFromTUN(packet []byte) []SwitchConn { _, dest, ok := parseIPAddrs(packet) if !ok { + log.Printf("[ROUTE] TUN->? parse failed size=%d", len(packet)) return nil } + proto := protoName(packet) if dest.IsGlobalUnicast() { if c := s.findByIP(dest); c != nil { + log.Printf("[ROUTE] TUN->client dst=%s proto=%s size=%d found=true", + dest, proto, len(packet)) return []SwitchConn{c} } + log.Printf("[ROUTE] TUN->client dst=%s proto=%s size=%d found=false", + dest, proto, len(packet)) return nil } - return s.allExcept(nil) + targets := s.allExcept(nil) + log.Printf("[ROUTE] TUN->broadcast dst=%s proto=%s size=%d targets=%d", + dest, proto, len(packet), len(targets)) + return targets } diff --git a/internal/vpn/tun_darwin.go b/internal/vpn/tun_darwin.go index d71d2b5..f2e677d 100644 --- a/internal/vpn/tun_darwin.go +++ b/internal/vpn/tun_darwin.go @@ -4,6 +4,7 @@ package vpn import ( "fmt" + "log" "net" ) @@ -32,5 +33,11 @@ func (t *TUNInterface) AddSubnetRoute(subnet *net.IPNet) error { } func (t *TUNInterface) SetMTU(mtu int) error { - return execCmd("ifconfig", t.Name(), "mtu", fmt.Sprintf("%d", mtu)) + err := execCmd("ifconfig", t.Name(), "mtu", fmt.Sprintf("%d", mtu)) + if err != nil { + log.Printf("[TUN] set MTU dev=%s mtu=%d ok=false err=%v", t.Name(), mtu, err) + } else { + log.Printf("[TUN] set MTU dev=%s mtu=%d ok=true", t.Name(), mtu) + } + return err } diff --git a/internal/vpn/tun_linux.go b/internal/vpn/tun_linux.go index 8fec5c2..c0c6c3b 100644 --- a/internal/vpn/tun_linux.go +++ b/internal/vpn/tun_linux.go @@ -4,6 +4,7 @@ package vpn import ( "fmt" + "log" "net" ) @@ -28,5 +29,11 @@ func (t *TUNInterface) AddSubnetRoute(subnet *net.IPNet) error { } func (t *TUNInterface) SetMTU(mtu int) error { - return execCmd("ip", "link", "set", "dev", t.Name(), "mtu", fmt.Sprintf("%d", mtu)) + err := execCmd("ip", "link", "set", "dev", t.Name(), "mtu", fmt.Sprintf("%d", mtu)) + if err != nil { + log.Printf("[TUN] set MTU dev=%s mtu=%d ok=false err=%v", t.Name(), mtu, err) + } else { + log.Printf("[TUN] set MTU dev=%s mtu=%d ok=true", t.Name(), mtu) + } + return err } diff --git a/internal/vpn/tunnel.go b/internal/vpn/tunnel.go index 621665b..d210b21 100644 --- a/internal/vpn/tunnel.go +++ b/internal/vpn/tunnel.go @@ -41,22 +41,42 @@ type tunnelConn struct { ready atomic.Bool rxBytes atomic.Int64 txBytes atomic.Int64 + rxPkts atomic.Int64 + txPkts atomic.Int64 } func (c *tunnelConn) AssignedIP() net.IP { return c.assignedIP } func (c *tunnelConn) AssignedIP6() net.IP { return c.assignedIP6 } +func (c *tunnelConn) label() string { + s := "user=" + c.user.Username + " ip=" + c.assignedIP.String() + if c.assignedIP6 != nil { + s += " ip6=" + c.assignedIP6.String() + } + return s +} + func (c *tunnelConn) WritePacket(data []byte) error { if !c.ready.Load() || len(data) == 0 { return nil } + lockStart := time.Now() c.writeMu.Lock() + lockWait := time.Since(lockStart) defer c.writeMu.Unlock() c.conn.SetWriteDeadline(time.Now().Add(writeTimeout)) - if err := c.conn.WriteMessage(websocket.BinaryMessage, data); err != nil { + writeStart := time.Now() + err := c.conn.WriteMessage(websocket.BinaryMessage, data) + writeDur := time.Since(writeStart) + if err != nil { + log.Printf("[WS-TX] %s size=%d lockWaitMs=%d writeMs=%d ok=false err=%v", + c.label(), len(data), lockWait.Milliseconds(), writeDur.Milliseconds(), err) return err } c.txBytes.Add(int64(len(data))) + c.txPkts.Add(1) + log.Printf("[WS-TX] %s size=%d lockWaitMs=%d writeMs=%d ok=true", + c.label(), len(data), lockWait.Milliseconds(), writeDur.Milliseconds()) return nil } @@ -150,13 +170,13 @@ func runTunnel(conn *websocket.Conn, user *model.User) { initMsg.ServerIP6 = VPN.ServerIP6().String() } if err := tc.writeControl(initMsg); err != nil { - log.Printf("用户 %s 发送 init 失败: %v", user.Username, err) + log.Printf("[CONN] init failed %s err=%v", tc.label(), err) return } - log.Printf("用户 %s 已连接,分配 IP %s", user.Username, ip4.String()) + log.Printf("[CONN] init sent %s mtu=%d prefix=%d server=%s", tc.label(), settings.MTU, VPN.Prefix(), VPN.ServerIP().String()) if ip6 != nil { - log.Printf(" IPv6: %s", ip6.String()) + log.Printf("[CONN] init6 sent %s ip6=%s prefix6=%d server6=%s", tc.label(), ip6.String(), VPN.Prefix6(), VPN.ServerIP6().String()) } conn.SetReadLimit(maxMessageSize) @@ -170,12 +190,32 @@ func runTunnel(conn *websocket.Conn, user *model.User) { tc.writeMu.Lock() if err := conn.WriteControl(websocket.PingMessage, nil, time.Now().Add(writeTimeout)); err != nil { tc.writeMu.Unlock() + log.Printf("[PING] failed %s err=%v", tc.label(), err) return } tc.writeMu.Unlock() } }() + go func() { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + var lastRx, lastTx int64 + for range ticker.C { + if !tc.ready.Load() { + return + } + rx := tc.rxBytes.Load() + tx := tc.txBytes.Load() + rxRate := rx - lastRx + txRate := tx - lastTx + log.Printf("[CONN] stats %s rxBytes=%d txBytes=%d rxRate=%dB/s txRate=%dB/s rxPkts=%d txPkts=%d", + tc.label(), rx, tx, rxRate/10, txRate/10, tc.rxPkts.Load(), tc.txPkts.Load()) + lastRx = rx + lastTx = tx + } + }() + conn.SetPongHandler(func(string) error { conn.SetReadDeadline(time.Now().Add(readTimeout)) return nil @@ -184,7 +224,10 @@ func runTunnel(conn *websocket.Conn, user *model.User) { for { messageType, data, err := conn.ReadMessage() if err != nil { - log.Printf("用户 %s 断开连接: %v", user.Username, err) + duration := time.Since(tc.connectedAt) + log.Printf("[CONN] disconnected %s duration=%s rxBytes=%d txBytes=%d rxPkts=%d txPkts=%d err=%v", + tc.label(), duration.Round(time.Second), tc.rxBytes.Load(), tc.txBytes.Load(), + tc.rxPkts.Load(), tc.txPkts.Load(), err) return } @@ -196,7 +239,7 @@ func runTunnel(conn *websocket.Conn, user *model.User) { if msg.Type == "ready" && !tc.ready.Load() { tc.ready.Store(true) conn.SetReadDeadline(time.Now().Add(readTimeout)) - log.Printf("用户 %s 就绪 (IP %s)", user.Username, ip4.String()) + log.Printf("[CONN] ready %s", tc.label()) } continue } @@ -207,7 +250,7 @@ func runTunnel(conn *websocket.Conn, user *model.User) { if !tc.ready.Load() { if time.Now().After(readyDeadline) { - log.Printf("用户 %s 等待 ready 超时", user.Username) + log.Printf("[CONN] ready timeout %s", tc.label()) return } conn.SetReadDeadline(readyDeadline) @@ -215,11 +258,20 @@ func runTunnel(conn *websocket.Conn, user *model.User) { } tc.rxBytes.Add(int64(len(data))) + tc.rxPkts.Add(1) + + srcIP, destIP, ok := parseIPAddrs(data) + if ok { + log.Printf("[WS-RX] %s size=%d proto=%s src=%s dst=%s", + tc.label(), len(data), protoName(data), srcIP, destIP) + } else { + log.Printf("[WS-RX] %s size=%d parse=fail", tc.label(), len(data)) + } targets := VPN.RouteFromClient(tc, data) if len(targets) == 0 { if err := VPN.WriteToTUN(data); err != nil { - log.Printf("用户 %s 写入 TUN 失败: %v", user.Username, err) + log.Printf("[WS-RX] %s action=tun err=%v size=%d", tc.label(), err, len(data)) } continue }