mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-08 22:49:10 +02:00
[client] pqkem: drain the data-path receive loop before clearing state
Transport.Run spawned an untracked goroutine, and Close only closed the socket without waiting for it. An in-flight onInbound could therefore derive and store a PSK after the manager cleared the peer, or call SetPresharedKey while the WireGuard interface was being torn down. Track the receive goroutine and make Close drain it (close the socket to unblock the read, then wait for the callback to return). Manager.Stop now closes and drains the transport before clearing the exchange and PSK maps, and the engine tears WireGuard down only after Stop returns, so no late callback repopulates cleared state or touches a closed device. Found in cubic review on #7098 (client/internal/pqkem_transport.go:61).
This commit is contained in:
@@ -57,7 +57,8 @@ type Transport interface {
|
||||
// Run starts delivering inbound datagrams as (source endpoint, msg) to onInbound
|
||||
// and returns immediately; it runs until Close.
|
||||
Run(onInbound func(src netip.AddrPort, msg []byte))
|
||||
// Close stops delivery and releases the socket.
|
||||
// Close stops delivery, drains any in-flight onInbound callback, and releases the
|
||||
// socket, so no callback is still running when Close returns.
|
||||
Close() error
|
||||
}
|
||||
|
||||
@@ -281,19 +282,25 @@ func (m *Manager) Stop() {
|
||||
// has cancelled here no new Add can race Wait.
|
||||
m.mu.Lock()
|
||||
m.rootCancel()
|
||||
m.mu.Unlock()
|
||||
m.wait.Wait()
|
||||
m.mu.Lock()
|
||||
t := m.transport
|
||||
m.transport = nil
|
||||
m.exchanges = make(map[RemoteID]*exchangeCtl)
|
||||
m.psks = make(map[RemoteID]PSK)
|
||||
m.mu.Unlock()
|
||||
|
||||
// Close and drain the data-path receive loop BEFORE clearing peer state: Close waits
|
||||
// for any in-flight onInbound to return, so no inbound message can derive a PSK into a
|
||||
// cleared map, and the engine only tears the WireGuard interface down after this Stop
|
||||
// returns, so no late SetPresharedKey hits a torn-down device.
|
||||
if t != nil {
|
||||
if err := t.Close(); err != nil {
|
||||
m.logger.Warn("pqkem: closing data-path transport", "err", err)
|
||||
}
|
||||
}
|
||||
m.wait.Wait()
|
||||
|
||||
m.mu.Lock()
|
||||
m.exchanges = make(map[RemoteID]*exchangeCtl)
|
||||
m.psks = make(map[RemoteID]PSK)
|
||||
m.mu.Unlock()
|
||||
}
|
||||
|
||||
// ---- Signalling channel (host-driven; rides the host's negotiation) ----
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"fmt"
|
||||
"net"
|
||||
"net/netip"
|
||||
"sync"
|
||||
|
||||
log "github.com/sirupsen/logrus"
|
||||
)
|
||||
@@ -20,6 +21,7 @@ const DefaultPort = 51833
|
||||
type pqTransport struct {
|
||||
conn *net.UDPConn
|
||||
port int
|
||||
wg sync.WaitGroup
|
||||
}
|
||||
|
||||
// newPQTransport binds a UDP socket on the WG overlay IP, preferring DefaultPort and
|
||||
@@ -58,7 +60,9 @@ func (t *pqTransport) LocalPort() int { return t.port }
|
||||
// Run implements pqkem.Transport: the receive loop, delivering each datagram as
|
||||
// (source endpoint, msg). Exits when the socket is closed.
|
||||
func (t *pqTransport) Run(onInbound func(src netip.AddrPort, msg []byte)) {
|
||||
t.wg.Add(1)
|
||||
go func() {
|
||||
defer t.wg.Done()
|
||||
buf := make([]byte, 2048)
|
||||
for {
|
||||
n, src, err := t.conn.ReadFromUDPAddrPort(buf)
|
||||
@@ -72,5 +76,10 @@ func (t *pqTransport) Run(onInbound func(src netip.AddrPort, msg []byte)) {
|
||||
}()
|
||||
}
|
||||
|
||||
// Close implements pqkem.Transport.
|
||||
func (t *pqTransport) Close() error { return t.conn.Close() }
|
||||
// Close implements pqkem.Transport. It closes the socket (unblocking the read) and drains
|
||||
// the receive loop, so no onInbound callback is still running when Close returns.
|
||||
func (t *pqTransport) Close() error {
|
||||
err := t.conn.Close()
|
||||
t.wg.Wait()
|
||||
return err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user