feat: add structured logging throughout audio pipeline
Add a lightweight log package (log/log.go) with debug/info/warn/error levels. Wire it up via a -log flag (default: warn). Logging covers: UDP transport: - Socket open/failure, CryptSetup receipt, crypto initialization - First UDP audio send confirmation - UDP read errors and packet counts - TCP fallback reason (no socket, waiting for crypto) Audio pipeline: - Source (mic) routine start with config details - Audio stream start/end per remote user - First audio stream creation for each user - Packet loss detection with sequence gaps and PLC generation - Decoder reset on sequence reordering Usage: barnard -log=debug -server=mumble.example.com Levels: debug, info, warn, error
This commit is contained in:
committed by
Brandon McGinty
parent
a218493eb0
commit
e483512f74
+21
-2
@@ -10,6 +10,7 @@ import (
|
||||
"time"
|
||||
|
||||
"git.stormux.org/storm/barnard/gumble/gumble/MumbleProto"
|
||||
"git.stormux.org/storm/barnard/log"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
@@ -172,7 +173,11 @@ func DialWithDialer(dialer *net.Dialer, config *Config, tlsConfig *tls.Config) (
|
||||
|
||||
// Start UDP transport for lower-latency audio. This is best-effort;
|
||||
// if UDP fails, audio falls back to TCP tunneling.
|
||||
client.startUDP()
|
||||
if err := client.startUDP(); err != nil {
|
||||
log.Warn("UDP setup failed, audio will use TCP tunnel: %v", err)
|
||||
} else if client.udpConn != nil {
|
||||
log.Info("UDP socket opened to %s, waiting for CryptSetup", client.udpConn.RemoteAddr())
|
||||
}
|
||||
|
||||
return client, nil
|
||||
}
|
||||
@@ -258,6 +263,7 @@ func (c *Client) readRoutine() {
|
||||
|
||||
// Clean up UDP connection
|
||||
if c.udpConn != nil {
|
||||
log.Debug("closing UDP connection")
|
||||
c.udpConn.Close()
|
||||
c.udpConn = nil
|
||||
}
|
||||
@@ -316,12 +322,25 @@ func (c *Client) EnableStereoEncoder() {
|
||||
|
||||
// WriteAudio writes an audio packet, preferring UDP when encryption is
|
||||
// set up. Falls back to TCP-tunneled audio when UDP is unavailable.
|
||||
var udpFallbackLogged bool
|
||||
|
||||
func (c *Client) WriteAudio(format, target byte, sequence int64, final bool, data []byte, X, Y, Z *float32) error {
|
||||
// Try UDP first
|
||||
if sent, err := c.WriteAudioUDP(format, target, sequence, final, data, X, Y, Z); sent {
|
||||
if err != nil {
|
||||
log.Error("UDP send error: %v", err)
|
||||
}
|
||||
return err
|
||||
}
|
||||
// Fall back to TCP tunnel
|
||||
// Fall back to TCP tunnel — log once per process
|
||||
if !udpFallbackLogged {
|
||||
udpFallbackLogged = true
|
||||
if c.udpConn == nil {
|
||||
log.Info("no UDP socket, audio using TCP tunnel")
|
||||
} else if !c.cryptOut.initialized {
|
||||
log.Info("waiting for CryptSetup, audio using TCP tunnel")
|
||||
}
|
||||
}
|
||||
return c.Conn.WriteAudio(format, target, sequence, final, data, X, Y, Z)
|
||||
}
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ import (
|
||||
"sync"
|
||||
|
||||
"git.stormux.org/storm/barnard/gumble/gumble/MumbleProto"
|
||||
"git.stormux.org/storm/barnard/log"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
@@ -289,14 +290,21 @@ func (c *Client) handleCryptSetup(buffer []byte) error {
|
||||
defer c.volatile.Unlock()
|
||||
|
||||
if packet.Key != nil && packet.ClientNonce != nil && packet.ServerNonce != nil {
|
||||
log.Info("received CryptSetup: key_len=%d client_nonce_len=%d server_nonce_len=%d",
|
||||
len(packet.Key), len(packet.ClientNonce), len(packet.ServerNonce))
|
||||
c.cryptOut.setup(packet.Key, packet.ClientNonce)
|
||||
c.cryptIn.setup(packet.Key, packet.ServerNonce)
|
||||
} else {
|
||||
log.Debug("received CryptSetup with incomplete fields")
|
||||
}
|
||||
|
||||
if c.cryptOut.initialized && c.udpConn != nil && !c.udpActive {
|
||||
c.udpActive = true
|
||||
log.Info("UDP crypto ready, starting UDP reader and pinger")
|
||||
go c.udpReadRoutine()
|
||||
go c.udpPingRoutine()
|
||||
} else if c.cryptOut.initialized && c.udpConn == nil {
|
||||
log.Warn("crypto ready but no UDP socket — audio will use TCP tunnel")
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"time"
|
||||
|
||||
"git.stormux.org/storm/barnard/gumble/gumble/MumbleProto"
|
||||
"git.stormux.org/storm/barnard/log"
|
||||
"git.stormux.org/storm/barnard/gumble/gumble/varint"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
@@ -125,11 +126,15 @@ func (c *Client) handleUDPTunnel(buffer []byte) error {
|
||||
if user.audioSequenceValid {
|
||||
gap := seq - user.audioSequence
|
||||
if gap > 1 && gap < 100 {
|
||||
log.Debug("audio seq gap for %s: %d -> %d (loss=%d), generating PLC",
|
||||
user.Name, user.audioSequence, seq, gap-1)
|
||||
// Lost packets detected; generate PLC frames for each.
|
||||
for i := int64(1); i < gap; i++ {
|
||||
c.dispatchPLC(user, audioTarget, decoder)
|
||||
}
|
||||
} else if gap < 0 && gap > -100 {
|
||||
log.Debug("audio seq reorder for %s: %d -> %d, resetting decoder",
|
||||
user.Name, user.audioSequence, seq)
|
||||
// Reordered packet — reset decoder to prevent corruption.
|
||||
decoder.Reset()
|
||||
}
|
||||
@@ -214,6 +219,7 @@ func (c *Client) dispatchAudio(user *User, packet *AudioPacket) {
|
||||
if ch == nil {
|
||||
ch = make(chan *AudioPacket)
|
||||
item.streams[user] = ch
|
||||
log.Debug("new audio stream from %s (session=%d)", user.Name, user.Session)
|
||||
streamEvent := AudioStreamEvent{
|
||||
Client: c,
|
||||
User: user,
|
||||
|
||||
+23
-3
@@ -8,6 +8,7 @@ import (
|
||||
"time"
|
||||
|
||||
"git.stormux.org/storm/barnard/gumble/gumble/varint"
|
||||
"git.stormux.org/storm/barnard/log"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -25,31 +26,41 @@ const (
|
||||
// audio packets. It should be called after the server address is known.
|
||||
func (c *Client) startUDP() error {
|
||||
addr := c.Conn.RemoteAddr()
|
||||
log.Debug("attempting UDP connection to %s", addr.String())
|
||||
|
||||
udpAddr, err := net.ResolveUDPAddr("udp", addr.String())
|
||||
if err != nil {
|
||||
log.Warn("failed to resolve UDP address %s: %v", addr.String(), err)
|
||||
return err
|
||||
}
|
||||
|
||||
conn, err := net.DialUDP("udp", nil, udpAddr)
|
||||
if err != nil {
|
||||
// UDP might not be available; fall back to TCP-only audio.
|
||||
// This is not an error — many Mumble servers work fine TCP-only.
|
||||
log.Warn("UDP dial failed (audio will use TCP tunnel): %v", err)
|
||||
return nil
|
||||
}
|
||||
|
||||
c.udpConn = conn
|
||||
log.Info("UDP socket connected to %s", conn.RemoteAddr())
|
||||
// The UDP reader and pinger will be started once CryptSetup is received.
|
||||
return nil
|
||||
}
|
||||
|
||||
// udpReadRoutine reads encrypted UDP audio packets from the server.
|
||||
func (c *Client) udpReadRoutine() {
|
||||
log.Info("UDP reader started")
|
||||
buf := make([]byte, maxUDPPacketSize)
|
||||
var packetCount uint64
|
||||
for {
|
||||
n, _, err := c.udpConn.ReadFromUDP(buf)
|
||||
n, addr, err := c.udpConn.ReadFromUDP(buf)
|
||||
if err != nil {
|
||||
log.Warn("UDP read error (stopping reader): %v", err)
|
||||
return
|
||||
}
|
||||
packetCount++
|
||||
if packetCount <= 3 || packetCount%100 == 0 {
|
||||
log.Debug("UDP recv #%d: %d bytes from %s", packetCount, n, addr)
|
||||
}
|
||||
packet := make([]byte, n)
|
||||
copy(packet, buf[:n])
|
||||
c.handleUDPPacket(packet)
|
||||
@@ -89,11 +100,20 @@ func (c *Client) sendUDPPing(counter uint32) {
|
||||
|
||||
// WriteAudioUDP writes an encrypted audio packet over UDP.
|
||||
// Returns true if the packet was sent over UDP, false if TCP should be used.
|
||||
var firstUDPSendLogged bool
|
||||
|
||||
func (c *Client) WriteAudioUDP(format, target byte, sequence int64, final bool, data []byte, X, Y, Z *float32) (bool, error) {
|
||||
if c.udpConn == nil || !c.cryptOut.initialized {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
if !firstUDPSendLogged {
|
||||
firstUDPSendLogged = true
|
||||
log.Info("first UDP audio packet sent — UDP transport active")
|
||||
} else {
|
||||
log.Debug("UDP send: seq=%d len=%d final=%v", sequence, len(data), final)
|
||||
}
|
||||
|
||||
// Build the unencrypted header
|
||||
var header [1 + varint.MaxVarintLen*2]byte
|
||||
header[0] = (format << 5) | target
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"time"
|
||||
|
||||
"git.stormux.org/storm/barnard/audio"
|
||||
"git.stormux.org/storm/barnard/log"
|
||||
"git.stormux.org/storm/barnard/gumble/go-openal/openal"
|
||||
"git.stormux.org/storm/barnard/gumble/gumble"
|
||||
"git.stormux.org/storm/barnard/noise"
|
||||
@@ -242,6 +243,7 @@ func (s *Stream) SetMicVolume(change float32, relative bool) {
|
||||
|
||||
func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
go func(e *gumble.AudioStreamEvent) {
|
||||
log.Info("audio stream started for user %s", e.User.Name)
|
||||
var source = openal.NewSource()
|
||||
e.User.SetAudioSource(&source)
|
||||
|
||||
@@ -359,6 +361,7 @@ func (s *Stream) OnAudioStream(e *gumble.AudioStreamEvent) {
|
||||
reclaim()
|
||||
emptyBufs.Delete()
|
||||
source.Delete()
|
||||
log.Debug("audio stream ended for user %s", e.User.Name)
|
||||
}(e)
|
||||
}
|
||||
|
||||
@@ -478,6 +481,8 @@ func (s *Stream) processAudioPacket(packet *gumble.AudioPacket, user *gumble.Use
|
||||
}
|
||||
|
||||
func (s *Stream) sourceRoutine(inputDevice *string) {
|
||||
log.Info("source routine started: interval=%v frameSize=%d channels=%d",
|
||||
s.client.Config.AudioInterval, s.client.Config.AudioFrameSize(), s.sourceChannels)
|
||||
interval := s.client.Config.AudioInterval
|
||||
frameSize := s.client.Config.AudioFrameSize()
|
||||
|
||||
|
||||
+116
@@ -0,0 +1,116 @@
|
||||
// Package log provides debug logging for the barnard audio pipeline.
|
||||
// Logging is disabled by default and can be enabled via SetLogger.
|
||||
package log
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Level represents logging severity.
|
||||
type Level int
|
||||
|
||||
const (
|
||||
LevelDebug Level = iota
|
||||
LevelInfo
|
||||
LevelWarn
|
||||
LevelError
|
||||
)
|
||||
|
||||
func (l Level) String() string {
|
||||
switch l {
|
||||
case LevelDebug:
|
||||
return "DEBUG"
|
||||
case LevelInfo:
|
||||
return "INFO"
|
||||
case LevelWarn:
|
||||
return "WARN"
|
||||
case LevelError:
|
||||
return "ERROR"
|
||||
default:
|
||||
return "???"
|
||||
}
|
||||
}
|
||||
|
||||
// Logger receives log messages. The default logger is a no-op.
|
||||
type Logger interface {
|
||||
Log(level Level, format string, args ...interface{})
|
||||
}
|
||||
|
||||
var (
|
||||
mu sync.Mutex
|
||||
logger Logger = &nopLogger{}
|
||||
)
|
||||
|
||||
type nopLogger struct{}
|
||||
|
||||
func (n *nopLogger) Log(level Level, format string, args ...interface{}) {}
|
||||
|
||||
// SetLogger sets the destination for log messages. Pass nil to disable.
|
||||
func SetLogger(l Logger) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if l == nil {
|
||||
logger = &nopLogger{}
|
||||
} else {
|
||||
logger = l
|
||||
}
|
||||
}
|
||||
|
||||
// WriterLogger is a simple Logger that writes to an io.Writer.
|
||||
type WriterLogger struct {
|
||||
mu sync.Mutex
|
||||
w io.Writer
|
||||
level Level
|
||||
buf []byte
|
||||
}
|
||||
|
||||
// NewWriterLogger creates a logger that writes to w, filtering below level.
|
||||
func NewWriterLogger(w io.Writer, level Level) *WriterLogger {
|
||||
if w == nil {
|
||||
w = os.Stderr
|
||||
}
|
||||
return &WriterLogger{w: w, level: level}
|
||||
}
|
||||
|
||||
func (wl *WriterLogger) Log(level Level, format string, args ...interface{}) {
|
||||
if level < wl.level {
|
||||
return
|
||||
}
|
||||
wl.mu.Lock()
|
||||
defer wl.mu.Unlock()
|
||||
now := time.Now().Format("15:04:05.000")
|
||||
msg := fmt.Sprintf(format, args...)
|
||||
fmt.Fprintf(wl.w, "%s [%-5s] %s\n", now, level.String(), msg)
|
||||
}
|
||||
|
||||
func Debug(format string, args ...interface{}) {
|
||||
mu.Lock()
|
||||
l := logger
|
||||
mu.Unlock()
|
||||
l.Log(LevelDebug, format, args...)
|
||||
}
|
||||
|
||||
func Info(format string, args ...interface{}) {
|
||||
mu.Lock()
|
||||
l := logger
|
||||
mu.Unlock()
|
||||
l.Log(LevelInfo, format, args...)
|
||||
}
|
||||
|
||||
func Warn(format string, args ...interface{}) {
|
||||
mu.Lock()
|
||||
l := logger
|
||||
mu.Unlock()
|
||||
l.Log(LevelWarn, format, args...)
|
||||
}
|
||||
|
||||
func Error(format string, args ...interface{}) {
|
||||
mu.Lock()
|
||||
l := logger
|
||||
mu.Unlock()
|
||||
l.Log(LevelError, format, args...)
|
||||
}
|
||||
@@ -15,6 +15,8 @@ import (
|
||||
"strings"
|
||||
"syscall"
|
||||
|
||||
barnlog "git.stormux.org/storm/barnard/log"
|
||||
|
||||
"git.stormux.org/storm/barnard/config"
|
||||
"git.stormux.org/storm/barnard/gumble/go-openal/openal"
|
||||
"git.stormux.org/storm/barnard/gumble/gumble"
|
||||
@@ -114,9 +116,26 @@ func main() {
|
||||
buffers := flag.Int("buffers", 16, "number of audio buffers to use")
|
||||
profile := flag.Bool("profile", false, "add http server to serve profiles")
|
||||
noiseSuppressionEnabled := flag.Bool("noise-suppression", false, "enable noise suppression for microphone input")
|
||||
logLevel := flag.String("log", "warn", "log level: debug, info, warn, error")
|
||||
|
||||
flag.Parse()
|
||||
|
||||
// Set up logging
|
||||
var level barnlog.Level
|
||||
switch strings.ToLower(*logLevel) {
|
||||
case "debug":
|
||||
level = barnlog.LevelDebug
|
||||
case "info":
|
||||
level = barnlog.LevelInfo
|
||||
case "warn":
|
||||
level = barnlog.LevelWarn
|
||||
case "error":
|
||||
level = barnlog.LevelError
|
||||
default:
|
||||
level = barnlog.LevelWarn
|
||||
}
|
||||
barnlog.SetLogger(barnlog.NewWriterLogger(os.Stderr, level))
|
||||
|
||||
if *profile == true {
|
||||
go func() {
|
||||
log.Println(http.ListenAndServe("localhost:6060", nil))
|
||||
|
||||
Reference in New Issue
Block a user