mirror of
https://github.com/slackhq/nebula.git
synced 2026-10-04 14:46:39 +02:00
Under multiport each port had exactly one socket, so each port had exactly one core. That hurt worst on the base port, which carries every handshake, every lighthouse and punch packet, every peer without multiport, and every tunnel whose lanes are down or firewalled. Bind multiport.ports consecutive ports and give each one a full group of `routines` sockets sharing it through SO_REUSEPORT, so the kernel hashes each arriving 4-tuple onto one of the group. `routines` is now per port and the total worker count is routines * multiport.ports, computed for you: with `routines: 8` and `multiport.ports: 4` you get 8 routines per port, 32 total. A routine still owns exactly one socket, which is what lets the read path own its state without locking. Sockets are laid out port-major, writers[s*routinesPerPort+r] being the r'th socket on port listen.port+s, so socket selection is a pure function of (queue, lane) with no borrowing: laneSock(q,s) = s*R + q%R, and base traffic takes laneSock(q,0). Sibling routines land on different sockets of the same port, so a lane's port is served by its whole group. multiport.ports is a new key and must be > 1; there is no default, since under these semantics defaulting it would silently multiply the worker count. The multiport decision moves above tun creation because tun queues are sized by the total. Also drop the base-port fallback for unanswered lane probes. A lane now aims only at its own paired peer port: sharing the base port's destination buys the lane a source port of its own but costs the peer the receive spread lanes exist to create, so a lane that cannot reach its port stays down and its flows ride the base tunnel.
1036 lines
36 KiB
Go
1036 lines
36 KiB
Go
package nebula
|
|
|
|
import (
|
|
"log/slog"
|
|
"net/netip"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/gaissmai/bart"
|
|
"github.com/slackhq/nebula/cert"
|
|
"github.com/slackhq/nebula/config"
|
|
"github.com/slackhq/nebula/firewall"
|
|
"github.com/slackhq/nebula/handshake"
|
|
"github.com/slackhq/nebula/header"
|
|
"github.com/slackhq/nebula/overlay/overlaytest"
|
|
"github.com/slackhq/nebula/overlay/tio"
|
|
"github.com/slackhq/nebula/test"
|
|
"github.com/slackhq/nebula/udp"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
)
|
|
|
|
var (
|
|
testMyAddr = netip.MustParseAddr("10.0.0.1")
|
|
testPeerAddr = netip.MustParseAddr("10.0.0.2")
|
|
)
|
|
|
|
// newTestLaneSet builds a lane set from a real handshake result so the sessions
|
|
// it derives hold usable keys.
|
|
func newTestLaneSet(t *testing.T, r *handshake.Result, myLanes int, peerPorts, peerBase, peerTxLanes uint32) *laneSet {
|
|
t.Helper()
|
|
r.PeerPortCount = peerPorts
|
|
r.PeerBasePort = peerBase
|
|
r.PeerTxLanes = peerTxLanes
|
|
return newLaneSet(r, myLanes, testMyAddr, testPeerAddr)
|
|
}
|
|
|
|
// laneSessionFor derives lane s's session and installs it, standing in for the
|
|
// data-plane call that would normally be the first to need it.
|
|
func laneSessionFor(t *testing.T, ls *laneSet, s int) *ConnectionState {
|
|
t.Helper()
|
|
cs, err := ls.session(s)
|
|
require.NoError(t, err)
|
|
require.NotNil(t, cs)
|
|
return cs
|
|
}
|
|
|
|
func newTestLaneHostInfo(t *testing.T, r *handshake.Result, ls *laneSet) *HostInfo {
|
|
t.Helper()
|
|
cs, err := newConnectionStateFromResult(r)
|
|
require.NoError(t, err)
|
|
hi := &HostInfo{
|
|
vpnAddrs: []netip.Addr{testPeerAddr},
|
|
localIndexId: 100,
|
|
remoteIndexId: 200,
|
|
remotes: NewRemoteList([]netip.Addr{testPeerAddr}, nil),
|
|
HandshakePacket: map[uint8][]byte{},
|
|
ConnectionState: cs,
|
|
lanes: ls,
|
|
}
|
|
hi.SetRemote(netip.MustParseAddrPort("192.0.2.1:4242"))
|
|
return hi
|
|
}
|
|
|
|
func TestLanePortOffset(t *testing.T) {
|
|
a := netip.MustParseAddr("10.0.0.1")
|
|
b := netip.MustParseAddr("10.0.0.2")
|
|
|
|
// Deterministic and in range.
|
|
for _, count := range []uint16{1, 2, 3, 4, 16, 256} {
|
|
o := lanePortOffset(a, b, count)
|
|
assert.Equal(t, o, lanePortOffset(a, b, count), "count %d not deterministic", count)
|
|
assert.Less(t, o, count, "count %d out of range", count)
|
|
}
|
|
assert.Equal(t, uint16(0), lanePortOffset(a, b, 0), "zero port count")
|
|
|
|
// The two sides' rotations cancel when port counts match, preserving the
|
|
// lane-i-reverses-lane-j conntrack pairing.
|
|
for _, count := range []uint16{2, 3, 4, 7, 16} {
|
|
for i := range 32 {
|
|
peer := netip.AddrFrom4([4]byte{192, 0, 2, byte(i)})
|
|
oA := lanePortOffset(a, peer, count)
|
|
oB := lanePortOffset(peer, a, count)
|
|
assert.Equal(t, uint16(0), (oA+oB)%count,
|
|
"offsets don't cancel for peer %s count %d", peer, count)
|
|
}
|
|
}
|
|
|
|
// Distinct small peers land on distinct rotations of a big peer's range,
|
|
// not all on the same first ports.
|
|
const bigPeerPorts = 16
|
|
distinct := map[uint16]struct{}{}
|
|
for i := range 64 {
|
|
client := netip.AddrFrom4([4]byte{192, 0, 2, byte(i)})
|
|
distinct[lanePortOffset(client, a, bigPeerPorts)] = struct{}{}
|
|
}
|
|
assert.GreaterOrEqual(t, len(distinct), 8, "64 clients only produced %d distinct offsets", len(distinct))
|
|
}
|
|
|
|
func TestLaneTargetPort(t *testing.T) {
|
|
ls := &laneSet{peerBasePort: 4242, peerPortCount: 4}
|
|
|
|
// No rotation: lane i targets base+i, wrapping past the peer's range.
|
|
for i, want := range map[int]uint16{1: 4243, 2: 4244, 3: 4245, 5: 4243} {
|
|
assert.Equal(t, want, ls.laneTargetPortLocked(i), "lane %d", i)
|
|
}
|
|
|
|
// Rotation shifts the whole mapping; the wrapped lane lands on the base
|
|
// port itself, which is a valid distinct 4-tuple (our source port differs).
|
|
ls.portOffset = 3
|
|
for i, want := range map[int]uint16{1: 4242, 2: 4243, 3: 4244} {
|
|
assert.Equal(t, want, ls.laneTargetPortLocked(i), "rotated lane %d", i)
|
|
}
|
|
|
|
// Fewer peer ports than local lanes: rotation still spreads across all of
|
|
// the peer's ports.
|
|
ls = &laneSet{peerBasePort: 4242, peerPortCount: 2, portOffset: 1}
|
|
assert.Equal(t, uint16(4242), ls.laneTargetPortLocked(1))
|
|
assert.Equal(t, uint16(4243), ls.laneTargetPortLocked(2))
|
|
}
|
|
|
|
// A lane's keys are derived, not negotiated, so the whole design rests on the
|
|
// two sides landing on the same pair without exchanging anything.
|
|
func TestLaneKeyDerivationSymmetry(t *testing.T) {
|
|
initR, respR := runTestHandshake(t)
|
|
|
|
initLS := newTestLaneSet(t, initR, 4, 4, 4242, 4)
|
|
respLS := newTestLaneSet(t, respR, 4, 4, 4242, 4)
|
|
require.Len(t, initLS.sessions, 4)
|
|
assert.Nil(t, initLS.sessions[0].Load(), "lane 0 is the base session, not a derived one")
|
|
|
|
nb := make([]byte, 12)
|
|
for s := 1; s < 4; s++ {
|
|
out := header.EncodeLane(make([]byte, 0, mtu), header.Version, header.Message, 0, 200, 1, uint8(s))
|
|
ct, err := laneSessionFor(t, initLS, s).eKey.EncryptDanger(out, out, []byte("lane payload"), 1, nb)
|
|
require.NoError(t, err)
|
|
|
|
pt, err := laneSessionFor(t, respLS, s).Decrypt(test.NewLogger(), 1, ct, nb)
|
|
require.NoError(t, err, "lane %d keys did not match", s)
|
|
assert.Equal(t, []byte("lane payload"), pt)
|
|
}
|
|
|
|
// Distinct lanes get distinct keys: lane 2's session must not open lane 1's
|
|
// ciphertext, or the header's lane index would be forgeable in effect.
|
|
out := header.EncodeLane(make([]byte, 0, mtu), header.Version, header.Message, 0, 200, 7, 1)
|
|
ct, err := laneSessionFor(t, initLS, 1).eKey.EncryptDanger(out, out, []byte("lane payload"), 7, nb)
|
|
require.NoError(t, err)
|
|
_, err = laneSessionFor(t, respLS, 2).Decrypt(test.NewLogger(), 7, ct, nb)
|
|
assert.Error(t, err)
|
|
}
|
|
|
|
func TestNewLaneSetSizing(t *testing.T) {
|
|
initR, _ := runTestHandshake(t)
|
|
|
|
// A peer with no multiport advert gets no lanes at all.
|
|
assert.Nil(t, newLaneSet(&handshake.Result{}, 4, testMyAddr, testPeerAddr))
|
|
|
|
// One lane means only the base tunnel, which is not a lane set.
|
|
assert.Nil(t, newLaneSet(&handshake.Result{PeerPortCount: 4, PeerTxLanes: 1}, 1, testMyAddr, testPeerAddr))
|
|
|
|
// Sessions cover both directions: enough for everything the peer may send,
|
|
// even though we may only send on a few.
|
|
ls := newTestLaneSet(t, initR, 2, 8, 4242, 6)
|
|
assert.Len(t, ls.sessions, 6, "sessions must cover the peer's tx lanes")
|
|
assert.Equal(t, 2, ls.txLanes, "we may only send on our own lanes")
|
|
|
|
// The peer's advert sizes the session table but nothing else. Its lane count
|
|
// is its own choice, so it must not be able to make us allocate per-lane tx
|
|
// state we will never use.
|
|
assert.Len(t, ls.txAddr, 2)
|
|
assert.Len(t, ls.demand, 2)
|
|
assert.Len(t, ls.probe, 2)
|
|
|
|
// Nothing is derived up front, for the same reason.
|
|
for s := range ls.sessions {
|
|
assert.Nil(t, ls.sessions[s].Load(), "lane %d derived before anything needed it", s)
|
|
}
|
|
|
|
// Every lane starts demanded so the first traffic tick probes them all: a
|
|
// lane that only comes up after its flows do is a lane the flows never move
|
|
// onto, because the kernel has pinned them to the queue they started on.
|
|
ls = newTestLaneSet(t, initR, 4, 4, 4242, 4)
|
|
for s := 1; s < ls.txLanes; s++ {
|
|
assert.True(t, ls.demand[s].Load(), "lane %d not demanded at creation", s)
|
|
}
|
|
|
|
// Our tx lanes are clamped to the ports the peer actually bound: a lane
|
|
// aimed past the peer's range would land on some unrelated socket.
|
|
ls = newTestLaneSet(t, initR, 8, 3, 4242, 8)
|
|
assert.Equal(t, 3, ls.txLanes)
|
|
assert.Len(t, ls.sessions, 8)
|
|
}
|
|
|
|
// The handshake log lines carry the negotiated lanes, so an operator can tell a
|
|
// peer that got none from one that was never asked.
|
|
func TestLaneLogAttr(t *testing.T) {
|
|
initR, _ := runTestHandshake(t)
|
|
|
|
// Not running multiport: an empty attr, which slog drops entirely.
|
|
assert.Equal(t, slog.Attr{}, laneLogAttr(0, nil))
|
|
|
|
// Running multiport against a peer that isn't: zeros, not silence.
|
|
assert.Equal(t, m{"tx": 0, "sessions": 0}, laneLogAttr(4, nil).Value.Any())
|
|
|
|
ls := newTestLaneSet(t, initR, 2, 8, 4242, 6)
|
|
attr := laneLogAttr(2, ls)
|
|
assert.Equal(t, "lanes", attr.Key)
|
|
assert.Equal(t, m{
|
|
"tx": 2,
|
|
"sessions": 6,
|
|
"peerBasePort": uint16(4242),
|
|
"peerPorts": uint16(8),
|
|
"portOffset": ls.portOffset,
|
|
}, attr.Value.Any())
|
|
}
|
|
|
|
// The RX path must not cache a session for a lane until a packet on it has
|
|
// actually decrypted, or a spoofer naming lanes at random could make us hold a
|
|
// replay window and two cipher states per lane without authenticating anything.
|
|
func TestLaneSessionRxDerivation(t *testing.T) {
|
|
initR, respR := runTestHandshake(t)
|
|
respLS := newTestLaneSet(t, respR, 4, 4, 4242, 4)
|
|
hi := newTestLaneHostInfo(t, respR, respLS)
|
|
|
|
ci, cached, err := hi.laneSession(2)
|
|
require.NoError(t, err)
|
|
require.NotNil(t, ci)
|
|
assert.False(t, cached, "the first packet on a lane derives, it does not hit")
|
|
assert.Nil(t, respLS.sessions[2].Load(), "an unauthenticated packet must not install a session")
|
|
|
|
// The lane the peer really is using decrypts, and that is what installs it.
|
|
nb := make([]byte, 12)
|
|
out := header.EncodeLane(make([]byte, 0, mtu), header.Version, header.Message, 0, 200, 1, 2)
|
|
ct, err := laneSessionFor(t, newTestLaneSet(t, initR, 4, 4, 4242, 4), 2).
|
|
eKey.EncryptDanger(out, out, []byte("real lane traffic"), 1, nb)
|
|
require.NoError(t, err)
|
|
pt, err := ci.Decrypt(test.NewLogger(), 1, ct, nb)
|
|
require.NoError(t, err)
|
|
assert.Equal(t, []byte("real lane traffic"), pt)
|
|
|
|
respLS.installSession(test.NewLogger(), 2, ci, 1)
|
|
assert.Same(t, ci, respLS.sessions[2].Load())
|
|
|
|
// Now it is a hit, and the replay window the decrypt above advanced is the
|
|
// one the next packet sees.
|
|
got, cached, err := hi.laneSession(2)
|
|
require.NoError(t, err)
|
|
assert.True(t, cached)
|
|
assert.Same(t, ci, got)
|
|
|
|
// A racing install loses rather than swapping the session out, which would
|
|
// throw away the replay window the live one has been accumulating. The loser's
|
|
// packet still has to be marked seen on the winner, or dropping its session
|
|
// would make that one counter replayable.
|
|
other, err := newLaneConnectionState(&respLS.material, 2)
|
|
require.NoError(t, err)
|
|
require.True(t, ci.window.Check(test.NewLogger(), 7), "counter 7 seen before the race")
|
|
respLS.installSession(test.NewLogger(), 2, other, 7)
|
|
assert.Same(t, ci, respLS.sessions[2].Load())
|
|
assert.False(t, ci.window.Check(test.NewLogger(), 7),
|
|
"the loser's counter was not carried to the surviving window")
|
|
}
|
|
|
|
func TestLaneTxGate(t *testing.T) {
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 4, 4, 4242, 4)
|
|
for s := range ls.demand {
|
|
// Creation seeds demand on every lane; start from cold so this test can
|
|
// tell what the TX path itself raises.
|
|
ls.demand[s].Store(false)
|
|
}
|
|
|
|
// A down lane hands back nothing and raises demand, which is what gets it
|
|
// re-probed after a demotion.
|
|
ci, addr := ls.txLane(1)
|
|
assert.Nil(t, ci)
|
|
assert.False(t, addr.IsValid())
|
|
assert.True(t, ls.demand[1].Load())
|
|
assert.False(t, ls.demand[2].Load(), "demand raised on an untouched lane")
|
|
|
|
// An up lane with no session derived yet cannot happen — the probe that
|
|
// promoted it derived one — but it must fall back rather than send in the
|
|
// clear if it ever does.
|
|
want := netip.MustParseAddrPort("192.0.2.1:4243")
|
|
ls.txAddr[1].Store(&want)
|
|
ci, _ = ls.txLane(1)
|
|
assert.Nil(t, ci)
|
|
|
|
// Promotion publishes the session and destination together.
|
|
sess := laneSessionFor(t, ls, 1)
|
|
ls.demand[1].Store(false)
|
|
ci, addr = ls.txLane(1)
|
|
assert.Same(t, sess, ci)
|
|
assert.Equal(t, want, addr)
|
|
assert.False(t, ls.demand[1].Load(), "a hit must not raise demand")
|
|
|
|
// Lane 0 is the base tunnel and lanes at or above txLanes are receive-only.
|
|
ci, _ = ls.txLane(0)
|
|
assert.Nil(t, ci)
|
|
ci, _ = ls.txLane(4)
|
|
assert.Nil(t, ci)
|
|
}
|
|
|
|
func TestLaneSessionLookup(t *testing.T) {
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 3, 4, 4242, 3)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
|
|
assertNoLane := func(s uint8, msg string) {
|
|
t.Helper()
|
|
ci, cached, err := hi.laneSession(s)
|
|
require.NoError(t, err)
|
|
assert.Nil(t, ci, msg)
|
|
assert.False(t, cached, msg)
|
|
}
|
|
|
|
assertNoLane(0, "lane 0 is the base session")
|
|
assertNoLane(3, "a lane beyond what this tunnel covers")
|
|
assertNoLane(255, "a lane beyond what this tunnel covers")
|
|
|
|
// An already-derived lane is returned as a hit.
|
|
sess := laneSessionFor(t, ls, 2)
|
|
ci, cached, err := hi.laneSession(2)
|
|
require.NoError(t, err)
|
|
assert.Same(t, sess, ci)
|
|
assert.True(t, cached)
|
|
|
|
// A peer without lanes answers nil for every lane rather than panicking.
|
|
bare := &HostInfo{}
|
|
ci, _, err = bare.laneSession(1)
|
|
require.NoError(t, err)
|
|
assert.Nil(t, ci)
|
|
}
|
|
|
|
func TestMaxMessageCounter(t *testing.T) {
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 3, 4, 4242, 3)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
|
|
hi.ConnectionState.messageCounter.Store(5)
|
|
assert.Equal(t, uint64(5), hi.maxMessageCounter())
|
|
|
|
// A lane past the base is what the rehandshake threshold has to notice: the
|
|
// base counter would sit still while the lane burns through its nonces.
|
|
laneSessionFor(t, ls, 2).messageCounter.Store(9000)
|
|
assert.Equal(t, uint64(9000), hi.maxMessageCounter())
|
|
|
|
assert.Equal(t, uint64(0), (&HostInfo{}).maxMessageCounter())
|
|
}
|
|
|
|
func TestLaneRetryDelay(t *testing.T) {
|
|
assert.Equal(t, laneRetryBase, laneRetryDelay(0))
|
|
assert.Equal(t, 2*laneRetryBase, laneRetryDelay(1))
|
|
assert.Equal(t, laneRetryMax, laneRetryDelay(laneMaxFails), "backoff must saturate")
|
|
}
|
|
|
|
func newLaneTestInterface(hostMap *HostMap) *Interface {
|
|
l := test.NewLogger()
|
|
lh := newTestLighthouse()
|
|
cs := &CertState{
|
|
initiatingVersion: cert.Version1,
|
|
privateKey: []byte{},
|
|
v1Cert: &dummyCert{version: cert.Version1},
|
|
v1Credential: nil,
|
|
}
|
|
ifce := &Interface{
|
|
hostMap: hostMap,
|
|
inside: &overlaytest.NoopTun{},
|
|
outside: &udp.NoopConn{},
|
|
firewall: &Firewall{},
|
|
lightHouse: lh,
|
|
pki: &PKI{},
|
|
handshakeManager: NewHandshakeManager(l, hostMap, lh, &udp.NoopConn{}, defaultHandshakeConfig),
|
|
myVpnNetworksTable: new(bart.Lite),
|
|
messageMetrics: newMessageMetricsOnlyRecvError(),
|
|
writers: []udp.Conn{&udp.NoopConn{}, &udp.NoopConn{}, &udp.NoopConn{}, &udp.NoopConn{}},
|
|
// One socket per port keeps writers indexed by lane, which is what most of
|
|
// these tests want; the tests that care about the group layout set their own.
|
|
routinesPerPort: 1,
|
|
l: l,
|
|
}
|
|
ifce.pki.cs.Store(cs)
|
|
|
|
conf := config.NewC(l)
|
|
punchy := NewPunchyFromConfig(l, conf, nil)
|
|
cm := newConnectionManagerFromConfig(l, conf, hostMap, punchy)
|
|
cm.intf = ifce
|
|
ifce.connectionManager = cm
|
|
ifce.handshakeManager.f = ifce
|
|
return ifce
|
|
}
|
|
|
|
// The full TX lifecycle of a lane: demand -> probe -> ack -> up, then an
|
|
// unanswered keepalive -> demoted.
|
|
func TestLaneProbeLifecycle(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
|
|
nb := make([]byte, 12)
|
|
out := make([]byte, mtu)
|
|
now := time.Now()
|
|
|
|
// No demand: nothing is probed, so a peer we barely talk to costs nothing
|
|
// beyond its base tunnel. (Creation seeds demand on every lane, so clear it
|
|
// to get at the no-demand case.)
|
|
for s := range ls.demand {
|
|
ls.demand[s].Store(false)
|
|
}
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
for s := 1; s < 4; s++ {
|
|
assert.True(t, ls.probe[s].sentAt.IsZero(), "lane %d probed without demand", s)
|
|
}
|
|
ls.mu.Unlock()
|
|
|
|
// Demand on lane 1 alone probes lane 1 alone, aimed at the peer's port for
|
|
// that lane.
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
require.False(t, ls.probe[1].sentAt.IsZero(), "demand did not produce a probe")
|
|
assert.True(t, ls.probe[2].sentAt.IsZero())
|
|
wantPort := ls.laneTargetPortLocked(1)
|
|
assert.Equal(t, netip.AddrPortFrom(netip.MustParseAddr("192.0.2.1"), wantPort), ls.probe[1].target)
|
|
gen := ls.probe[1].gen
|
|
assert.False(t, ls.demand[1].Load(), "the probe did not consume the demand")
|
|
ls.mu.Unlock()
|
|
|
|
// The lane stays down until the ack lands, and a stale generation cannot
|
|
// bring it up.
|
|
assert.Nil(t, ls.txAddr[1].Load())
|
|
_, promoted := ls.noteAck(1, gen+1, now)
|
|
assert.False(t, promoted, "an ack for a superseded probe promoted the lane")
|
|
assert.Nil(t, ls.txAddr[1].Load())
|
|
|
|
// The matching ack promotes it, and the destination is the probed target.
|
|
ackTarget, promoted := ls.noteAck(1, gen, now)
|
|
assert.True(t, promoted)
|
|
assert.Equal(t, netip.AddrPortFrom(netip.MustParseAddr("192.0.2.1"), wantPort), ackTarget)
|
|
addr := ls.txAddr[1].Load()
|
|
require.NotNil(t, addr)
|
|
assert.Equal(t, netip.AddrPortFrom(netip.MustParseAddr("192.0.2.1"), wantPort), *addr)
|
|
|
|
// A second ack is a keepalive, not a promotion.
|
|
ls.mu.Lock()
|
|
ls.probe[1].sentAt = now
|
|
ls.mu.Unlock()
|
|
_, promoted = ls.noteAck(1, gen, now)
|
|
assert.False(t, promoted)
|
|
assert.NotNil(t, ls.txAddr[1].Load())
|
|
|
|
// An up lane is left alone until the keepalive comes due.
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
assert.True(t, ls.probe[1].sentAt.IsZero(), "an up lane was re-probed early")
|
|
ls.mu.Unlock()
|
|
|
|
// Past the keepalive it re-proves its path...
|
|
now = now.Add(laneKeepalive + time.Second)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
require.False(t, ls.probe[1].sentAt.IsZero(), "keepalive did not probe")
|
|
ls.mu.Unlock()
|
|
assert.NotNil(t, ls.txAddr[1].Load(), "lane demoted before its probe aged out")
|
|
|
|
// ...and an unanswered keepalive demotes it with backoff, so the routine
|
|
// falls back to the base tunnel.
|
|
now = now.Add(laneProbeTimeout + time.Second)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
assert.Nil(t, ls.txAddr[1].Load(), "unanswered keepalive did not demote the lane")
|
|
ls.mu.Lock()
|
|
assert.Equal(t, uint8(1), ls.probe[1].fails)
|
|
assert.True(t, ls.probe[1].retryAt.After(now))
|
|
ls.mu.Unlock()
|
|
|
|
// The backoff holds even with fresh demand.
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
assert.True(t, ls.probe[1].sentAt.IsZero(), "backoff was ignored")
|
|
ls.mu.Unlock()
|
|
}
|
|
|
|
// A lane aims at its own peer port and nothing else. There is no fallback to the
|
|
// peer's base port: sharing that destination would cost the peer the receive
|
|
// spread lanes exist to create, so an unreachable lane port just means this lane
|
|
// stays down and its flows ride the base tunnel.
|
|
func TestLaneProbeAlwaysTargetsItsOwnPort(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
|
|
nb := make([]byte, 12)
|
|
out := make([]byte, mtu)
|
|
now := time.Now()
|
|
|
|
ls.mu.Lock()
|
|
lanePort := ls.laneTargetPortLocked(1)
|
|
ls.mu.Unlock()
|
|
require.NotEqual(t, ls.peerBasePort, lanePort, "this test needs a lane that isn't aimed at the base port")
|
|
|
|
// Repeated failures never move the target off the lane's own port.
|
|
for i := 0; i < 4; i++ {
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
require.False(t, ls.probe[1].sentAt.IsZero(), "the lane was not probed")
|
|
assert.Equal(t, lanePort, ls.probe[1].target.Port(), "probe %d left the lane port", i)
|
|
gen := ls.probe[1].gen
|
|
ls.mu.Unlock()
|
|
|
|
// Age the probe out, then wait out the backoff for the next attempt.
|
|
now = now.Add(laneProbeTimeout + time.Second)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
require.Nil(t, ls.txAddr[1].Load())
|
|
assert.False(t, promotedByAck(ls, 1, gen, now), "an aged-out probe was still promotable")
|
|
now = now.Add(laneRetryMax + time.Second)
|
|
}
|
|
|
|
// It comes up on that port, and the keepalive re-proves the same one.
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
gen := ls.probe[1].gen
|
|
ls.mu.Unlock()
|
|
target, promoted := ls.noteAck(1, gen, now)
|
|
require.True(t, promoted)
|
|
assert.Equal(t, lanePort, target.Port())
|
|
|
|
now = now.Add(laneKeepalive + time.Second)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
require.False(t, ls.probe[1].sentAt.IsZero(), "the keepalive did not probe")
|
|
assert.Equal(t, lanePort, ls.probe[1].target.Port(), "the keepalive left the lane port")
|
|
ls.mu.Unlock()
|
|
|
|
// And a demotion returns to probing that same port.
|
|
now = now.Add(laneProbeTimeout + time.Second)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
require.Nil(t, ls.txAddr[1].Load(), "the unanswered keepalive did not demote the lane")
|
|
now = now.Add(laneRetryMax + time.Second)
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
require.False(t, ls.probe[1].sentAt.IsZero(), "the lane was not retried")
|
|
assert.Equal(t, lanePort, ls.probe[1].target.Port())
|
|
ls.mu.Unlock()
|
|
}
|
|
|
|
// promotedByAck is noteAck's boolean alone, for assertions that only care whether
|
|
// an ack could bring the lane up.
|
|
func promotedByAck(ls *laneSet, s int, gen uint8, now time.Time) bool {
|
|
_, promoted := ls.noteAck(s, gen, now)
|
|
return promoted
|
|
}
|
|
|
|
// A roam is a new path with no derivable relationship to the old lane ports, so
|
|
// every lane has to be rebuilt rather than moved.
|
|
func TestLaneProbeRoamResets(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
|
|
now := time.Now()
|
|
nb := make([]byte, 12)
|
|
out := make([]byte, mtu)
|
|
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
gen := ls.probe[1].gen
|
|
ls.mu.Unlock()
|
|
_, promoted := ls.noteAck(1, gen, now)
|
|
require.True(t, promoted)
|
|
require.NotNil(t, ls.txAddr[1].Load())
|
|
|
|
// New remote address: the lane is taken down, not retargeted.
|
|
hi.SetRemote(netip.MustParseAddrPort("198.51.100.7:4242"))
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
assert.Nil(t, ls.txAddr[1].Load(), "lane survived a roam")
|
|
ls.mu.Lock()
|
|
assert.Zero(t, ls.probe[1].fails, "a roam is not a lane failure")
|
|
assert.Equal(t, netip.MustParseAddr("198.51.100.7"), ls.peerAddr)
|
|
ls.mu.Unlock()
|
|
|
|
// Relayed (no direct remote) means no lanes at all.
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
ls.mu.Lock()
|
|
gen = ls.probe[1].gen
|
|
ls.mu.Unlock()
|
|
_, promoted = ls.noteAck(1, gen, now)
|
|
require.True(t, promoted)
|
|
hi.SetRemote(netip.AddrPort{})
|
|
ifce.probeLanes(hi, now, nb, out)
|
|
assert.Nil(t, ls.txAddr[1].Load(), "lane survived losing the direct path")
|
|
}
|
|
|
|
// A lane probe is answered on the base tunnel, echoing the header's lane, so a
|
|
// peer cannot get us to vouch for a lane it never probed.
|
|
func TestHandleLaneProbe(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
|
|
initR, respR := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, respR, 4, 4, 4242, 4)
|
|
hi := newTestLaneHostInfo(t, respR, ls)
|
|
|
|
rxc := &rxContext{nb: make([]byte, 12), scratch: make([]byte, mtu)}
|
|
sent := &recordingUdpConn{}
|
|
ifce.writers = []udp.Conn{sent, &udp.NoopConn{}, &udp.NoopConn{}, &udp.NoopConn{}}
|
|
|
|
// The payload's lane is ignored in favour of the header's.
|
|
ifce.handleLaneProbe(hi, 2, []byte{3, 42}, rxc)
|
|
require.Len(t, sent.bufs, 1, "the ack must ride the base tunnel's socket")
|
|
|
|
h := &header.H{}
|
|
require.NoError(t, h.Parse(sent.bufs[0]))
|
|
assert.Equal(t, header.Test, h.Type)
|
|
assert.Equal(t, header.LaneProbeAck, h.Subtype)
|
|
assert.Equal(t, uint8(0), h.Lane(), "the ack is base-tunnel traffic")
|
|
|
|
pt, err := laneSessionFor(t, newTestLaneSet(t, initR, 4, 4, 4242, 4), 1).dKey.DecryptDanger(
|
|
nil, sent.bufs[0][:header.Len], sent.bufs[0][header.Len:], h.MessageCounter, make([]byte, 12))
|
|
_ = pt
|
|
assert.Error(t, err, "the ack must not be readable with a lane key")
|
|
|
|
// A probe claiming lane 0 or with a truncated payload is answered with
|
|
// nothing at all.
|
|
sent.bufs = nil
|
|
ifce.handleLaneProbe(hi, 0, []byte{1, 2}, rxc)
|
|
ifce.handleLaneProbe(hi, 1, []byte{1}, rxc)
|
|
assert.Empty(t, sent.bufs)
|
|
}
|
|
|
|
func TestHandleLaneProbeAck(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
|
|
now := time.Now()
|
|
ls.demand[1].Store(true)
|
|
ifce.probeLanes(hi, now, make([]byte, 12), make([]byte, mtu))
|
|
ls.mu.Lock()
|
|
gen := ls.probe[1].gen
|
|
ls.mu.Unlock()
|
|
|
|
// A short or out-of-range ack is ignored, and a peer without lanes does not
|
|
// panic the handler.
|
|
ifce.handleLaneProbeAck(hi, []byte{1})
|
|
ifce.handleLaneProbeAck(hi, []byte{99, gen})
|
|
ifce.handleLaneProbeAck(&HostInfo{}, []byte{1, gen})
|
|
assert.Nil(t, ls.txAddr[1].Load())
|
|
|
|
ifce.handleLaneProbeAck(hi, []byte{1, gen})
|
|
assert.NotNil(t, ls.txAddr[1].Load())
|
|
}
|
|
|
|
// recordingUdpConn records what was written to it, one datagram per write.
|
|
type recordingUdpConn struct {
|
|
udp.NoopConn
|
|
bufs [][]byte
|
|
dsts []netip.AddrPort
|
|
}
|
|
|
|
func (c *recordingUdpConn) WriteTo(b []byte, addr netip.AddrPort) error {
|
|
c.bufs = append(c.bufs, append([]byte(nil), b...))
|
|
c.dsts = append(c.dsts, addr)
|
|
return nil
|
|
}
|
|
|
|
// recordingBatchConn is a udp.Conn that records the batches flushed to it, so a
|
|
// test can see which socket a packet left on.
|
|
type recordingBatchConn struct {
|
|
udp.NoopConn
|
|
bufs [][]byte
|
|
dsts []netip.AddrPort
|
|
}
|
|
|
|
func (c *recordingBatchConn) WriteBatch(bufs [][]byte, addrs []netip.AddrPort) (int, error) {
|
|
for i := range bufs {
|
|
c.bufs = append(c.bufs, append([]byte(nil), bufs[i]...))
|
|
c.dsts = append(c.dsts, addrs[i])
|
|
}
|
|
return len(bufs), nil
|
|
}
|
|
|
|
// laneOfFlow is what the TX path computes to pick a lane, without needing a
|
|
// promoted lane to observe it.
|
|
func laneOfFlow(ls *laneSet, p *firewall.Packet) int {
|
|
return int((laneFlowHash(p) + uint32(ls.laneBias)) % uint32(ls.txLanes))
|
|
}
|
|
|
|
func testFlow(port uint16) *firewall.Packet {
|
|
return &firewall.Packet{
|
|
LocalAddr: testMyAddr,
|
|
RemoteAddr: testPeerAddr,
|
|
LocalPort: port,
|
|
RemotePort: 443,
|
|
Protocol: 6,
|
|
}
|
|
}
|
|
|
|
// flowForLane finds a flow that hashes onto lane s, so a test can aim traffic at
|
|
// a specific lane.
|
|
func flowForLane(t *testing.T, ls *laneSet, s int) *firewall.Packet {
|
|
t.Helper()
|
|
for port := uint16(1024); port < 4096; port++ {
|
|
p := testFlow(port)
|
|
if laneOfFlow(ls, p) == s {
|
|
return p
|
|
}
|
|
}
|
|
t.Fatalf("no flow hashes onto lane %d", s)
|
|
return nil
|
|
}
|
|
|
|
func TestSendInsideMessageLaneSwap(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
ifce.multiport = true
|
|
ifce.laneCount = 4
|
|
|
|
conns := []*recordingBatchConn{{}, {}, {}, {}}
|
|
ifce.writers = []udp.Conn{conns[0], conns[1], conns[2], conns[3]}
|
|
|
|
initR, respR := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
peerLS := newTestLaneSet(t, respR, 4, 4, 4242, 4)
|
|
|
|
tx := ifce.newTxQueue(1)
|
|
pkt := tio.Packet{Bytes: []byte{0x45, 0, 0, 4, 1, 2, 3, 4}}
|
|
nb := make([]byte, 12)
|
|
flow1 := flowForLane(t, ls, 1)
|
|
for s := range ls.demand {
|
|
ls.demand[s].Store(false)
|
|
}
|
|
|
|
// A down lane rides the base tunnel and asks for a probe.
|
|
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx)
|
|
tx.flush(ifce)
|
|
require.Len(t, conns[0].bufs, 1)
|
|
assert.Empty(t, conns[1].bufs)
|
|
assert.True(t, ls.demand[1].Load(), "a miss did not raise demand")
|
|
assert.False(t, ls.demand[2].Load(), "demand raised on a lane nothing hashed onto")
|
|
|
|
h := &header.H{}
|
|
require.NoError(t, h.Parse(conns[0].bufs[0]))
|
|
assert.Equal(t, uint8(0), h.Lane())
|
|
|
|
// Once the lane is up, a flow hashing onto it rides the lane session, the
|
|
// lane socket and the lane's destination, tagged with the lane index.
|
|
// Promotion normally happens on the ack of a probe, which is also what
|
|
// derived the session, so stand both up here.
|
|
laneRemote := netip.MustParseAddrPort("192.0.2.1:5354")
|
|
laneSessionFor(t, ls, 1)
|
|
ls.txAddr[1].Store(&laneRemote)
|
|
ls.demand[1].Store(false)
|
|
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx)
|
|
tx.flush(ifce)
|
|
require.Len(t, conns[1].bufs, 1)
|
|
assert.Equal(t, laneRemote, conns[1].dsts[0])
|
|
assert.False(t, ls.demand[1].Load(), "a hit raised demand")
|
|
|
|
require.NoError(t, h.Parse(conns[1].bufs[0]))
|
|
assert.Equal(t, uint8(1), h.Lane())
|
|
assert.Equal(t, hi.remoteIndexId, h.RemoteIndex)
|
|
|
|
// And the peer's derived lane-1 session is what opens it.
|
|
pt, err := laneSessionFor(t, peerLS, 1).Decrypt(test.NewLogger(), h.MessageCounter, conns[1].bufs[0], nb)
|
|
require.NoError(t, err)
|
|
assert.Equal(t, pkt.Bytes, pt)
|
|
|
|
// Another routine sending the same flow rides the same lane socket and
|
|
// session: the lane follows the flow, not the routine.
|
|
tx2 := ifce.newTxQueue(2)
|
|
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx2)
|
|
tx2.flush(ifce)
|
|
require.Len(t, conns[1].bufs, 2)
|
|
require.NoError(t, h.Parse(conns[1].bufs[1]))
|
|
assert.Equal(t, uint8(1), h.Lane())
|
|
|
|
// A flow that hashes onto lane 0 stays on the base tunnel even with lane 1
|
|
// up: a full share of flows belongs there.
|
|
ifce.sendInsideMessage(hi, pkt, flowForLane(t, ls, 0), nb, tx)
|
|
tx.flush(ifce)
|
|
require.Len(t, conns[0].bufs, 2)
|
|
assert.Equal(t, hi.GetRemote(), conns[0].dsts[1])
|
|
require.NoError(t, h.Parse(conns[0].bufs[1]))
|
|
assert.Equal(t, uint8(0), h.Lane())
|
|
|
|
// Demotion falls back to the base tunnel on the very next packet.
|
|
ls.txAddr[1].Store(nil)
|
|
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx)
|
|
tx.flush(ifce)
|
|
require.Len(t, conns[0].bufs, 3)
|
|
require.Len(t, conns[1].bufs, 2)
|
|
}
|
|
|
|
// Without multiport nothing about the send path changes: every socket shares one
|
|
// port under SO_REUSEPORT, so a routine keeps writing to its own, and there are
|
|
// no lane batches at all. A peer that negotiated no lanes is the same story on a
|
|
// node that does run multiport, except that base traffic owes its 4-tuple to
|
|
// socket 0.
|
|
func TestTxQueueWithoutLanes(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
conns := []*recordingBatchConn{{}, {}, {}, {}}
|
|
ifce.writers = []udp.Conn{conns[0], conns[1], conns[2], conns[3]}
|
|
|
|
initR, _ := runTestHandshake(t)
|
|
hi := newTestLaneHostInfo(t, initR, nil)
|
|
pkt := tio.Packet{Bytes: []byte{0x45, 0, 0, 4, 1, 2, 3, 4}}
|
|
nb := make([]byte, 12)
|
|
flow := testFlow(1234)
|
|
|
|
tx := ifce.newTxQueue(2)
|
|
assert.Empty(t, tx.lane, "a queue with no lanes must not hold lane batches")
|
|
require.Len(t, tx.live, 1)
|
|
|
|
ifce.sendInsideMessage(hi, pkt, flow, nb, tx)
|
|
tx.flush(ifce)
|
|
assert.Len(t, conns[2].bufs, 1, "routine 2 did not send on its own socket")
|
|
assert.Empty(t, conns[0].bufs)
|
|
|
|
// Multiport on, but this peer has no lanes: socket 0, lane 0, one batch.
|
|
ifce.multiport = true
|
|
ifce.laneCount = 4
|
|
tx = ifce.newTxQueue(2)
|
|
require.Len(t, tx.live, 1)
|
|
ifce.sendInsideMessage(hi, pkt, flow, nb, tx)
|
|
tx.flush(ifce)
|
|
require.Len(t, conns[0].bufs, 1, "base traffic left a socket other than 0")
|
|
assert.Len(t, conns[2].bufs, 1)
|
|
|
|
h := &header.H{}
|
|
require.NoError(t, h.Parse(conns[0].bufs[0]))
|
|
assert.Equal(t, uint8(0), h.Lane())
|
|
assert.Len(t, tx.live, 1, "a lane batch was built for a peer with no lanes")
|
|
}
|
|
|
|
// `routines` is per port, so a port's sockets are a group sharing it through
|
|
// SO_REUSEPORT and every routine has a sibling of its own on every port. Pin the
|
|
// arithmetic that turns (queue, lane) into one of them.
|
|
func TestLaneSockGroupLayout(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
|
|
// Two routines per port over three ports: writers[s*2+r] is the r'th socket
|
|
// on port listen.port+s.
|
|
ifce.routinesPerPort = 2
|
|
ifce.routines = 6
|
|
ifce.laneCount = 3
|
|
|
|
// Off, every socket is on the one port, so a routine keeps writing to its own
|
|
// and any of them can carry any lane.
|
|
for q := 0; q < ifce.routines; q++ {
|
|
assert.Equal(t, q, ifce.egressSock(q))
|
|
for s := 0; s < ifce.laneCount; s++ {
|
|
assert.Equal(t, q, ifce.laneSock(q, s))
|
|
}
|
|
}
|
|
|
|
ifce.multiport = true
|
|
for q := 0; q < ifce.routines; q++ {
|
|
// Base traffic is lane 0's port, and each routine has its own socket
|
|
// there rather than sharing one with the other five.
|
|
assert.Equal(t, q%2, ifce.egressSock(q), "queue %d", q)
|
|
for s := 0; s < ifce.laneCount; s++ {
|
|
assert.Equal(t, s*2+q%2, ifce.laneSock(q, s), "queue %d lane %d", q, s)
|
|
}
|
|
}
|
|
|
|
// Sibling routines land on different sockets of the same port, so a lane's
|
|
// port is served by its whole group and not by one socket.
|
|
assert.NotEqual(t, ifce.laneSock(0, 2), ifce.laneSock(1, 2))
|
|
assert.Equal(t, ifce.laneSock(0, 2), ifce.laneSock(2, 2), "queues 0 and 2 share a group position")
|
|
}
|
|
|
|
// The lane batches a routine builds must flush to that routine's own sockets, not
|
|
// to whatever socket happens to be indexed by the lane.
|
|
func TestTxQueueLaneBatchesFollowTheGroup(t *testing.T) {
|
|
hostMap := newHostMap(test.NewLogger())
|
|
ifce := newLaneTestInterface(hostMap)
|
|
ifce.multiport = true
|
|
ifce.routinesPerPort = 2
|
|
ifce.routines = 4
|
|
ifce.laneCount = 2
|
|
|
|
conns := []*recordingBatchConn{{}, {}, {}, {}}
|
|
ifce.writers = []udp.Conn{conns[0], conns[1], conns[2], conns[3]}
|
|
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 2, 2, 5353, 2)
|
|
hi := newTestLaneHostInfo(t, initR, ls)
|
|
laneSessionFor(t, ls, 1)
|
|
laneRemote := netip.MustParseAddrPort("192.0.2.1:5354")
|
|
ls.txAddr[1].Store(&laneRemote)
|
|
|
|
pkt := tio.Packet{Bytes: []byte{0x45, 0, 0, 4, 1, 2, 3, 4}}
|
|
nb := make([]byte, 12)
|
|
flow1 := flowForLane(t, ls, 1)
|
|
|
|
// Routine 1 is the second socket in each group, so its lane 1 traffic leaves
|
|
// writers[1*2+1].
|
|
tx := ifce.newTxQueue(1)
|
|
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx)
|
|
tx.flush(ifce)
|
|
require.Len(t, conns[3].bufs, 1, "lane 1 did not leave routine 1's socket on the lane port")
|
|
assert.Equal(t, laneRemote, conns[3].dsts[0])
|
|
assert.Empty(t, conns[2].bufs)
|
|
|
|
// Routine 0 sends the same flow out the other socket on that same port.
|
|
tx0 := ifce.newTxQueue(0)
|
|
ifce.sendInsideMessage(hi, pkt, flow1, nb, tx0)
|
|
tx0.flush(ifce)
|
|
require.Len(t, conns[2].bufs, 1, "sibling routine shared a socket instead of its own")
|
|
assert.Len(t, conns[3].bufs, 1)
|
|
}
|
|
|
|
// One routine has to be able to reach every lane, since the kernel decides which
|
|
// routine sees a flow and it may well hand one routine everything.
|
|
func TestTxLaneForFlowSpread(t *testing.T) {
|
|
initR, _ := runTestHandshake(t)
|
|
ls := newTestLaneSet(t, initR, 4, 4, 5353, 4)
|
|
|
|
seen := map[int]int{}
|
|
for port := uint16(1024); port < 1224; port++ {
|
|
s := laneOfFlow(ls, testFlow(port))
|
|
require.Less(t, s, ls.txLanes)
|
|
seen[s]++
|
|
}
|
|
assert.Len(t, seen, 4, "200 flows did not reach all four lanes: %v", seen)
|
|
for s, n := range seen {
|
|
assert.Greater(t, n, 10, "lane %d got a negligible share of flows", s)
|
|
}
|
|
|
|
// A flow's lane is a function of its 5-tuple alone, so it never moves
|
|
// mid-flow.
|
|
p := testFlow(1234)
|
|
assert.Equal(t, laneOfFlow(ls, p), laneOfFlow(ls, p))
|
|
|
|
// A tunnel with nothing but the base lane always answers lane 0.
|
|
one := newTestLaneSet(t, initR, 1, 1, 5353, 4)
|
|
s, ci, addr := one.txLaneForFlow(p)
|
|
assert.Zero(t, s)
|
|
assert.Nil(t, ci)
|
|
assert.False(t, addr.IsValid())
|
|
s, ci, _ = (*laneSet)(nil).txLaneForFlow(p)
|
|
assert.Zero(t, s)
|
|
assert.Nil(t, ci)
|
|
}
|
|
|
|
// The two ends of a flow have to pick partner lanes, or the flow's two
|
|
// directions use unrelated 4-tuples instead of exact reverses and neither side's
|
|
// traffic arrives through the conntrack entry the other's probe opened.
|
|
func TestLaneFlowSymmetry(t *testing.T) {
|
|
initR, respR := runTestHandshake(t)
|
|
|
|
// Our side is the low address; the peer builds the same pair reversed.
|
|
const ourBase, peerBase = 4242, 5353
|
|
ourLS := newTestLaneSet(t, initR, 4, 4, peerBase, 4)
|
|
respR.PeerPortCount, respR.PeerBasePort, respR.PeerTxLanes = 4, ourBase, 4
|
|
peerLS := newLaneSet(respR, 4, testPeerAddr, testMyAddr)
|
|
require.NotNil(t, peerLS)
|
|
require.Equal(t, uint16(0), ourLS.laneBias, "the low address hashes straight")
|
|
|
|
partnered := 0
|
|
for port := uint16(1024); port < 1124; port++ {
|
|
p := testFlow(port)
|
|
reverse := &firewall.Packet{
|
|
LocalAddr: p.RemoteAddr, RemoteAddr: p.LocalAddr,
|
|
LocalPort: p.RemotePort, RemotePort: p.LocalPort,
|
|
Protocol: p.Protocol,
|
|
}
|
|
|
|
ours := laneOfFlow(ourLS, p)
|
|
theirs := laneOfFlow(peerLS, reverse)
|
|
require.Equal(t, (ours+int(ourLS.portOffset))%4, theirs,
|
|
"flow on port %d: our lane %d is not partnered with their lane %d", port, ours, theirs)
|
|
|
|
if ours == 0 || theirs == 0 {
|
|
// The lane that lands on the other side's base port has no partner
|
|
// lane; that flow rides one side's base tunnel. See lanePortOffset.
|
|
continue
|
|
}
|
|
partnered++
|
|
|
|
// Our source and destination ports are the reverse of theirs.
|
|
ourLS.mu.Lock()
|
|
peerLS.mu.Lock()
|
|
assert.Equal(t, uint16(peerBase+theirs), ourLS.laneTargetPortLocked(ours), "flow on port %d", port)
|
|
assert.Equal(t, uint16(ourBase+ours), peerLS.laneTargetPortLocked(theirs), "flow on port %d", port)
|
|
peerLS.mu.Unlock()
|
|
ourLS.mu.Unlock()
|
|
}
|
|
assert.Greater(t, partnered, 0, "no flow landed on a partnered lane")
|
|
}
|
|
|
|
// Regression: deleting a hostinfo whose pending entry is NOT the one recorded
|
|
// in vpnIps must not evict a concurrently pending handshake for that address.
|
|
func TestHandshakeManagerVpnIpsIdentityDelete(t *testing.T) {
|
|
l := test.NewLogger()
|
|
hostMap := newHostMap(l)
|
|
lh := newTestLighthouse()
|
|
hm := NewHandshakeManager(l, hostMap, lh, &udp.NoopConn{}, defaultHandshakeConfig)
|
|
|
|
vpnIp := netip.MustParseAddr("172.1.1.4")
|
|
pendingBase := hm.StartHandshake(vpnIp, nil)
|
|
require.NotNil(t, pendingBase)
|
|
|
|
other := &HostInfo{vpnAddrs: []netip.Addr{vpnIp}, localIndexId: 999}
|
|
hm.DeleteHostInfo(other)
|
|
|
|
// The pending base handshake must still be tracked.
|
|
assert.Equal(t, pendingBase, hm.QueryVpnAddr(vpnIp))
|
|
|
|
// And deleting the actual owner still works.
|
|
hm.DeleteHostInfo(pendingBase)
|
|
assert.Nil(t, hm.QueryVpnAddr(vpnIp))
|
|
}
|