mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-07 22:19:08 +02:00
Merge branch 'main' into feature/shared-service-config-loader
This commit is contained in:
@@ -1,88 +0,0 @@
|
|||||||
name: Mobile
|
|
||||||
|
|
||||||
on:
|
|
||||||
push:
|
|
||||||
branches:
|
|
||||||
- main
|
|
||||||
- "release-*"
|
|
||||||
pull_request:
|
|
||||||
|
|
||||||
concurrency:
|
|
||||||
group: ${{ github.workflow }}-${{ github.ref }}-${{ github.head_ref || github.actor_id }}
|
|
||||||
cancel-in-progress: true
|
|
||||||
|
|
||||||
jobs:
|
|
||||||
android_build:
|
|
||||||
name: "Android / Build"
|
|
||||||
runs-on: ubuntu-latest
|
|
||||||
steps:
|
|
||||||
- name: Checkout repository
|
|
||||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
|
|
||||||
with:
|
|
||||||
persist-credentials: false
|
|
||||||
- name: Install Go
|
|
||||||
uses: actions/setup-go@924ae3a1cded613372ab5595356fb5720e22ba16 # v6.5.0
|
|
||||||
with:
|
|
||||||
go-version-file: "go.mod"
|
|
||||||
- name: Setup Android SDK
|
|
||||||
uses: android-actions/setup-android@40fd30fb8d7440372e1316f5d1809ec01dcd3699 # v4.0.1
|
|
||||||
with:
|
|
||||||
cmdline-tools-version: 8512546
|
|
||||||
- name: Setup Java
|
|
||||||
uses: actions/setup-java@1bcf9fb12cf4aa7d266a90ae39939e61372fe520
|
|
||||||
with:
|
|
||||||
java-version: "11"
|
|
||||||
distribution: "adopt"
|
|
||||||
- name: NDK Cache
|
|
||||||
id: ndk-cache
|
|
||||||
uses: actions/cache@2c8a9bd7457de244a408f35966fab2fb45fda9c8 # v6.0.0
|
|
||||||
with:
|
|
||||||
path: /usr/local/lib/android/sdk/ndk
|
|
||||||
key: ndk-cache-23.1.7779620
|
|
||||||
- name: Setup NDK
|
|
||||||
run: /usr/local/lib/android/sdk/cmdline-tools/7.0/bin/sdkmanager --install "ndk;23.1.7779620"
|
|
||||||
- name: install gomobile
|
|
||||||
run: go install golang.org/x/mobile/cmd/gomobile@v0.0.0-20251113184115-a159579294ab
|
|
||||||
# `gomobile init` re-installs gobind from golang.org/x/mobile@latest
|
|
||||||
# regardless of the pin above (cmd/gomobile/init.go: "Make sure gobind is
|
|
||||||
# up to date"), so this step resolves a version nobody chose, on every run.
|
|
||||||
#
|
|
||||||
# setup-go sets GOTOOLCHAIN=local, so that install fails outright once
|
|
||||||
# x/mobile@latest declares a newer Go than go.mod does — which it did on
|
|
||||||
# 2026-08-21, breaking both jobs on every branch at once. GOTOOLCHAIN=auto
|
|
||||||
# lets this one install fetch the toolchain it asks for. Scoped to the
|
|
||||||
# step: the repo's own Go version, and every build below, is unaffected.
|
|
||||||
- name: gomobile init
|
|
||||||
run: gomobile init
|
|
||||||
env:
|
|
||||||
GOTOOLCHAIN: auto
|
|
||||||
- name: build android netbird lib
|
|
||||||
run: PATH=$PATH:$(go env GOPATH) gomobile bind -o $GITHUB_WORKSPACE/netbird.aar -javapkg=io.netbird.gomobile -ldflags="-checklinkname=0 -X golang.zx2c4.com/wireguard/ipc.socketDirectory=/data/data/io.netbird.client/cache/wireguard -X github.com/netbirdio/netbird/version.version=buildtest" $GITHUB_WORKSPACE/client/android
|
|
||||||
env:
|
|
||||||
CGO_ENABLED: 0
|
|
||||||
ANDROID_NDK_HOME: /usr/local/lib/android/sdk/ndk/23.1.7779620
|
|
||||||
ios_build:
|
|
||||||
name: "iOS / Build"
|
|
||||||
runs-on: macos-latest
|
|
||||||
steps:
|
|
||||||
- name: Checkout repository
|
|
||||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7.0.0
|
|
||||||
with:
|
|
||||||
persist-credentials: false
|
|
||||||
- name: Install Go
|
|
||||||
uses: actions/setup-go@924ae3a1cded613372ab5595356fb5720e22ba16 # v6.5.0
|
|
||||||
with:
|
|
||||||
go-version-file: "go.mod"
|
|
||||||
- name: install gomobile
|
|
||||||
run: go install golang.org/x/mobile/cmd/gomobile@v0.0.0-20251113184115-a159579294ab
|
|
||||||
# See the Android job: `gomobile init` re-installs gobind from
|
|
||||||
# golang.org/x/mobile@latest regardless of the pin above, and needs a
|
|
||||||
# toolchain it may pick newer than go.mod's.
|
|
||||||
- name: gomobile init
|
|
||||||
run: gomobile init
|
|
||||||
env:
|
|
||||||
GOTOOLCHAIN: auto
|
|
||||||
- name: build iOS netbird lib
|
|
||||||
run: PATH=$PATH:$(go env GOPATH) gomobile bind -target=ios -bundleid=io.netbird.framework -ldflags="-X github.com/netbirdio/netbird/version.version=buildtest" -o ./NetBirdSDK.xcframework ./client/ios/NetBirdSDK
|
|
||||||
env:
|
|
||||||
CGO_ENABLED: 0
|
|
||||||
@@ -2,21 +2,17 @@ package ebpf
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
_ "embed"
|
_ "embed"
|
||||||
"fmt"
|
|
||||||
"net"
|
"net"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
"github.com/cilium/ebpf/link"
|
"github.com/cilium/ebpf/link"
|
||||||
"github.com/cilium/ebpf/rlimit"
|
"github.com/cilium/ebpf/rlimit"
|
||||||
log "github.com/sirupsen/logrus"
|
log "github.com/sirupsen/logrus"
|
||||||
"golang.org/x/sys/unix"
|
|
||||||
|
|
||||||
"github.com/netbirdio/netbird/client/internal/ebpf/manager"
|
"github.com/netbirdio/netbird/client/internal/ebpf/manager"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
xdpProgName = "nb_xdp_prog"
|
|
||||||
|
|
||||||
mapKeyFeatures uint32 = 0
|
mapKeyFeatures uint32 = 0
|
||||||
|
|
||||||
featureFlagWGProxy = 0b00000001
|
featureFlagWGProxy = 0b00000001
|
||||||
@@ -72,50 +68,21 @@ func (tf *GeneralManager) loadXdp() error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
// lo has no native XDP, so the program runs in generic mode. Unless it
|
// load pre-compiled programs into the kernel.
|
||||||
// declares multi-buffer support the kernel must linearize every non-linear
|
err = loadBpfObjects(&tf.bpfObjs, nil)
|
||||||
// skb before running it. Loopback packets are up to 64 KB, so that is a
|
|
||||||
// contiguous GFP_ATOMIC allocation per packet, and when it fails the packet
|
|
||||||
// is dropped before the program runs, stalling local TCP connections.
|
|
||||||
// Multi-buffer XDP in generic mode requires kernel 6.3, so fall back to a
|
|
||||||
// plain attach when the kernel rejects it.
|
|
||||||
err = tf.attachXdp(iFace.Index, true)
|
|
||||||
if err == nil {
|
|
||||||
return nil
|
|
||||||
}
|
|
||||||
log.Debugf("failed to attach multi-buffer xdp program, retrying without it: %s", err)
|
|
||||||
|
|
||||||
return tf.attachXdp(iFace.Index, false)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (tf *GeneralManager) attachXdp(iFaceIndex int, multiBuffer bool) error {
|
|
||||||
spec, err := loadBpf()
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("load bpf spec: %w", err)
|
return err
|
||||||
}
|
|
||||||
|
|
||||||
if multiBuffer {
|
|
||||||
prog, ok := spec.Programs[xdpProgName]
|
|
||||||
if !ok {
|
|
||||||
return fmt.Errorf("program %s not found in bpf spec", xdpProgName)
|
|
||||||
}
|
|
||||||
prog.Flags |= unix.BPF_F_XDP_HAS_FRAGS
|
|
||||||
}
|
|
||||||
|
|
||||||
if err := spec.LoadAndAssign(&tf.bpfObjs, nil); err != nil {
|
|
||||||
return fmt.Errorf("load bpf objects: %w", err)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
tf.link, err = link.AttachXDP(link.XDPOptions{
|
tf.link, err = link.AttachXDP(link.XDPOptions{
|
||||||
Program: tf.bpfObjs.NbXdpProg,
|
Program: tf.bpfObjs.NbXdpProg,
|
||||||
Interface: iFaceIndex,
|
Interface: iFace.Index,
|
||||||
})
|
})
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if closeErr := tf.bpfObjs.Close(); closeErr != nil {
|
_ = tf.bpfObjs.Close()
|
||||||
log.Debugf("failed to close bpf objects after xdp attach error: %s", closeErr)
|
|
||||||
}
|
|
||||||
tf.link = nil
|
tf.link = nil
|
||||||
return fmt.Errorf("attach xdp: %w", err)
|
return err
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -81,14 +81,19 @@ type Handshaker struct {
|
|||||||
|
|
||||||
func NewHandshaker(log *log.Entry, config ConnConfig, signaler *Signaler, ice *WorkerICE, relay *WorkerRelay, metricsStages *MetricsStages) *Handshaker {
|
func NewHandshaker(log *log.Entry, config ConnConfig, signaler *Signaler, ice *WorkerICE, relay *WorkerRelay, metricsStages *MetricsStages) *Handshaker {
|
||||||
h := &Handshaker{
|
h := &Handshaker{
|
||||||
log: log,
|
log: log,
|
||||||
config: config,
|
config: config,
|
||||||
signaler: signaler,
|
signaler: signaler,
|
||||||
ice: ice,
|
ice: ice,
|
||||||
relay: relay,
|
relay: relay,
|
||||||
metricsStages: metricsStages,
|
metricsStages: metricsStages,
|
||||||
remoteOffersCh: make(chan OfferAnswer),
|
// Buffered by one so an offer or answer that arrives between Open launching
|
||||||
remoteAnswerCh: make(chan OfferAnswer),
|
// the Listen goroutine and it reaching its receive is held rather than
|
||||||
|
// dropped. A peer activated by an incoming signal receives the remote's
|
||||||
|
// message in that window; an unbuffered channel skips it as "receiver not
|
||||||
|
// ready", and the connection cannot proceed until the remote re-sends.
|
||||||
|
remoteOffersCh: make(chan OfferAnswer, 1),
|
||||||
|
remoteAnswerCh: make(chan OfferAnswer, 1),
|
||||||
}
|
}
|
||||||
// assume remote supports ICE until we learn otherwise from received offers
|
// assume remote supports ICE until we learn otherwise from received offers
|
||||||
h.remoteICESupported.Store(ice != nil)
|
h.remoteICESupported.Store(ice != nil)
|
||||||
@@ -162,29 +167,38 @@ func (h *Handshaker) SendOffer() error {
|
|||||||
return h.sendOffer()
|
return h.sendOffer()
|
||||||
}
|
}
|
||||||
|
|
||||||
// OnRemoteOffer handles an offer from the remote peer and returns true if the message was accepted, false otherwise
|
// OnRemoteOffer hands an offer to Listen without blocking, keeping only the most
|
||||||
// doesn't block, discards the message if connection wasn't ready
|
// recent one if several arrive before Listen reads them.
|
||||||
func (h *Handshaker) OnRemoteOffer(offer OfferAnswer) {
|
func (h *Handshaker) OnRemoteOffer(offer OfferAnswer) {
|
||||||
select {
|
enqueueLatest(h.remoteOffersCh, offer)
|
||||||
case h.remoteOffersCh <- offer:
|
|
||||||
return
|
|
||||||
default:
|
|
||||||
h.log.Warnf("skipping remote offer message because receiver not ready")
|
|
||||||
// connection might not be ready yet to receive so we ignore the message
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// OnRemoteAnswer handles an offer from the remote peer and returns true if the message was accepted, false otherwise
|
// OnRemoteAnswer hands an answer to Listen without blocking, keeping only the most
|
||||||
// doesn't block, discards the message if connection wasn't ready
|
// recent one if several arrive before Listen reads them.
|
||||||
func (h *Handshaker) OnRemoteAnswer(answer OfferAnswer) {
|
func (h *Handshaker) OnRemoteAnswer(answer OfferAnswer) {
|
||||||
|
enqueueLatest(h.remoteAnswerCh, answer)
|
||||||
|
}
|
||||||
|
|
||||||
|
// enqueueLatest delivers msg on a one-slot channel without blocking. When the slot
|
||||||
|
// already holds an unread message the older one is discarded in favor of msg, so a
|
||||||
|
// message arriving before Listen starts reading is held rather than dropped, and
|
||||||
|
// the newest wins if several arrive first. Safe because there is a single producer
|
||||||
|
// (the engine loop): after draining the stale value the send always has room.
|
||||||
|
func enqueueLatest(ch chan OfferAnswer, msg OfferAnswer) {
|
||||||
select {
|
select {
|
||||||
case h.remoteAnswerCh <- answer:
|
case ch <- msg:
|
||||||
return
|
return
|
||||||
default:
|
default:
|
||||||
// connection might not be ready yet to receive so we ignore the message
|
}
|
||||||
h.log.Warnf("skipping remote answer message because receiver not ready")
|
|
||||||
return
|
select {
|
||||||
|
case <-ch:
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case ch <- msg:
|
||||||
|
default:
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,63 @@
|
|||||||
|
package peer
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
log "github.com/sirupsen/logrus"
|
||||||
|
"github.com/stretchr/testify/assert"
|
||||||
|
)
|
||||||
|
|
||||||
|
func newTestHandshaker(t *testing.T) *Handshaker {
|
||||||
|
t.Helper()
|
||||||
|
// The tests exercise the answer path, whose Listen branch dispatches to the
|
||||||
|
// relay listener without sending an answer, so no signaler/ICE/relay is needed.
|
||||||
|
return NewHandshaker(log.WithField("test", t.Name()), ConnConfig{}, nil, nil, nil, nil)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHandshakerHoldsSignalArrivingBeforeListen covers the case where a peer is
|
||||||
|
// activated by an incoming signal: the remote's offer/answer arrives in the same
|
||||||
|
// step that opens the connection, before the Listen loop starts reading. The
|
||||||
|
// message must be held rather than dropped, or the connection cannot proceed until
|
||||||
|
// the remote re-sends. This is the path taken when an eager peer connects to a
|
||||||
|
// lazily-managed one.
|
||||||
|
func TestHandshakerHoldsSignalArrivingBeforeListen(t *testing.T) {
|
||||||
|
h := newTestHandshaker(t)
|
||||||
|
|
||||||
|
processed := make(chan *OfferAnswer, 4)
|
||||||
|
h.AddRelayListener(func(o *OfferAnswer) { processed <- o })
|
||||||
|
|
||||||
|
// Delivered before Listen is reading, as when the peer is woken by the remote's
|
||||||
|
// signal and the message is delivered right after Open.
|
||||||
|
h.OnRemoteAnswer(OfferAnswer{WgListenPort: 51820})
|
||||||
|
|
||||||
|
go h.Listen(t.Context())
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-processed:
|
||||||
|
case <-time.After(2 * time.Second):
|
||||||
|
assert.Fail(t, "remote-answer dispatch: signal delivered before Listen was ready was dropped")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestHandshakerKeepsLatestSignalBeforeListen covers several signals arriving
|
||||||
|
// before Listen reads: the newest must win (matching the latest-offer contract),
|
||||||
|
// rather than the first being kept and later ones discarded.
|
||||||
|
func TestHandshakerKeepsLatestSignalBeforeListen(t *testing.T) {
|
||||||
|
h := newTestHandshaker(t)
|
||||||
|
|
||||||
|
processed := make(chan *OfferAnswer, 4)
|
||||||
|
h.AddRelayListener(func(o *OfferAnswer) { processed <- o })
|
||||||
|
|
||||||
|
h.OnRemoteAnswer(OfferAnswer{WgListenPort: 1111})
|
||||||
|
h.OnRemoteAnswer(OfferAnswer{WgListenPort: 2222})
|
||||||
|
|
||||||
|
go h.Listen(t.Context())
|
||||||
|
|
||||||
|
select {
|
||||||
|
case got := <-processed:
|
||||||
|
assert.Equal(t, 2222, got.WgListenPort, "remote-answer dispatch: the latest queued signal should be processed")
|
||||||
|
case <-time.After(2 * time.Second):
|
||||||
|
assert.Fail(t, "remote-answer dispatch: queued signal was dropped")
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user