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.
This commit is contained in:
riccardom
2026-07-28 22:43:07 +02:00
parent e447011dc3
commit c9a66a7fbc
4 changed files with 112 additions and 8 deletions

View File

@@ -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

View File

@@ -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 }

View File

@@ -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

View File

@@ -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)
}