From c9a66a7fbc5c39a220337b86da5dd9ec63fd9fb7 Mon Sep 17 00:00:00 2001 From: riccardom Date: Tue, 28 Jul 2026 22:43:07 +0200 Subject: [PATCH] Adds log tracepoints - Add a trace slog level (NB_PQ_MLKEM_LOG_LEVEL=trace) and move the verbose per-exchange lifecycle logs (offer/answer/PSK/ack/rotation) to it, so debug stays quiet and troubleshooting is opt-in. - Stop logging the raw preshared key; drop the temporary pqkem-dbg OnRemoteOffer/ OnRemoteAnswer probes. - Demote the per-handshake conn log to trace. --- client/internal/pqkem/convergence.go | 38 +++++++++++++++ client/internal/pqkem/env.go | 69 +++++++++++++++++++++++++--- client/internal/pqkem/manager.go | 11 ++++- client/internal/pqkem_adapter.go | 2 +- 4 files changed, 112 insertions(+), 8 deletions(-) diff --git a/client/internal/pqkem/convergence.go b/client/internal/pqkem/convergence.go index 63e1d61e9..deaf68cba 100644 --- a/client/internal/pqkem/convergence.go +++ b/client/internal/pqkem/convergence.go @@ -2,9 +2,22 @@ package pqkem import ( "context" + "crypto/sha256" + "encoding/hex" "time" ) +// idHex renders an exchange ID for logs. +func idHex(id ExchangeID) string { return hex.EncodeToString(id[:]) } + +// pskFingerprint is a short, non-secret digest of a derived PSK: identical on both +// peers iff they derived the same key. Logged instead of the raw PSK so debug logs +// never carry the actual WireGuard preshared key. +func pskFingerprint(psk PSK) string { + sum := sha256.Sum256(psk[:]) + return hex.EncodeToString(sum[:8]) +} + // startExchange creates a fresh initiator exchange (acknowledging ackID, zero for a // bootstrap) and returns the framed offer for the caller to send — pushed over the // data path for a chained rekey, or handed to the host for signalling when viaSignal @@ -41,6 +54,12 @@ func (m *Manager) startExchange(remoteID RemoteID, viaSignal bool, ackID Exchang m.wait.Add(1) go m.initiatorLoop(ctx, remoteID, id) + + via := "data-path" + if viaSignal { + via = "signal" + } + m.trace("pqkem: offer sent", "peer", remoteID, "exchange", idHex(id), "acks", idHex(ackID), "via", via) return raw, nil } @@ -49,6 +68,8 @@ func (m *Manager) startExchange(remoteID RemoteID, viaSignal bool, ackID Exchang // then derives the PSK for the new offer, commits it optimistically, and returns the // framed answer. A duplicate offer returns the cached answer without re-deriving. func (m *Manager) processOffer(remoteID RemoteID, o *OfferMsg) ([]byte, error) { + m.trace("pqkem: offer received", "peer", remoteID, "exchange", idHex(o.ExchangeID), "acks", idHex(o.AckID)) + if o.AckID != (ExchangeID{}) { m.ackConverged(remoteID, o.AckID) } @@ -60,6 +81,7 @@ func (m *Manager) processOffer(remoteID RemoteID, o *OfferMsg) ([]byte, error) { if state == stateReserved { return nil, nil } + m.trace("pqkem: duplicate offer, resending cached answer", "peer", remoteID, "exchange", idHex(o.ExchangeID)) return last, nil } // Reserve the slot so a concurrent duplicate offer bails. @@ -79,6 +101,7 @@ func (m *Manager) processOffer(remoteID RemoteID, o *OfferMsg) ([]byte, error) { ex := m.exchanges[remoteID] if ex == nil || ex.id != o.ExchangeID { m.mu.Unlock() + m.trace("pqkem: exchange superseded during respond, dropping answer", "peer", remoteID, "exchange", idHex(o.ExchangeID)) return nil, nil } ex.state = stateAwaitingAck @@ -87,10 +110,13 @@ func (m *Manager) processOffer(remoteID RemoteID, o *OfferMsg) ([]byte, error) { m.psks[remoteID] = psk m.mu.Unlock() + m.trace("pqkem: new PSK derived", "peer", remoteID, "exchange", idHex(o.ExchangeID), "role", "responder", "psk_fp", pskFingerprint(psk)) + // Commit optimistically so our data path can rekey to the new PSK. if err := m.cbHandler.OnNewPSKReady(remoteID, psk); err != nil { return nil, err } + m.trace("pqkem: answer sent", "peer", remoteID, "exchange", idHex(o.ExchangeID)) return raw, nil } @@ -102,7 +128,12 @@ func (m *Manager) processAnswer(remoteID RemoteID, a *AnswerMsg) error { m.mu.Lock() ex := m.exchanges[remoteID] if ex == nil || ex.id != a.ExchangeID || ex.state != stateAwaitingAnswer { + haveID := "none" + if ex != nil { + haveID = idHex(ex.id) + } m.mu.Unlock() + m.trace("pqkem: unexpected answer dropped (inconsistency)", "peer", remoteID, "answer_for", idHex(a.ExchangeID), "have_exchange", haveID) return nil } ex.state = stateAwaitingRekey @@ -110,6 +141,8 @@ func (m *Manager) processAnswer(remoteID RemoteID, a *AnswerMsg) error { ex.initiator = nil m.mu.Unlock() + m.trace("pqkem: answer received", "peer", remoteID, "exchange", idHex(a.ExchangeID)) + psk, err := init.Finish(a.KEMAnswer, m.binding(remoteID)) if err != nil { return err @@ -122,6 +155,8 @@ func (m *Manager) processAnswer(remoteID RemoteID, a *AnswerMsg) error { m.psks[remoteID] = psk m.mu.Unlock() + m.trace("pqkem: new PSK derived", "peer", remoteID, "exchange", idHex(a.ExchangeID), "role", "initiator", "psk_fp", pskFingerprint(psk)) + return m.cbHandler.OnNewPSKReady(remoteID, psk) } @@ -133,6 +168,7 @@ func (m *Manager) ackConverged(remoteID RemoteID, ackID ExchangeID) { ex := m.exchanges[remoteID] if ex == nil || ex.id != ackID || ex.state != stateAwaitingAck { m.mu.Unlock() + m.trace("pqkem: ack for unknown/mismatched exchange, ignored (inconsistency)", "peer", remoteID, "acks", idHex(ackID)) return } delete(m.exchanges, remoteID) @@ -140,6 +176,8 @@ func (m *Manager) ackConverged(remoteID RemoteID, ackID ExchangeID) { m.failures[remoteID] = 0 _ = time.Since(ex.startedAt) // convergence latency (metrics hook, later step) m.mu.Unlock() + + m.trace("pqkem: previous exchange confirmed by ack", "peer", remoteID, "exchange", idHex(ackID)) } // initiatorLoop enforces the offer->answer convergence deadline and retransmits the diff --git a/client/internal/pqkem/env.go b/client/internal/pqkem/env.go index 8eb402672..f299d5fce 100644 --- a/client/internal/pqkem/env.go +++ b/client/internal/pqkem/env.go @@ -1,6 +1,7 @@ package pqkem import ( + "context" "log/slog" "os" "strconv" @@ -34,19 +35,26 @@ func Enabled() bool { return enabled } -// EnvLogLevel overrides the ML-KEM manager's slog level (debug/info/warn/error). -// Defaults to info. +// EnvLogLevel overrides the ML-KEM manager's slog level (trace/debug/info/warn/error). +// Defaults to info. The verbose per-exchange lifecycle logs are emitted at trace. const EnvLogLevel = "NB_PQ_MLKEM_LOG_LEVEL" -// NewLogger builds the slog logger for the ML-KEM manager: a text handler to stdout -// at the level from EnvLogLevel. Mirrors the Rosenpass manager's logger setup so PQ -// components log consistently. +// LevelTrace is a custom slog level below Debug for the verbose per-exchange lifecycle +// logs, so they stay off unless NB_PQ_MLKEM_LOG_LEVEL=trace (and the daemon log level +// is trace, since the records are forwarded to logrus). +const LevelTrace = slog.LevelDebug - 4 + +// NewLogger builds the slog logger for the ML-KEM manager. It forwards records to +// logrus so PQ logs land in the same sink as the rest of the daemon (console + +// client.log) rather than stdout. Verbosity is gated by EnvLogLevel. func NewLogger() *slog.Logger { - return slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: logLevel()})) + return slog.New(slogToLogrus{}) } func logLevel() slog.Level { switch strings.ToLower(strings.TrimSpace(os.Getenv(EnvLogLevel))) { + case "trace": + return LevelTrace case "debug": return slog.LevelDebug case "warn": @@ -57,3 +65,52 @@ func logLevel() slog.Level { return slog.LevelInfo } } + +// slogToLogrus is a slog.Handler that forwards records to logrus, so the ML-KEM +// manager's logs go wherever the daemon's logrus is configured (console + client.log) +// instead of stdout. Verbosity is gated by EnvLogLevel via logLevel(). +type slogToLogrus struct { + fields log.Fields +} + +func (h slogToLogrus) Enabled(_ context.Context, level slog.Level) bool { + return level >= logLevel() +} + +func (h slogToLogrus) Handle(_ context.Context, r slog.Record) error { + fields := make(log.Fields, len(h.fields)+r.NumAttrs()) + for k, v := range h.fields { + fields[k] = v + } + r.Attrs(func(a slog.Attr) bool { + fields[a.Key] = a.Value.Any() + return true + }) + entry := log.WithFields(fields) + switch { + case r.Level >= slog.LevelError: + entry.Error(r.Message) + case r.Level >= slog.LevelWarn: + entry.Warn(r.Message) + case r.Level >= slog.LevelInfo: + entry.Info(r.Message) + case r.Level >= slog.LevelDebug: + entry.Debug(r.Message) + default: + entry.Trace(r.Message) + } + return nil +} + +func (h slogToLogrus) WithAttrs(attrs []slog.Attr) slog.Handler { + fields := make(log.Fields, len(h.fields)+len(attrs)) + for k, v := range h.fields { + fields[k] = v + } + for _, a := range attrs { + fields[a.Key] = a.Value.Any() + } + return slogToLogrus{fields: fields} +} + +func (h slogToLogrus) WithGroup(_ string) slog.Handler { return h } diff --git a/client/internal/pqkem/manager.go b/client/internal/pqkem/manager.go index 9cc06a00a..9cf4a189b 100644 --- a/client/internal/pqkem/manager.go +++ b/client/internal/pqkem/manager.go @@ -169,6 +169,12 @@ func (m *Manager) PSK(remoteID RemoteID) (PSK, bool) { return psk, ok } +// trace logs at LevelTrace, the verbose per-exchange lifecycle level gated by +// NB_PQ_MLKEM_LOG_LEVEL=trace. +func (m *Manager) trace(msg string, args ...any) { + m.logger.Log(context.Background(), LevelTrace, msg, args...) +} + // AddPeer registers where a peer's data-path messages are sent and received: its // overlay endpoint (IP:port). Re-adding updates the endpoint. func (m *Manager) AddPeer(remoteID RemoteID, endpoint netip.AddrPort) { @@ -280,7 +286,7 @@ func (m *Manager) onDataPathInbound(src netip.AddrPort, msg []byte) { return } if err := m.OnDataPathMessage(remoteID, msg); err != nil { - m.logger.Debug("pqkem: inbound", "peer", remoteID, "err", err) + m.trace("pqkem: inbound", "peer", remoteID, "err", err) } } @@ -323,6 +329,7 @@ func (m *Manager) OnDataPathRekeyed(remoteID RemoteID) { } m.mu.Unlock() + m.trace("pqkem: data-path rekey signal", "peer", remoteID, "chaining", chain) if !chain { return } @@ -333,7 +340,9 @@ func (m *Manager) OnDataPathRekeyed(remoteID RemoteID) { } if err := m.pushDataPath(remoteID, offer); err != nil { m.logger.Warn("pqkem: send chain offer failed", "peer", remoteID, "err", err) + return } + m.trace("pqkem: chain offer sent over data path", "peer", remoteID) } // OnDataPathDown notifies that the peer's data path went down. Rotations resume once diff --git a/client/internal/pqkem_adapter.go b/client/internal/pqkem_adapter.go index 94ce8ccf2..074039b9e 100644 --- a/client/internal/pqkem_adapter.go +++ b/client/internal/pqkem_adapter.go @@ -27,7 +27,7 @@ func (h pqCallbackHandler) OnNewPSKReady(remoteID pqkem.RemoteID, psk pqkem.PSK) // updateOnly: applies to an already-configured peer (rotation). At bootstrap the // peer is not configured yet, so this is a no-op there and the PSK is instead // pulled at peer-config time (pqHandshaker.PSK / conn.presharedKey). - log.Debugf("pqkem: programming PSK for peer %s", remoteID) + log.Tracef("pqkem: programming PSK for peer %s", remoteID) return h.wg.SetPresharedKey(string(remoteID), wgtypes.Key(psk), true) }