mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-11 16:09:07 +02:00
The nonce secret came from the server's WireGuard key, which is generated afresh in every process and never persisted. A nonce carries no state, so the only thing that lets one instance verify what another issued is deriving the same secret — and that premise, written in the comment above the challenger, was not met: every instance had its own key. A peer reconnecting after a restart therefore presented a nonce minted under the previous secret, verification failed with a mismatch, its whole proof set was rejected and the certificates stored for it were dropped until it signed again. Reproduced three times on the lab, each one logging "nonce was not issued to this peer", which only a changed secret produces. On a single instance it costs seconds of lost policy access per restart; across instances it is not transient at all, because every reconnect that lands elsewhere is rejected the same way. Derive from the data store encryption key instead: it is generated once, written back to the configuration and read by every instance, so it survives restarts and is shared. Where none is configured the secret falls back to the WireGuard key with a warning — degraded but still unpredictable, which is the property that matters most: a peer able to guess it could mint the nonces of future windows, sign them while its key is present and keep passing after it is gone. The challenger is now built once and passed to the two places that need it, rather than re-derived per message.
141 lines
4.6 KiB
Go
141 lines
4.6 KiB
Go
package grpc
|
|
|
|
import (
|
|
"context"
|
|
"time"
|
|
|
|
"github.com/netbirdio/netbird/encryption"
|
|
"github.com/netbirdio/netbird/management/internals/controllers/network_map"
|
|
"github.com/netbirdio/netbird/management/server/telemetry"
|
|
"github.com/netbirdio/netbird/shared/management/certposture"
|
|
"github.com/netbirdio/netbird/shared/management/proto"
|
|
log "github.com/sirupsen/logrus"
|
|
"golang.zx2c4.com/wireguard/wgctrl/wgtypes"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/status"
|
|
)
|
|
|
|
func PeerUpdateHandlerFactory(
|
|
peerKey wgtypes.Key,
|
|
updates chan *network_map.UpdateMessage,
|
|
secretsManager SecretsManager,
|
|
challenger *certposture.Challenger,
|
|
srv proto.ManagementService_SyncServer,
|
|
cleanupfunc func()) *PeerUpdateHandler {
|
|
return &PeerUpdateHandler{
|
|
peerKey: peerKey,
|
|
updates: updates,
|
|
secretsManager: secretsManager,
|
|
challenger: challenger,
|
|
srv: srv,
|
|
encrypter: encryption.DefaultEncrypter{},
|
|
debouncer: NewUpdateDebouncer(1000 * time.Millisecond),
|
|
cleanupFunc: cleanupfunc,
|
|
}
|
|
}
|
|
|
|
// PeerUpdateHandler sends updates to the connected peer until the updates channel is closed.
|
|
// It implements a backpressure mechanism that sends the first update immediately,
|
|
// then debounces subsequent rapid updates, ensuring only the latest update is sent
|
|
// after a quiet period.
|
|
type PeerUpdateHandler struct {
|
|
peerKey wgtypes.Key
|
|
updates chan *network_map.UpdateMessage
|
|
appMetrics telemetry.AppMetrics
|
|
secretsManager SecretsManager
|
|
challenger *certposture.Challenger
|
|
srv syncSender
|
|
encrypter encryption.Encrypter
|
|
debouncer Debouncer
|
|
cleanupFunc func()
|
|
}
|
|
|
|
func (pu *PeerUpdateHandler) WithMetrics(appMetrics telemetry.AppMetrics) *PeerUpdateHandler {
|
|
pu.appMetrics = appMetrics
|
|
return pu
|
|
}
|
|
|
|
//go:generate go tool mockgen -source=./peer_update_handler.go -destination=./sync_sender_mock.go -package=grpc
|
|
type syncSender interface {
|
|
Send(*proto.EncryptedMessage) error
|
|
Context() context.Context
|
|
}
|
|
|
|
func (pu *PeerUpdateHandler) HandleUpdates(ctx context.Context) error {
|
|
log.WithContext(ctx).Tracef("starting to handle updates for peer %s", pu.peerKey.String())
|
|
|
|
defer pu.debouncer.Stop()
|
|
|
|
for {
|
|
select {
|
|
// condition when there are some updates
|
|
// todo set the updates channel size to 1
|
|
case update, open := <-pu.updates:
|
|
if pu.appMetrics != nil {
|
|
pu.appMetrics.GRPCMetrics().UpdateChannelQueueLength(len(pu.updates) + 1)
|
|
}
|
|
|
|
if !open {
|
|
log.WithContext(ctx).Debugf("updates channel for peer %s was closed", pu.peerKey.String())
|
|
pu.cleanupFunc()
|
|
return nil
|
|
}
|
|
|
|
log.WithContext(ctx).Tracef("received an update for peer %s", pu.peerKey.String())
|
|
if pu.debouncer.ProcessUpdate(update) {
|
|
// Send immediately (first update or after quiet period)
|
|
if err := pu.SendUpdate(ctx, update); err != nil {
|
|
log.WithContext(ctx).Debugf("error while sending an update to peer %s: %v", pu.peerKey.String(), err)
|
|
return err
|
|
}
|
|
}
|
|
|
|
// Timer expired - quiet period reached, send pending updates if any
|
|
case <-pu.debouncer.TimerChannel():
|
|
pendingUpdates := pu.debouncer.GetPendingUpdates()
|
|
if len(pendingUpdates) == 0 {
|
|
continue
|
|
}
|
|
log.WithContext(ctx).Debugf("sending %d debounced update(s) for peer %s", len(pendingUpdates), pu.peerKey.String())
|
|
for _, pendingUpdate := range pendingUpdates {
|
|
if err := pu.SendUpdate(ctx, pendingUpdate); err != nil {
|
|
log.WithContext(ctx).Debugf("error while sending an update to peer %s: %v", pu.peerKey.String(), err)
|
|
return err
|
|
}
|
|
}
|
|
|
|
// condition when client <-> server connection has been terminated
|
|
case <-pu.srv.Context().Done():
|
|
// happens when connection drops, e.g. client disconnects
|
|
log.WithContext(ctx).Debugf("stream of peer %s has been closed", pu.peerKey.String())
|
|
pu.cleanupFunc()
|
|
return pu.srv.Context().Err()
|
|
}
|
|
}
|
|
}
|
|
|
|
func (pu *PeerUpdateHandler) SendUpdate(ctx context.Context, update *network_map.UpdateMessage) error {
|
|
key, err := pu.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
pu.cleanupFunc()
|
|
return status.Errorf(codes.Internal, "failed processing update message")
|
|
}
|
|
|
|
stampCertificateChallenges(update.Update.GetChecks(), pu.challenger, pu.peerKey)
|
|
encryptedResp, err := pu.encrypter.EncryptMessage(pu.peerKey, key, update.Update)
|
|
if err != nil {
|
|
pu.cleanupFunc()
|
|
return status.Errorf(codes.Internal, "failed processing update message")
|
|
}
|
|
err = pu.srv.Send(&proto.EncryptedMessage{
|
|
WgPubKey: key.PublicKey().String(),
|
|
Body: encryptedResp,
|
|
})
|
|
if err != nil {
|
|
pu.cleanupFunc()
|
|
return status.Errorf(codes.Internal, "failed sending update message")
|
|
}
|
|
log.WithContext(ctx).Tracef("sent an update to peer %s", pu.peerKey.String())
|
|
return nil
|
|
}
|