Compare commits

..

1 Commits

Author SHA1 Message Date
Viktor Liu
b3ead5ee7e Localize daemon notifications via stable message keys 2026-08-05 17:53:02 +02:00
142 changed files with 3557 additions and 5446 deletions

View File

@@ -16,8 +16,6 @@ jobs:
const allowedTags = [
'management',
'client',
'android',
'ios',
'signal',
'proxy',
'relay',

View File

@@ -93,9 +93,7 @@ nfpms:
- src: client/ui/build/appicon.png
dst: /usr/share/pixmaps/netbird.png
dependencies:
- netbird (>= 0.75.0)
- libgtk-4-1 (>= 4.14)
- libwebkitgtk-6.0-4
- netbird
- maintainer: Netbird <dev@netbird.io>
description: Netbird client UI.
@@ -116,9 +114,7 @@ nfpms:
- src: client/ui/build/appicon.png
dst: /usr/share/pixmaps/netbird.png
dependencies:
- netbird >= 0.75.0
- (gtk4 >= 4.14 or libgtk-4-1 >= 4.14)
- (webkitgtk6.0 or libwebkitgtk-6_0-4)
- netbird
rpm:
signature:

View File

@@ -57,12 +57,6 @@ type DnsReadyListener interface {
dns.ReadyListener
}
// TunSettings is a snapshot of the settings the TUN device is rebuilt with
type TunSettings struct {
Routes string
SearchDomains string
}
func init() {
formatter.SetLogcatFormatter(log.StandardLogger())
}
@@ -82,8 +76,6 @@ type Client struct {
connectClient *internal.ConnectClient
config *profilemanager.Config
cacheDir string
// Identifies the running profile for the SSO login hint; see profile_state.go.
cfgPath string
stateChangeMu sync.Mutex
stateChangeSubID string
@@ -104,12 +96,11 @@ type Client struct {
extendCancel context.CancelFunc
}
func (c *Client) setState(cfg *profilemanager.Config, cacheDir string, cfgPath string, cc *internal.ConnectClient) {
func (c *Client) setState(cfg *profilemanager.Config, cacheDir string, cc *internal.ConnectClient) {
c.stateMu.Lock()
defer c.stateMu.Unlock()
c.config = cfg
c.cacheDir = cacheDir
c.cfgPath = cfgPath
c.connectClient = cc
}
@@ -119,16 +110,6 @@ func (c *Client) stateSnapshot() (*profilemanager.Config, string, *internal.Conn
return c.config, c.cacheDir, c.connectClient
}
// authSnapshot returns the config together with the path it was loaded from, in
// one lock: the path identifies the profile whose account email backs the login
// hint, so reading it separately could pair one profile's config with another's
// hint when a profile switch lands in between.
func (c *Client) authSnapshot() (*profilemanager.Config, string, *internal.ConnectClient) {
c.stateMu.RLock()
defer c.stateMu.RUnlock()
return c.config, c.cfgPath, c.connectClient
}
func (c *Client) getConnectClient() *internal.ConnectClient {
c.stateMu.RLock()
defer c.stateMu.RUnlock()
@@ -181,7 +162,7 @@ func (c *Client) Run(platformFiles PlatformFiles, urlOpener URLOpener, isAndroid
defer c.ctxCancel()
c.ctxCancelLock.Unlock()
auth := NewAuthWithConfig(ctx, cfg, cfgFile)
auth := NewAuthWithConfig(ctx, cfg)
err = auth.login(urlOpener, isAndroidTV)
if err != nil {
return err
@@ -189,7 +170,7 @@ func (c *Client) Run(platformFiles PlatformFiles, urlOpener URLOpener, isAndroid
// todo do not throw error in case of cancelled context
ctx = internal.CtxInitState(ctx)
connectClient := internal.NewConnectClient(ctx, cfg, c.recorder)
c.setState(cfg, cacheDir, cfgFile, connectClient)
c.setState(cfg, cacheDir, connectClient)
// This path runs the interactive SSO flow, so reaching here means the peer
// is authenticated again — release the latch Status() reports from. Clear
// only once the fresh connect client is installed: until then Status()
@@ -230,7 +211,7 @@ func (c *Client) RunWithoutLogin(platformFiles PlatformFiles, dns *DNSList, dnsR
// todo do not throw error in case of cancelled context
ctx = internal.CtxInitState(ctx)
connectClient := internal.NewConnectClient(ctx, cfg, c.recorder)
c.setState(cfg, cacheDir, cfgFile, connectClient)
c.setState(cfg, cacheDir, connectClient)
return connectClient.RunOnAndroid(c.tunAdapter, c.iFaceDiscover, c.networkChangeListener, slices.Clone(dns.items), dnsReadyListener, stateFile, cacheDir)
}
@@ -259,24 +240,6 @@ func (c *Client) RenewTun(fd int) error {
return e.RenewTun(fd)
}
func (c *Client) GetTunSettings() (*TunSettings, error) {
cc := c.getConnectClient()
if cc == nil {
return nil, fmt.Errorf("engine not running")
}
e := cc.Engine()
if e == nil {
return nil, fmt.Errorf("engine not initialized")
}
routes, searchDomains := e.TunSettings()
return &TunSettings{
Routes: strings.Join(routes, ";"),
SearchDomains: strings.Join(searchDomains, ";"),
}, nil
}
// DebugBundle generates a debug bundle, uploads it, and returns the upload key.
// It works both with and without a running engine.
func (c *Client) DebugBundle(platformFiles PlatformFiles, anonymize bool) (string, error) {

View File

@@ -4,8 +4,6 @@ import (
"context"
"fmt"
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/internal/auth"
"github.com/netbirdio/netbird/client/internal/profilemanager"
"github.com/netbirdio/netbird/client/system"
@@ -38,20 +36,12 @@ type Auth struct {
}
// NewAuth instantiate Auth struct and validate the management URL
//
// The configuration at cfgPath is reused when one is already there, and only created when it is
// not. Building a fresh in-memory config unconditionally gives the client a new WireGuard key on
// every call: the peer registers under that key, the key is written out, and any peer registered by
// an earlier call is orphaned on the server. It also breaks a client that enrols and then runs from
// the persisted config, because the identity it registered is not the one it runs with — the
// management stream rejects it with "no peer auth method provided".
func NewAuth(cfgPath string, mgmURL string) (*Auth, error) {
inputCfg := profilemanager.ConfigInput{
ConfigPath: cfgPath,
ManagementURL: mgmURL,
}
cfg, err := profilemanager.UpdateOrCreateConfig(inputCfg)
cfg, err := profilemanager.CreateInMemoryConfig(inputCfg)
if err != nil {
return nil, err
}
@@ -63,14 +53,11 @@ func NewAuth(cfgPath string, mgmURL string) (*Auth, error) {
}, nil
}
// NewAuthWithConfig instantiate Auth based on existing config. cfgPath is the
// file the config was loaded from; it identifies the profile whose account email
// backs the login_hint.
func NewAuthWithConfig(ctx context.Context, config *profilemanager.Config, cfgPath string) *Auth {
// NewAuthWithConfig instantiate Auth based on existing config
func NewAuthWithConfig(ctx context.Context, config *profilemanager.Config) *Auth {
return &Auth{
ctx: ctx,
config: config,
cfgPath: cfgPath,
ctx: ctx,
config: config,
}
}
@@ -163,14 +150,12 @@ func (a *Auth) login(urlOpener URLOpener, isAndroidTV bool) error {
}
jwtToken := ""
email := ""
if needsLogin {
tokenInfo, err := a.foregroundGetTokenInfo(authClient, urlOpener, isAndroidTV)
if err != nil {
return fmt.Errorf("interactive sso login failed: %v", err)
}
jwtToken = tokenInfo.GetTokenToUse()
email = tokenInfo.Email
}
err, _ = authClient.Login(a.ctx, "", jwtToken)
@@ -178,42 +163,17 @@ func (a *Auth) login(urlOpener URLOpener, isAndroidTV bool) error {
return fmt.Errorf("login failed: %v", err)
}
// Stored after Login, not before: a rejected token must not leave a hint
// pointing at an account that cannot be used.
if email != "" && a.cfgPath != "" {
if err := writeProfileEmail(a.cfgPath, email); err != nil {
log.Warnf("failed to store profile account email: %v", err)
}
}
go urlOpener.OnLoginSuccess()
return nil
}
// loginHintSetter is implemented by both concrete flows (PKCE and device code)
// but absent from the OAuthFlow interface, hence the assertion below — the same
// way internal/auth wires it in authenticateWithPKCEFlow.
type loginHintSetter interface {
SetLoginHint(hint string)
}
func (a *Auth) foregroundGetTokenInfo(authClient *auth.Auth, urlOpener URLOpener, isAndroidTV bool) (*auth.TokenInfo, error) {
oAuthFlow, err := authClient.GetOAuthFlow(a.ctx, isAndroidTV)
if err != nil {
return nil, fmt.Errorf("failed to get OAuth flow: %v", err)
}
// An empty hint is deliberate, not a fallback: a fresh or logged-out profile
// leaves the choice to the IdP, which is how accounts get switched.
if a.cfgPath != "" {
if hint := readProfileEmail(a.cfgPath); hint != "" {
if setter, ok := oAuthFlow.(loginHintSetter); ok {
setter.SetLoginHint(hint)
}
}
}
flowInfo, err := oAuthFlow.RequestAuthInfo(context.TODO())
if err != nil {
return nil, fmt.Errorf("getting a request OAuth flow info failed: %v", err)

View File

@@ -1,51 +0,0 @@
package android
import (
"path/filepath"
"testing"
)
// NewAuth must reuse the configuration already at cfgPath rather than building a fresh one.
//
// Creating a new in-memory config on every call gives the client a new WireGuard private key each
// time. The peer registers under that key and the key is written out, so a peer registered by an
// earlier call is orphaned on the server — a client that enrols twice leaves two entries and owns
// neither. It also breaks enrol-then-run: RunWithoutLogin reloads the configuration from disk, so
// the identity that registered is not the identity that runs, and the management stream rejects it
// with "no peer auth method provided, please use a setup key or interactive SSO login".
func TestNewAuth_ReusesPersistedIdentity(t *testing.T) {
cfgPath := filepath.Join(t.TempDir(), "config.json")
first, err := NewAuth(cfgPath, "https://api.example.com:443")
if err != nil {
t.Fatalf("first NewAuth: %v", err)
}
if first.config.PrivateKey == "" {
t.Fatal("first NewAuth produced no private key")
}
second, err := NewAuth(cfgPath, "https://api.example.com:443")
if err != nil {
t.Fatalf("second NewAuth: %v", err)
}
if second.config.PrivateKey != first.config.PrivateKey {
t.Errorf("private key changed between calls: a second enrolment would orphan the peer registered by the first")
}
}
// A missing configuration is still created, so a first enrolment works unchanged.
func TestNewAuth_CreatesConfigWhenAbsent(t *testing.T) {
cfgPath := filepath.Join(t.TempDir(), "config.json")
auth, err := NewAuth(cfgPath, "https://api.example.com:443")
if err != nil {
t.Fatalf("NewAuth: %v", err)
}
if auth.config == nil || auth.config.PrivateKey == "" {
t.Fatal("NewAuth did not create a usable configuration")
}
if auth.cfgPath != cfgPath {
t.Errorf("cfgPath = %q, want %q", auth.cfgPath, cfgPath)
}
}

View File

@@ -13,17 +13,18 @@ import (
)
const (
// Android-specific config filename (different from desktop default.json)
defaultConfigFilename = "netbird.cfg"
// Subdirectory for non-default profiles (must match Java Preferences.java)
profilesSubdir = "profiles"
// Android uses a single user context per app (non-empty username required by ServiceManager)
androidUsername = "android"
)
// Profile represents a profile for gomobile
type Profile struct {
ID string
Name string
// Email is the account this profile last logged in with, "" if it never
// completed an SSO login or was logged out. See profile_state.go.
Email string
ID string
Name string
IsActive bool
}
@@ -100,7 +101,6 @@ func (pm *ProfileManager) ListProfiles() (*ProfileArray, error) {
profiles = append(profiles, &Profile{
ID: p.ID.String(),
Name: p.Name,
Email: pm.profileEmail(p.ID.String()),
IsActive: p.IsActive,
})
}
@@ -123,22 +123,7 @@ func (pm *ProfileManager) GetActiveProfile() (*Profile, error) {
if err != nil {
return nil, fmt.Errorf("failed to resolve active profile %q: %w", activeState.ID, err)
}
return &Profile{
ID: prof.ID.String(),
Name: prof.Name,
Email: pm.profileEmail(prof.ID.String()),
IsActive: true,
}, nil
}
// profileEmail returns the account email recorded for a profile. Display-only, so
// an unresolvable path degrades to "" rather than an error.
func (pm *ProfileManager) profileEmail(id string) string {
configPath, err := pm.getProfileConfigPath(id)
if err != nil {
return ""
}
return readProfileEmail(configPath)
return &Profile{ID: prof.ID.String(), Name: prof.Name, IsActive: true}, nil
}
// SwitchProfile switches to a different profile
@@ -200,11 +185,6 @@ func (pm *ProfileManager) LogoutProfile(id string) error {
return fmt.Errorf("failed to save config: %w", err)
}
// Not fatal: a stale hint costs an account switch, not the logout itself.
if err := removeProfileEmail(configPath); err != nil {
log.Warnf("failed to clear stored account email for profile %s: %v", id, err)
}
log.Infof("logged out from profile: %s", id)
return nil
}

View File

@@ -1,108 +0,0 @@
package android
import (
"context"
"fmt"
"os"
"path/filepath"
"strings"
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/internal/profilemanager"
"github.com/netbirdio/netbird/util"
)
const (
// Android-specific config filename (different from desktop default.json)
defaultConfigFilename = "netbird.cfg"
// Subdirectory for non-default profiles (must match Java Preferences.java)
profilesSubdir = "profiles"
// profileAccountSuffix names the file holding the profile's account email.
// Deliberately not ".state.json", which desktop uses for the same data:
// there the email and the engine's state manager live in different
// directories, but on Android both resolve under files/, so sharing the name
// would have the two overwrite each other — the state manager rewrites the
// whole file from its own keys (see statemanager.Manager.PersistState), and
// this package's writer does the same in reverse.
profileAccountSuffix = ".account.json"
)
// profileAccountPathFor derives the account file path from a profile's config
// path: netbird.cfg -> netbird.account.json, <id>.json -> <id>.account.json.
//
// Deriving from the config path rather than resolving the active profile keeps
// the write on the profile the login actually ran for: Auth.login runs in a
// goroutine, so the active profile can change under a flow already in flight.
func profileAccountPathFor(configPath string) (string, error) {
if configPath == "" {
return "", fmt.Errorf("empty config path")
}
base := filepath.Base(configPath)
stem := strings.TrimSuffix(base, filepath.Ext(base))
if stem == "" || stem == "." {
return "", fmt.Errorf("config path %q has no filename stem", configPath)
}
return filepath.Join(filepath.Dir(configPath), stem+profileAccountSuffix), nil
}
// readProfileEmail returns the account email stored for the profile whose config
// lives at configPath. A missing or unreadable file yields "", which leaves the
// account choice to the IdP.
func readProfileEmail(configPath string) string {
accountPath, err := profileAccountPathFor(configPath)
if err != nil {
log.Debugf("no profile account path for login hint: %v", err)
return ""
}
var state profilemanager.ProfileState
if _, err := util.ReadJson(accountPath, &state); err != nil {
if !os.IsNotExist(err) {
log.Debugf("failed to read profile account for login hint: %v", err)
}
return ""
}
return state.Email
}
// writeProfileEmail records the account email for the profile whose config lives
// at configPath, so later logins can pass it as an OIDC login_hint. An empty
// email is ignored rather than blanking what is already stored.
func writeProfileEmail(configPath string, email string) error {
if email == "" {
return nil
}
accountPath, err := profileAccountPathFor(configPath)
if err != nil {
return fmt.Errorf("resolve profile account path: %w", err)
}
state := profilemanager.ProfileState{Email: email}
if err := util.WriteJsonWithRestrictedPermission(context.Background(), accountPath, state); err != nil {
return fmt.Errorf("write profile account: %w", err)
}
return nil
}
// removeProfileEmail drops the stored account email. Called on logout: while the
// email is on disk it goes out as a login_hint, which would steer the next login
// straight back into the account just logged out of. Mirrors the desktop UI's
// RemoveProfileState call.
func removeProfileEmail(configPath string) error {
accountPath, err := profileAccountPathFor(configPath)
if err != nil {
return fmt.Errorf("resolve profile account path: %w", err)
}
if err := os.Remove(accountPath); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove profile account: %w", err)
}
return nil
}

View File

@@ -1,161 +0,0 @@
package android
import (
"os"
"path/filepath"
"testing"
)
func TestProfileAccountPathFor(t *testing.T) {
tests := []struct {
name string
configPath string
want string
wantErr bool
}{
{
name: "default profile",
configPath: "/data/data/io.netbird.client/files/netbird.cfg",
want: filepath.FromSlash("/data/data/io.netbird.client/files/netbird.account.json"),
},
{
name: "id profile",
configPath: "/data/data/io.netbird.client/files/profiles/4c5f5c8198c3989cffb5b5394f5a7ae0.json",
want: filepath.FromSlash("/data/data/io.netbird.client/files/profiles/4c5f5c8198c3989cffb5b5394f5a7ae0.account.json"),
},
{
name: "legacy name-keyed profile is handled the same way",
configPath: "/data/data/io.netbird.client/files/profiles/work.json",
want: filepath.FromSlash("/data/data/io.netbird.client/files/profiles/work.account.json"),
},
{
name: "empty path is rejected",
configPath: "",
wantErr: true,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := profileAccountPathFor(tt.configPath)
if tt.wantErr {
if err == nil {
t.Fatalf("expected an error, got path %q", got)
}
return
}
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if got != tt.want {
t.Errorf("got %q, want %q", got, tt.want)
}
})
}
}
func TestProfileAccountPathForDefaultDoesNotCollide(t *testing.T) {
root := "/data/data/io.netbird.client/files"
defaultAccount, err := profileAccountPathFor(filepath.Join(root, defaultConfigFilename))
if err != nil {
t.Fatalf("default profile: %v", err)
}
idAccount, err := profileAccountPathFor(filepath.Join(root, profilesSubdir, "abc123.json"))
if err != nil {
t.Fatalf("id profile: %v", err)
}
if defaultAccount == idAccount {
t.Fatalf("default and id profile share an account file: %q", defaultAccount)
}
}
// The account file must never land on the engine state file: on Android both
// resolve under files/, and the state manager rewrites the whole file from its
// own keys, so sharing a path would have the two overwrite each other. The
// expected names here mirror ProfileManager.GetStateFilePath.
func TestProfileAccountPathAvoidsEngineStateFile(t *testing.T) {
root := "/data/data/io.netbird.client/files"
cases := []struct {
configPath string
engineState string
}{
{
configPath: filepath.Join(root, defaultConfigFilename),
engineState: filepath.Join(root, "state.json"),
},
{
configPath: filepath.Join(root, profilesSubdir, "abc123.json"),
engineState: filepath.Join(root, profilesSubdir, "abc123.state.json"),
},
}
for _, c := range cases {
account, err := profileAccountPathFor(c.configPath)
if err != nil {
t.Fatalf("%s: %v", c.configPath, err)
}
if account == c.engineState {
t.Errorf("account file collides with the engine state file: %q", account)
}
}
}
func TestWriteThenReadProfileEmail(t *testing.T) {
configPath := filepath.Join(t.TempDir(), "profiles", "abc123.json")
if err := ensureDirFor(t, configPath); err != nil {
t.Fatalf("prepare dir: %v", err)
}
if got := readProfileEmail(configPath); got != "" {
t.Errorf("expected no email before a login, got %q", got)
}
const email = "user@example.com"
if err := writeProfileEmail(configPath, email); err != nil {
t.Fatalf("write: %v", err)
}
if got := readProfileEmail(configPath); got != email {
t.Errorf("got %q, want %q", got, email)
}
if err := removeProfileEmail(configPath); err != nil {
t.Fatalf("remove: %v", err)
}
if got := readProfileEmail(configPath); got != "" {
t.Errorf("expected no email after logout, got %q", got)
}
// Logout may run on a never-logged-in profile, so a second remove must pass.
if err := removeProfileEmail(configPath); err != nil {
t.Fatalf("second remove should be a no-op: %v", err)
}
}
func TestWriteProfileEmailIgnoresEmpty(t *testing.T) {
configPath := filepath.Join(t.TempDir(), "profiles", "abc123.json")
if err := ensureDirFor(t, configPath); err != nil {
t.Fatalf("prepare dir: %v", err)
}
const email = "user@example.com"
if err := writeProfileEmail(configPath, email); err != nil {
t.Fatalf("write: %v", err)
}
if err := writeProfileEmail(configPath, ""); err != nil {
t.Fatalf("write empty: %v", err)
}
if got := readProfileEmail(configPath); got != email {
t.Errorf("empty write clobbered the stored email: got %q, want %q", got, email)
}
}
func ensureDirFor(t *testing.T, path string) error {
t.Helper()
return os.MkdirAll(filepath.Dir(path), 0o700)
}

View File

@@ -278,7 +278,7 @@ func (c *Client) endExtend() {
}
func (c *Client) extendAuthSession(ctx context.Context, urlOpener URLOpener, isAndroidTV bool) error {
cfg, cfgPath, cc := c.authSnapshot()
cfg, _, cc := c.stateSnapshot()
if cfg == nil || cc == nil {
return fmt.Errorf("engine is not running")
}
@@ -293,10 +293,7 @@ func (c *Client) extendAuthSession(ctx context.Context, urlOpener URLOpener, isA
}
defer authClient.Close()
// Passing the config path makes the flow pick up the login_hint: an extend
// renews the session of the account already signed in, so it must not stop to
// offer a choice.
a := NewAuthWithConfig(ctx, cfg, cfgPath)
a := &Auth{ctx: ctx, config: cfg}
tokenInfo, err := a.foregroundGetTokenInfo(authClient, urlOpener, isAndroidTV)
if err != nil {
return fmt.Errorf("interactive sso login failed: %v", err)

View File

@@ -11,10 +11,11 @@ import (
// emits.
// Metadata keys attached by the daemon to session-warning SystemEvents.
// The UI tray reads these to build a locale-aware notification without
// relying on the daemon's locale-less UserMessage string, and to
// disambiguate the T-WarningLead notification from the T-FinalWarningLead
// fallback that auto-opens the SessionAboutToExpire dialog.
// The notification text itself travels as a message key (see
// proto.UserMsgSessionExpiresIn); these keys carry the structured deadline
// the UI needs for its own countdown label, and disambiguate the
// T-WarningLead notification from the T-FinalWarningLead fallback that
// auto-opens the SessionAboutToExpire dialog.
const (
// MetaSessionWarning is set to "true" on both warning events (T-10 and
// T-2) so the UI can detect a session-warning SystemEvent without
@@ -36,10 +37,9 @@ const (
// MetaSessionDeadlineRejected is attached to the ERROR/AUTHENTICATION
// SystemEvent the daemon emits when it discards a deadline from the
// management server (pre-epoch, too far in the future, or past the
// clock-skew tolerance). The value is the rejection reason string.
// userMessage is left empty; the UI detects the event via this key
// and builds a localized notification — same pattern as the session
// warnings above.
// clock-skew tolerance). The value is the rejection reason string,
// which is diagnostic only: the user-facing text travels as
// proto.UserMsgSessionDeadlineReject.
MetaSessionDeadlineRejected = "session_deadline_rejected"
)

View File

@@ -21,6 +21,7 @@ import (
log "github.com/sirupsen/logrus"
cProto "github.com/netbirdio/netbird/client/proto"
nbstatus "github.com/netbirdio/netbird/client/status"
)
const (
@@ -80,7 +81,7 @@ type StatusRecorder interface {
severity cProto.SystemEvent_Severity,
category cProto.SystemEvent_Category,
message string,
userMessage string,
userMessage *cProto.UserMessage,
metadata map[string]string,
)
}
@@ -376,7 +377,22 @@ func publishWarning(recorder StatusRecorder, deadline time.Time, final bool) {
cProto.SystemEvent_CRITICAL,
cProto.SystemEvent_AUTHENTICATION,
message,
"",
warningUserMessage(deadline),
meta,
)
}
// warningUserMessage builds the localizable body for a session warning. The
// remaining time is rendered here rather than in the UI so every consumer of the
// event agrees on it; a deadline that is already gone (a warning delivered late)
// drops to the variant without a countdown.
func warningUserMessage(deadline time.Time) *cProto.UserMessage {
remaining := time.Until(deadline)
if remaining <= 0 {
return cProto.NewUserMessage(cProto.UserMsgSessionExpiresSoon).
WithTitle(cProto.TitleSessionWarning)
}
return cProto.NewUserMessage(cProto.UserMsgSessionExpiresIn,
cProto.ArgRemaining, nbstatus.HumaniseDuration(remaining)).
WithTitle(cProto.TitleSessionWarning)
}

View File

@@ -34,6 +34,9 @@ type event struct {
severity cProto.SystemEvent_Severity
category cProto.SystemEvent_Category
message string
msgKey cProto.UserMessageKey
titleKey cProto.UserMessageKey
msgArgs map[string]string
meta map[string]string
}
@@ -62,7 +65,7 @@ func (r *fakeRecorder) PublishEvent(
severity cProto.SystemEvent_Severity,
category cProto.SystemEvent_Category,
message string,
_ string,
userMessage *cProto.UserMessage,
metadata map[string]string,
) {
r.mu.Lock()
@@ -72,6 +75,9 @@ func (r *fakeRecorder) PublishEvent(
severity: severity,
category: category,
message: message,
msgKey: userMessage.Key(),
titleKey: userMessage.TitleKey(),
msgArgs: userMessage.Args(),
meta: metadata,
})
}
@@ -186,6 +192,33 @@ func TestWarningFiresOnceWithinLeadWindow(t *testing.T) {
}
}
// The UI localizes the warning body from the key rather than from the daemon's
// English text, so a warning that ships no key would silently regress to English.
func TestWarningCarriesLocalizableMessage(t *testing.T) {
r := &fakeRecorder{}
w := newWatcher(50*time.Millisecond, r)
defer w.Close()
_ = w.Update(time.Now().Add(80 * time.Millisecond))
events := waitForEvents(t, r, 2)
warning := events[1]
if !warning.isWarning() {
t.Fatalf("event[1] should be a warning publish, got %+v", warning)
}
if warning.msgKey != cProto.UserMsgSessionExpiresIn {
t.Errorf("warning message key = %q, want %q", warning.msgKey, cProto.UserMsgSessionExpiresIn)
}
if warning.titleKey != cProto.TitleSessionWarning {
t.Errorf("warning title key = %q, want %q", warning.titleKey, cProto.TitleSessionWarning)
}
// The remaining time is rendered at publish time so every consumer of the
// event agrees on it; the exact value depends on timer slack.
if remaining := warning.msgArgs[cProto.ArgRemaining]; remaining == "" {
t.Errorf("warning is missing the %q argument, args=%v", cProto.ArgRemaining, warning.msgArgs)
}
}
func TestWarningFiresImmediatelyWhenAlreadyInsideWindow(t *testing.T) {
r := &fakeRecorder{}
w := newWatcher(time.Hour, r) // lead > delta => fire immediately

View File

@@ -113,14 +113,11 @@ func (c *ConnectClient) RunOnAndroid(
stateFilePath string,
cacheDir string,
) error {
notifier := tunnelnotifier.New(networkChangeListener, nil)
defer notifier.Close()
// in case of non Android os these variables will be nil
mobileDependency := MobileDependency{
TunAdapter: tunAdapter,
IFaceDiscover: iFaceDiscover,
NetworkChangeListener: notifier,
NetworkChangeListener: networkChangeListener,
HostDNSAddresses: dnsAddresses,
DnsReadyListener: dnsReadyListener,
StateFilePath: stateFilePath,
@@ -166,7 +163,7 @@ func (c *ConnectClient) run(mobileDependency MobileDependency, runningChan chan
rec.PublishEvent(
cProto.SystemEvent_CRITICAL, cProto.SystemEvent_SYSTEM,
"panic occurred",
"The Netbird service panicked. Please restart the service and submit a bug report with the client logs.",
cProto.NewUserMessage(cProto.UserMsgPanic),
nil,
)
}

View File

@@ -51,5 +51,7 @@ func (n *notifier) notify() {
return
}
n.listener.OnNetworkChanged("")
go func(l listener.NetworkChangeListener) {
l.OnNetworkChanged("")
}(n.listener)
}

View File

@@ -252,7 +252,7 @@ func NewDefaultServerPermanentUpstream(
ds.hostsDNSHolder.set(hostsDnsList)
ds.permanent = true
ds.currentConfig = dnsConfigToHostDNSConfig(config, ds.service.RuntimeIP(), ds.service.RuntimePort())
ds.searchDomainNotifier = newNotifier(ds.searchDomains())
ds.searchDomainNotifier = newNotifier(ds.SearchDomains())
ds.searchDomainNotifier.setListener(listener)
setServerDns(ds)
return ds
@@ -602,12 +602,6 @@ func (s *DefaultServer) UpdateDNSServer(serial uint64, update nbdns.Config) erro
}
func (s *DefaultServer) SearchDomains() []string {
s.mux.Lock()
defer s.mux.Unlock()
return s.searchDomains()
}
func (s *DefaultServer) searchDomains() []string {
var searchDomains []string
for _, dConf := range s.currentConfig.Domains {
@@ -692,7 +686,7 @@ func (s *DefaultServer) applyConfiguration(update nbdns.Config) error {
}()
if s.searchDomainNotifier != nil {
s.searchDomainNotifier.onNewSearchDomains(s.searchDomains())
s.searchDomainNotifier.onNewSearchDomains(s.SearchDomains())
}
s.updateNSGroupStates(update.NameServerGroups)
@@ -1140,7 +1134,7 @@ func (s *DefaultServer) projectHealthy(p *nsGroupProj, servers []netip.AddrPort)
proto.SystemEvent_INFO,
proto.SystemEvent_DNS,
"Nameserver group recovered",
"DNS servers are reachable again.",
proto.NewUserMessage(proto.UserMsgDNSRecovered),
map[string]string{"upstreams": joinAddrPorts(servers)},
)
p.warningActive = false
@@ -1163,7 +1157,7 @@ func (s *DefaultServer) projectUnhealthy(p *nsGroupProj, servers []netip.AddrPor
proto.SystemEvent_WARNING,
proto.SystemEvent_DNS,
"Nameserver group unreachable",
"Unable to reach one or more DNS servers. This might affect your ability to connect to some services.",
proto.NewUserMessage(proto.UserMsgDNSUnreachable),
map[string]string{"upstreams": joinAddrPorts(servers)},
)
p.warningActive = true

View File

@@ -48,7 +48,6 @@ import (
"github.com/netbirdio/netbird/client/internal/peer"
"github.com/netbirdio/netbird/client/internal/peer/guard"
icemaker "github.com/netbirdio/netbird/client/internal/peer/ice"
"github.com/netbirdio/netbird/client/internal/peer/signaling"
"github.com/netbirdio/netbird/client/internal/peerstore"
"github.com/netbirdio/netbird/client/internal/portforward"
"github.com/netbirdio/netbird/client/internal/profilemanager"
@@ -187,7 +186,7 @@ type EngineServices struct {
type Engine struct {
// signal is a Signal Service client
signal signal.Client
signaler *signaling.Signaler
signaler *peer.Signaler
// mgmClient is a Management Service client
mgmClient mgm.Client
// peerConns is a map that holds all the peers that are known to this peer
@@ -330,7 +329,7 @@ func NewEngine(
ctx: ctx,
cancel: cancel,
signal: services.SignalClient,
signaler: signaling.NewSignaler(services.SignalClient, config.WgPrivateKey),
signaler: peer.NewSignaler(services.SignalClient, config.WgPrivateKey),
mgmClient: services.MgmClient,
relayManager: services.RelayManager,
peerStore: peerstore.NewConnStore(),
@@ -573,7 +572,12 @@ func (e *Engine) Start(netbirdConfig *mgmProto.NetbirdConfig, mgmtURL *url.URL)
}
e.stateManager.Start()
dnsServer, err := e.newDnsServer()
initialRoutes, dnsConfig, dnsFeatureFlag, err := e.readInitialSettings()
if err != nil {
return fmt.Errorf("read initial settings: %w", err)
}
dnsServer, err := e.newDnsServer(dnsConfig)
if err != nil {
return fmt.Errorf("create dns server: %w", err)
}
@@ -591,8 +595,10 @@ func (e *Engine) Start(netbirdConfig *mgmProto.NetbirdConfig, mgmtURL *url.URL)
WGInterface: e.wgInterface,
StatusRecorder: e.statusRecorder,
RelayManager: e.relayManager,
InitialRoutes: initialRoutes,
StateManager: e.stateManager,
DNSServer: dnsServer,
DNSFeatureFlag: dnsFeatureFlag,
PeerStore: e.peerStore,
DisableClientRoutes: e.config.DisableClientRoutes,
DisableServerRoutes: e.config.DisableServerRoutes,
@@ -1068,7 +1074,7 @@ func (e *Engine) handleSync(update *mgmProto.SyncResponse) error {
return err
}
e.statusRecorder.PublishEvent(cProto.SystemEvent_INFO, cProto.SystemEvent_SYSTEM, "Network map updated", "", nil)
e.statusRecorder.PublishEvent(cProto.SystemEvent_INFO, cProto.SystemEvent_SYSTEM, "Network map updated", nil, nil)
return nil
}
@@ -2096,6 +2102,42 @@ func (e *Engine) close() {
}
}
func (e *Engine) readInitialSettings() ([]*route.Route, *nbdns.Config, bool, error) {
if runtime.GOOS != "android" {
// nolint:nilnil
return nil, nil, false, nil
}
info := system.GetInfo(e.ctx)
info.SetFlags(
e.config.RosenpassEnabled,
e.config.RosenpassPermissive,
&e.config.ServerSSHAllowed,
e.config.DisableClientRoutes,
e.config.DisableServerRoutes,
e.config.DisableDNS,
e.config.DisableFirewall,
e.config.BlockLANAccess,
e.config.BlockInbound,
e.config.DisableIPv6,
e.config.SyncMessageVersion,
e.config.EnableSSHRoot,
e.config.EnableSSHSFTP,
e.config.EnableSSHLocalPortForwarding,
e.config.EnableSSHRemotePortForwarding,
e.config.DisableSSHAuth,
)
netMap, err := e.mgmClient.GetNetworkMap(info)
if err != nil {
return nil, nil, false, err
}
routes := toRoutes(netMap.GetRoutes())
dnsCfg := toDNSConfig(netMap.GetDNSConfig(), e.wgInterface.Address())
dnsFeatureFlag := toDNSFeatureFlag(netMap)
return routes, &dnsCfg, dnsFeatureFlag, nil
}
func (e *Engine) newWgIface() (*iface.WGIface, error) {
transportNet, err := e.newStdNet()
if err != nil {
@@ -2130,7 +2172,7 @@ func (e *Engine) newWgIface() (*iface.WGIface, error) {
func (e *Engine) wgInterfaceCreate() (err error) {
switch runtime.GOOS {
case "android":
err = e.wgInterface.CreateOnAndroid(e.routeManager.CurrentRouteRange(), e.dnsServer.DnsIP().String(), e.dnsServer.SearchDomains())
err = e.wgInterface.CreateOnAndroid(e.routeManager.InitialRouteRange(), e.dnsServer.DnsIP().String(), e.dnsServer.SearchDomains())
case "ios":
e.mobileDep.NetworkChangeListener.SetInterfaceIP(e.config.WgAddr.String())
if e.config.WgAddr.HasIPv6() {
@@ -2143,7 +2185,7 @@ func (e *Engine) wgInterfaceCreate() (err error) {
return err
}
func (e *Engine) newDnsServer() (dns.Server, error) {
func (e *Engine) newDnsServer(dnsConfig *nbdns.Config) (dns.Server, error) {
// due to tests where we are using a mocked version of the DNS server
if e.dnsServer != nil {
return e.dnsServer, nil
@@ -2155,7 +2197,7 @@ func (e *Engine) newDnsServer() (dns.Server, error) {
e.ctx,
e.wgInterface,
e.mobileDep.HostDNSAddresses,
nbdns.Config{},
*dnsConfig,
e.mobileDep.NetworkChangeListener,
e.statusRecorder,
e.config.DisableDNS,
@@ -2816,7 +2858,7 @@ func createFile(path string) error {
return file.Close()
}
func convertToOfferAnswer(msg *sProto.Message) (*signaling.OfferAnswer, error) {
func convertToOfferAnswer(msg *sProto.Message) (*peer.OfferAnswer, error) {
remoteCred, err := signal.UnMarshalCredential(msg)
if err != nil {
return nil, err
@@ -2832,9 +2874,9 @@ func convertToOfferAnswer(msg *sProto.Message) (*signaling.OfferAnswer, error) {
}
// Handle optional SessionID
var sessionID *icemaker.SessionID
var sessionID *peer.ICESessionID
if sessionBytes := msg.GetBody().GetSessionId(); sessionBytes != nil {
if id, err := icemaker.SessionIDFromBytes(sessionBytes); err != nil {
if id, err := peer.ICESessionIDFromBytes(sessionBytes); err != nil {
log.Warnf("Invalid session ID in message: %v", err)
sessionID = nil // Set to nil if conversion fails
} else {
@@ -2844,8 +2886,8 @@ func convertToOfferAnswer(msg *sProto.Message) (*signaling.OfferAnswer, error) {
relayIP := decodeRelayIP(msg.GetBody().GetRelayServerIP())
offerAnswer := signaling.OfferAnswer{
IceCredentials: signaling.IceCredentials{
offerAnswer := peer.OfferAnswer{
IceCredentials: peer.IceCredentials{
UFrag: remoteCred.UFrag,
Pwd: remoteCred.Pwd,
},

View File

@@ -57,7 +57,8 @@ func (e *Engine) ApplySessionDeadline(ts *timestamppb.Timestamp) {
cProto.SystemEvent_ERROR,
cProto.SystemEvent_AUTHENTICATION,
"session deadline rejected",
"",
cProto.NewUserMessage(cProto.UserMsgSessionDeadlineReject).
WithTitle(cProto.TitleSessionDeadlineReject),
map[string]string{sessionwatch.MetaSessionDeadlineRejected: err.Error()},
)
}

View File

@@ -1,20 +0,0 @@
package internal
func (e *Engine) TunSettings() ([]string, []string) {
e.syncMsgMux.Lock()
routeManager := e.routeManager
dnsServer := e.dnsServer
e.syncMsgMux.Unlock()
var routes []string
if routeManager != nil {
routes = routeManager.CurrentRouteRange()
}
var searchDomains []string
if dnsServer != nil {
searchDomains = dnsServer.SearchDomains()
}
return routes, searchDomains
}

File diff suppressed because it is too large Load Diff

View File

@@ -1,5 +1,18 @@
package peer
import (
log "github.com/sirupsen/logrus"
)
const (
// StatusIdle indicate the peer is in disconnected state
StatusIdle ConnStatus = iota
// StatusConnecting indicate the peer is in connecting state
StatusConnecting
// StatusConnected indicate the peer is in connected state
StatusConnected
)
// connStatusInputs is the primitive-valued snapshot of the state that drives the
// tri-state connection classification. Extracted so the decision logic can be unit-tested
// without constructing full Worker/Handshaker objects.
@@ -8,7 +21,24 @@ type connStatusInputs struct {
peerUsesRelay bool // remote peer advertises relay support AND local has relay
relayConnected bool // statusRelay reports Connected (independent of whether peer uses relay)
remoteSupportsICE bool // remote peer sent ICE credentials
iceWorkerCreated bool // local ICE worker exists (false in force-relay mode)
iceWorkerCreated bool // local WorkerICE exists (false in force-relay mode)
iceStatusConnecting bool // statusICE is anything other than Disconnected
iceInProgress bool // a negotiation is currently in flight
}
// ConnStatus describe the status of a peer's connection
type ConnStatus int32
func (s ConnStatus) String() string {
switch s {
case StatusConnecting:
return "Connecting"
case StatusConnected:
return "Connected"
case StatusIdle:
return "Idle"
default:
log.Errorf("unknown status: %d", s)
return "INVALID_PEER_CONNECTION_STATUS"
}
}

View File

@@ -1,4 +1,4 @@
package status
package peer
import (
"testing"

View File

@@ -3,33 +3,28 @@ package peer
import (
"context"
"fmt"
"net/netip"
"os"
"testing"
"time"
log "github.com/sirupsen/logrus"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/client/iface"
"github.com/netbirdio/netbird/client/internal/peer/dispatcher"
"github.com/netbirdio/netbird/client/internal/peer/guard"
"github.com/netbirdio/netbird/client/internal/peer/ice"
"github.com/netbirdio/netbird/client/internal/peer/metricsstages"
"github.com/netbirdio/netbird/client/internal/peer/signaling"
"github.com/netbirdio/netbird/client/internal/peer/status"
"github.com/netbirdio/netbird/client/internal/stdnet"
"github.com/netbirdio/netbird/util"
)
var testDispatcher = dispatcher.NewConnectionDispatcher()
var connConf = ConnConfig{
Key: "LLHf3Ma6z6mdLbriAJbqhX7+nM/B71lgw2+91q3LfhU=",
LocalKey: "RRHf3Ma6z6mdLbriAJbqhX7+nM/B71lgw2+91q3LfhU=",
Timeout: time.Second,
LocalWgPort: 51820,
WgConfig: WgConfig{
AllowedIps: []netip.Prefix{netip.MustParsePrefix("100.64.0.1/32")},
},
ICEConfig: ice.Config{
InterfaceBlackList: nil,
},
@@ -57,37 +52,92 @@ func TestConn_GetKey(t *testing.T) {
swWatcher := guard.NewSRWatcher(nil, nil, nil, connConf.ICEConfig)
sd := ServiceDependencies{
SrWatcher: swWatcher,
SrWatcher: swWatcher,
PeerConnDispatcher: testDispatcher,
}
conn, err := NewConn(connConf, sd)
require.NoError(t, err)
if err != nil {
return
}
got := conn.GetKey()
assert.Equal(t, got, connConf.Key, "they should be equal")
}
// TestConn_DiscardMessagesWhenNotOpened: signal messages posted to a not yet
// opened connection must be discarded without blocking or panicking.
func TestConn_DiscardMessagesWhenNotOpened(t *testing.T) {
func TestConn_OnRemoteOffer(t *testing.T) {
swWatcher := guard.NewSRWatcher(nil, nil, nil, connConf.ICEConfig)
sd := ServiceDependencies{
StatusRecorder: status.NewRecorder("https://mgm"),
SrWatcher: swWatcher,
StatusRecorder: NewRecorder("https://mgm"),
SrWatcher: swWatcher,
PeerConnDispatcher: testDispatcher,
}
conn, err := NewConn(connConf, sd)
require.NoError(t, err)
if err != nil {
return
}
offerAnswer := signaling.OfferAnswer{
IceCredentials: signaling.IceCredentials{
onNewOfferChan := make(chan struct{})
conn.handshaker.AddRelayListener(func(remoteOfferAnswer *OfferAnswer) {
onNewOfferChan <- struct{}{}
})
conn.OnRemoteOffer(OfferAnswer{
IceCredentials: IceCredentials{
UFrag: "test",
Pwd: "test",
},
WgListenPort: 0,
Version: "",
})
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
select {
case <-onNewOfferChan:
// success
case <-ctx.Done():
t.Error("expected to receive a new offer notification, but timed out")
}
}
func TestConn_OnRemoteAnswer(t *testing.T) {
swWatcher := guard.NewSRWatcher(nil, nil, nil, connConf.ICEConfig)
sd := ServiceDependencies{
StatusRecorder: NewRecorder("https://mgm"),
SrWatcher: swWatcher,
PeerConnDispatcher: testDispatcher,
}
conn, err := NewConn(connConf, sd)
if err != nil {
return
}
onNewOfferChan := make(chan struct{})
conn.handshaker.AddRelayListener(func(remoteOfferAnswer *OfferAnswer) {
onNewOfferChan <- struct{}{}
})
conn.OnRemoteAnswer(OfferAnswer{
IceCredentials: IceCredentials{
UFrag: "test",
Pwd: "test",
},
WgListenPort: 0,
Version: "",
})
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
select {
case <-onNewOfferChan:
// success
case <-ctx.Done():
t.Error("expected to receive a new offer notification, but timed out")
}
conn.OnRemoteOffer(offerAnswer)
conn.OnRemoteAnswer(offerAnswer)
conn.OnRemoteCandidate(nil, nil)
conn.Close(false)
}
func TestConn_presharedKey(t *testing.T) {
@@ -270,7 +320,7 @@ func newWGTimeoutTestConn(rosenpassEnabled bool, disconnected *[]string) *Conn {
ctx: context.Background(),
config: cfg,
Log: log.WithField("peer", cfg.Key),
metricsStages: &metricsstages.MetricsStages{},
metricsStages: &MetricsStages{},
}
conn.SetOnDisconnected(func(remotePeer string) {
*disconnected = append(*disconnected, remotePeer)
@@ -289,20 +339,20 @@ func TestConn_onWGDisconnected_EscalatesToRosenpassReset(t *testing.T) {
conn := newWGTimeoutTestConn(true, &disconnected)
for i := 0; i < wgTimeoutEscalationThreshold-1; i++ {
conn.handleWGTimeout()
conn.onWGDisconnected(conn.ctx)
}
assert.Empty(t, disconnected, "escalation must not fire below the threshold")
conn.handleWGTimeout()
conn.onWGDisconnected(conn.ctx)
assert.Equal(t, []string{conn.config.WgConfig.RemoteKey}, disconnected,
"reaching the threshold must report the peer disconnected once")
for i := 0; i < wgTimeoutEscalationThreshold-1; i++ {
conn.handleWGTimeout()
conn.onWGDisconnected(conn.ctx)
}
assert.Len(t, disconnected, 1, "escalation must restart counting after firing")
conn.handleWGTimeout()
conn.onWGDisconnected(conn.ctx)
assert.Len(t, disconnected, 2, "continued timeouts must escalate again")
}
@@ -314,12 +364,12 @@ func TestConn_onWGDisconnected_CheckSuccessResetsEscalation(t *testing.T) {
conn := newWGTimeoutTestConn(true, &disconnected)
for i := 0; i < wgTimeoutEscalationThreshold-1; i++ {
conn.handleWGTimeout()
conn.onWGDisconnected(conn.ctx)
}
conn.handleWGCheckSuccess()
conn.onWGCheckSuccess()
for i := 0; i < wgTimeoutEscalationThreshold-1; i++ {
conn.handleWGTimeout()
conn.onWGDisconnected(conn.ctx)
}
assert.Empty(t, disconnected, "handshake success must reset the timeout count")
}
@@ -327,30 +377,12 @@ func TestConn_onWGDisconnected_CheckSuccessResetsEscalation(t *testing.T) {
// TestConn_onWGDisconnected_NoEscalationWithoutRosenpass: without rosenpass
// there is no per-peer key state to reset; repeated timeouts must not report
// disconnects.
func TestConn_handleEvent_DropsStaleWGTimeout(t *testing.T) {
var disconnected []string
conn := newWGTimeoutTestConn(true, &disconnected)
staleCtx, cancel := context.WithCancel(context.Background())
cancel()
for i := 0; i < wgTimeoutEscalationThreshold; i++ {
conn.handleEvent(evWGTimeout{ctx: staleCtx})
}
assert.Empty(t, disconnected, "timeouts from a cancelled watcher must be dropped")
liveCtx := context.Background()
for i := 0; i < wgTimeoutEscalationThreshold; i++ {
conn.handleEvent(evWGTimeout{ctx: liveCtx})
}
assert.Len(t, disconnected, 1, "timeouts from the live watcher must be dispatched")
}
func TestConn_onWGDisconnected_NoEscalationWithoutRosenpass(t *testing.T) {
var disconnected []string
conn := newWGTimeoutTestConn(false, &disconnected)
for i := 0; i < wgTimeoutEscalationThreshold*3; i++ {
conn.handleWGTimeout()
conn.onWGDisconnected(conn.ctx)
}
assert.Empty(t, disconnected, "escalation must be limited to rosenpass connections")
}

View File

@@ -1,4 +1,4 @@
package worker
package conntype
import (
"fmt"

View File

@@ -0,0 +1,52 @@
package dispatcher
import (
"sync"
"github.com/netbirdio/netbird/client/internal/peer/id"
)
type ConnectionListener struct {
OnConnected func(peerID id.ConnID)
OnDisconnected func(peerID id.ConnID)
}
type ConnectionDispatcher struct {
listeners map[*ConnectionListener]struct{}
mu sync.Mutex
}
func NewConnectionDispatcher() *ConnectionDispatcher {
return &ConnectionDispatcher{
listeners: make(map[*ConnectionListener]struct{}),
}
}
func (e *ConnectionDispatcher) AddListener(listener *ConnectionListener) {
e.mu.Lock()
defer e.mu.Unlock()
e.listeners[listener] = struct{}{}
}
func (e *ConnectionDispatcher) RemoveListener(listener *ConnectionListener) {
e.mu.Lock()
defer e.mu.Unlock()
delete(e.listeners, listener)
}
func (e *ConnectionDispatcher) NotifyConnected(peerConnID id.ConnID) {
e.mu.Lock()
defer e.mu.Unlock()
for listener := range e.listeners {
listener.OnConnected(peerConnID)
}
}
func (e *ConnectionDispatcher) NotifyDisconnected(peerConnID id.ConnID) {
e.mu.Lock()
defer e.mu.Unlock()
for listener := range e.listeners {
listener.OnDisconnected(peerConnID)
}
}

View File

@@ -1,83 +0,0 @@
package peer
import (
"context"
"time"
"github.com/pion/ice/v4"
"github.com/netbirdio/netbird/client/internal/peer/signaling"
"github.com/netbirdio/netbird/client/internal/peer/worker"
"github.com/netbirdio/netbird/route"
)
// event is a message processed by the Conn event loop. All mutable Conn state
// is owned by that loop; producers deliver events through the mailbox and
// never mutate Conn state directly.
type event any
// staleableEvent is implemented by events tied to the lifetime of a transport
// component (WG watcher, ICE agent, relay connection). Each such component runs
// under its own context, cancelled when the component is superseded; an event
// carrying a cancelled context is dropped at dispatch time. A cancel performed
// by an earlier event in the same drained batch already suppresses it.
type staleableEvent interface {
isStale() bool
}
// evClose asks the event loop to tear down the connection. done is closed
// once the teardown finished.
type evClose struct {
signalToRemote bool
done chan struct{}
}
type evRemoteOffer struct {
offer signaling.OfferAnswer
}
type evRemoteAnswer struct {
answer signaling.OfferAnswer
}
type evRemoteCandidate struct {
candidate ice.Candidate
haRoutes route.HAMap
}
type evICEReady struct {
priority worker.ConnPriority
info worker.ICEConnInfo
}
type evICEDown struct {
sessionChanged bool
}
type evRelayReady struct {
info worker.RelayConnInfo
}
type evRelayDown struct{}
// evRelayDialDone reports that the relay dial helper goroutine finished,
// successfully or not, so the loop may dispatch a pending offer.
type evRelayDialDone struct{}
type evWGTimeout struct {
ctx context.Context
}
func (e evWGTimeout) isStale() bool { return e.ctx.Err() != nil }
// evWGHandshake reports the first WireGuard handshake of the current watcher run.
type evWGHandshake struct {
when time.Time
}
// evWGCheckOK reports a watcher check that observed a fresh handshake,
// including handshakes of connections that were already up.
type evWGCheckOK struct{}
// evGuardTick asks the loop to send a new offer to restore connectivity.
type evGuardTick struct{}

View File

@@ -21,6 +21,8 @@ const (
)
type ICEMonitor struct {
ReconnectCh chan struct{}
iFaceDiscover stdnet.ExternalIFaceDiscover
iceConfig icemaker.Config
tickerPeriod time.Duration
@@ -32,6 +34,7 @@ type ICEMonitor struct {
func NewICEMonitor(iFaceDiscover stdnet.ExternalIFaceDiscover, config icemaker.Config, period time.Duration) *ICEMonitor {
log.Debugf("prepare ICE monitor with period: %s", period)
cm := &ICEMonitor{
ReconnectCh: make(chan struct{}, 1),
iFaceDiscover: iFaceDiscover,
iceConfig: config,
tickerPeriod: period,

View File

@@ -0,0 +1,246 @@
package peer
import (
"context"
"errors"
"net/netip"
"sync"
"sync/atomic"
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/version"
)
var (
ErrSignalIsNotReady = errors.New("signal is not ready")
)
// IceCredentials ICE protocol credentials struct
type IceCredentials struct {
UFrag string
Pwd string
}
// OfferAnswer represents a session establishment offer or answer
type OfferAnswer struct {
IceCredentials IceCredentials
// WgListenPort is a remote WireGuard listen port.
// This field is used when establishing a direct WireGuard connection without any proxy.
// We can set the remote peer's endpoint with this port.
WgListenPort int
// Version of NetBird Agent
Version string
// RosenpassPubKey is the Rosenpass public key of the remote peer when receiving this message
// This value is the local Rosenpass server public key when sending the message
RosenpassPubKey []byte
// RosenpassAddr is the Rosenpass server address (IP:port) of the remote peer when receiving this message
// This value is the local Rosenpass server address when sending the message
RosenpassAddr string
// relay server address
RelaySrvAddress string
// RelaySrvIP is the IP the remote peer is connected to on its
// relay server. Used as a dial target if DNS for RelaySrvAddress
// fails. Zero value if the peer did not advertise an IP.
RelaySrvIP netip.Addr
// SessionID is the unique identifier of the session, used to discard old messages
SessionID *ICESessionID
}
func (o *OfferAnswer) hasICECredentials() bool {
return o.IceCredentials.UFrag != "" && o.IceCredentials.Pwd != ""
}
type Handshaker struct {
mu sync.Mutex
log *log.Entry
config ConnConfig
signaler *Signaler
ice *WorkerICE
relay *WorkerRelay
metricsStages *MetricsStages
// relayListener is not blocking because the listener is using a goroutine to process the messages
// and it will only keep the latest message if multiple offers are received in a short time
// this is to avoid blocking the handshaker if the listener is doing some heavy processing
// and also to avoid processing old offers if multiple offers are received in a short time
// the listener will always process the latest offer
relayListener *AsyncOfferListener
iceListener func(remoteOfferAnswer *OfferAnswer)
// remoteICESupported tracks whether the remote peer includes ICE credentials in its offers/answers.
// When false, the local side skips ICE listener dispatch and suppresses ICE credentials in responses.
remoteICESupported atomic.Bool
// remoteOffersCh is a channel used to wait for remote credentials to proceed with the connection
remoteOffersCh chan OfferAnswer
// remoteAnswerCh is a channel used to wait for remote credentials answer (confirmation of our offer) to proceed with the connection
remoteAnswerCh chan OfferAnswer
}
func NewHandshaker(log *log.Entry, config ConnConfig, signaler *Signaler, ice *WorkerICE, relay *WorkerRelay, metricsStages *MetricsStages) *Handshaker {
h := &Handshaker{
log: log,
config: config,
signaler: signaler,
ice: ice,
relay: relay,
metricsStages: metricsStages,
remoteOffersCh: make(chan OfferAnswer),
remoteAnswerCh: make(chan OfferAnswer),
}
// assume remote supports ICE until we learn otherwise from received offers
h.remoteICESupported.Store(ice != nil)
return h
}
func (h *Handshaker) RemoteICESupported() bool {
return h.remoteICESupported.Load()
}
func (h *Handshaker) AddRelayListener(offer func(remoteOfferAnswer *OfferAnswer)) {
h.relayListener = NewAsyncOfferListener(offer)
}
func (h *Handshaker) AddICEListener(offer func(remoteOfferAnswer *OfferAnswer)) {
h.iceListener = offer
}
func (h *Handshaker) Listen(ctx context.Context) {
for {
select {
case remoteOfferAnswer := <-h.remoteOffersCh:
h.log.Infof("received offer, running version %s, remote WireGuard listen port %d, session id: %s, remote ICE supported: %t", remoteOfferAnswer.Version, remoteOfferAnswer.WgListenPort, remoteOfferAnswer.SessionIDString(), remoteOfferAnswer.hasICECredentials())
// Record signaling received for reconnection attempts
if h.metricsStages != nil {
h.metricsStages.RecordSignalingReceived()
}
h.updateRemoteICEState(&remoteOfferAnswer)
if h.relayListener != nil {
h.relayListener.Notify(&remoteOfferAnswer)
}
if h.iceListener != nil && h.RemoteICESupported() {
h.iceListener(&remoteOfferAnswer)
}
if err := h.sendAnswer(); err != nil {
h.log.Errorf("failed to send remote offer confirmation: %s", err)
continue
}
case remoteOfferAnswer := <-h.remoteAnswerCh:
h.log.Infof("received answer, running version %s, remote WireGuard listen port %d, session id: %s, remote ICE supported: %t", remoteOfferAnswer.Version, remoteOfferAnswer.WgListenPort, remoteOfferAnswer.SessionIDString(), remoteOfferAnswer.hasICECredentials())
// Record signaling received for reconnection attempts
if h.metricsStages != nil {
h.metricsStages.RecordSignalingReceived()
}
h.updateRemoteICEState(&remoteOfferAnswer)
if h.relayListener != nil {
h.relayListener.Notify(&remoteOfferAnswer)
}
if h.iceListener != nil && h.RemoteICESupported() {
h.iceListener(&remoteOfferAnswer)
}
case <-ctx.Done():
h.log.Infof("stop listening for remote offers and answers")
return
}
}
}
func (h *Handshaker) SendOffer() error {
h.mu.Lock()
defer h.mu.Unlock()
return h.sendOffer()
}
// OnRemoteOffer handles an offer from the remote peer and returns true if the message was accepted, false otherwise
// doesn't block, discards the message if connection wasn't ready
func (h *Handshaker) OnRemoteOffer(offer OfferAnswer) {
select {
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
// doesn't block, discards the message if connection wasn't ready
func (h *Handshaker) OnRemoteAnswer(answer OfferAnswer) {
select {
case h.remoteAnswerCh <- answer:
return
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
}
}
// sendOffer prepares local user credentials and signals them to the remote peer
func (h *Handshaker) sendOffer() error {
if !h.signaler.Ready() {
return ErrSignalIsNotReady
}
offer := h.buildOfferAnswer()
h.log.Debugf("sending offer with serial: %s", offer.SessionIDString())
return h.signaler.SignalOffer(offer, h.config.Key)
}
func (h *Handshaker) sendAnswer() error {
answer := h.buildOfferAnswer()
h.log.Debugf("sending answer with serial: %s", answer.SessionIDString())
return h.signaler.SignalAnswer(answer, h.config.Key)
}
func (h *Handshaker) buildOfferAnswer() OfferAnswer {
answer := OfferAnswer{
WgListenPort: h.config.LocalWgPort,
Version: version.NetbirdVersion(),
RosenpassPubKey: h.config.RosenpassConfig.PubKey,
RosenpassAddr: h.config.RosenpassConfig.Addr,
}
if h.ice != nil && h.RemoteICESupported() {
uFrag, pwd := h.ice.GetLocalUserCredentials()
sid := h.ice.SessionID()
answer.IceCredentials = IceCredentials{uFrag, pwd}
answer.SessionID = &sid
}
if addr, ip, err := h.relay.RelayInstanceAddress(); err == nil {
answer.RelaySrvAddress = addr
answer.RelaySrvIP = ip
}
return answer
}
func (h *Handshaker) updateRemoteICEState(offer *OfferAnswer) {
hasICE := offer.hasICECredentials()
prev := h.remoteICESupported.Swap(hasICE)
if prev != hasICE {
if hasICE {
h.log.Infof("remote peer started sending ICE credentials")
} else {
h.log.Infof("remote peer stopped sending ICE credentials")
if h.ice != nil {
h.ice.Close()
}
}
}
}

View File

@@ -0,0 +1,62 @@
package peer
import (
"sync"
)
type callbackFunc func(remoteOfferAnswer *OfferAnswer)
func (oa *OfferAnswer) SessionIDString() string {
if oa.SessionID == nil {
return "unknown"
}
return oa.SessionID.String()
}
type AsyncOfferListener struct {
fn callbackFunc
running bool
latest *OfferAnswer
mu sync.Mutex
}
func NewAsyncOfferListener(fn callbackFunc) *AsyncOfferListener {
return &AsyncOfferListener{
fn: fn,
}
}
func (o *AsyncOfferListener) Notify(remoteOfferAnswer *OfferAnswer) {
o.mu.Lock()
defer o.mu.Unlock()
// Store the latest offer
o.latest = remoteOfferAnswer
// If already running, the running goroutine will pick up this latest value
if o.running {
return
}
// Start processing
o.running = true
// Process in a goroutine to avoid blocking the caller
go func(remoteOfferAnswer *OfferAnswer) {
for {
o.fn(remoteOfferAnswer)
o.mu.Lock()
if o.latest == nil {
// No more work to do
o.running = false
o.mu.Unlock()
return
}
remoteOfferAnswer = o.latest
// Clear the latest to mark it as being processed
o.latest = nil
o.mu.Unlock()
}
}(remoteOfferAnswer)
}

View File

@@ -0,0 +1,39 @@
package peer
import (
"testing"
"time"
)
func Test_newOfferListener(t *testing.T) {
dummyOfferAnswer := &OfferAnswer{}
runChan := make(chan struct{}, 10)
longRunningFn := func(remoteOfferAnswer *OfferAnswer) {
time.Sleep(1 * time.Second)
runChan <- struct{}{}
}
hl := NewAsyncOfferListener(longRunningFn)
hl.Notify(dummyOfferAnswer)
hl.Notify(dummyOfferAnswer)
hl.Notify(dummyOfferAnswer)
// Wait for exactly 2 callbacks
for i := 0; i < 2; i++ {
select {
case <-runChan:
case <-time.After(3 * time.Second):
t.Fatal("Timeout waiting for callback")
}
}
// Verify no additional callbacks happen
select {
case <-runChan:
t.Fatal("Unexpected additional callback")
case <-time.After(100 * time.Millisecond):
t.Log("Correctly received exactly 2 callbacks")
}
}

View File

@@ -0,0 +1,22 @@
package peer
import (
"net"
"net/netip"
"time"
"golang.zx2c4.com/wireguard/wgctrl/wgtypes"
"github.com/netbirdio/netbird/client/iface/configurer"
"github.com/netbirdio/netbird/client/iface/wgaddr"
"github.com/netbirdio/netbird/client/iface/wgproxy"
)
type WGIface interface {
UpdatePeer(peerKey string, allowedIps []netip.Prefix, keepAlive time.Duration, endpoint *net.UDPAddr, preSharedKey *wgtypes.Key) error
RemovePeer(peerKey string) error
GetStats() (map[string]configurer.WGStats, error)
GetProxy() wgproxy.Proxy
Address() wgaddr.Address
RemoveEndpointAddress(key string) error
}

View File

@@ -0,0 +1,11 @@
package peer
// Listener is a callback type about the NetBird network connection state
type Listener interface {
OnConnected()
OnDisconnected()
OnConnecting()
OnDisconnecting()
OnAddressChanged(string, string)
OnPeersListChanged(int)
}

View File

@@ -1,116 +0,0 @@
package peer
import (
"sync"
)
// maxQueuedCandidates bounds the remote candidate queue; on overflow the
// oldest candidate is dropped. Lost candidates are recovered by the next
// offer exchange triggered by the guard.
const maxQueuedCandidates = 128
// mailbox is the coalescing inbox of the Conn event loop. Posting never
// blocks. Per message kind either the latest value wins (offer, answer,
// guard tick), the values queue in bounded FIFO order (candidates) or in
// unbounded FIFO order (lifecycle and transport state changes, which are
// low-volume and must not be lost). A new offer flushes the queued
// candidates because they belong to the superseded session.
type mailbox struct {
mu sync.Mutex
closed bool
lifecycle []event
transport []event
offer *evRemoteOffer
answer *evRemoteAnswer
candidates []evRemoteCandidate
guardTick bool
wake chan struct{}
}
func newMailbox() *mailbox {
return &mailbox{
wake: make(chan struct{}, 1),
}
}
// post stores the event and wakes the loop. It reports false if the mailbox
// is already closed and the event was not accepted.
func (m *mailbox) post(ev event) bool {
m.mu.Lock()
if m.closed {
m.mu.Unlock()
return false
}
switch e := ev.(type) {
case evClose:
m.lifecycle = append(m.lifecycle, e)
case evRemoteOffer:
m.offer = &e
m.candidates = nil
case evRemoteAnswer:
m.answer = &e
case evRemoteCandidate:
if len(m.candidates) >= maxQueuedCandidates {
m.candidates = m.candidates[1:]
}
m.candidates = append(m.candidates, e)
case evGuardTick:
m.guardTick = true
default:
m.transport = append(m.transport, ev)
}
m.mu.Unlock()
select {
case m.wake <- struct{}{}:
default:
}
return true
}
// drain returns the pending events in processing order: lifecycle first,
// then transport state changes, the coalesced offer and answer, the queued
// candidates and finally the guard tick.
func (m *mailbox) drain() []event {
m.mu.Lock()
defer m.mu.Unlock()
return m.drainLocked()
}
// closeAndDrain marks the mailbox closed so further posts are rejected and
// returns the events that were still pending.
func (m *mailbox) closeAndDrain() []event {
m.mu.Lock()
defer m.mu.Unlock()
m.closed = true
return m.drainLocked()
}
func (m *mailbox) drainLocked() []event {
evs := make([]event, 0, len(m.lifecycle)+len(m.transport)+len(m.candidates)+3)
evs = append(evs, m.lifecycle...)
evs = append(evs, m.transport...)
if m.offer != nil {
evs = append(evs, *m.offer)
}
if m.answer != nil {
evs = append(evs, *m.answer)
}
for _, c := range m.candidates {
evs = append(evs, c)
}
if m.guardTick {
evs = append(evs, evGuardTick{})
}
m.lifecycle = nil
m.transport = nil
m.offer = nil
m.answer = nil
m.candidates = nil
m.guardTick = false
return evs
}

View File

@@ -1,128 +0,0 @@
package peer
import (
"testing"
"github.com/netbirdio/netbird/client/internal/peer/signaling"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
func TestMailbox_OfferCoalescing(t *testing.T) {
mb := newMailbox()
require.True(t, mb.post(evRemoteOffer{offer: signaling.OfferAnswer{WgListenPort: 1}}))
require.True(t, mb.post(evRemoteOffer{offer: signaling.OfferAnswer{WgListenPort: 2}}))
require.True(t, mb.post(evRemoteOffer{offer: signaling.OfferAnswer{WgListenPort: 3}}))
evs := mb.drain()
require.Len(t, evs, 1, "consecutive offers must coalesce to a single event")
offer, ok := evs[0].(evRemoteOffer)
require.True(t, ok, "coalesced event must be an offer")
assert.Equal(t, 3, offer.offer.WgListenPort, "the newest offer must win")
}
func TestMailbox_OfferFlushesCandidates(t *testing.T) {
mb := newMailbox()
require.True(t, mb.post(evRemoteCandidate{}))
require.True(t, mb.post(evRemoteCandidate{}))
require.True(t, mb.post(evRemoteOffer{offer: signaling.OfferAnswer{}}))
evs := mb.drain()
require.Len(t, evs, 1, "candidates of the superseded session must be flushed")
_, ok := evs[0].(evRemoteOffer)
assert.True(t, ok, "only the offer must remain after the flush")
}
func TestMailbox_CandidatesKeepOrderAfterOffer(t *testing.T) {
mb := newMailbox()
require.True(t, mb.post(evRemoteOffer{offer: signaling.OfferAnswer{}}))
require.True(t, mb.post(evRemoteCandidate{haRoutes: nil}))
require.True(t, mb.post(evRemoteCandidate{haRoutes: nil}))
evs := mb.drain()
require.Len(t, evs, 3)
_, ok := evs[0].(evRemoteOffer)
assert.True(t, ok, "offer must be processed before the candidates")
for _, ev := range evs[1:] {
_, ok := ev.(evRemoteCandidate)
assert.True(t, ok, "candidates posted after the offer must survive")
}
}
func TestMailbox_CandidateQueueBounded(t *testing.T) {
mb := newMailbox()
for i := 0; i < maxQueuedCandidates+10; i++ {
require.True(t, mb.post(evRemoteCandidate{}))
}
evs := mb.drain()
assert.Len(t, evs, maxQueuedCandidates, "candidate queue must stay bounded")
}
func TestMailbox_DrainOrder(t *testing.T) {
mb := newMailbox()
require.True(t, mb.post(evGuardTick{}))
require.True(t, mb.post(evRemoteAnswer{answer: signaling.OfferAnswer{}}))
require.True(t, mb.post(evRemoteOffer{offer: signaling.OfferAnswer{}}))
require.True(t, mb.post(evRelayDown{}))
require.True(t, mb.post(evICEDown{sessionChanged: true}))
require.True(t, mb.post(evClose{}))
evs := mb.drain()
require.Len(t, evs, 6)
_, ok := evs[0].(evClose)
assert.True(t, ok, "lifecycle events must come first")
_, ok = evs[1].(evRelayDown)
assert.True(t, ok, "transport events must keep FIFO order")
_, ok = evs[2].(evICEDown)
assert.True(t, ok, "transport events must keep FIFO order")
_, ok = evs[3].(evRemoteOffer)
assert.True(t, ok, "offer must come after transport events")
_, ok = evs[4].(evRemoteAnswer)
assert.True(t, ok, "answer must come after the offer")
_, ok = evs[5].(evGuardTick)
assert.True(t, ok, "guard tick must come last")
}
func TestMailbox_GuardTickCoalesced(t *testing.T) {
mb := newMailbox()
require.True(t, mb.post(evGuardTick{}))
require.True(t, mb.post(evGuardTick{}))
require.True(t, mb.post(evGuardTick{}))
evs := mb.drain()
assert.Len(t, evs, 1, "guard ticks must coalesce to a single event")
}
func TestMailbox_PostAfterCloseRejected(t *testing.T) {
mb := newMailbox()
require.True(t, mb.post(evRelayDown{}))
leftovers := mb.closeAndDrain()
assert.Len(t, leftovers, 1, "pending events must be returned on close")
assert.False(t, mb.post(evRelayDown{}), "posts must be rejected after close")
assert.Empty(t, mb.drain(), "no events must remain after close")
}
func TestMailbox_WakeSignal(t *testing.T) {
mb := newMailbox()
require.True(t, mb.post(evRelayDown{}))
require.True(t, mb.post(evGuardTick{}))
select {
case <-mb.wake:
default:
t.Fatal("wake signal must be pending after posts")
}
assert.Len(t, mb.drain(), 2, "a single wake must deliver all pending events")
}

View File

@@ -1,4 +1,4 @@
package metricsstages
package peer
import (
"sync"

View File

@@ -1,4 +1,4 @@
package metricsstages
package peer
import (
"testing"

View File

@@ -1,4 +1,4 @@
package status
package peer
import (
"sync"
@@ -11,16 +11,6 @@ const (
stateDisconnecting
)
// Listener is a callback type about the NetBird network connection state
type Listener interface {
OnConnected()
OnDisconnected()
OnConnecting()
OnDisconnecting()
OnAddressChanged(string, string)
OnPeersListChanged(int)
}
type notifier struct {
serverStateLock sync.Mutex
listenersLock sync.Mutex

View File

@@ -1,4 +1,4 @@
package status
package peer
import (
"sync"

View File

@@ -1,4 +1,4 @@
package status
package peer
import (
"net/netip"

View File

@@ -1,4 +1,4 @@
package ice
package peer
import (
"crypto/rand"
@@ -9,26 +9,26 @@ import (
const sessionIDSize = 5
type SessionID string
type ICESessionID string
// NewSessionID generates a new session ID for distinguishing sessions
func NewSessionID() (SessionID, error) {
// NewICESessionID generates a new session ID for distinguishing sessions
func NewICESessionID() (ICESessionID, error) {
b := make([]byte, sessionIDSize)
if _, err := io.ReadFull(rand.Reader, b); err != nil {
return "", fmt.Errorf("failed to generate session ID: %w", err)
}
return SessionID(hex.EncodeToString(b)), nil
return ICESessionID(hex.EncodeToString(b)), nil
}
func SessionIDFromBytes(b []byte) (SessionID, error) {
func ICESessionIDFromBytes(b []byte) (ICESessionID, error) {
if len(b) != sessionIDSize {
return "", fmt.Errorf("invalid session ID length: %d", len(b))
}
return SessionID(hex.EncodeToString(b)), nil
return ICESessionID(hex.EncodeToString(b)), nil
}
// Bytes returns the raw bytes of the session ID for protobuf serialization
func (id SessionID) Bytes() ([]byte, error) {
func (id ICESessionID) Bytes() ([]byte, error) {
if len(id) == 0 {
return nil, fmt.Errorf("ICE session ID is empty")
}
@@ -42,6 +42,6 @@ func (id SessionID) Bytes() ([]byte, error) {
return b, nil
}
func (id SessionID) String() string {
func (id ICESessionID) String() string {
return string(id)
}

View File

@@ -1,4 +1,4 @@
package signaling
package peer
import (
"github.com/pion/ice/v4"

View File

@@ -1,189 +0,0 @@
package signaling
import (
"errors"
"net/netip"
"sync"
"sync/atomic"
log "github.com/sirupsen/logrus"
icemaker "github.com/netbirdio/netbird/client/internal/peer/ice"
relayClient "github.com/netbirdio/netbird/shared/relay/client"
"github.com/netbirdio/netbird/version"
)
var (
ErrSignalIsNotReady = errors.New("signal is not ready")
)
// IceCredentials ICE protocol credentials struct
type IceCredentials struct {
UFrag string
Pwd string
}
// OfferAnswer represents a session establishment offer or answer
type OfferAnswer struct {
IceCredentials IceCredentials
// WgListenPort is a remote WireGuard listen port.
// This field is used when establishing a direct WireGuard connection without any proxy.
// We can set the remote peer's endpoint with this port.
WgListenPort int
// Version of NetBird Agent
Version string
// RosenpassPubKey is the Rosenpass public key of the remote peer when receiving this message
// This value is the local Rosenpass server public key when sending the message
RosenpassPubKey []byte
// RosenpassAddr is the Rosenpass server address (IP:port) of the remote peer when receiving this message
// This value is the local Rosenpass server address when sending the message
RosenpassAddr string
// relay server address
RelaySrvAddress string
// RelaySrvIP is the IP the remote peer is connected to on its
// relay server. Used as a dial target if DNS for RelaySrvAddress
// fails. Zero value if the peer did not advertise an IP.
RelaySrvIP netip.Addr
// SessionID is the unique identifier of the session, used to discard old messages
SessionID *icemaker.SessionID
}
func (o *OfferAnswer) HasICECredentials() bool {
return o.IceCredentials.UFrag != "" && o.IceCredentials.Pwd != ""
}
func (o *OfferAnswer) SessionIDString() string {
if o.SessionID == nil {
return "unknown"
}
return o.SessionID.String()
}
// Config carries the peer-specific values the Handshaker embeds into offers
// and answers.
type Config struct {
Key string
LocalWgPort int
RosenpassPubKey []byte
RosenpassAddr string
}
// Credentials are the local ICE credentials and session id the Handshaker embeds in offers.
type Credentials struct {
UFrag string
Pwd string
SessionID icemaker.SessionID
}
// ICEWorker is the subset of the ICE worker the Handshaker needs to build offers.
type ICEWorker interface {
Credentials() Credentials
Close()
}
// Handshaker keeps the signaling protocol logic: building and sending offers
// and answers and tracking whether the remote peer supports ICE. Incoming
// message processing is driven by the Conn event loop.
type Handshaker struct {
mu sync.Mutex
log *log.Entry
config Config
signaler *Signaler
ice ICEWorker
relayManager *relayClient.Manager
// remoteICESupported tracks whether the remote peer includes ICE credentials in its offers/answers.
// When false, the local side skips ICE dispatch and suppresses ICE credentials in responses.
remoteICESupported atomic.Bool
}
func NewHandshaker(log *log.Entry, config Config, signaler *Signaler, ice ICEWorker, relayManager *relayClient.Manager) *Handshaker {
h := &Handshaker{
log: log,
config: config,
signaler: signaler,
ice: ice,
relayManager: relayManager,
}
// assume remote supports ICE until we learn otherwise from received offers
h.remoteICESupported.Store(ice != nil)
return h
}
func (h *Handshaker) RemoteICESupported() bool {
return h.remoteICESupported.Load()
}
func (h *Handshaker) SendOffer() error {
h.mu.Lock()
defer h.mu.Unlock()
return h.sendOffer()
}
func (h *Handshaker) SendAnswer() error {
h.mu.Lock()
defer h.mu.Unlock()
return h.sendAnswer()
}
// sendOffer prepares local user credentials and signals them to the remote peer
func (h *Handshaker) sendOffer() error {
if !h.signaler.Ready() {
return ErrSignalIsNotReady
}
offer := h.buildOfferAnswer()
h.log.Debugf("sending offer with serial: %s", offer.SessionIDString())
return h.signaler.SignalOffer(offer, h.config.Key)
}
func (h *Handshaker) sendAnswer() error {
answer := h.buildOfferAnswer()
h.log.Debugf("sending answer with serial: %s", answer.SessionIDString())
return h.signaler.SignalAnswer(answer, h.config.Key)
}
func (h *Handshaker) buildOfferAnswer() OfferAnswer {
answer := OfferAnswer{
WgListenPort: h.config.LocalWgPort,
Version: version.NetbirdVersion(),
RosenpassPubKey: h.config.RosenpassPubKey,
RosenpassAddr: h.config.RosenpassAddr,
}
if h.ice != nil && h.RemoteICESupported() {
creds := h.ice.Credentials()
answer.IceCredentials = IceCredentials{creds.UFrag, creds.Pwd}
sid := creds.SessionID
answer.SessionID = &sid
}
if addr, ip, err := h.relayManager.RelayInstanceAddress(); err == nil {
answer.RelaySrvAddress = addr
answer.RelaySrvIP = ip
}
return answer
}
// UpdateRemoteICEState refreshes the remote ICE support flag from a received
// offer or answer and closes the ICE worker when the remote peer stopped
// sending ICE credentials. Runs on the Conn event loop.
func (h *Handshaker) UpdateRemoteICEState(offer *OfferAnswer) {
hasICE := offer.HasICECredentials()
prev := h.remoteICESupported.Swap(hasICE)
if prev != hasICE {
if hasICE {
h.log.Infof("remote peer started sending ICE credentials")
} else {
h.log.Infof("remote peer stopped sending ICE credentials")
if h.ice != nil {
h.ice.Close()
}
}
}
}

View File

@@ -1,4 +1,4 @@
package state_dump
package peer
import (
"context"
@@ -6,13 +6,11 @@ import (
"time"
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/internal/peer/status"
)
type StateDump struct {
type stateDump struct {
log *log.Entry
status *status.Recorder
status *Status
key string
sentOffer int
@@ -28,15 +26,15 @@ type StateDump struct {
mu sync.Mutex
}
func NewStateDump(key string, log *log.Entry, statusRecorder *status.Recorder) *StateDump {
return &StateDump{
func newStateDump(key string, log *log.Entry, statusRecorder *Status) *stateDump {
return &stateDump{
log: log,
status: statusRecorder,
key: key,
}
}
func (s *StateDump) Start(ctx context.Context) {
func (s *stateDump) Start(ctx context.Context) {
ticker := time.NewTicker(10 * time.Minute)
defer ticker.Stop()
@@ -50,25 +48,25 @@ func (s *StateDump) Start(ctx context.Context) {
}
}
func (s *StateDump) RemoteOffer() {
func (s *stateDump) RemoteOffer() {
s.mu.Lock()
defer s.mu.Unlock()
s.remoteOffer++
}
func (s *StateDump) RemoteCandidate() {
func (s *stateDump) RemoteCandidate() {
s.mu.Lock()
defer s.mu.Unlock()
s.remoteCandidate++
}
func (s *StateDump) SendOffer() {
func (s *stateDump) SendOffer() {
s.mu.Lock()
defer s.mu.Unlock()
s.sentOffer++
}
func (s *StateDump) dumpState() {
func (s *stateDump) dumpState() {
s.mu.Lock()
defer s.mu.Unlock()
@@ -82,41 +80,41 @@ func (s *StateDump) dumpState() {
status, s.sentOffer, s.remoteOffer, s.remoteAnswer, s.remoteCandidate, s.p2pConnected, s.switchToRelay, s.wgCheckSuccess, s.relayConnected, s.localProxies)
}
func (s *StateDump) RemoteAnswer() {
func (s *stateDump) RemoteAnswer() {
s.mu.Lock()
defer s.mu.Unlock()
s.remoteAnswer++
}
func (s *StateDump) P2PConnected() {
func (s *stateDump) P2PConnected() {
s.mu.Lock()
defer s.mu.Unlock()
s.p2pConnected++
}
func (s *StateDump) SwitchToRelay() {
func (s *stateDump) SwitchToRelay() {
s.mu.Lock()
defer s.mu.Unlock()
s.switchToRelay++
}
func (s *StateDump) WGcheckSuccess() {
func (s *stateDump) WGcheckSuccess() {
s.mu.Lock()
defer s.mu.Unlock()
s.wgCheckSuccess++
}
func (s *StateDump) RelayConnected() {
func (s *stateDump) RelayConnected() {
s.mu.Lock()
defer s.mu.Unlock()
s.relayConnected++
}
func (s *StateDump) NewLocalProxy() {
func (s *stateDump) NewLocalProxy() {
s.mu.Lock()
defer s.mu.Unlock()

View File

@@ -1,31 +0,0 @@
package status
import (
log "github.com/sirupsen/logrus"
)
const (
// StatusIdle indicate the peer is in disconnected state
StatusIdle ConnStatus = iota
// StatusConnecting indicate the peer is in connecting state
StatusConnecting
// StatusConnected indicate the peer is in connected state
StatusConnected
)
// ConnStatus describe the status of a peer's connection
type ConnStatus int32
func (s ConnStatus) String() string {
switch s {
case StatusConnecting:
return "Connecting"
case StatusConnected:
return "Connected"
case StatusIdle:
return "Idle"
default:
log.Errorf("unknown status: %d", s)
return "INVALID_PEER_CONNECTION_STATUS"
}
}

View File

@@ -1,48 +0,0 @@
package status
import (
"slices"
"sync"
"github.com/netbirdio/netbird/client/proto"
)
type EventQueue struct {
maxSize int
events []*proto.SystemEvent
mutex sync.RWMutex
}
func NewEventQueue(size int) *EventQueue {
return &EventQueue{
maxSize: size,
events: make([]*proto.SystemEvent, 0, size),
}
}
func (q *EventQueue) Add(event *proto.SystemEvent) {
q.mutex.Lock()
defer q.mutex.Unlock()
q.events = append(q.events, event)
if len(q.events) > q.maxSize {
q.events = q.events[len(q.events)-q.maxSize:]
}
}
func (q *EventQueue) GetAll() []*proto.SystemEvent {
q.mutex.RLock()
defer q.mutex.RUnlock()
return slices.Clone(q.events)
}
type EventSubscription struct {
id string
events chan *proto.SystemEvent
}
func (s *EventSubscription) Events() <-chan *proto.SystemEvent {
return s.events
}

View File

@@ -1,122 +0,0 @@
package status
import (
"golang.org/x/exp/maps"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
"github.com/netbirdio/netbird/client/internal/relay"
"github.com/netbirdio/netbird/client/proto"
)
// FullStatus contains the full state held by the Recorder instance
type FullStatus struct {
Peers []State
ManagementState ManagementState
SignalState SignalState
LocalPeerState LocalPeerState
RosenpassState RosenpassState
Relays []relay.ProbeResult
NSGroupStates []NSGroupState
NumOfForwardingRules int
LazyConnectionEnabled bool
Events []*proto.SystemEvent
}
// ToProto converts FullStatus to proto.FullStatus.
func (fs FullStatus) ToProto() *proto.FullStatus {
pbFullStatus := proto.FullStatus{
ManagementState: &proto.ManagementState{},
SignalState: &proto.SignalState{},
LocalPeerState: &proto.LocalPeerState{},
Peers: []*proto.PeerState{},
}
pbFullStatus.ManagementState.URL = fs.ManagementState.URL
pbFullStatus.ManagementState.Connected = fs.ManagementState.Connected
if err := fs.ManagementState.Error; err != nil {
pbFullStatus.ManagementState.Error = err.Error()
}
pbFullStatus.SignalState.URL = fs.SignalState.URL
pbFullStatus.SignalState.Connected = fs.SignalState.Connected
if err := fs.SignalState.Error; err != nil {
pbFullStatus.SignalState.Error = err.Error()
}
pbFullStatus.LocalPeerState.IP = fs.LocalPeerState.IP
pbFullStatus.LocalPeerState.Ipv6 = fs.LocalPeerState.IPv6
pbFullStatus.LocalPeerState.PubKey = fs.LocalPeerState.PubKey
pbFullStatus.LocalPeerState.KernelInterface = fs.LocalPeerState.KernelInterface
pbFullStatus.LocalPeerState.Fqdn = fs.LocalPeerState.FQDN
pbFullStatus.LocalPeerState.WgPort = int32(fs.LocalPeerState.WgPort)
pbFullStatus.LocalPeerState.RosenpassPermissive = fs.RosenpassState.Permissive
pbFullStatus.LocalPeerState.RosenpassEnabled = fs.RosenpassState.Enabled
pbFullStatus.NumberOfForwardingRules = int32(fs.NumOfForwardingRules)
pbFullStatus.LazyConnectionEnabled = fs.LazyConnectionEnabled
pbFullStatus.LocalPeerState.Networks = maps.Keys(fs.LocalPeerState.Routes)
for _, peerState := range fs.Peers {
networks := maps.Keys(peerState.GetRoutes())
pbPeerState := &proto.PeerState{
IP: peerState.IP,
Ipv6: peerState.IPv6,
PubKey: peerState.PubKey,
ConnStatus: peerState.ConnStatus.String(),
ConnStatusUpdate: timestamppb.New(peerState.ConnStatusUpdate),
Relayed: peerState.Relayed,
LocalIceCandidateType: peerState.LocalIceCandidateType,
RemoteIceCandidateType: peerState.RemoteIceCandidateType,
LocalIceCandidateEndpoint: peerState.LocalIceCandidateEndpoint,
RemoteIceCandidateEndpoint: peerState.RemoteIceCandidateEndpoint,
RelayAddress: peerState.RelayServerAddress,
Fqdn: peerState.FQDN,
LastWireguardHandshake: timestamppb.New(peerState.LastWireguardHandshake),
BytesRx: peerState.BytesRx,
BytesTx: peerState.BytesTx,
RosenpassEnabled: peerState.RosenpassEnabled,
Networks: networks,
Latency: durationpb.New(peerState.Latency),
SshHostKey: peerState.SSHHostKey,
}
pbFullStatus.Peers = append(pbFullStatus.Peers, pbPeerState)
}
for _, relayState := range fs.Relays {
pbRelayState := &proto.RelayState{
URI: relayState.URI,
Available: relayState.Err == nil,
Transport: relayState.Transport,
}
if err := relayState.Err; err != nil {
pbRelayState.Error = err.Error()
}
pbFullStatus.Relays = append(pbFullStatus.Relays, pbRelayState)
}
for _, dnsState := range fs.NSGroupStates {
var err string
if dnsState.Error != nil {
err = dnsState.Error.Error()
}
var servers []string
for _, server := range dnsState.Servers {
servers = append(servers, server.String())
}
pbDnsState := &proto.NSGroupState{
Servers: servers,
Domains: dnsState.Domains,
Enabled: dnsState.Enabled,
Error: err,
}
pbFullStatus.DnsServers = append(pbFullStatus.DnsServers, pbDnsState)
}
pbFullStatus.Events = fs.Events
return &pbFullStatus
}

View File

@@ -1,63 +0,0 @@
package status
import (
"sync"
"time"
"golang.org/x/exp/maps"
)
// State contains the latest state of a peer
type State struct {
Mux *sync.RWMutex
IP string
IPv6 string
PubKey string
FQDN string
ConnStatus ConnStatus
ConnStatusUpdate time.Time
Relayed bool
LocalIceCandidateType string
RemoteIceCandidateType string
LocalIceCandidateEndpoint string
RemoteIceCandidateEndpoint string
RelayServerAddress string
LastWireguardHandshake time.Time
BytesTx int64
BytesRx int64
Latency time.Duration
RosenpassEnabled bool
SSHHostKey []byte
routes map[string]struct{}
}
// AddRoute add a single route to routes map
func (s *State) AddRoute(network string) {
s.Mux.Lock()
defer s.Mux.Unlock()
if s.routes == nil {
s.routes = make(map[string]struct{})
}
s.routes[network] = struct{}{}
}
// SetRoutes set state routes
func (s *State) SetRoutes(routes map[string]struct{}) {
s.Mux.Lock()
defer s.Mux.Unlock()
s.routes = routes
}
// DeleteRoute removes a route from the network amp
func (s *State) DeleteRoute(network string) {
s.Mux.Lock()
defer s.Mux.Unlock()
delete(s.routes, network)
}
// GetRoutes return routes map
func (s *State) GetRoutes() map[string]struct{} {
s.Mux.RLock()
defer s.Mux.RUnlock()
return maps.Clone(s.routes)
}

View File

@@ -1,36 +0,0 @@
package peer
import "github.com/netbirdio/netbird/client/internal/peer/status"
// Transitional aliases re-exporting the peer status recorder from its own
// package. Callers are being migrated to reference the status package
// directly; these aliases will be removed once the migration completes.
type (
Status = status.Recorder
State = status.State
ConnStatus = status.ConnStatus
FullStatus = status.FullStatus
RouterState = status.RouterState
LocalPeerState = status.LocalPeerState
SignalState = status.SignalState
ManagementState = status.ManagementState
RosenpassState = status.RosenpassState
NSGroupState = status.NSGroupState
ResolvedDomainInfo = status.ResolvedDomainInfo
StatusChangeSubscription = status.StatusChangeSubscription
EventQueue = status.EventQueue
EventSubscription = status.EventSubscription
WGIfaceStatus = status.WGIfaceStatus
Listener = status.Listener
EventListener = status.EventListener
)
const (
StatusIdle = status.StatusIdle
StatusConnecting = status.StatusConnecting
StatusConnected = status.StatusConnected
)
var (
NewRecorder = status.NewRecorder
)

View File

@@ -1,4 +1,4 @@
package status
package peer
import (
"context"

View File

@@ -1,4 +1,4 @@
package wg_watcher
package peer
import (
"context"
@@ -8,7 +8,6 @@ import (
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/iface/configurer"
"github.com/netbirdio/netbird/client/internal/peer/state_dump"
)
const (
@@ -30,7 +29,7 @@ type WGWatcher struct {
log *log.Entry
wgIfaceStater WGInterfaceStater
peerKey string
stateDump *state_dump.StateDump
stateDump *stateDump
// initialHandshake is not thread-safe; never call PrepareInitialHandshake and EnableWgWatcher concurrently.
initialHandshake time.Time
@@ -38,7 +37,7 @@ type WGWatcher struct {
resetCh chan struct{}
}
func NewWGWatcher(log *log.Entry, wgIfaceStater WGInterfaceStater, peerKey string, stateDump *state_dump.StateDump) *WGWatcher {
func NewWGWatcher(log *log.Entry, wgIfaceStater WGInterfaceStater, peerKey string, stateDump *stateDump) *WGWatcher {
return &WGWatcher{
log: log,
wgIfaceStater: wgIfaceStater,

View File

@@ -1,4 +1,4 @@
package wg_watcher
package peer
import (
"context"
@@ -9,8 +9,6 @@ import (
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/iface/configurer"
"github.com/netbirdio/netbird/client/internal/peer/state_dump"
"github.com/netbirdio/netbird/client/internal/peer/status"
)
type MocWgIface struct {
@@ -58,7 +56,7 @@ func TestWGWatcher_CheckSuccessCallback(t *testing.T) {
// platforms with coarse clock resolution (Windows), where two time.Now() calls
// microseconds apart can return the same instant and read as a timed-out handshake.
stats := &mockHandshakeStats{handshake: time.Now().Add(-time.Hour)}
watcher := NewWGWatcher(mlog, stats, "", state_dump.NewStateDump("peer", mlog, &status.Recorder{}))
watcher := NewWGWatcher(mlog, stats, "", newStateDump("peer", mlog, &Status{}))
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
@@ -67,18 +65,14 @@ func TestWGWatcher_CheckSuccessCallback(t *testing.T) {
firstHandshake := make(chan struct{}, 1)
checkSuccess := make(chan struct{}, 1)
watcherDone := make(chan struct{})
go func() {
defer close(watcherDone)
watcher.EnableWgWatcher(ctx, time.Now(), func() {}, func(when time.Time) {
firstHandshake <- struct{}{}
}, func() {
select {
case checkSuccess <- struct{}{}:
default:
}
})
}()
go watcher.EnableWgWatcher(ctx, time.Now(), func() {}, func(when time.Time) {
firstHandshake <- struct{}{}
}, func() {
select {
case checkSuccess <- struct{}{}:
default:
}
})
stats.advance()
@@ -93,11 +87,6 @@ func TestWGWatcher_CheckSuccessCallback(t *testing.T) {
t.Errorf("first-handshake callback must not fire for a non-zero baseline")
default:
}
// Wait for the watcher goroutine to exit so it cannot race with other
// tests mutating the package-level check timing variables.
cancel()
<-watcherDone
}
func TestWGWatcher_EnableWgWatcher(t *testing.T) {
@@ -106,7 +95,7 @@ func TestWGWatcher_EnableWgWatcher(t *testing.T) {
mlog := log.WithField("peer", "tet")
mocWgIface := &MocWgIface{}
watcher := NewWGWatcher(mlog, mocWgIface, "", state_dump.NewStateDump("peer", mlog, &status.Recorder{}))
watcher := NewWGWatcher(mlog, mocWgIface, "", newStateDump("peer", mlog, &Status{}))
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
@@ -138,7 +127,7 @@ func TestWGWatcher_ReEnable(t *testing.T) {
mlog := log.WithField("peer", "tet")
mocWgIface := &MocWgIface{}
watcher := NewWGWatcher(mlog, mocWgIface, "", state_dump.NewStateDump("peer", mlog, &status.Recorder{}))
watcher := NewWGWatcher(mlog, mocWgIface, "", newStateDump("peer", mlog, &Status{}))
ctx, cancel := context.WithCancel(context.Background())
watcher.PrepareInitialHandshake()

View File

@@ -1,4 +1,4 @@
package peer
package worker
import (
"sync/atomic"
@@ -7,17 +7,17 @@ import (
)
const (
WorkerStatusDisconnected WorkerStatus = iota
WorkerStatusConnected
StatusDisconnected Status = iota
StatusConnected
)
type WorkerStatus int32
type Status int32
func (s WorkerStatus) String() string {
func (s Status) String() string {
switch s {
case WorkerStatusDisconnected:
case StatusDisconnected:
return "Disconnected"
case WorkerStatusConnected:
case StatusConnected:
return "Connected"
default:
log.Errorf("unknown status: %d", s)
@@ -37,16 +37,16 @@ func NewAtomicStatus() *AtomicWorkerStatus {
}
// Get returns the current connection status
func (acs *AtomicWorkerStatus) Get() WorkerStatus {
return WorkerStatus(acs.status.Load())
func (acs *AtomicWorkerStatus) Get() Status {
return Status(acs.status.Load())
}
func (acs *AtomicWorkerStatus) SetConnected() {
acs.status.Store(int32(WorkerStatusConnected))
acs.status.Store(int32(StatusConnected))
}
func (acs *AtomicWorkerStatus) SetDisconnected() {
acs.status.Store(int32(WorkerStatusDisconnected))
acs.status.Store(int32(StatusDisconnected))
}
// String returns the string representation of the current status

View File

@@ -1,4 +1,4 @@
package worker
package peer
import (
"context"
@@ -13,9 +13,8 @@ import (
"github.com/netbirdio/netbird/client/iface"
"github.com/netbirdio/netbird/client/iface/udpmux"
"github.com/netbirdio/netbird/client/internal/peer/conntype"
icemaker "github.com/netbirdio/netbird/client/internal/peer/ice"
"github.com/netbirdio/netbird/client/internal/peer/signaling"
"github.com/netbirdio/netbird/client/internal/peer/status"
"github.com/netbirdio/netbird/client/internal/portforward"
"github.com/netbirdio/netbird/client/internal/stdnet"
"github.com/netbirdio/netbird/route"
@@ -33,68 +32,57 @@ type ICEConnInfo struct {
RelayedOnLocal bool
}
type ICEDependencies struct {
Signaler *signaling.Signaler
IFaceDiscover stdnet.ExternalIFaceDiscover
StatusRecorder *status.Recorder
PortForwardManager *portforward.Manager
}
type ICE struct {
log *log.Entry
key string
iceConfig icemaker.Config
isController bool
onConnReady func(priority ConnPriority, iceConnInfo ICEConnInfo)
onStatusDisconnect func(sessionChanged bool)
signaler *signaling.Signaler
iFaceDiscover stdnet.ExternalIFaceDiscover
statusRecorder *status.Recorder
portForwardManager *portforward.Manager
hasRelayOnLocally bool
type WorkerICE struct {
ctx context.Context
log *log.Entry
config ConnConfig
conn *Conn
signaler *Signaler
iFaceDiscover stdnet.ExternalIFaceDiscover
statusRecorder *Status
hasRelayOnLocally bool
agent *icemaker.ThreadSafeAgent
agentDialerCancel context.CancelFunc
agentConnecting bool // while it is true, drop all incoming offers
lastSuccess time.Time // with this avoid the too frequent ICE agent recreation
// connectedAgent is the agent whose connection was last reported ready; guarded by muxAgent
connectedAgent *icemaker.ThreadSafeAgent
// remoteSessionID represents the peer's session identifier from the latest remote offer.
remoteSessionID icemaker.SessionID
remoteSessionID ICESessionID
// sessionID is used to track the current session ID of the ICE agent
// increase by one when disconnecting the agent
// with it the remote peer can discard the already deprecated offer/answer
// Without it the remote peer may recreate a workable ICE connection
sessionID icemaker.SessionID
sessionID ICESessionID
remoteSessionChanged bool
muxAgent sync.Mutex
localUfrag string
localPwd string
// we record the last known state of the ICE agent to avoid duplicate on disconnected events
lastKnownState ice.ConnectionState
// portForwardAttempted tracks if we've already tried port forwarding this session
portForwardAttempted bool
}
func NewICE(log *log.Entry, key string, iceConfig icemaker.Config, isController bool, onConnReady func(ConnPriority, ICEConnInfo), onStatusDisconnect func(bool), services ICEDependencies, hasRelayOnLocally bool) (*ICE, error) {
sessionID, err := icemaker.NewSessionID()
func NewWorkerICE(ctx context.Context, log *log.Entry, config ConnConfig, conn *Conn, signaler *Signaler, ifaceDiscover stdnet.ExternalIFaceDiscover, statusRecorder *Status, hasRelayOnLocally bool) (*WorkerICE, error) {
sessionID, err := NewICESessionID()
if err != nil {
return nil, err
}
w := &ICE{
log: log,
key: key,
iceConfig: iceConfig,
isController: isController,
onConnReady: onConnReady,
onStatusDisconnect: onStatusDisconnect,
signaler: services.Signaler,
iFaceDiscover: services.IFaceDiscover,
statusRecorder: services.StatusRecorder,
portForwardManager: services.PortForwardManager,
hasRelayOnLocally: hasRelayOnLocally,
sessionID: sessionID,
w := &WorkerICE{
ctx: ctx,
log: log,
config: config,
conn: conn,
signaler: signaler,
iFaceDiscover: ifaceDiscover,
statusRecorder: statusRecorder,
hasRelayOnLocally: hasRelayOnLocally,
lastKnownState: ice.ConnectionStateDisconnected,
sessionID: sessionID,
}
localUfrag, localPwd, err := icemaker.GenerateICECredentials()
@@ -106,7 +94,7 @@ func NewICE(log *log.Entry, key string, iceConfig icemaker.Config, isController
return w, nil
}
func (w *ICE) OnNewOffer(ctx context.Context, remoteOfferAnswer *signaling.OfferAnswer) {
func (w *WorkerICE) OnNewOffer(remoteOfferAnswer *OfferAnswer) {
w.log.Debugf("OnNewOffer for ICE, serial: %s", remoteOfferAnswer.SessionIDString())
w.muxAgent.Lock()
defer w.muxAgent.Unlock()
@@ -130,7 +118,7 @@ func (w *ICE) OnNewOffer(ctx context.Context, remoteOfferAnswer *signaling.Offer
}
}
sessionID, err := icemaker.NewSessionID()
sessionID, err := NewICESessionID()
if err != nil {
w.log.Errorf("failed to create new session ID: %s", err)
}
@@ -148,8 +136,8 @@ func (w *ICE) OnNewOffer(ctx context.Context, remoteOfferAnswer *signaling.Offer
if remoteOfferAnswer.SessionID != nil {
w.log.Debugf("recreate ICE agent: %s / %s", w.sessionID, *remoteOfferAnswer.SessionID)
}
dialerCtx, dialerCancel := context.WithCancel(ctx)
agent, err := w.reCreateAgent(ctx, dialerCancel, preferredCandidateTypes)
dialerCtx, dialerCancel := context.WithCancel(w.ctx)
agent, err := w.reCreateAgent(dialerCancel, preferredCandidateTypes)
if err != nil {
w.log.Errorf("failed to recreate ICE Agent: %s", err)
return
@@ -163,14 +151,14 @@ func (w *ICE) OnNewOffer(ctx context.Context, remoteOfferAnswer *signaling.Offer
w.remoteSessionID = ""
}
go w.connect(dialerCtx, dialerCancel, agent, remoteOfferAnswer)
go w.connect(dialerCtx, agent, remoteOfferAnswer)
}
// OnRemoteCandidate Handles ICE connection Candidate provided by the remote peer.
func (w *ICE) OnRemoteCandidate(candidate ice.Candidate, haRoutes route.HAMap) {
func (w *WorkerICE) OnRemoteCandidate(candidate ice.Candidate, haRoutes route.HAMap) {
w.muxAgent.Lock()
defer w.muxAgent.Unlock()
w.log.Debugf("OnRemoteCandidate from peer %s -> %s", w.key, candidate.String())
w.log.Debugf("OnRemoteCandidate from peer %s -> %s", w.config.Key, candidate.String())
if w.agent == nil {
w.log.Warnf("ICE Agent is not initialized yet")
return
@@ -197,24 +185,18 @@ func (w *ICE) OnRemoteCandidate(candidate ice.Candidate, haRoutes route.HAMap) {
}
}
func (w *ICE) Credentials() signaling.Credentials {
w.muxAgent.Lock()
defer w.muxAgent.Unlock()
return signaling.Credentials{
UFrag: w.localUfrag,
Pwd: w.localPwd,
SessionID: w.sessionID,
}
func (w *WorkerICE) GetLocalUserCredentials() (frag string, pwd string) {
return w.localUfrag, w.localPwd
}
func (w *ICE) InProgress() bool {
func (w *WorkerICE) InProgress() bool {
w.muxAgent.Lock()
defer w.muxAgent.Unlock()
return w.agentConnecting
}
func (w *ICE) Close() {
func (w *WorkerICE) Close() {
w.muxAgent.Lock()
defer w.muxAgent.Unlock()
@@ -230,10 +212,10 @@ func (w *ICE) Close() {
w.agent = nil
}
func (w *ICE) reCreateAgent(ctx context.Context, dialerCancel context.CancelFunc, candidates []ice.CandidateType) (*icemaker.ThreadSafeAgent, error) {
func (w *WorkerICE) reCreateAgent(dialerCancel context.CancelFunc, candidates []ice.CandidateType) (*icemaker.ThreadSafeAgent, error) {
w.portForwardAttempted = false
agent, err := icemaker.NewAgent(ctx, w.iFaceDiscover, w.iceConfig, candidates, w.localUfrag, w.localPwd)
agent, err := icemaker.NewAgent(w.ctx, w.iFaceDiscover, w.config.ICEConfig, candidates, w.localUfrag, w.localPwd)
if err != nil {
return nil, fmt.Errorf("create agent: %w", err)
}
@@ -255,7 +237,7 @@ func (w *ICE) reCreateAgent(ctx context.Context, dialerCancel context.CancelFunc
return agent, nil
}
func (w *ICE) getSessionID() icemaker.SessionID {
func (w *WorkerICE) SessionID() ICESessionID {
w.muxAgent.Lock()
defer w.muxAgent.Unlock()
@@ -265,11 +247,11 @@ func (w *ICE) getSessionID() icemaker.SessionID {
// will block until connection succeeded
// but it won't release if ICE Agent went into Disconnected or Failed state,
// so we have to cancel it with the provided context once agent detected a broken connection
func (w *ICE) connect(ctx context.Context, dialerCancel context.CancelFunc, agent *icemaker.ThreadSafeAgent, remoteOfferAnswer *signaling.OfferAnswer) {
func (w *WorkerICE) connect(ctx context.Context, agent *icemaker.ThreadSafeAgent, remoteOfferAnswer *OfferAnswer) {
w.log.Debugf("gather candidates")
if err := agent.GatherCandidates(); err != nil {
w.log.Warnf("failed to gather candidates: %s", err)
w.closeAgent(agent, dialerCancel)
w.closeAgent(agent, w.agentDialerCancel)
return
}
@@ -277,19 +259,19 @@ func (w *ICE) connect(ctx context.Context, dialerCancel context.CancelFunc, agen
remoteConn, err := w.turnAgentDial(ctx, agent, remoteOfferAnswer)
if err != nil {
w.log.Debugf("failed to dial the remote peer: %s", err)
w.closeAgent(agent, dialerCancel)
w.closeAgent(agent, w.agentDialerCancel)
return
}
w.log.Debugf("agent dial succeeded")
pair, err := agent.GetSelectedCandidatePair()
if err != nil {
w.closeAgent(agent, dialerCancel)
w.closeAgent(agent, w.agentDialerCancel)
return
}
if pair == nil {
w.log.Warnf("selected candidate pair is nil, cannot proceed")
w.closeAgent(agent, dialerCancel)
w.closeAgent(agent, w.agentDialerCancel)
return
}
@@ -317,22 +299,17 @@ func (w *ICE) connect(ctx context.Context, dialerCancel context.CancelFunc, agen
}
w.log.Debugf("on ICE conn is ready to use")
w.log.Infof("connection succeeded with offer session: %s", remoteOfferAnswer.SessionIDString())
w.muxAgent.Lock()
if w.agent != agent {
w.muxAgent.Unlock()
w.log.Debugf("agent has been replaced during connect, dropping obsolete connection")
return
}
w.agentConnecting = false
w.lastSuccess = time.Now()
w.connectedAgent = agent
w.muxAgent.Unlock()
w.log.Infof("connection succeeded with offer session: %s", remoteOfferAnswer.SessionIDString())
w.onConnReady(selectedPriority(pair), ci)
// todo: the potential problem is a race between the onConnectionStateChange
w.conn.onICEConnectionIsReady(selectedPriority(pair), ci)
}
func (w *ICE) closeAgent(agent *icemaker.ThreadSafeAgent, cancel context.CancelFunc) bool {
func (w *WorkerICE) closeAgent(agent *icemaker.ThreadSafeAgent, cancel context.CancelFunc) bool {
cancel()
if err := agent.Close(); err != nil {
w.log.Warnf("failed to close ICE agent: %s", err)
@@ -346,7 +323,7 @@ func (w *ICE) closeAgent(agent *icemaker.ThreadSafeAgent, cancel context.CancelF
if w.agent == agent {
// consider to remove from here and move to the OnNewOffer
sessionID, err := icemaker.NewSessionID()
sessionID, err := NewICESessionID()
if err != nil {
w.log.Errorf("failed to create new session ID: %s", err)
}
@@ -358,7 +335,7 @@ func (w *ICE) closeAgent(agent *icemaker.ThreadSafeAgent, cancel context.CancelF
return sessionChanged
}
func (w *ICE) punchRemoteWGPort(pair *ice.CandidatePair, remoteWgPort int) {
func (w *WorkerICE) punchRemoteWGPort(pair *ice.CandidatePair, remoteWgPort int) {
// wait local endpoint configuration
time.Sleep(time.Second)
addr, err := net.ResolveUDPAddr("udp", net.JoinHostPort(pair.Remote.Address(), strconv.Itoa(remoteWgPort)))
@@ -367,7 +344,7 @@ func (w *ICE) punchRemoteWGPort(pair *ice.CandidatePair, remoteWgPort int) {
return
}
mux, ok := w.iceConfig.UDPMuxSrflx.(*udpmux.UniversalUDPMuxDefault)
mux, ok := w.config.ICEConfig.UDPMuxSrflx.(*udpmux.UniversalUDPMuxDefault)
if !ok {
w.log.Warn("invalid udp mux conversion")
return
@@ -380,7 +357,7 @@ func (w *ICE) punchRemoteWGPort(pair *ice.CandidatePair, remoteWgPort int) {
// onICECandidate is a callback attached to an ICE Agent to receive new local connection candidates
// and then signals them to the remote peer
func (w *ICE) onICECandidate(candidate ice.Candidate) {
func (w *WorkerICE) onICECandidate(candidate ice.Candidate) {
// nil means candidate gathering has been ended
if candidate == nil {
return
@@ -389,9 +366,9 @@ func (w *ICE) onICECandidate(candidate ice.Candidate) {
// TODO: reported port is incorrect for CandidateTypeHost, makes understanding ICE use via logs confusing as port is ignored
w.log.Debugf("discovered local candidate %s", candidate.String())
go func() {
err := w.signaler.SignalICECandidate(candidate, w.key)
err := w.signaler.SignalICECandidate(candidate, w.config.Key)
if err != nil {
w.log.Errorf("failed signaling candidate to the remote peer %s %s", w.key, err)
w.log.Errorf("failed signaling candidate to the remote peer %s %s", w.config.Key, err)
}
}()
@@ -401,8 +378,8 @@ func (w *ICE) onICECandidate(candidate ice.Candidate) {
}
// injectPortForwardedCandidate signals an additional candidate using the pre-created port mapping.
func (w *ICE) injectPortForwardedCandidate(srflxCandidate ice.Candidate) {
pfManager := w.portForwardManager
func (w *WorkerICE) injectPortForwardedCandidate(srflxCandidate ice.Candidate) {
pfManager := w.conn.portForwardManager
if pfManager == nil {
return
}
@@ -430,7 +407,7 @@ func (w *ICE) injectPortForwardedCandidate(srflxCandidate ice.Candidate) {
forwardedCandidate.String(), mapping.InternalPort, mapping.ExternalPort, mapping.NATType, forwardedCandidate.Priority())
go func() {
if err := w.signaler.SignalICECandidate(forwardedCandidate, w.key); err != nil {
if err := w.signaler.SignalICECandidate(forwardedCandidate, w.config.Key); err != nil {
w.log.Errorf("signal port-forwarded candidate: %v", err)
}
}()
@@ -438,7 +415,7 @@ func (w *ICE) injectPortForwardedCandidate(srflxCandidate ice.Candidate) {
// createForwardedCandidate creates a new server reflexive candidate with the forwarded port.
// It uses the NAT gateway's external IP with the forwarded port.
func (w *ICE) createForwardedCandidate(srflxCandidate ice.Candidate, mapping *portforward.Mapping) (ice.Candidate, error) {
func (w *WorkerICE) createForwardedCandidate(srflxCandidate ice.Candidate, mapping *portforward.Mapping) (ice.Candidate, error) {
var externalIP string
if mapping.ExternalIP != nil && !mapping.ExternalIP.IsUnspecified() {
externalIP = mapping.ExternalIP.String()
@@ -483,9 +460,9 @@ func (w *ICE) createForwardedCandidate(srflxCandidate ice.Candidate, mapping *po
return candidate, nil
}
func (w *ICE) onICESelectedCandidatePair(agent *icemaker.ThreadSafeAgent, c1, c2 ice.Candidate) {
func (w *WorkerICE) onICESelectedCandidatePair(agent *icemaker.ThreadSafeAgent, c1, c2 ice.Candidate) {
w.log.Debugf("selected candidate pair [local <-> remote] -> [%s <-> %s], peer %s", c1.String(), c2.String(),
w.key)
w.config.Key)
pairStat, ok := agent.GetSelectedCandidatePairStats()
if !ok {
@@ -494,14 +471,14 @@ func (w *ICE) onICESelectedCandidatePair(agent *icemaker.ThreadSafeAgent, c1, c2
}
duration := time.Duration(pairStat.CurrentRoundTripTime * float64(time.Second))
if err := w.statusRecorder.UpdateLatency(w.key, duration); err != nil {
if err := w.statusRecorder.UpdateLatency(w.config.Key, duration); err != nil {
w.log.Debugf("failed to update latency for peer: %s", err)
return
}
}
func (w *ICE) logSuccessfulPaths(agent *icemaker.ThreadSafeAgent) {
sessionID := w.getSessionID()
func (w *WorkerICE) logSuccessfulPaths(agent *icemaker.ThreadSafeAgent) {
sessionID := w.SessionID()
stats := agent.GetCandidatePairsStats()
localCandidates, _ := agent.GetLocalCandidates()
remoteCandidates, _ := agent.GetRemoteCandidates()
@@ -531,44 +508,32 @@ func (w *ICE) logSuccessfulPaths(agent *icemaker.ThreadSafeAgent) {
}
}
func (w *ICE) onConnectionStateChange(agent *icemaker.ThreadSafeAgent, dialerCancel context.CancelFunc) func(ice.ConnectionState) {
// per-agent state; pion delivers callbacks of one agent sequentially
var connected bool
func (w *WorkerICE) onConnectionStateChange(agent *icemaker.ThreadSafeAgent, dialerCancel context.CancelFunc) func(ice.ConnectionState) {
return func(state ice.ConnectionState) {
w.log.Debugf("ICE ConnectionState has changed to %s", state.String())
switch state {
case ice.ConnectionStateConnected:
connected = true
w.lastKnownState = ice.ConnectionStateConnected
w.logSuccessfulPaths(agent)
return
case ice.ConnectionStateFailed, ice.ConnectionStateDisconnected, ice.ConnectionStateClosed:
// ice.ConnectionStateClosed happens when we recreate the agent. For the P2P to TURN switch important to
// notify the conn.onICEStateDisconnected changes to update the current used priority
sessionChanged := w.closeAgent(agent, dialerCancel)
if !connected {
return
if w.lastKnownState == ice.ConnectionStateConnected {
w.lastKnownState = ice.ConnectionStateDisconnected
w.conn.onICEStateDisconnected(sessionChanged)
}
connected = false
w.muxAgent.Lock()
stale := w.connectedAgent != agent
if !stale {
w.connectedAgent = nil
}
w.muxAgent.Unlock()
if stale {
w.log.Debugf("suppress disconnected event of replaced ICE agent")
return
}
w.onStatusDisconnect(sessionChanged)
default:
return
}
}
}
func (w *ICE) turnAgentDial(ctx context.Context, agent *icemaker.ThreadSafeAgent, remoteOfferAnswer *signaling.OfferAnswer) (*ice.Conn, error) {
if w.isController {
func (w *WorkerICE) turnAgentDial(ctx context.Context, agent *icemaker.ThreadSafeAgent, remoteOfferAnswer *OfferAnswer) (*ice.Conn, error) {
if isController(w.config) {
return agent.Dial(ctx, remoteOfferAnswer.IceCredentials.UFrag, remoteOfferAnswer.IceCredentials.Pwd)
} else {
return agent.Accept(ctx, remoteOfferAnswer.IceCredentials.UFrag, remoteOfferAnswer.IceCredentials.Pwd)
@@ -630,10 +595,10 @@ func isRelayed(pair *ice.CandidatePair) bool {
return false
}
func selectedPriority(pair *ice.CandidatePair) ConnPriority {
func selectedPriority(pair *ice.CandidatePair) conntype.ConnPriority {
if isRelayed(pair) {
return ICETurn
return conntype.ICETurn
} else {
return ICEP2P
return conntype.ICEP2P
}
}

View File

@@ -1,4 +1,4 @@
package worker
package peer
import (
"context"
@@ -10,23 +10,22 @@ import (
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/internal/peer/signaling"
relayClient "github.com/netbirdio/netbird/shared/relay/client"
)
type RelayConnInfo struct {
RelayedConn net.Conn
RosenpassPubKey []byte
RosenpassAddr string
relayedConn net.Conn
rosenpassPubKey []byte
rosenpassAddr string
}
type WorkerRelay struct {
log *log.Entry
key string
isController bool
onConnReady func(RelayConnInfo)
onDisconnected func()
relayManager *relayClient.Manager
peerCtx context.Context
log *log.Entry
isController bool
config ConnConfig
conn *Conn
relayManager *relayClient.Manager
relayedConn net.Conn
relayLock sync.Mutex
@@ -34,19 +33,19 @@ type WorkerRelay struct {
relaySupportedOnRemotePeer atomic.Bool
}
func NewWorkerRelay(log *log.Entry, key string, isController bool, onConnReady func(RelayConnInfo), onDisconnected func(), relayManager *relayClient.Manager) *WorkerRelay {
func NewWorkerRelay(ctx context.Context, log *log.Entry, ctrl bool, config ConnConfig, conn *Conn, relayManager *relayClient.Manager) *WorkerRelay {
r := &WorkerRelay{
log: log,
key: key,
isController: isController,
onConnReady: onConnReady,
onDisconnected: onDisconnected,
relayManager: relayManager,
peerCtx: ctx,
log: log,
isController: ctrl,
config: config,
conn: conn,
relayManager: relayManager,
}
return r
}
func (w *WorkerRelay) OnNewOffer(ctx context.Context, remoteOfferAnswer *signaling.OfferAnswer) {
func (w *WorkerRelay) OnNewOffer(remoteOfferAnswer *OfferAnswer) {
if !w.isRelaySupported(remoteOfferAnswer) {
w.log.Infof("Relay is not supported by remote peer")
w.relaySupportedOnRemotePeer.Store(false)
@@ -67,7 +66,7 @@ func (w *WorkerRelay) OnNewOffer(ctx context.Context, remoteOfferAnswer *signali
serverIP = remoteOfferAnswer.RelaySrvIP
}
relayedConn, err := w.relayManager.OpenConn(ctx, srv, w.key, serverIP)
relayedConn, err := w.relayManager.OpenConn(w.peerCtx, srv, w.config.Key, serverIP)
if err != nil {
if errors.Is(err, relayClient.ErrConnAlreadyExists) {
w.log.Debugf("handled offer by reusing existing relay connection")
@@ -89,10 +88,10 @@ func (w *WorkerRelay) OnNewOffer(ctx context.Context, remoteOfferAnswer *signali
}
w.log.Debugf("peer conn opened via Relay: %s", srv)
w.onConnReady(RelayConnInfo{
RelayedConn: relayedConn,
RosenpassPubKey: remoteOfferAnswer.RosenpassPubKey,
RosenpassAddr: remoteOfferAnswer.RosenpassAddr,
go w.conn.onRelayConnectionIsReady(RelayConnInfo{
relayedConn: relayedConn,
rosenpassPubKey: remoteOfferAnswer.RosenpassPubKey,
rosenpassAddr: remoteOfferAnswer.RosenpassAddr,
})
}
@@ -120,7 +119,7 @@ func (w *WorkerRelay) CloseConn() {
}
}
func (w *WorkerRelay) isRelaySupported(answer *signaling.OfferAnswer) bool {
func (w *WorkerRelay) isRelaySupported(answer *OfferAnswer) bool {
if !w.relayManager.HasRelayAddress() {
return false
}
@@ -135,5 +134,5 @@ func (w *WorkerRelay) preferredRelayServer(myRelayAddress, remoteRelayAddress st
}
func (w *WorkerRelay) onRelayClientDisconnected() {
w.onDisconnected()
go w.conn.onRelayDisconnected()
}

View File

@@ -45,35 +45,12 @@ func (pm *ProfileManager) GetProfileState(id ID) (*ProfileState, error) {
return &state, nil
}
// SetProfileState writes the state file of the profile identified by id. Prefer
// it over SetActiveProfileState whenever the caller knows which profile the data
// belongs to: an SSO login spans seconds of user interaction, and the active
// profile can change during it, which would file the account email under
// whichever profile happened to be active when the flow returned.
func (pm *ProfileManager) SetProfileState(id ID, state *ProfileState) error {
func (pm *ProfileManager) SetActiveProfileState(state *ProfileState) error {
configDir, err := getConfigDir()
if err != nil {
return fmt.Errorf("get config directory: %w", err)
}
if id == "" {
return fmt.Errorf("empty profile ID")
}
if id != defaultProfileName && !IsValidProfileFilenameStem(id) {
return fmt.Errorf("invalid profile ID: %q", id)
}
stateFile := filepath.Join(configDir, id.String()+".state.json")
if err := util.WriteJsonWithRestrictedPermission(context.Background(), stateFile, state); err != nil {
return fmt.Errorf("write profile state: %w", err)
}
return nil
}
// SetActiveProfileState writes the state file of whichever profile is active at
// call time. Use SetProfileState when the target profile is known.
func (pm *ProfileManager) SetActiveProfileState(state *ProfileState) error {
activeProf, err := pm.GetActiveProfile()
if err != nil {
if errors.Is(err, ErrNoActiveProfile) {
@@ -82,7 +59,18 @@ func (pm *ProfileManager) SetActiveProfileState(state *ProfileState) error {
return fmt.Errorf("get active profile: %w", err)
}
return pm.SetProfileState(activeProf.ID, state)
id := activeProf.ID
if id != defaultProfileName && !IsValidProfileFilenameStem(id) {
return fmt.Errorf("invalid active profile ID: %q", id)
}
stateFile := filepath.Join(configDir, id.String()+".state.json")
err = util.WriteJsonWithRestrictedPermission(context.Background(), stateFile, state)
if err != nil {
return fmt.Errorf("write profile state: %w", err)
}
return nil
}
// RemoveProfileState deletes the per-profile state file (which holds the

View File

@@ -403,7 +403,7 @@ func (w *Watcher) connectEvent(route *route.Route) {
proto.SystemEvent_INFO,
proto.SystemEvent_NETWORK,
"Default route added",
"Exit node connected.",
proto.NewUserMessage(proto.UserMsgExitNodeConnected),
meta,
)
}
@@ -423,7 +423,7 @@ func (w *Watcher) disconnectEvent(route *route.Route, rsn reason) {
var severity proto.SystemEvent_Severity
var message string
var userMessage string
var userMessage *proto.UserMessage
meta := make(map[string]string)
if route != nil {
@@ -435,22 +435,22 @@ func (w *Watcher) disconnectEvent(route *route.Route, rsn reason) {
case reasonShutdown:
severity = proto.SystemEvent_INFO
message = "Default route removed"
userMessage = "Exit node disconnected."
userMessage = proto.NewUserMessage(proto.UserMsgExitNodeDisconnected)
case reasonRouteUpdate:
severity = proto.SystemEvent_INFO
message = "Default route updated due to configuration change"
case reasonPeerUpdate:
severity = proto.SystemEvent_WARNING
message = "Default route disconnected due to peer unreachability"
userMessage = "Exit node connection lost. Your internet access might be affected."
userMessage = proto.NewUserMessage(proto.UserMsgExitNodeConnectionLost)
case reasonHA:
severity = proto.SystemEvent_INFO
message = "Default route disconnected due to high availability change"
userMessage = "Exit node disconnected due to high availability change."
userMessage = proto.NewUserMessage(proto.UserMsgExitNodeHAChange)
default:
severity = proto.SystemEvent_ERROR
message = "Default route disconnected for unknown reasons"
userMessage = "Exit node disconnected for unknown reasons."
userMessage = proto.NewUserMessage(proto.UserMsgExitNodeDisconnectedUnknown)
}
w.statusRecorder.PublishEvent(

View File

@@ -479,7 +479,7 @@ func (d *DnsInterceptor) removeDNATMappings(realPrefixes []netip.Prefix, logger
// internalDnatFw checks if the firewall supports internal DNAT
func (d *DnsInterceptor) internalDnatFw() (internalDNATer, bool) {
if d.firewall == nil || d.fakeIPManager == nil || runtime.GOOS != "android" {
if d.firewall == nil || runtime.GOOS != "android" {
return nil, false
}
fw, ok := d.firewall.(internalDNATer)

View File

@@ -8,13 +8,14 @@ import (
"net/netip"
"net/url"
"runtime"
"sort"
"slices"
"strings"
"sync"
"sync/atomic"
"syscall"
"time"
"github.com/google/uuid"
"github.com/hashicorp/go-multierror"
log "github.com/sirupsen/logrus"
"golang.org/x/exp/maps"
@@ -61,7 +62,7 @@ type Manager interface {
GetActiveClientRoutes() route.HAMap
GetClientRoutesWithNetID() map[route.NetID][]*route.Route
SetRouteChangeListener(listener listener.NetworkChangeListener)
CurrentRouteRange() []string
InitialRouteRange() []string
SetFirewall(firewall.Manager) error
SetDNSForwarderPort(port uint16)
ReconcilePeerAllowedIPs(peerKey string) error
@@ -75,8 +76,10 @@ type ManagerConfig struct {
WGInterface iface.WGIface
StatusRecorder *peer.Status
RelayManager *relayClient.Manager
InitialRoutes []*route.Route
StateManager *statemanager.Manager
DNSServer dns.Server
DNSFeatureFlag bool
PeerStore *peerstore.Store
DisableClientRoutes bool
DisableServerRoutes bool
@@ -146,12 +149,45 @@ func NewManager(config ManagerConfig) *DefaultManager {
useNoop := netstack.IsEnabled() || config.DisableClientRoutes
dm.setupRefCounters(useNoop)
// don't proceed with client routes if it is disabled
if config.DisableClientRoutes {
return dm
}
if runtime.GOOS == "android" {
dm.setupAndroidRoutes(config)
}
return dm
}
func (m *DefaultManager) setupAndroidRoutes(config ManagerConfig) {
cr := m.initialClientRoutes(config.InitialRoutes)
func (m *DefaultManager) enableFakeIPRoutes() {
m.fakeIPManager = fakeip.NewManager()
m.notifier.NotifyRouteChange()
routesForComparison := slices.Clone(cr)
if config.DNSFeatureFlag {
m.fakeIPManager = fakeip.NewManager()
v4ID := uuid.NewString()
fakeIPRoute := &route.Route{
ID: route.ID(v4ID),
Network: m.fakeIPManager.GetFakeIPBlock(),
NetID: route.NetID(v4ID),
Peer: m.pubKey,
NetworkType: route.IPv4Network,
}
v6ID := uuid.NewString()
fakeIPv6Route := &route.Route{
ID: route.ID(v6ID),
Network: m.fakeIPManager.GetFakeIPv6Block(),
NetID: route.NetID(v6ID),
Peer: m.pubKey,
NetworkType: route.IPv6Network,
}
cr = append(cr, fakeIPRoute, fakeIPv6Route)
m.notifier.SetFakeIPRoutes([]*route.Route{fakeIPRoute, fakeIPv6Route})
}
m.notifier.SetInitialClientRoutes(cr, routesForComparison)
}
func (m *DefaultManager) setupRefCounters(useNoop bool) {
@@ -428,9 +464,6 @@ func (m *DefaultManager) UpdateRoutes(
var merr *multierror.Error
if !m.disableClientRoutes {
if runtime.GOOS == "android" && useNewDNSRoute && m.fakeIPManager == nil {
m.enableFakeIPRoutes()
}
// Update route selector based on management server's isSelected status
m.updateRouteSelectorFromManagement(clientRoutes)
@@ -467,32 +500,9 @@ func (m *DefaultManager) SetRouteChangeListener(listener listener.NetworkChangeL
m.notifier.SetListener(listener)
}
// CurrentRouteRange returns the current TUN route list. It is used by mobile systems
func (m *DefaultManager) CurrentRouteRange() []string {
m.mux.Lock()
defer m.mux.Unlock()
if m.disableClientRoutes {
return nil
}
filtered := m.routeSelector.FilterSelectedExitNodes(m.clientRoutes)
var nets []string
for _, routes := range filtered {
for _, r := range routes {
if r.IsDynamic() {
continue
}
nets = append(nets, r.NetString())
}
}
if m.fakeIPManager != nil {
nets = append(nets, m.fakeIPManager.GetFakeIPBlock().String(), m.fakeIPManager.GetFakeIPv6Block().String())
}
sort.Strings(nets)
return nets
// InitialRouteRange return the list of initial routes. It used by mobile systems
func (m *DefaultManager) InitialRouteRange() []string {
return m.notifier.GetInitialRouteRanges()
}
// GetRouteSelector returns the route selector
@@ -690,6 +700,16 @@ func (m *DefaultManager) ClassifyRoutes(newRoutes []*route.Route) (map[route.ID]
return newServerRoutesMap, newClientRoutesIDMap
}
func (m *DefaultManager) initialClientRoutes(initialRoutes []*route.Route) []*route.Route {
_, crMap := m.ClassifyRoutes(initialRoutes)
rs := make([]*route.Route, 0, len(crMap))
for _, routes := range crMap {
rs = append(rs, routes...)
}
return rs
}
func isRouteSupported(route *route.Route) bool {
if netstack.IsEnabled() || !nbnet.CustomRoutingDisabled() || route.IsDynamic() {
return true

View File

@@ -30,8 +30,8 @@ func (m *MockManager) Init() error {
return nil
}
// CurrentRouteRange mock implementation of CurrentRouteRange from Manager interface
func (m *MockManager) CurrentRouteRange() []string {
// InitialRouteRange mock implementation of InitialRouteRange from Manager interface
func (m *MockManager) InitialRouteRange() []string {
return nil
}

View File

@@ -6,6 +6,7 @@ import (
"net/netip"
"slices"
"sort"
"strings"
"sync"
"github.com/netbirdio/netbird/client/internal/listener"
@@ -13,15 +14,12 @@ import (
)
type Notifier struct {
mu sync.Mutex
// currentRoutes is the last announced route set. It exists only to
// suppress noise: without it every network map sync would trigger the
// Java side, even when the routes did not change. The actual TUN route
// state is owned by the route manager and pulled from there.
initialRoutes []*route.Route
currentRoutes []*route.Route
fakeIPRoutes []*route.Route
listener listener.NetworkChangeListener
listener listener.NetworkChangeListener
listenerMux sync.Mutex
}
func NewNotifier() *Notifier {
@@ -29,15 +27,20 @@ func NewNotifier() *Notifier {
}
func (n *Notifier) SetListener(listener listener.NetworkChangeListener) {
n.mu.Lock()
defer n.mu.Unlock()
n.listenerMux.Lock()
defer n.listenerMux.Unlock()
n.listener = listener
}
func (n *Notifier) NotifyRouteChange() {
n.mu.Lock()
defer n.mu.Unlock()
n.notifyLocked()
// SetInitialClientRoutes stores the initial route sets for TUN configuration.
func (n *Notifier) SetInitialClientRoutes(initialRoutes []*route.Route, routesForComparison []*route.Route) {
n.initialRoutes = filterStatic(initialRoutes)
n.currentRoutes = filterStatic(routesForComparison)
}
// SetFakeIPRoutes stores the fake IP routes to be included in every TUN rebuild.
func (n *Notifier) SetFakeIPRoutes(routes []*route.Route) {
n.fakeIPRoutes = routes
}
func (n *Notifier) OnNewRoutes(idMap route.HAMap) {
@@ -51,32 +54,46 @@ func (n *Notifier) OnNewRoutes(idMap route.HAMap) {
}
}
n.mu.Lock()
defer n.mu.Unlock()
if !hasRouteDiff(n.currentRoutes, newRoutes) {
if !n.hasRouteDiff(n.currentRoutes, newRoutes) {
return
}
n.currentRoutes = newRoutes
n.notifyLocked()
n.notify()
}
func (n *Notifier) OnNewPrefixes([]netip.Prefix) {
// Not used on Android
}
func (n *Notifier) notifyLocked() {
func (n *Notifier) notify() {
n.listenerMux.Lock()
defer n.listenerMux.Unlock()
if n.listener == nil {
return
}
n.listener.OnNetworkChanged("")
allRoutes := slices.Clone(n.currentRoutes)
allRoutes = append(allRoutes, n.fakeIPRoutes...)
routeStrings := n.routesToStrings(allRoutes)
sort.Strings(routeStrings)
go func(l listener.NetworkChangeListener) {
l.OnNetworkChanged(strings.Join(routeStrings, ","))
}(n.listener)
}
func (n *Notifier) Close() {
// unused
func filterStatic(routes []*route.Route) []*route.Route {
out := make([]*route.Route, 0, len(routes))
for _, r := range routes {
if !r.IsDynamic() {
out = append(out, r)
}
}
return out
}
func routesToStrings(routes []*route.Route) []string {
func (n *Notifier) routesToStrings(routes []*route.Route) []string {
nets := make([]string, 0, len(routes))
for _, r := range routes {
nets = append(nets, r.NetString())
@@ -84,10 +101,25 @@ func routesToStrings(routes []*route.Route) []string {
return nets
}
func hasRouteDiff(a []*route.Route, b []*route.Route) bool {
as := routesToStrings(a)
bs := routesToStrings(b)
sort.Strings(as)
sort.Strings(bs)
return !slices.Equal(as, bs)
func (n *Notifier) hasRouteDiff(a []*route.Route, b []*route.Route) bool {
slices.SortFunc(a, func(x, y *route.Route) int {
return strings.Compare(x.NetString(), y.NetString())
})
slices.SortFunc(b, func(x, y *route.Route) int {
return strings.Compare(x.NetString(), y.NetString())
})
return !slices.EqualFunc(a, b, func(x, y *route.Route) bool {
return x.NetString() == y.NetString()
})
}
func (n *Notifier) GetInitialRouteRanges() []string {
initialStrings := n.routesToStrings(n.initialRoutes)
sort.Strings(initialStrings)
return initialStrings
}
func (n *Notifier) Close() {
// unused
}

View File

@@ -29,7 +29,11 @@ func (n *Notifier) SetListener(listener listener.NetworkChangeListener) {
n.listener = listener
}
func (n *Notifier) NotifyRouteChange() {
func (n *Notifier) SetInitialClientRoutes([]*route.Route, []*route.Route) {
// iOS doesn't care about initial routes
}
func (n *Notifier) SetFakeIPRoutes([]*route.Route) {
// Not used on iOS
}

View File

@@ -19,7 +19,11 @@ func (n *Notifier) SetListener(listener listener.NetworkChangeListener) {
// Not used on non-mobile platforms
}
func (n *Notifier) NotifyRouteChange() {
func (n *Notifier) SetInitialClientRoutes([]*route.Route, []*route.Route) {
// Not used on non-mobile platforms
}
func (n *Notifier) SetFakeIPRoutes([]*route.Route) {
// Not used on non-mobile platforms
}
@@ -31,6 +35,10 @@ func (n *Notifier) OnNewPrefixes(prefixes []netip.Prefix) {
// Not used on non-mobile platforms
}
func (n *Notifier) GetInitialRouteRanges() []string {
return []string{}
}
func (n *Notifier) Close() {
// unused
}

View File

@@ -98,44 +98,47 @@ func (u *Installer) startDaemon(daemonFolder string) error {
func (u *Installer) startUIAsUser() error {
log.Infof("starting netbird-ui: %s", uiBinary)
username, err := consoleUser()
// Get the current console user
cmd := exec.Command("stat", "-f", "%Su", "/dev/console")
output, err := cmd.Output()
if err != nil {
return err
return fmt.Errorf("failed to get console user: %w", err)
}
username := strings.TrimSpace(string(output))
if username == "" || username == "root" {
return fmt.Errorf("no active user session found")
}
log.Infof("starting UI for user: %s", username)
// Get user's UID
userInfo, err := user.Lookup(username)
if err != nil {
return fmt.Errorf("lookup user %s: %w", username, err)
return fmt.Errorf("failed to lookup user %s: %w", username, err)
}
log.Infof("starting UI for user: %s (uid %s)", username, userInfo.Uid)
launchCmd := exec.Command("launchctl", "asuser", userInfo.Uid, "sudo", "-u", username, "-H", "open", "-a", uiBinary)
// Start the UI process as the console user using launchctl
// This ensures the app runs in the user's context with proper GUI access
launchCmd := exec.Command("launchctl", "asuser", userInfo.Uid, "open", "-a", uiBinary)
log.Infof("launchCmd: %s", launchCmd.String())
// Set the user's home directory for proper macOS app behavior
launchCmd.Env = append(os.Environ(), "HOME="+userInfo.HomeDir)
log.Infof("set HOME environment variable: %s", userInfo.HomeDir)
if err := launchCmd.Run(); err != nil {
return fmt.Errorf("run UI launch: %w", err)
if err := launchCmd.Start(); err != nil {
return fmt.Errorf("failed to start UI process: %w", err)
}
// Release the process so it can run independently
if err := launchCmd.Process.Release(); err != nil {
log.Warnf("failed to release UI process: %v", err)
}
log.Infof("netbird-ui started successfully for user %s", username)
return nil
}
func consoleUser() (string, error) {
output, err := exec.Command("stat", "-f", "%Su", "/dev/console").Output()
if err != nil {
return "", fmt.Errorf("get console user: %w", err)
}
username := strings.TrimSpace(string(output))
switch username {
case "", "root", "loginwindow", "_mbsetupuser":
return "", fmt.Errorf("no active GUI user session, console user: %q", username)
}
return username, nil
}
func (u *Installer) installPkgFile(ctx context.Context, path string) error {
log.Infof("installing pkg file: %s", path)

View File

@@ -94,7 +94,7 @@ func (m *Manager) CheckUpdateSuccess(ctx context.Context) {
cProto.SystemEvent_ERROR,
cProto.SystemEvent_SYSTEM,
"Auto-update failed",
fmt.Sprintf("Auto-update failed: %s", reason),
cProto.NewUserMessage(cProto.UserMsgUpdateFailed, cProto.ArgReason, reason),
nil,
)
}
@@ -115,7 +115,7 @@ func (m *Manager) CheckUpdateSuccess(ctx context.Context) {
cProto.SystemEvent_INFO,
cProto.SystemEvent_SYSTEM,
"Auto-update completed",
fmt.Sprintf("Your NetBird Client was auto-updated to version %s", m.currentVersion),
cProto.NewUserMessage(cProto.UserMsgUpdateCompleted, cProto.ArgVersion, m.currentVersion),
nil,
)
return
@@ -272,7 +272,7 @@ func (m *Manager) NotifyUI() {
cProto.SystemEvent_INFO,
cProto.SystemEvent_SYSTEM,
"New version available",
"",
nil,
map[string]string{"new_version_available": latestVersion.String()},
)
return
@@ -283,7 +283,7 @@ func (m *Manager) NotifyUI() {
cProto.SystemEvent_INFO,
cProto.SystemEvent_SYSTEM,
"New version available",
"",
nil,
map[string]string{"new_version_available": pendingVersion.String(), "enforced": "true"},
)
}
@@ -384,7 +384,7 @@ func (m *Manager) handleUpdate(ctx context.Context) {
cProto.SystemEvent_INFO,
cProto.SystemEvent_SYSTEM,
"New version available",
"",
nil,
map[string]string{"new_version_available": updateVersion.String()},
)
return
@@ -401,7 +401,7 @@ func (m *Manager) handleUpdate(ctx context.Context) {
cProto.SystemEvent_INFO,
cProto.SystemEvent_SYSTEM,
"New version available",
"",
nil,
map[string]string{"new_version_available": updateVersion.String(), "enforced": "true"},
)
}
@@ -411,14 +411,14 @@ func (m *Manager) install(ctx context.Context, pendingVersion *v.Version) error
cProto.SystemEvent_CRITICAL,
cProto.SystemEvent_SYSTEM,
"Updating client",
"Installing update now.",
cProto.NewUserMessage(cProto.UserMsgUpdateInstalling),
nil,
)
m.statusRecorder.PublishEvent(
cProto.SystemEvent_CRITICAL,
cProto.SystemEvent_SYSTEM,
"",
"",
nil,
map[string]string{"progress_window": "show", "version": pendingVersion.String()},
)
@@ -441,7 +441,7 @@ func (m *Manager) install(ctx context.Context, pendingVersion *v.Version) error
cProto.SystemEvent_ERROR,
cProto.SystemEvent_SYSTEM,
"Auto-update failed",
fmt.Sprintf("Auto-update failed: %v", err),
cProto.NewUserMessage(cProto.UserMsgUpdateFailed, cProto.ArgReason, err.Error()),
nil,
)
return err

View File

@@ -3909,14 +3909,30 @@ func (*SubscribeRequest) Descriptor() ([]byte, []int) {
}
type SystemEvent struct {
state protoimpl.MessageState `protogen:"open.v1"`
Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"`
Severity SystemEvent_Severity `protobuf:"varint,2,opt,name=severity,proto3,enum=daemon.SystemEvent_Severity" json:"severity,omitempty"`
Category SystemEvent_Category `protobuf:"varint,3,opt,name=category,proto3,enum=daemon.SystemEvent_Category" json:"category,omitempty"`
Message string `protobuf:"bytes,4,opt,name=message,proto3" json:"message,omitempty"`
UserMessage string `protobuf:"bytes,5,opt,name=userMessage,proto3" json:"userMessage,omitempty"`
Timestamp *timestamppb.Timestamp `protobuf:"bytes,6,opt,name=timestamp,proto3" json:"timestamp,omitempty"`
Metadata map[string]string `protobuf:"bytes,7,rep,name=metadata,proto3" json:"metadata,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"`
state protoimpl.MessageState `protogen:"open.v1"`
Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"`
Severity SystemEvent_Severity `protobuf:"varint,2,opt,name=severity,proto3,enum=daemon.SystemEvent_Severity" json:"severity,omitempty"`
Category SystemEvent_Category `protobuf:"varint,3,opt,name=category,proto3,enum=daemon.SystemEvent_Category" json:"category,omitempty"`
Message string `protobuf:"bytes,4,opt,name=message,proto3" json:"message,omitempty"`
// userMessage is the daemon's English rendering of messageKey, kept for the
// CLI and for UIs that predate messageKey. UIs that localise read messageKey
// and treat this as the fallback.
UserMessage string `protobuf:"bytes,5,opt,name=userMessage,proto3" json:"userMessage,omitempty"`
Timestamp *timestamppb.Timestamp `protobuf:"bytes,6,opt,name=timestamp,proto3" json:"timestamp,omitempty"`
Metadata map[string]string `protobuf:"bytes,7,rep,name=metadata,proto3" json:"metadata,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"`
// messageKey names a user-facing message in a stable, locale-independent
// form. A UI resolves it against its own translation bundle and substitutes
// messageArgs; an unrecognised key falls back to userMessage. Empty on events
// that carry no user-facing text.
MessageKey string `protobuf:"bytes,8,opt,name=messageKey,proto3" json:"messageKey,omitempty"`
// messageArgs holds the placeholder name/value pairs for messageKey, e.g.
// {"version": "0.60.1"} for a "{version}" template. Values are data the
// daemon cannot localise (versions, addresses, error text).
MessageArgs map[string]string `protobuf:"bytes,9,rep,name=messageArgs,proto3" json:"messageArgs,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"`
// titleKey names the notification title in the same way as messageKey. Empty
// when the event has no dedicated title, in which case a UI composes one from
// severity and category.
TitleKey string `protobuf:"bytes,10,opt,name=titleKey,proto3" json:"titleKey,omitempty"`
unknownFields protoimpl.UnknownFields
sizeCache protoimpl.SizeCache
}
@@ -4000,6 +4016,27 @@ func (x *SystemEvent) GetMetadata() map[string]string {
return nil
}
func (x *SystemEvent) GetMessageKey() string {
if x != nil {
return x.MessageKey
}
return ""
}
func (x *SystemEvent) GetMessageArgs() map[string]string {
if x != nil {
return x.MessageArgs
}
return nil
}
func (x *SystemEvent) GetTitleKey() string {
if x != nil {
return x.TitleKey
}
return ""
}
type GetEventsRequest struct {
state protoimpl.MessageState `protogen:"open.v1"`
unknownFields protoimpl.UnknownFields
@@ -7332,7 +7369,7 @@ const file_daemon_proto_rawDesc = "" +
"\x13TracePacketResponse\x12*\n" +
"\x06stages\x18\x01 \x03(\v2\x12.daemon.TraceStageR\x06stages\x12+\n" +
"\x11final_disposition\x18\x02 \x01(\bR\x10finalDisposition\"\x12\n" +
"\x10SubscribeRequest\"\x93\x04\n" +
"\x10SubscribeRequest\"\xd7\x05\n" +
"\vSystemEvent\x12\x0e\n" +
"\x02id\x18\x01 \x01(\tR\x02id\x128\n" +
"\bseverity\x18\x02 \x01(\x0e2\x1c.daemon.SystemEvent.SeverityR\bseverity\x128\n" +
@@ -7340,9 +7377,18 @@ const file_daemon_proto_rawDesc = "" +
"\amessage\x18\x04 \x01(\tR\amessage\x12 \n" +
"\vuserMessage\x18\x05 \x01(\tR\vuserMessage\x128\n" +
"\ttimestamp\x18\x06 \x01(\v2\x1a.google.protobuf.TimestampR\ttimestamp\x12=\n" +
"\bmetadata\x18\a \x03(\v2!.daemon.SystemEvent.MetadataEntryR\bmetadata\x1a;\n" +
"\bmetadata\x18\a \x03(\v2!.daemon.SystemEvent.MetadataEntryR\bmetadata\x12\x1e\n" +
"\n" +
"messageKey\x18\b \x01(\tR\n" +
"messageKey\x12F\n" +
"\vmessageArgs\x18\t \x03(\v2$.daemon.SystemEvent.MessageArgsEntryR\vmessageArgs\x12\x1a\n" +
"\btitleKey\x18\n" +
" \x01(\tR\btitleKey\x1a;\n" +
"\rMetadataEntry\x12\x10\n" +
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\x1a>\n" +
"\x10MessageArgsEntry\x12\x10\n" +
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\":\n" +
"\bSeverity\x12\b\n" +
"\x04INFO\x10\x00\x12\v\n" +
@@ -7660,7 +7706,7 @@ func file_daemon_proto_rawDescGZIP() []byte {
}
var file_daemon_proto_enumTypes = make([]protoimpl.EnumInfo, 4)
var file_daemon_proto_msgTypes = make([]protoimpl.MessageInfo, 110)
var file_daemon_proto_msgTypes = make([]protoimpl.MessageInfo, 111)
var file_daemon_proto_goTypes = []any{
(LogLevel)(0), // 0: daemon.LogLevel
(ExposeProtocol)(0), // 1: daemon.ExposeProtocol
@@ -7776,16 +7822,17 @@ var file_daemon_proto_goTypes = []any{
nil, // 111: daemon.Network.ResolvedIPsEntry
(*PortInfo_Range)(nil), // 112: daemon.PortInfo.Range
nil, // 113: daemon.SystemEvent.MetadataEntry
(*durationpb.Duration)(nil), // 114: google.protobuf.Duration
(*timestamppb.Timestamp)(nil), // 115: google.protobuf.Timestamp
nil, // 114: daemon.SystemEvent.MessageArgsEntry
(*durationpb.Duration)(nil), // 115: google.protobuf.Duration
(*timestamppb.Timestamp)(nil), // 116: google.protobuf.Timestamp
}
var file_daemon_proto_depIdxs = []int32{
114, // 0: daemon.LoginRequest.dnsRouteInterval:type_name -> google.protobuf.Duration
115, // 0: daemon.LoginRequest.dnsRouteInterval:type_name -> google.protobuf.Duration
25, // 1: daemon.StatusResponse.fullStatus:type_name -> daemon.FullStatus
115, // 2: daemon.StatusResponse.sessionExpiresAt:type_name -> google.protobuf.Timestamp
115, // 3: daemon.PeerState.connStatusUpdate:type_name -> google.protobuf.Timestamp
115, // 4: daemon.PeerState.lastWireguardHandshake:type_name -> google.protobuf.Timestamp
114, // 5: daemon.PeerState.latency:type_name -> google.protobuf.Duration
116, // 2: daemon.StatusResponse.sessionExpiresAt:type_name -> google.protobuf.Timestamp
116, // 3: daemon.PeerState.connStatusUpdate:type_name -> google.protobuf.Timestamp
116, // 4: daemon.PeerState.lastWireguardHandshake:type_name -> google.protobuf.Timestamp
115, // 5: daemon.PeerState.latency:type_name -> google.protobuf.Duration
23, // 6: daemon.SSHServerState.sessions:type_name -> daemon.SSHSessionInfo
20, // 7: daemon.FullStatus.managementState:type_name -> daemon.ManagementState
19, // 8: daemon.FullStatus.signalState:type_name -> daemon.SignalState
@@ -7808,114 +7855,115 @@ var file_daemon_proto_depIdxs = []int32{
54, // 25: daemon.TracePacketResponse.stages:type_name -> daemon.TraceStage
2, // 26: daemon.SystemEvent.severity:type_name -> daemon.SystemEvent.Severity
3, // 27: daemon.SystemEvent.category:type_name -> daemon.SystemEvent.Category
115, // 28: daemon.SystemEvent.timestamp:type_name -> google.protobuf.Timestamp
116, // 28: daemon.SystemEvent.timestamp:type_name -> google.protobuf.Timestamp
113, // 29: daemon.SystemEvent.metadata:type_name -> daemon.SystemEvent.MetadataEntry
57, // 30: daemon.GetEventsResponse.events:type_name -> daemon.SystemEvent
114, // 31: daemon.SetConfigRequest.dnsRouteInterval:type_name -> google.protobuf.Duration
72, // 32: daemon.ListProfilesResponse.profiles:type_name -> daemon.Profile
115, // 33: daemon.WaitExtendAuthSessionResponse.sessionExpiresAt:type_name -> google.protobuf.Timestamp
1, // 34: daemon.ExposeServiceRequest.protocol:type_name -> daemon.ExposeProtocol
104, // 35: daemon.ExposeServiceEvent.ready:type_name -> daemon.ExposeServiceReady
114, // 36: daemon.StartCaptureRequest.duration:type_name -> google.protobuf.Duration
114, // 37: daemon.StartBundleCaptureRequest.timeout:type_name -> google.protobuf.Duration
30, // 38: daemon.Network.ResolvedIPsEntry.value:type_name -> daemon.IPList
5, // 39: daemon.DaemonService.Login:input_type -> daemon.LoginRequest
7, // 40: daemon.DaemonService.WaitSSOLogin:input_type -> daemon.WaitSSOLoginRequest
9, // 41: daemon.DaemonService.Up:input_type -> daemon.UpRequest
11, // 42: daemon.DaemonService.Status:input_type -> daemon.StatusRequest
11, // 43: daemon.DaemonService.SubscribeStatus:input_type -> daemon.StatusRequest
13, // 44: daemon.DaemonService.Down:input_type -> daemon.DownRequest
15, // 45: daemon.DaemonService.GetConfig:input_type -> daemon.GetConfigRequest
26, // 46: daemon.DaemonService.ListNetworks:input_type -> daemon.ListNetworksRequest
28, // 47: daemon.DaemonService.SelectNetworks:input_type -> daemon.SelectNetworksRequest
28, // 48: daemon.DaemonService.DeselectNetworks:input_type -> daemon.SelectNetworksRequest
4, // 49: daemon.DaemonService.ForwardingRules:input_type -> daemon.EmptyRequest
35, // 50: daemon.DaemonService.DebugBundle:input_type -> daemon.DebugBundleRequest
37, // 51: daemon.DaemonService.GetLogLevel:input_type -> daemon.GetLogLevelRequest
39, // 52: daemon.DaemonService.SetLogLevel:input_type -> daemon.SetLogLevelRequest
44, // 53: daemon.DaemonService.ListStates:input_type -> daemon.ListStatesRequest
46, // 54: daemon.DaemonService.CleanState:input_type -> daemon.CleanStateRequest
48, // 55: daemon.DaemonService.DeleteState:input_type -> daemon.DeleteStateRequest
50, // 56: daemon.DaemonService.SetSyncResponsePersistence:input_type -> daemon.SetSyncResponsePersistenceRequest
53, // 57: daemon.DaemonService.TracePacket:input_type -> daemon.TracePacketRequest
105, // 58: daemon.DaemonService.StartCapture:input_type -> daemon.StartCaptureRequest
107, // 59: daemon.DaemonService.StartBundleCapture:input_type -> daemon.StartBundleCaptureRequest
109, // 60: daemon.DaemonService.StopBundleCapture:input_type -> daemon.StopBundleCaptureRequest
56, // 61: daemon.DaemonService.SubscribeEvents:input_type -> daemon.SubscribeRequest
58, // 62: daemon.DaemonService.GetEvents:input_type -> daemon.GetEventsRequest
41, // 63: daemon.DaemonService.RegisterUILog:input_type -> daemon.RegisterUILogRequest
60, // 64: daemon.DaemonService.SwitchProfile:input_type -> daemon.SwitchProfileRequest
62, // 65: daemon.DaemonService.SetConfig:input_type -> daemon.SetConfigRequest
64, // 66: daemon.DaemonService.AddProfile:input_type -> daemon.AddProfileRequest
66, // 67: daemon.DaemonService.RenameProfile:input_type -> daemon.RenameProfileRequest
68, // 68: daemon.DaemonService.RemoveProfile:input_type -> daemon.RemoveProfileRequest
70, // 69: daemon.DaemonService.ListProfiles:input_type -> daemon.ListProfilesRequest
73, // 70: daemon.DaemonService.GetActiveProfile:input_type -> daemon.GetActiveProfileRequest
75, // 71: daemon.DaemonService.Logout:input_type -> daemon.LogoutRequest
79, // 72: daemon.DaemonService.GetFeatures:input_type -> daemon.GetFeaturesRequest
82, // 73: daemon.DaemonService.TriggerUpdate:input_type -> daemon.TriggerUpdateRequest
84, // 74: daemon.DaemonService.GetPeerSSHHostKey:input_type -> daemon.GetPeerSSHHostKeyRequest
86, // 75: daemon.DaemonService.RequestJWTAuth:input_type -> daemon.RequestJWTAuthRequest
88, // 76: daemon.DaemonService.WaitJWTToken:input_type -> daemon.WaitJWTTokenRequest
90, // 77: daemon.DaemonService.RequestExtendAuthSession:input_type -> daemon.RequestExtendAuthSessionRequest
92, // 78: daemon.DaemonService.WaitExtendAuthSession:input_type -> daemon.WaitExtendAuthSessionRequest
94, // 79: daemon.DaemonService.DismissSessionWarning:input_type -> daemon.DismissSessionWarningRequest
96, // 80: daemon.DaemonService.StartCPUProfile:input_type -> daemon.StartCPUProfileRequest
98, // 81: daemon.DaemonService.StopCPUProfile:input_type -> daemon.StopCPUProfileRequest
100, // 82: daemon.DaemonService.GetInstallerResult:input_type -> daemon.InstallerResultRequest
102, // 83: daemon.DaemonService.ExposeService:input_type -> daemon.ExposeServiceRequest
77, // 84: daemon.DaemonService.WailsUIReady:input_type -> daemon.WailsUIReadyRequest
6, // 85: daemon.DaemonService.Login:output_type -> daemon.LoginResponse
8, // 86: daemon.DaemonService.WaitSSOLogin:output_type -> daemon.WaitSSOLoginResponse
10, // 87: daemon.DaemonService.Up:output_type -> daemon.UpResponse
12, // 88: daemon.DaemonService.Status:output_type -> daemon.StatusResponse
12, // 89: daemon.DaemonService.SubscribeStatus:output_type -> daemon.StatusResponse
14, // 90: daemon.DaemonService.Down:output_type -> daemon.DownResponse
16, // 91: daemon.DaemonService.GetConfig:output_type -> daemon.GetConfigResponse
27, // 92: daemon.DaemonService.ListNetworks:output_type -> daemon.ListNetworksResponse
29, // 93: daemon.DaemonService.SelectNetworks:output_type -> daemon.SelectNetworksResponse
29, // 94: daemon.DaemonService.DeselectNetworks:output_type -> daemon.SelectNetworksResponse
34, // 95: daemon.DaemonService.ForwardingRules:output_type -> daemon.ForwardingRulesResponse
36, // 96: daemon.DaemonService.DebugBundle:output_type -> daemon.DebugBundleResponse
38, // 97: daemon.DaemonService.GetLogLevel:output_type -> daemon.GetLogLevelResponse
40, // 98: daemon.DaemonService.SetLogLevel:output_type -> daemon.SetLogLevelResponse
45, // 99: daemon.DaemonService.ListStates:output_type -> daemon.ListStatesResponse
47, // 100: daemon.DaemonService.CleanState:output_type -> daemon.CleanStateResponse
49, // 101: daemon.DaemonService.DeleteState:output_type -> daemon.DeleteStateResponse
51, // 102: daemon.DaemonService.SetSyncResponsePersistence:output_type -> daemon.SetSyncResponsePersistenceResponse
55, // 103: daemon.DaemonService.TracePacket:output_type -> daemon.TracePacketResponse
106, // 104: daemon.DaemonService.StartCapture:output_type -> daemon.CapturePacket
108, // 105: daemon.DaemonService.StartBundleCapture:output_type -> daemon.StartBundleCaptureResponse
110, // 106: daemon.DaemonService.StopBundleCapture:output_type -> daemon.StopBundleCaptureResponse
57, // 107: daemon.DaemonService.SubscribeEvents:output_type -> daemon.SystemEvent
59, // 108: daemon.DaemonService.GetEvents:output_type -> daemon.GetEventsResponse
42, // 109: daemon.DaemonService.RegisterUILog:output_type -> daemon.RegisterUILogResponse
61, // 110: daemon.DaemonService.SwitchProfile:output_type -> daemon.SwitchProfileResponse
63, // 111: daemon.DaemonService.SetConfig:output_type -> daemon.SetConfigResponse
65, // 112: daemon.DaemonService.AddProfile:output_type -> daemon.AddProfileResponse
67, // 113: daemon.DaemonService.RenameProfile:output_type -> daemon.RenameProfileResponse
69, // 114: daemon.DaemonService.RemoveProfile:output_type -> daemon.RemoveProfileResponse
71, // 115: daemon.DaemonService.ListProfiles:output_type -> daemon.ListProfilesResponse
74, // 116: daemon.DaemonService.GetActiveProfile:output_type -> daemon.GetActiveProfileResponse
76, // 117: daemon.DaemonService.Logout:output_type -> daemon.LogoutResponse
80, // 118: daemon.DaemonService.GetFeatures:output_type -> daemon.GetFeaturesResponse
83, // 119: daemon.DaemonService.TriggerUpdate:output_type -> daemon.TriggerUpdateResponse
85, // 120: daemon.DaemonService.GetPeerSSHHostKey:output_type -> daemon.GetPeerSSHHostKeyResponse
87, // 121: daemon.DaemonService.RequestJWTAuth:output_type -> daemon.RequestJWTAuthResponse
89, // 122: daemon.DaemonService.WaitJWTToken:output_type -> daemon.WaitJWTTokenResponse
91, // 123: daemon.DaemonService.RequestExtendAuthSession:output_type -> daemon.RequestExtendAuthSessionResponse
93, // 124: daemon.DaemonService.WaitExtendAuthSession:output_type -> daemon.WaitExtendAuthSessionResponse
95, // 125: daemon.DaemonService.DismissSessionWarning:output_type -> daemon.DismissSessionWarningResponse
97, // 126: daemon.DaemonService.StartCPUProfile:output_type -> daemon.StartCPUProfileResponse
99, // 127: daemon.DaemonService.StopCPUProfile:output_type -> daemon.StopCPUProfileResponse
101, // 128: daemon.DaemonService.GetInstallerResult:output_type -> daemon.InstallerResultResponse
103, // 129: daemon.DaemonService.ExposeService:output_type -> daemon.ExposeServiceEvent
78, // 130: daemon.DaemonService.WailsUIReady:output_type -> daemon.WailsUIReadyResponse
85, // [85:131] is the sub-list for method output_type
39, // [39:85] is the sub-list for method input_type
39, // [39:39] is the sub-list for extension type_name
39, // [39:39] is the sub-list for extension extendee
0, // [0:39] is the sub-list for field type_name
114, // 30: daemon.SystemEvent.messageArgs:type_name -> daemon.SystemEvent.MessageArgsEntry
57, // 31: daemon.GetEventsResponse.events:type_name -> daemon.SystemEvent
115, // 32: daemon.SetConfigRequest.dnsRouteInterval:type_name -> google.protobuf.Duration
72, // 33: daemon.ListProfilesResponse.profiles:type_name -> daemon.Profile
116, // 34: daemon.WaitExtendAuthSessionResponse.sessionExpiresAt:type_name -> google.protobuf.Timestamp
1, // 35: daemon.ExposeServiceRequest.protocol:type_name -> daemon.ExposeProtocol
104, // 36: daemon.ExposeServiceEvent.ready:type_name -> daemon.ExposeServiceReady
115, // 37: daemon.StartCaptureRequest.duration:type_name -> google.protobuf.Duration
115, // 38: daemon.StartBundleCaptureRequest.timeout:type_name -> google.protobuf.Duration
30, // 39: daemon.Network.ResolvedIPsEntry.value:type_name -> daemon.IPList
5, // 40: daemon.DaemonService.Login:input_type -> daemon.LoginRequest
7, // 41: daemon.DaemonService.WaitSSOLogin:input_type -> daemon.WaitSSOLoginRequest
9, // 42: daemon.DaemonService.Up:input_type -> daemon.UpRequest
11, // 43: daemon.DaemonService.Status:input_type -> daemon.StatusRequest
11, // 44: daemon.DaemonService.SubscribeStatus:input_type -> daemon.StatusRequest
13, // 45: daemon.DaemonService.Down:input_type -> daemon.DownRequest
15, // 46: daemon.DaemonService.GetConfig:input_type -> daemon.GetConfigRequest
26, // 47: daemon.DaemonService.ListNetworks:input_type -> daemon.ListNetworksRequest
28, // 48: daemon.DaemonService.SelectNetworks:input_type -> daemon.SelectNetworksRequest
28, // 49: daemon.DaemonService.DeselectNetworks:input_type -> daemon.SelectNetworksRequest
4, // 50: daemon.DaemonService.ForwardingRules:input_type -> daemon.EmptyRequest
35, // 51: daemon.DaemonService.DebugBundle:input_type -> daemon.DebugBundleRequest
37, // 52: daemon.DaemonService.GetLogLevel:input_type -> daemon.GetLogLevelRequest
39, // 53: daemon.DaemonService.SetLogLevel:input_type -> daemon.SetLogLevelRequest
44, // 54: daemon.DaemonService.ListStates:input_type -> daemon.ListStatesRequest
46, // 55: daemon.DaemonService.CleanState:input_type -> daemon.CleanStateRequest
48, // 56: daemon.DaemonService.DeleteState:input_type -> daemon.DeleteStateRequest
50, // 57: daemon.DaemonService.SetSyncResponsePersistence:input_type -> daemon.SetSyncResponsePersistenceRequest
53, // 58: daemon.DaemonService.TracePacket:input_type -> daemon.TracePacketRequest
105, // 59: daemon.DaemonService.StartCapture:input_type -> daemon.StartCaptureRequest
107, // 60: daemon.DaemonService.StartBundleCapture:input_type -> daemon.StartBundleCaptureRequest
109, // 61: daemon.DaemonService.StopBundleCapture:input_type -> daemon.StopBundleCaptureRequest
56, // 62: daemon.DaemonService.SubscribeEvents:input_type -> daemon.SubscribeRequest
58, // 63: daemon.DaemonService.GetEvents:input_type -> daemon.GetEventsRequest
41, // 64: daemon.DaemonService.RegisterUILog:input_type -> daemon.RegisterUILogRequest
60, // 65: daemon.DaemonService.SwitchProfile:input_type -> daemon.SwitchProfileRequest
62, // 66: daemon.DaemonService.SetConfig:input_type -> daemon.SetConfigRequest
64, // 67: daemon.DaemonService.AddProfile:input_type -> daemon.AddProfileRequest
66, // 68: daemon.DaemonService.RenameProfile:input_type -> daemon.RenameProfileRequest
68, // 69: daemon.DaemonService.RemoveProfile:input_type -> daemon.RemoveProfileRequest
70, // 70: daemon.DaemonService.ListProfiles:input_type -> daemon.ListProfilesRequest
73, // 71: daemon.DaemonService.GetActiveProfile:input_type -> daemon.GetActiveProfileRequest
75, // 72: daemon.DaemonService.Logout:input_type -> daemon.LogoutRequest
79, // 73: daemon.DaemonService.GetFeatures:input_type -> daemon.GetFeaturesRequest
82, // 74: daemon.DaemonService.TriggerUpdate:input_type -> daemon.TriggerUpdateRequest
84, // 75: daemon.DaemonService.GetPeerSSHHostKey:input_type -> daemon.GetPeerSSHHostKeyRequest
86, // 76: daemon.DaemonService.RequestJWTAuth:input_type -> daemon.RequestJWTAuthRequest
88, // 77: daemon.DaemonService.WaitJWTToken:input_type -> daemon.WaitJWTTokenRequest
90, // 78: daemon.DaemonService.RequestExtendAuthSession:input_type -> daemon.RequestExtendAuthSessionRequest
92, // 79: daemon.DaemonService.WaitExtendAuthSession:input_type -> daemon.WaitExtendAuthSessionRequest
94, // 80: daemon.DaemonService.DismissSessionWarning:input_type -> daemon.DismissSessionWarningRequest
96, // 81: daemon.DaemonService.StartCPUProfile:input_type -> daemon.StartCPUProfileRequest
98, // 82: daemon.DaemonService.StopCPUProfile:input_type -> daemon.StopCPUProfileRequest
100, // 83: daemon.DaemonService.GetInstallerResult:input_type -> daemon.InstallerResultRequest
102, // 84: daemon.DaemonService.ExposeService:input_type -> daemon.ExposeServiceRequest
77, // 85: daemon.DaemonService.WailsUIReady:input_type -> daemon.WailsUIReadyRequest
6, // 86: daemon.DaemonService.Login:output_type -> daemon.LoginResponse
8, // 87: daemon.DaemonService.WaitSSOLogin:output_type -> daemon.WaitSSOLoginResponse
10, // 88: daemon.DaemonService.Up:output_type -> daemon.UpResponse
12, // 89: daemon.DaemonService.Status:output_type -> daemon.StatusResponse
12, // 90: daemon.DaemonService.SubscribeStatus:output_type -> daemon.StatusResponse
14, // 91: daemon.DaemonService.Down:output_type -> daemon.DownResponse
16, // 92: daemon.DaemonService.GetConfig:output_type -> daemon.GetConfigResponse
27, // 93: daemon.DaemonService.ListNetworks:output_type -> daemon.ListNetworksResponse
29, // 94: daemon.DaemonService.SelectNetworks:output_type -> daemon.SelectNetworksResponse
29, // 95: daemon.DaemonService.DeselectNetworks:output_type -> daemon.SelectNetworksResponse
34, // 96: daemon.DaemonService.ForwardingRules:output_type -> daemon.ForwardingRulesResponse
36, // 97: daemon.DaemonService.DebugBundle:output_type -> daemon.DebugBundleResponse
38, // 98: daemon.DaemonService.GetLogLevel:output_type -> daemon.GetLogLevelResponse
40, // 99: daemon.DaemonService.SetLogLevel:output_type -> daemon.SetLogLevelResponse
45, // 100: daemon.DaemonService.ListStates:output_type -> daemon.ListStatesResponse
47, // 101: daemon.DaemonService.CleanState:output_type -> daemon.CleanStateResponse
49, // 102: daemon.DaemonService.DeleteState:output_type -> daemon.DeleteStateResponse
51, // 103: daemon.DaemonService.SetSyncResponsePersistence:output_type -> daemon.SetSyncResponsePersistenceResponse
55, // 104: daemon.DaemonService.TracePacket:output_type -> daemon.TracePacketResponse
106, // 105: daemon.DaemonService.StartCapture:output_type -> daemon.CapturePacket
108, // 106: daemon.DaemonService.StartBundleCapture:output_type -> daemon.StartBundleCaptureResponse
110, // 107: daemon.DaemonService.StopBundleCapture:output_type -> daemon.StopBundleCaptureResponse
57, // 108: daemon.DaemonService.SubscribeEvents:output_type -> daemon.SystemEvent
59, // 109: daemon.DaemonService.GetEvents:output_type -> daemon.GetEventsResponse
42, // 110: daemon.DaemonService.RegisterUILog:output_type -> daemon.RegisterUILogResponse
61, // 111: daemon.DaemonService.SwitchProfile:output_type -> daemon.SwitchProfileResponse
63, // 112: daemon.DaemonService.SetConfig:output_type -> daemon.SetConfigResponse
65, // 113: daemon.DaemonService.AddProfile:output_type -> daemon.AddProfileResponse
67, // 114: daemon.DaemonService.RenameProfile:output_type -> daemon.RenameProfileResponse
69, // 115: daemon.DaemonService.RemoveProfile:output_type -> daemon.RemoveProfileResponse
71, // 116: daemon.DaemonService.ListProfiles:output_type -> daemon.ListProfilesResponse
74, // 117: daemon.DaemonService.GetActiveProfile:output_type -> daemon.GetActiveProfileResponse
76, // 118: daemon.DaemonService.Logout:output_type -> daemon.LogoutResponse
80, // 119: daemon.DaemonService.GetFeatures:output_type -> daemon.GetFeaturesResponse
83, // 120: daemon.DaemonService.TriggerUpdate:output_type -> daemon.TriggerUpdateResponse
85, // 121: daemon.DaemonService.GetPeerSSHHostKey:output_type -> daemon.GetPeerSSHHostKeyResponse
87, // 122: daemon.DaemonService.RequestJWTAuth:output_type -> daemon.RequestJWTAuthResponse
89, // 123: daemon.DaemonService.WaitJWTToken:output_type -> daemon.WaitJWTTokenResponse
91, // 124: daemon.DaemonService.RequestExtendAuthSession:output_type -> daemon.RequestExtendAuthSessionResponse
93, // 125: daemon.DaemonService.WaitExtendAuthSession:output_type -> daemon.WaitExtendAuthSessionResponse
95, // 126: daemon.DaemonService.DismissSessionWarning:output_type -> daemon.DismissSessionWarningResponse
97, // 127: daemon.DaemonService.StartCPUProfile:output_type -> daemon.StartCPUProfileResponse
99, // 128: daemon.DaemonService.StopCPUProfile:output_type -> daemon.StopCPUProfileResponse
101, // 129: daemon.DaemonService.GetInstallerResult:output_type -> daemon.InstallerResultResponse
103, // 130: daemon.DaemonService.ExposeService:output_type -> daemon.ExposeServiceEvent
78, // 131: daemon.DaemonService.WailsUIReady:output_type -> daemon.WailsUIReadyResponse
86, // [86:132] is the sub-list for method output_type
40, // [40:86] is the sub-list for method input_type
40, // [40:40] is the sub-list for extension type_name
40, // [40:40] is the sub-list for extension extendee
0, // [0:40] is the sub-list for field type_name
}
func init() { file_daemon_proto_init() }
@@ -7947,7 +7995,7 @@ func file_daemon_proto_init() {
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
RawDescriptor: unsafe.Slice(unsafe.StringData(file_daemon_proto_rawDesc), len(file_daemon_proto_rawDesc)),
NumEnums: 4,
NumMessages: 110,
NumMessages: 111,
NumExtensions: 0,
NumServices: 1,
},

View File

@@ -677,9 +677,25 @@ message SystemEvent {
Severity severity = 2;
Category category = 3;
string message = 4;
// userMessage is the daemon's English rendering of messageKey, kept for the
// CLI and for UIs that predate messageKey. UIs that localise read messageKey
// and treat this as the fallback.
string userMessage = 5;
google.protobuf.Timestamp timestamp = 6;
map<string, string> metadata = 7;
// messageKey names a user-facing message in a stable, locale-independent
// form. A UI resolves it against its own translation bundle and substitutes
// messageArgs; an unrecognised key falls back to userMessage. Empty on events
// that carry no user-facing text.
string messageKey = 8;
// messageArgs holds the placeholder name/value pairs for messageKey, e.g.
// {"version": "0.60.1"} for a "{version}" template. Values are data the
// daemon cannot localise (versions, addresses, error text).
map<string, string> messageArgs = 9;
// titleKey names the notification title in the same way as messageKey. Empty
// when the event has no dedicated title, in which case a UI composes one from
// severity and category.
string titleKey = 10;
}
message GetEventsRequest {}

View File

@@ -43,10 +43,10 @@ const (
// UIs to re-fetch their cached config + features. UserMessage is empty so
// the change is silent; the source is carried in MetadataSourceKey.
MetadataTypeConfigChanged = "config_changed"
// MetadataTypePolicyApplied marks an MDM-policy-driven config change. The
// daemon stamps it with a (non-localised) UserMessage; the UI suppresses
// that and builds its own localised toast off the paired config_changed
// event instead.
// MetadataTypePolicyApplied marks an MDM-policy-driven config change. It is
// the user-facing half of the pair: the daemon stamps it with the message
// and title keys a UI localises, while the paired config_changed event stays
// silent and only drives the cache refresh.
MetadataTypePolicyApplied = "policy_applied"
// MetadataSourceKey is the SystemEvent.metadata key carrying what

183
client/proto/usermsg.go Normal file
View File

@@ -0,0 +1,183 @@
package proto
import (
"strings"
log "github.com/sirupsen/logrus"
)
// UserMessageKey names a piece of user-facing SystemEvent text in a
// locale-independent form. The daemon has no notion of the user's language, so it
// publishes the key and lets each UI resolve it.
//
// The value doubles as the lookup key in the desktop UI's translation bundles
// (client/ui/i18n/locales/<code>/common.json), so renaming a constant here means
// renaming the key in every bundle. The tests in client/ui/i18n lock the two
// sides together. Keys that predate this mechanism keep their original bundle
// names so the shipped translations still apply.
type UserMessageKey string
// Message-body keys published by the daemon. Every key needs an entry in
// UserMessageTexts below and a translation in the UI bundles.
const (
// UserMsgPanic is the CRITICAL event published from the recover() guard in
// the connect loop.
UserMsgPanic UserMessageKey = "event.panic"
// UserMsgDNSRecovered and UserMsgDNSUnreachable bracket a nameserver
// group's health transitions.
UserMsgDNSRecovered UserMessageKey = "event.dns.recovered"
UserMsgDNSUnreachable UserMessageKey = "event.dns.unreachable"
// Exit-node (default route) transitions.
UserMsgExitNodeConnected UserMessageKey = "event.exitNode.connected"
UserMsgExitNodeDisconnected UserMessageKey = "event.exitNode.disconnected"
UserMsgExitNodeConnectionLost UserMessageKey = "event.exitNode.connectionLost"
UserMsgExitNodeHAChange UserMessageKey = "event.exitNode.haChange"
UserMsgExitNodeDisconnectedUnknown UserMessageKey = "event.exitNode.disconnectedUnknown"
// Auto-update lifecycle. UserMsgUpdateFailed takes a {reason} argument and
// UserMsgUpdateCompleted a {version}.
UserMsgUpdateInstalling UserMessageKey = "event.update.installing"
UserMsgUpdateCompleted UserMessageKey = "event.update.completed"
UserMsgUpdateFailed UserMessageKey = "event.update.failed"
// UserMsgMDMPolicyApplied reports that an MDM policy replaced the config.
UserMsgMDMPolicyApplied UserMessageKey = "notify.mdm.policyApplied.body"
// Session-expiry events. UserMsgSessionExpiresIn takes a {remaining}
// argument; the "soon" variant is published when the deadline has already
// passed by the time the warning fires.
UserMsgSessionExpiresIn UserMessageKey = "notify.sessionWarning.body"
UserMsgSessionExpiresSoon UserMessageKey = "notify.sessionWarning.bodyGeneric"
UserMsgSessionDeadlineReject UserMessageKey = "notify.sessionDeadlineRejected.body"
)
// Notification-title keys. An event without one leaves titleKey empty and the UI
// composes a title from severity and category, which is what every event did
// before message keys existed. These carry no English fallback: a UI old enough
// to ignore titleKey already builds its own title.
const (
TitleMDMPolicyApplied UserMessageKey = "notify.mdm.policyApplied.title"
TitleSessionWarning UserMessageKey = "notify.sessionWarning.title"
TitleSessionDeadlineReject UserMessageKey = "notify.sessionDeadlineRejected.title"
)
// UserMessageTitleKeys lists every title key a daemon can publish. It sits next
// to the constants above because the two must grow together: a title key missing
// from this slice ships untranslated, and the i18n test that would have caught it
// reads this list.
var UserMessageTitleKeys = []UserMessageKey{
TitleMDMPolicyApplied,
TitleSessionWarning,
TitleSessionDeadlineReject,
}
// ArgReason, ArgVersion and ArgRemaining are the placeholder names used by the
// templates below. Producers pass them to NewUserMessage; the UI bundles use the
// same names inside {}.
const (
ArgReason = "reason"
ArgVersion = "version"
ArgRemaining = "remaining"
)
// UserMessageTexts is the English rendering of every body key, and the source of
// the SystemEvent.userMessage fallback that the CLI and pre-messageKey UIs read.
// It must match the en bundle, which the i18n test enforces.
var UserMessageTexts = map[UserMessageKey]string{
UserMsgPanic: "The NetBird service panicked. Please restart the service and submit a bug report with the client logs.",
UserMsgDNSRecovered: "DNS servers are reachable again.",
UserMsgDNSUnreachable: "Unable to reach one or more DNS servers. This might affect your ability to connect to some services.",
UserMsgExitNodeConnected: "Exit node connected.",
UserMsgExitNodeDisconnected: "Exit node disconnected.",
UserMsgExitNodeConnectionLost: "Exit node connection lost. Your internet access might be affected.",
UserMsgExitNodeHAChange: "Exit node disconnected due to high availability change.",
UserMsgExitNodeDisconnectedUnknown: "Exit node disconnected for unknown reasons.",
UserMsgUpdateInstalling: "Installing update now.",
UserMsgUpdateCompleted: "Your NetBird client was auto-updated to version {version}.",
UserMsgUpdateFailed: "Auto-update failed: {reason}",
UserMsgMDMPolicyApplied: "Your NetBird configuration was updated by your IT policy.",
UserMsgSessionExpiresIn: "Your NetBird session expires in {remaining}. Click Extend now to renew.",
UserMsgSessionExpiresSoon: "Your NetBird session is about to expire. Click Extend now to renew.",
UserMsgSessionDeadlineReject: "The server sent an invalid session deadline. Please sign in again.",
}
// UserMessage is a localizable user-facing event message: a stable body key, an
// optional title key, and the placeholder values to substitute. A nil
// *UserMessage carries no user-facing text, which is how internal control events
// are published.
type UserMessage struct {
key UserMessageKey
title UserMessageKey
args map[string]string
}
// NewUserMessage builds a UserMessage from key and flat placeholder name/value
// pairs, e.g. NewUserMessage(UserMsgUpdateCompleted, ArgVersion, "0.60.1"). An
// unpaired trailing argument is dropped.
func NewUserMessage(key UserMessageKey, args ...string) *UserMessage {
m := &UserMessage{key: key}
if len(args)%2 != 0 {
log.Debugf("user message %q: placeholder args not paired: %d items, last dropped", key, len(args))
args = args[:len(args)-1]
}
if len(args) > 0 {
m.args = make(map[string]string, len(args)/2)
for i := 0; i < len(args); i += 2 {
m.args[args[i]] = args[i+1]
}
}
return m
}
// WithTitle attaches a notification title key and returns m, so it chains onto
// NewUserMessage.
func (m *UserMessage) WithTitle(key UserMessageKey) *UserMessage {
m.title = key
return m
}
// Key returns the body key, or the empty key for a nil message.
func (m *UserMessage) Key() UserMessageKey {
if m == nil {
return ""
}
return m.key
}
// TitleKey returns the title key, which is empty unless WithTitle was called.
func (m *UserMessage) TitleKey() UserMessageKey {
if m == nil {
return ""
}
return m.title
}
// Args returns the placeholder values, or nil for a nil message. The map is the
// message's own and must not be mutated by the caller; it is handed straight to
// SystemEvent.messageArgs.
func (m *UserMessage) Args() map[string]string {
if m == nil {
return nil
}
return m.args
}
// Text renders the English fallback for m, with placeholders substituted. A nil
// message, or a key with no registered template, renders empty so a UI treats
// the event as carrying no user-facing text.
func (m *UserMessage) Text() string {
if m == nil {
return ""
}
tmpl, ok := UserMessageTexts[m.key]
if !ok {
log.Warnf("no English template for user message key %q", m.key)
return ""
}
for name, value := range m.args {
tmpl = strings.ReplaceAll(tmpl, "{"+name+"}", value)
}
return tmpl
}

View File

@@ -0,0 +1,91 @@
package proto
import (
"testing"
"github.com/stretchr/testify/assert"
)
func TestUserMessageText(t *testing.T) {
tests := []struct {
name string
msg *UserMessage
want string
}{
{
name: "no placeholders",
msg: NewUserMessage(UserMsgExitNodeConnected),
want: "Exit node connected.",
},
{
name: "placeholder substituted",
msg: NewUserMessage(UserMsgUpdateCompleted, ArgVersion, "0.60.1"),
want: "Your NetBird client was auto-updated to version 0.60.1.",
},
{
// A dangling arg is a caller mistake; dropping it must still yield a
// readable sentence rather than panicking on the wire path.
name: "unpaired trailing arg dropped",
msg: NewUserMessage(UserMsgUpdateFailed, ArgReason, "disk full", "extra"),
want: "Auto-update failed: disk full",
},
{
name: "unknown key renders empty so the UI treats the event as silent",
msg: NewUserMessage("event.notRegistered"),
want: "",
},
{
name: "nil message carries no text",
msg: nil,
want: "",
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
assert.Equal(t, tc.want, tc.msg.Text())
})
}
}
// A nil *UserMessage is how every internal control event is published, so the
// accessors must stay usable without a nil check at each of the ~25 call sites.
func TestNilUserMessageAccessors(t *testing.T) {
var msg *UserMessage
assert.Empty(t, msg.Key(), "nil message must have no body key")
assert.Empty(t, msg.TitleKey(), "nil message must have no title key")
assert.Nil(t, msg.Args(), "nil message must have no args")
assert.Empty(t, msg.Text(), "nil message must render empty")
}
func TestUserMessageArgsAndTitle(t *testing.T) {
msg := NewUserMessage(UserMsgSessionExpiresIn, ArgRemaining, "10m").
WithTitle(TitleSessionWarning)
assert.Equal(t, UserMsgSessionExpiresIn, msg.Key())
assert.Equal(t, TitleSessionWarning, msg.TitleKey())
assert.Equal(t, map[string]string{ArgRemaining: "10m"}, msg.Args())
assert.Equal(t, "Your NetBird session expires in 10m. Click Extend now to renew.", msg.Text())
}
// Every registered template must be reachable: a key whose text is empty would
// publish an event with a key but no fallback for the CLI and older UIs.
func TestUserMessageTextsAreNonEmpty(t *testing.T) {
texts := UserMessageTexts
assert.NotEmpty(t, texts, "the catalog must not be empty")
for key, text := range texts {
assert.NotEmpty(t, text, "message key %q has an empty template", key)
assert.NotEmpty(t, key, "the catalog must not contain an empty key")
}
}
func TestUserMessageTitleKeysAreUnique(t *testing.T) {
seen := make(map[UserMessageKey]struct{})
for _, key := range UserMessageTitleKeys {
assert.NotEmpty(t, key, "title keys must not be empty")
_, dup := seen[key]
assert.False(t, dup, "title key %q listed twice", key)
seen[key] = struct{}{}
}
}

View File

@@ -1,110 +0,0 @@
package server
import (
"context"
"encoding/json"
"errors"
"os"
"testing"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/client/internal"
"github.com/netbirdio/netbird/client/proto"
)
// A login that never reached Management is not a decision about the peer's
// credentials, so it must come back as a retryable error rather than an SSO
// prompt: the user cannot finish a browser login while Management is down, and
// the CLI's own backoff resolves the outage on its own once the daemon reports
// the failure. Reproduces `netbird down; netbird up` printing a device-code URL
// because Management happened to be restarting when the daemon dialed it.
func TestLogin_ManagementUnreachableIsReturnedInsteadOfDemandingSSO(t *testing.T) {
s, _, _, username, _ := setupServerWithProfile(t)
s.rootCtx = internal.CtxInitState(context.Background())
unreachable := errors.New("create connection: dial context: context deadline exceeded")
attempts := 0
s.isLoginRequiredFn = func(context.Context) (bool, error) {
attempts++
return false, unreachable
}
resp, err := s.Login(userCtx(), &proto.LoginRequest{Username: &username})
require.Error(t, err)
require.ErrorIs(t, err, unreachable, "the transport failure was replaced by something else")
require.Nil(t, resp, "a failed login must not answer with a login response")
require.Equal(t, 1, attempts)
require.Nil(t, s.oauthAuthFlow.flow, "the daemon started an SSO flow for a peer whose login was never decided")
status, err := internal.CtxGetState(s.rootCtx).Status()
require.NoError(t, err)
require.Equal(t, internal.StatusLoginFailed, status,
"a peer that could not reach Management is not waiting on a login")
}
// The counterpart: Management refusing the peer's credentials is a decision, and
// the SSO flow still has to start for it. The profile carries an unusable
// private key so the flow setup fails immediately instead of dialing, which is
// enough to show the branch was entered — the refusal itself is never what comes
// back out.
func TestLogin_AuthRefusalStartsSSOFlow(t *testing.T) {
s, _, _, username, cfgPath := setupServerWithProfile(t)
s.rootCtx = internal.CtxInitState(context.Background())
breakProfilePrivateKey(t, cfgPath)
s.isLoginRequiredFn = func(context.Context) (bool, error) {
return true, nil
}
_, err := s.Login(userCtx(), &proto.LoginRequest{Username: &username})
require.Error(t, err)
status, stateErr := internal.CtxGetState(s.rootCtx).Status()
require.NoError(t, stateErr)
require.Equal(t, internal.StatusLoginFailed, status,
"the SSO flow setup was never reached with the broken key")
}
func TestLogin_SetupKeyStillRunsWhenPeerNeedsLogin(t *testing.T) {
s, _, _, username, _ := setupServerWithProfile(t)
s.rootCtx = internal.CtxInitState(context.Background())
s.isLoginRequiredFn = func(context.Context) (bool, error) {
return true, nil
}
var keysTried []string
s.loginAttemptFn = func(_ context.Context, setupKey, _ string) (internal.StatusType, error) {
keysTried = append(keysTried, setupKey)
return "", nil
}
setupKey := "A2C8E32F-AEB2-4B45-8FD3-8A0C1B2D3E4F"
resp, err := s.Login(userCtx(), &proto.LoginRequest{Username: &username, SetupKey: setupKey})
require.NoError(t, err, "the probe's outcome leaked out as the login result")
require.NotNil(t, resp)
require.Equal(t, []string{setupKey}, keysTried, "the setup key never reached the login attempt")
require.Nil(t, s.oauthAuthFlow.flow, "a setup-key login started an SSO flow")
status, err := internal.CtxGetState(s.rootCtx).Status()
require.NoError(t, err)
require.Equal(t, internal.StatusIdle, status)
}
// breakProfilePrivateKey replaces the profile's private key with an unparseable
// one, which makes any attempt to build a Management client fail on the spot.
func breakProfilePrivateKey(t *testing.T, cfgPath string) {
t.Helper()
raw, err := os.ReadFile(cfgPath)
require.NoError(t, err)
var cfg map[string]any
require.NoError(t, json.Unmarshal(raw, &cfg))
cfg["PrivateKey"] = "not-a-key"
patched, err := json.Marshal(cfg)
require.NoError(t, err)
require.NoError(t, os.WriteFile(cfgPath, patched, 0o600))
}

View File

@@ -94,12 +94,12 @@ func (s *Server) onMDMPolicyChange(_, _ *mdm.Policy) error {
// publishConfigChangedEvent has already fired inside
// restartEngineForMDMLocked with source="mdm". Emit an MDM-specific
// user-visible toast so the operator knows their IT policy was
// applied (UserMessage != "" triggers the GUI notifier).
// applied; the message and title keys let the GUI localise it.
s.statusRecorder.PublishEvent(
proto.SystemEvent_INFO,
proto.SystemEvent_SYSTEM,
"MDM policy applied",
"NetBird configuration was updated by your IT policy.",
proto.NewUserMessage(proto.UserMsgMDMPolicyApplied).WithTitle(proto.TitleMDMPolicyApplied),
map[string]string{
proto.MetadataSourceKey: proto.MetadataSourceMDM,
proto.MetadataTypeKey: proto.MetadataTypePolicyApplied,
@@ -126,7 +126,7 @@ func (s *Server) publishConfigChangedEvent(source string) {
proto.SystemEvent_INFO,
proto.SystemEvent_SYSTEM,
fmt.Sprintf("daemon config changed (source=%s)", source),
"",
nil,
map[string]string{
proto.MetadataSourceKey: source,
proto.MetadataTypeKey: proto.MetadataTypeConfigChanged,

View File

@@ -170,7 +170,7 @@ func (s *Server) SelectNetworks(_ context.Context, req *proto.SelectNetworksRequ
proto.SystemEvent_INFO,
proto.SystemEvent_SYSTEM,
"Network selection changed",
"",
nil,
map[string]string{
"networks": strings.Join(req.GetNetworkIDs(), ", "),
"append": fmt.Sprint(req.GetAppend()),
@@ -214,7 +214,7 @@ func (s *Server) DeselectNetworks(_ context.Context, req *proto.SelectNetworksRe
proto.SystemEvent_INFO,
proto.SystemEvent_SYSTEM,
"Network deselection changed",
"",
nil,
map[string]string{
"networks": strings.Join(req.GetNetworkIDs(), ", "),
"append": fmt.Sprint(req.GetAppend()),
@@ -232,4 +232,3 @@ func toNetIDs(routes []string) []route.NetID {
}
return netIDs
}

View File

@@ -135,13 +135,6 @@ type Server struct {
updateManager *updater.Manager
jwtCache *jwtCache
// loginAttemptFn stands in for the Management login round trip. Tests set
// it to drive the login outcomes that need a server on the other end;
// production leaves it nil, and every login goes through loginAttempt.
loginAttemptFn func(ctx context.Context, setupKey, jwtToken string) (internal.StatusType, error)
isLoginRequiredFn func(ctx context.Context) (bool, error)
}
type oauthAuthFlow struct {
@@ -377,34 +370,7 @@ func (s *Server) connectionGoroutineRunning() bool {
}
}
// attemptLogin runs a login round trip against Management, or the stand-in a
// test installed in place of it.
func (s *Server) attemptLogin(ctx context.Context, setupKey, jwtToken string) (internal.StatusType, error) {
if s.loginAttemptFn != nil {
return s.loginAttemptFn(ctx, setupKey, jwtToken)
}
return s.loginAttempt(ctx, setupKey, jwtToken)
}
func (s *Server) isLoginRequired(ctx context.Context) (bool, error) {
if s.isLoginRequiredFn != nil {
return s.isLoginRequiredFn(ctx)
}
authClient, err := auth.NewAuth(ctx, s.config.PrivateKey, s.config.ManagementURL, s.config)
if err != nil {
log.Errorf("failed to create auth client: %v", err)
return false, err
}
defer authClient.Close()
return authClient.IsLoginRequired(ctx)
}
// loginAttempt attempts to login using the provided information. It returns
// StatusNeedsLogin when Management refused the peer's credentials and
// StatusLoginFailed for every other failure, so callers can tell an
// authentication decision apart from a login that never got made.
// loginAttempt attempts to login using the provided information. it returns a status in case something fails
func (s *Server) loginAttempt(ctx context.Context, setupKey, jwtToken string) (internal.StatusType, error) {
authClient, err := auth.NewAuth(ctx, s.config.PrivateKey, s.config.ManagementURL, s.config)
if err != nil {
@@ -657,19 +623,7 @@ func (s *Server) Login(callerCtx context.Context, msg *proto.LoginRequest) (*pro
s.config = config
s.mutex.Unlock()
// A probe that errors leaves the login undecided: Management unreachable, a
// restart mid-request, an internal error. Those are returned for the caller
// to retry, because turning them into an SSO prompt asks the user to solve
// something that is not theirs to solve, and a browser login cannot succeed
// while Management is unreachable anyway. Only Management refusing the
// peer's key is a decision, and IsLoginRequired reports that as
// needsLogin=true rather than an error.
needsLogin, err := s.isLoginRequired(ctx)
if err != nil {
state.Set(internal.StatusLoginFailed)
return nil, err
}
if !needsLogin {
if _, err := s.loginAttempt(ctx, "", ""); err == nil {
state.Set(internal.StatusIdle)
return &proto.LoginResponse{}, nil
}
@@ -730,7 +684,7 @@ func (s *Server) Login(callerCtx context.Context, msg *proto.LoginRequest) (*pro
// which returns NeedsLogin and parks on the browser leg.
state.Set(internal.StatusConnecting)
if loginStatus, err := s.attemptLogin(ctx, msg.SetupKey, ""); err != nil {
if loginStatus, err := s.loginAttempt(ctx, msg.SetupKey, ""); err != nil {
state.Set(loginStatus)
return nil, err
}
@@ -885,7 +839,7 @@ func (s *Server) WaitSSOLogin(callerCtx context.Context, msg *proto.WaitSSOLogin
s.oauthAuthFlow.expiresAt = time.Now()
s.mutex.Unlock()
if loginStatus, err := s.attemptLogin(ctx, "", tokenInfo.GetTokenToUse()); err != nil {
if loginStatus, err := s.loginAttempt(ctx, "", tokenInfo.GetTokenToUse()); err != nil {
state.Set(loginStatus)
return nil, err
}
@@ -1815,9 +1769,6 @@ func (s *Server) RequestExtendAuthSession(
if connectClient == nil {
return nil, gstatus.Errorf(codes.FailedPrecondition, "client is not running")
}
if connectClient.Engine() == nil {
return nil, gstatus.Errorf(codes.FailedPrecondition, "session can no longer be extended, log in again to reconnect")
}
hint := ""
if msg.Hint != nil {
@@ -2230,7 +2181,7 @@ func (s *Server) publishProfileListChanged(profileName string) {
proto.SystemEvent_INFO,
proto.SystemEvent_SYSTEM,
"Profile list changed",
"",
nil,
map[string]string{proto.MetadataKindKey: proto.MetadataKindProfileListChanged, proto.MetadataProfileKey: profileName},
)
}
@@ -2249,7 +2200,7 @@ func (s *Server) publishLogLevelChanged(level string) {
proto.SystemEvent_INFO,
proto.SystemEvent_SYSTEM,
"Log level changed",
"",
nil,
map[string]string{proto.MetadataKindKey: proto.MetadataKindLogLevelChanged, proto.MetadataLevelKey: level},
)
}

View File

@@ -26,17 +26,17 @@ contents:
# Default dependencies for the GTK4 + WebKitGTK 6.0 stack (Ubuntu 24.04+ / Debian 13+)
depends:
- libgtk-4-1 (>= 4.14)
- libgtk-4-1
- libwebkitgtk-6.0-4
- xdg-utils
# Distribution-specific overrides for different package formats
overrides:
# RPM packages for Fedora / RHEL / AlmaLinux / Rocky Linux / openSUSE
# RPM packages for Fedora / RHEL / AlmaLinux / Rocky Linux
rpm:
depends:
- (gtk4 >= 4.14 or libgtk-4-1 >= 4.14)
- (webkitgtk6.0 or libwebkitgtk-6_0-4)
- gtk4
- webkitgtk6.0
- xdg-utils
# Arch Linux packages

View File

@@ -43,12 +43,7 @@ function buildSsoCancelPromise(state: SsoState, signal?: AbortSignal): Promise<v
}
async function runSsoLogin(
result: {
verificationUri: string;
verificationUriComplete: string;
userCode: string;
profileId: string;
},
result: { verificationUri: string; verificationUriComplete: string; userCode: string },
state: SsoState,
signal?: AbortSignal,
): Promise<void> {
@@ -61,7 +56,7 @@ async function runSsoLogin(
// suspended, so a frontend-driven Up (a promise continuation) would not
// fire until the user woke the window (e.g. hovering the tray icon).
const waitPromise = Connection.WaitSSOLoginAndUp(
{ userCode: result.userCode, hostname: "", profileId: result.profileId },
{ userCode: result.userCode, hostname: "" },
{ profileName: "", username: "" },
);

View File

@@ -11,7 +11,7 @@ import { DialogHeading } from "@/components/dialog/DialogHeading";
import { SquareIcon } from "@/components/SquareIcon";
import { Connection, Profiles as ProfilesSvc, Session, WindowManager } from "@bindings/services";
import { useAutoSizeWindow } from "@/hooks/useAutoSizeWindow";
import { EVENT_BROWSER_LOGIN_CANCEL, EVENT_TRIGGER_LOGIN } from "@/lib/connection";
import { EVENT_BROWSER_LOGIN_CANCEL } from "@/lib/connection";
import { errorDialog, formatErrorMessage } from "@/lib/errors.ts";
import { formatRemaining } from "@/lib/formatters";
@@ -131,21 +131,6 @@ export default function SessionExpirationDialog() {
}
}, [busy, t]);
const authenticate = useCallback(async () => {
if (busy) return;
setBusy(true);
try {
await Events.Emit(EVENT_TRIGGER_LOGIN);
await WindowManager.CloseSessionExpiration();
} catch (e) {
setBusy(false);
await errorDialog({
Title: t("connect.error.loginTitle"),
Message: formatErrorMessage(e),
});
}
}, [busy, t]);
const logout = useCallback(async () => {
if (busy) return;
setBusy(true);
@@ -200,7 +185,7 @@ export default function SessionExpirationDialog() {
variant={"primary"}
size={"md"}
className={"w-full"}
onClick={expired ? authenticate : stay}
onClick={stay}
disabled={busy}
>
{expired ? t("sessionExpiration.authenticate") : t("sessionExpiration.stay")}

View File

@@ -40,6 +40,12 @@ i18n/locales/<code>/common.json a target — message only
Chrome-extension JSON, each key → `{ "message", "description" }`. You translate the **`message`**.
The `event.*` keys are a special group: the background service names them when it
publishes a notification, and the app looks them up here. Their names are part of
a Go↔JSON contract (`client/proto/usermsg.go`), so they are even less renameable
than the rest — and a missing one shows the user English. Tests fail the build if
any locale drops one.
| ✅ Do | ❌ Don't |
|---|---|
| Keep **every key** from `en`, in the same order | Translate, rename, reorder, drop, or add keys (they're identifiers; the set grows over time) |

View File

@@ -126,18 +126,29 @@ func (b *Bundle) BundleFor(code LanguageCode) (map[string]string, error) {
// pairs ("version", "1.2.3" replaces "{version}"). Unknown keys fall back to
// the default language, then to the key itself so a miss is visible in the UI.
func (b *Bundle) Translate(lang LanguageCode, key string, args ...string) string {
if v, ok := b.Lookup(lang, key, args...); ok {
return v
}
return key
}
// Lookup resolves key like Translate but reports whether it was found in the
// requested or the default bundle. Callers holding a better fallback than the
// raw key — a daemon-supplied English string for a key this build predates —
// use this to tell a miss from a hit.
func (b *Bundle) Lookup(lang LanguageCode, key string, args ...string) (string, bool) {
b.mu.RLock()
defer b.mu.RUnlock()
if v, ok := b.bundles[lang][key]; ok {
return applyPlaceholders(v, args)
return applyPlaceholders(v, args), true
}
if lang != DefaultLanguage {
if v, ok := b.bundles[DefaultLanguage][key]; ok {
return applyPlaceholders(v, args)
return applyPlaceholders(v, args), true
}
}
return key
return "", false
}
// applyPlaceholders substitutes {name} in s using args as flat name/value

View File

@@ -0,0 +1,126 @@
//go:build !android && !ios && !freebsd && !js
package i18n
import (
"os"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/client/proto"
)
// shippedBundle loads the real locale tree rather than the fstest fixture the
// other tests use: these checks exist to catch a key the daemon publishes but no
// bundle translates, which only the shipped files can prove.
func shippedBundle(t *testing.T) *Bundle {
t.Helper()
b, err := NewBundle(os.DirFS("locales"))
require.NoError(t, err, "the shipped locale tree must load")
return b
}
// The daemon publishes a message key and each UI resolves it locally, so a key
// with no en entry degrades to the daemon's English fallback and silently stops
// being translatable. Fail the build instead.
func TestUserMessageKeysExistInEnglishBundle(t *testing.T) {
b := shippedBundle(t)
for key, text := range proto.UserMessageTexts {
got, ok := b.Lookup(DefaultLanguage, string(key))
if !assert.True(t, ok, "message key %q has no en translation", key) {
continue
}
assert.Equal(t, text, got,
"en translation of %q must match the daemon's English fallback", key)
}
}
func TestUserMessageTitleKeysExistInEnglishBundle(t *testing.T) {
b := shippedBundle(t)
for _, key := range proto.UserMessageTitleKeys {
_, ok := b.Lookup(DefaultLanguage, string(key))
assert.True(t, ok, "title key %q has no en translation", key)
}
}
// Every shipped locale must translate the daemon's keys, not just en. A missing
// one still renders (Lookup falls back to en) but the notification would show up
// in English for that user, which is the bug this whole mechanism exists to fix.
func TestUserMessageKeysTranslatedInEveryLanguage(t *testing.T) {
b := shippedBundle(t)
keys := make([]proto.UserMessageKey, 0, len(proto.UserMessageTexts))
for key := range proto.UserMessageTexts {
keys = append(keys, key)
}
keys = append(keys, proto.UserMessageTitleKeys...)
for _, lang := range b.Languages() {
bundle, err := b.BundleFor(lang.Code)
require.NoError(t, err, "BundleFor(%q)", lang.Code)
for _, key := range keys {
text, ok := bundle[string(key)]
if !assert.True(t, ok, "locale %q is missing key %q", lang.Code, key) {
continue
}
assert.NotEmpty(t, text, "locale %q has an empty message for %q", lang.Code, key)
}
}
}
// The tray composes a title from these when an event carries no title key, so a
// gap here would render "event.severity.warning: DNS" to the user.
func TestEventTitleKeysTranslatedInEveryLanguage(t *testing.T) {
b := shippedBundle(t)
keys := []string{
"event.title",
"event.severity.info", "event.severity.warning",
"event.severity.error", "event.severity.critical",
"event.category.network", "event.category.dns",
"event.category.authentication", "event.category.connectivity",
"event.category.system",
}
for _, lang := range b.Languages() {
bundle, err := b.BundleFor(lang.Code)
require.NoError(t, err, "BundleFor(%q)", lang.Code)
for _, key := range keys {
text, ok := bundle[key]
if !assert.True(t, ok, "locale %q is missing key %q", lang.Code, key) {
continue
}
assert.NotEmpty(t, text, "locale %q has an empty message for %q", lang.Code, key)
}
}
// The composed title is useless without both slots.
title, ok := b.Lookup(DefaultLanguage, "event.title", "severity", "Warning", "category", "DNS")
require.True(t, ok)
assert.Equal(t, "Warning: DNS", title, "event.title must substitute both placeholders")
}
func TestBundleLookupReportsMisses(t *testing.T) {
b, err := NewBundle(fakeLocales())
require.NoError(t, err)
got, ok := b.Lookup("en", "tray.menu.connect")
assert.True(t, ok)
assert.Equal(t, "Connect", got)
// An absent key must report a miss rather than echo the key, so callers can
// substitute their own fallback.
got, ok = b.Lookup("en", "tray.missing")
assert.False(t, ok, "unknown key must report a miss")
assert.Empty(t, got, "a miss must not return the key")
// Empty keys reach Lookup from events that carry no title key at all.
_, ok = b.Lookup("en", "")
assert.False(t, ok, "empty key must report a miss")
}

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "Ihre NetBird-Konfiguration wurde durch Ihre IT-Richtlinie aktualisiert."
},
"event.title": {
"message": "{severity}: {category}"
},
"event.severity.info": {
"message": "Info"
},
"event.severity.warning": {
"message": "Warnung"
},
"event.severity.error": {
"message": "Fehler"
},
"event.severity.critical": {
"message": "Kritisch"
},
"event.category.network": {
"message": "Netzwerk"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "Authentifizierung"
},
"event.category.connectivity": {
"message": "Konnektivität"
},
"event.category.system": {
"message": "System"
},
"event.panic": {
"message": "Der NetBird-Dienst ist abgestürzt. Bitte starten Sie den Dienst neu und senden Sie einen Fehlerbericht mit den Client-Protokollen."
},
"event.dns.recovered": {
"message": "DNS-Server sind wieder erreichbar."
},
"event.dns.unreachable": {
"message": "Ein oder mehrere DNS-Server sind nicht erreichbar. Das kann die Verbindung zu einigen Diensten beeinträchtigen."
},
"event.exitNode.connected": {
"message": "Exit Node verbunden."
},
"event.exitNode.disconnected": {
"message": "Exit Node getrennt."
},
"event.exitNode.connectionLost": {
"message": "Verbindung zum Exit Node verloren. Ihr Internetzugang kann beeinträchtigt sein."
},
"event.exitNode.haChange": {
"message": "Exit Node aufgrund einer Änderung der Hochverfügbarkeit getrennt."
},
"event.exitNode.disconnectedUnknown": {
"message": "Exit Node aus unbekannten Gründen getrennt."
},
"event.update.installing": {
"message": "Update wird jetzt installiert."
},
"event.update.completed": {
"message": "Ihr NetBird-Client wurde automatisch auf Version {version} aktualisiert."
},
"event.update.failed": {
"message": "Automatisches Update fehlgeschlagen: {reason}"
},
"common.cancel": {
"message": "Abbrechen"
},

View File

@@ -239,6 +239,90 @@
"message": "Your NetBird configuration was updated by your IT policy.",
"description": "Body of the MDM policy-applied notification, telling the user their settings were changed by their organization's device-management policy."
},
"event.title": {
"message": "{severity}: {category}",
"description": "Notification title for a daemon event, composed from severity and category, e.g. \"Warning: DNS\". Keep both placeholders; use your locale's colon spacing."
},
"event.severity.info": {
"message": "Info",
"description": "Severity label used in the {severity} slot of event.title. Keep it short."
},
"event.severity.warning": {
"message": "Warning",
"description": "Severity label used in the {severity} slot of event.title. Keep it short."
},
"event.severity.error": {
"message": "Error",
"description": "Severity label used in the {severity} slot of event.title. Keep it short."
},
"event.severity.critical": {
"message": "Critical",
"description": "Severity label used in the {severity} slot of event.title. Keep it short."
},
"event.category.network": {
"message": "Network",
"description": "Event category label used in the {category} slot of event.title. Refers to the overlay network."
},
"event.category.dns": {
"message": "DNS",
"description": "Event category label used in the {category} slot of event.title. Acronym, do not translate."
},
"event.category.authentication": {
"message": "Authentication",
"description": "Event category label used in the {category} slot of event.title. Refers to signing in to the management server."
},
"event.category.connectivity": {
"message": "Connectivity",
"description": "Event category label used in the {category} slot of event.title. Refers to reaching peers."
},
"event.category.system": {
"message": "System",
"description": "Event category label used in the {category} slot of event.title. Refers to the local machine and the NetBird service."
},
"event.panic": {
"message": "The NetBird service panicked. Please restart the service and submit a bug report with the client logs.",
"description": "Notification body after the NetBird background service crashed. \"Service\" is the daemon, not a remote service."
},
"event.dns.recovered": {
"message": "DNS servers are reachable again.",
"description": "Notification body when previously unreachable upstream DNS servers respond again."
},
"event.dns.unreachable": {
"message": "Unable to reach one or more DNS servers. This might affect your ability to connect to some services.",
"description": "Notification body when one or more upstream DNS servers stop responding."
},
"event.exitNode.connected": {
"message": "Exit node connected.",
"description": "Notification body when a full-tunnel exit node becomes active."
},
"event.exitNode.disconnected": {
"message": "Exit node disconnected.",
"description": "Notification body when the user or the client shuts the exit node down deliberately."
},
"event.exitNode.connectionLost": {
"message": "Exit node connection lost. Your internet access might be affected.",
"description": "Notification body when the exit node peer became unreachable. \"Internet access\" means the user's own browsing."
},
"event.exitNode.haChange": {
"message": "Exit node disconnected due to high availability change.",
"description": "Notification body when a high-availability group switched away from this exit node. High availability is the standard IT term."
},
"event.exitNode.disconnectedUnknown": {
"message": "Exit node disconnected for unknown reasons.",
"description": "Notification body when the exit node dropped for a reason the client could not classify."
},
"event.update.installing": {
"message": "Installing update now.",
"description": "Notification body shown as an automatic client update starts installing."
},
"event.update.completed": {
"message": "Your NetBird client was auto-updated to version {version}.",
"description": "Notification body after an automatic client update succeeded. {version} is a version number, keep verbatim."
},
"event.update.failed": {
"message": "Auto-update failed: {reason}",
"description": "Notification body when an automatic client update failed. {reason} is an untranslated technical error string; keep your locale's colon spacing."
},
"common.cancel": {
"message": "Cancel",
"description": "Generic Cancel button label, reused across dialogs. Keep short."

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "Su configuración de NetBird fue actualizada por su política de TI."
},
"event.title": {
"message": "{severity}: {category}"
},
"event.severity.info": {
"message": "Información"
},
"event.severity.warning": {
"message": "Advertencia"
},
"event.severity.error": {
"message": "Error"
},
"event.severity.critical": {
"message": "Crítico"
},
"event.category.network": {
"message": "Red"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "Autenticación"
},
"event.category.connectivity": {
"message": "Conectividad"
},
"event.category.system": {
"message": "Sistema"
},
"event.panic": {
"message": "El servicio de NetBird falló de forma inesperada. Reinicie el servicio y envíe un informe de error con los registros del cliente."
},
"event.dns.recovered": {
"message": "Los servidores DNS vuelven a estar accesibles."
},
"event.dns.unreachable": {
"message": "No se puede acceder a uno o más servidores DNS. Esto puede afectar la conexión a algunos servicios."
},
"event.exitNode.connected": {
"message": "Nodo de salida conectado."
},
"event.exitNode.disconnected": {
"message": "Nodo de salida desconectado."
},
"event.exitNode.connectionLost": {
"message": "Se perdió la conexión con el nodo de salida. Su acceso a Internet puede verse afectado."
},
"event.exitNode.haChange": {
"message": "Nodo de salida desconectado por un cambio de alta disponibilidad."
},
"event.exitNode.disconnectedUnknown": {
"message": "Nodo de salida desconectado por motivos desconocidos."
},
"event.update.installing": {
"message": "Instalando la actualización ahora."
},
"event.update.completed": {
"message": "Su cliente de NetBird se actualizó automáticamente a la versión {version}."
},
"event.update.failed": {
"message": "La actualización automática falló: {reason}"
},
"common.cancel": {
"message": "Cancelar"
},

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "Votre configuration NetBird a été mise à jour par votre politique informatique."
},
"event.title": {
"message": "{severity} : {category}"
},
"event.severity.info": {
"message": "Info"
},
"event.severity.warning": {
"message": "Avertissement"
},
"event.severity.error": {
"message": "Erreur"
},
"event.severity.critical": {
"message": "Critique"
},
"event.category.network": {
"message": "Réseau"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "Authentification"
},
"event.category.connectivity": {
"message": "Connectivité"
},
"event.category.system": {
"message": "Système"
},
"event.panic": {
"message": "Le service NetBird s'est arrêté brutalement. Veuillez redémarrer le service et envoyer un rapport de bug avec les journaux du client."
},
"event.dns.recovered": {
"message": "Les serveurs DNS sont de nouveau joignables."
},
"event.dns.unreachable": {
"message": "Impossible de joindre un ou plusieurs serveurs DNS. Cela peut affecter la connexion à certains services."
},
"event.exitNode.connected": {
"message": "Nœud de sortie connecté."
},
"event.exitNode.disconnected": {
"message": "Nœud de sortie déconnecté."
},
"event.exitNode.connectionLost": {
"message": "Connexion au nœud de sortie perdue. Votre accès à Internet peut être affecté."
},
"event.exitNode.haChange": {
"message": "Nœud de sortie déconnecté suite à un changement de haute disponibilité."
},
"event.exitNode.disconnectedUnknown": {
"message": "Nœud de sortie déconnecté pour une raison inconnue."
},
"event.update.installing": {
"message": "Installation de la mise à jour en cours."
},
"event.update.completed": {
"message": "Votre client NetBird a été mis à jour automatiquement vers la version {version}."
},
"event.update.failed": {
"message": "Échec de la mise à jour automatique : {reason}"
},
"common.cancel": {
"message": "Annuler"
},

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "A NetBird konfigurációt az IT-szabályzat frissítette."
},
"event.title": {
"message": "{severity}: {category}"
},
"event.severity.info": {
"message": "Információ"
},
"event.severity.warning": {
"message": "Figyelmeztetés"
},
"event.severity.error": {
"message": "Hiba"
},
"event.severity.critical": {
"message": "Kritikus"
},
"event.category.network": {
"message": "Hálózat"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "Hitelesítés"
},
"event.category.connectivity": {
"message": "Kapcsolat"
},
"event.category.system": {
"message": "Rendszer"
},
"event.panic": {
"message": "A NetBird szolgáltatás összeomlott. Kérjük, indítsa újra a szolgáltatást, és küldjön hibajelentést a kliens naplóival."
},
"event.dns.recovered": {
"message": "A DNS-kiszolgálók ismét elérhetők."
},
"event.dns.unreachable": {
"message": "Egy vagy több DNS-kiszolgáló nem érhető el. Ez befolyásolhatja egyes szolgáltatások elérését."
},
"event.exitNode.connected": {
"message": "Exit Node csatlakoztatva."
},
"event.exitNode.disconnected": {
"message": "Exit Node leválasztva."
},
"event.exitNode.connectionLost": {
"message": "Megszakadt a kapcsolat az Exit Node-dal. Ez érintheti az internetelérést."
},
"event.exitNode.haChange": {
"message": "Az Exit Node leválasztva a magas rendelkezésre állás változása miatt."
},
"event.exitNode.disconnectedUnknown": {
"message": "Az Exit Node ismeretlen okból leválasztva."
},
"event.update.installing": {
"message": "A frissítés telepítése folyamatban."
},
"event.update.completed": {
"message": "A NetBird kliens automatikusan a {version} verzióra frissült."
},
"event.update.failed": {
"message": "Az automatikus frissítés sikertelen: {reason}"
},
"common.cancel": {
"message": "Mégse"
},

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "La configurazione di NetBird è stata aggiornata dalla policy IT."
},
"event.title": {
"message": "{severity}: {category}"
},
"event.severity.info": {
"message": "Info"
},
"event.severity.warning": {
"message": "Avviso"
},
"event.severity.error": {
"message": "Errore"
},
"event.severity.critical": {
"message": "Critico"
},
"event.category.network": {
"message": "Rete"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "Autenticazione"
},
"event.category.connectivity": {
"message": "Connettività"
},
"event.category.system": {
"message": "Sistema"
},
"event.panic": {
"message": "Il servizio NetBird si è arrestato in modo anomalo. Riavvii il servizio e invii una segnalazione di bug con i log del client."
},
"event.dns.recovered": {
"message": "I server DNS sono di nuovo raggiungibili."
},
"event.dns.unreachable": {
"message": "Impossibile raggiungere uno o più server DNS. Questo potrebbe influire sulla connessione ad alcuni servizi."
},
"event.exitNode.connected": {
"message": "Nodo di uscita connesso."
},
"event.exitNode.disconnected": {
"message": "Nodo di uscita disconnesso."
},
"event.exitNode.connectionLost": {
"message": "Connessione al nodo di uscita perduta. L'accesso a Internet potrebbe essere compromesso."
},
"event.exitNode.haChange": {
"message": "Nodo di uscita disconnesso a causa di una modifica dell'alta disponibilità."
},
"event.exitNode.disconnectedUnknown": {
"message": "Nodo di uscita disconnesso per motivi sconosciuti."
},
"event.update.installing": {
"message": "Installazione dell'aggiornamento in corso."
},
"event.update.completed": {
"message": "Il client NetBird è stato aggiornato automaticamente alla versione {version}."
},
"event.update.failed": {
"message": "Aggiornamento automatico non riuscito: {reason}"
},
"common.cancel": {
"message": "Annulla"
},

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "NetBird の構成が IT ポリシーによって更新されました。"
},
"event.title": {
"message": "{severity}: {category}"
},
"event.severity.info": {
"message": "情報"
},
"event.severity.warning": {
"message": "警告"
},
"event.severity.error": {
"message": "エラー"
},
"event.severity.critical": {
"message": "重大"
},
"event.category.network": {
"message": "ネットワーク"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "認証"
},
"event.category.connectivity": {
"message": "接続"
},
"event.category.system": {
"message": "システム"
},
"event.panic": {
"message": "NetBird サービスがクラッシュしました。サービスを再起動し、クライアントログを添えてバグを報告してください。"
},
"event.dns.recovered": {
"message": "DNS サーバーに再び到達できるようになりました。"
},
"event.dns.unreachable": {
"message": "1 つ以上の DNS サーバーに到達できません。一部のサービスへの接続に影響する可能性があります。"
},
"event.exitNode.connected": {
"message": "出口ノードに接続しました。"
},
"event.exitNode.disconnected": {
"message": "出口ノードの接続を解除しました。"
},
"event.exitNode.connectionLost": {
"message": "出口ノードとの接続が失われました。インターネット接続に影響する可能性があります。"
},
"event.exitNode.haChange": {
"message": "高可用性の変更により出口ノードの接続が解除されました。"
},
"event.exitNode.disconnectedUnknown": {
"message": "不明な理由により出口ノードの接続が解除されました。"
},
"event.update.installing": {
"message": "更新をインストールしています。"
},
"event.update.completed": {
"message": "NetBird クライアントがバージョン {version} に自動更新されました。"
},
"event.update.failed": {
"message": "自動更新に失敗しました: {reason}"
},
"common.cancel": {
"message": "キャンセル"
},

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "A sua configuração do NetBird foi atualizada pela política de TI."
},
"event.title": {
"message": "{severity}: {category}"
},
"event.severity.info": {
"message": "Informação"
},
"event.severity.warning": {
"message": "Aviso"
},
"event.severity.error": {
"message": "Erro"
},
"event.severity.critical": {
"message": "Crítico"
},
"event.category.network": {
"message": "Rede"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "Autenticação"
},
"event.category.connectivity": {
"message": "Conectividade"
},
"event.category.system": {
"message": "Sistema"
},
"event.panic": {
"message": "O serviço NetBird falhou de forma inesperada. Reinicie o serviço e envie um relatório de erro com os registros do cliente."
},
"event.dns.recovered": {
"message": "Os servidores DNS estão novamente acessíveis."
},
"event.dns.unreachable": {
"message": "Não é possível acessar um ou mais servidores DNS. Isto pode afetar a conexão a alguns serviços."
},
"event.exitNode.connected": {
"message": "Nó de saída conectado."
},
"event.exitNode.disconnected": {
"message": "Nó de saída desconectado."
},
"event.exitNode.connectionLost": {
"message": "Conexão com o nó de saída perdida. O seu acesso à Internet pode ser afetado."
},
"event.exitNode.haChange": {
"message": "Nó de saída desconectado devido a uma alteração de alta disponibilidade."
},
"event.exitNode.disconnectedUnknown": {
"message": "Nó de saída desconectado por motivos desconhecidos."
},
"event.update.installing": {
"message": "Instalando a atualização agora."
},
"event.update.completed": {
"message": "O seu cliente NetBird foi atualizado automaticamente para a versão {version}."
},
"event.update.failed": {
"message": "Falha na atualização automática: {reason}"
},
"common.cancel": {
"message": "Cancelar"
},

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "Конфигурация NetBird была обновлена в соответствии с вашей ИТ-политикой."
},
"event.title": {
"message": "{severity}: {category}"
},
"event.severity.info": {
"message": "Информация"
},
"event.severity.warning": {
"message": "Предупреждение"
},
"event.severity.error": {
"message": "Ошибка"
},
"event.severity.critical": {
"message": "Критично"
},
"event.category.network": {
"message": "Сеть"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "Аутентификация"
},
"event.category.connectivity": {
"message": "Связь"
},
"event.category.system": {
"message": "Система"
},
"event.panic": {
"message": "Служба NetBird аварийно завершилась. Перезапустите службу и отправьте отчёт об ошибке с журналами клиента."
},
"event.dns.recovered": {
"message": "DNS-серверы снова доступны."
},
"event.dns.unreachable": {
"message": "Не удалось связаться с одним или несколькими DNS-серверами. Это может повлиять на подключение к некоторым сервисам."
},
"event.exitNode.connected": {
"message": "Выходной узел подключён."
},
"event.exitNode.disconnected": {
"message": "Выходной узел отключён."
},
"event.exitNode.connectionLost": {
"message": "Соединение с выходным узлом потеряно. Доступ в интернет может быть нарушен."
},
"event.exitNode.haChange": {
"message": "Выходной узел отключён из-за изменения конфигурации высокой доступности."
},
"event.exitNode.disconnectedUnknown": {
"message": "Выходной узел отключён по неизвестной причине."
},
"event.update.installing": {
"message": "Устанавливается обновление."
},
"event.update.completed": {
"message": "Клиент NetBird автоматически обновлён до версии {version}."
},
"event.update.failed": {
"message": "Не удалось выполнить автоматическое обновление: {reason}"
},
"common.cancel": {
"message": "Отмена"
},

View File

@@ -179,6 +179,69 @@
"notify.mdm.policyApplied.body": {
"message": "您的 NetBird 配置已根据 IT 策略更新。"
},
"event.title": {
"message": "{severity}{category}"
},
"event.severity.info": {
"message": "信息"
},
"event.severity.warning": {
"message": "警告"
},
"event.severity.error": {
"message": "错误"
},
"event.severity.critical": {
"message": "严重"
},
"event.category.network": {
"message": "网络"
},
"event.category.dns": {
"message": "DNS"
},
"event.category.authentication": {
"message": "身份验证"
},
"event.category.connectivity": {
"message": "连接"
},
"event.category.system": {
"message": "系统"
},
"event.panic": {
"message": "NetBird 服务发生崩溃。请重启该服务,并附上客户端日志提交错误报告。"
},
"event.dns.recovered": {
"message": "DNS 服务器已恢复可访问。"
},
"event.dns.unreachable": {
"message": "无法访问一个或多个 DNS 服务器。这可能影响您连接部分服务。"
},
"event.exitNode.connected": {
"message": "出口节点已连接。"
},
"event.exitNode.disconnected": {
"message": "出口节点已断开。"
},
"event.exitNode.connectionLost": {
"message": "与出口节点的连接已丢失。您的互联网访问可能受到影响。"
},
"event.exitNode.haChange": {
"message": "由于高可用性变更,出口节点已断开。"
},
"event.exitNode.disconnectedUnknown": {
"message": "出口节点因未知原因已断开。"
},
"event.update.installing": {
"message": "正在安装更新。"
},
"event.update.completed": {
"message": "NetBird 客户端已自动更新到版本 {version}。"
},
"event.update.failed": {
"message": "自动更新失败:{reason}"
},
"common.cancel": {
"message": "取消"
},

View File

@@ -63,6 +63,20 @@ func (l *Localizer) T(key string, args ...string) string {
return l.bundle.Translate(lang, key, args...)
}
// Lookup resolves a key supplied at runtime by the daemon, substituting args as
// {placeholder}/value pairs. It reports false when the key is in no bundle, so
// the caller can fall back to the daemon's own English text instead of showing a
// bare key. An empty key never resolves.
func (l *Localizer) Lookup(key string, args map[string]string) (string, bool) {
if l == nil || l.bundle == nil || key == "" {
return "", false
}
l.mu.RLock()
lang := l.lang
l.mu.RUnlock()
return l.bundle.Lookup(lang, key, flattenArgs(args)...)
}
// Watch invokes cb on each language change, after the cached language is
// updated so cb may call l.T with the new locale. Replaces any prior subscription.
func (l *Localizer) Watch(cb func(lang i18n.LanguageCode)) {
@@ -128,3 +142,17 @@ func (l *Localizer) StatusLabel(status string) string {
}
return status
}
// flattenArgs turns a placeholder map into the flat name/value slice the bundle
// takes. Iteration order is irrelevant: each pair substitutes an independent
// {name}.
func flattenArgs(args map[string]string) []string {
if len(args) == 0 {
return nil
}
out := make([]string, 0, len(args)*2)
for name, value := range args {
out = append(out, name, value)
}
return out
}

View File

@@ -14,6 +14,7 @@ import (
"github.com/sirupsen/logrus"
"github.com/wailsapp/wails/v3/pkg/application"
"github.com/wailsapp/wails/v3/pkg/events"
"github.com/wailsapp/wails/v3/pkg/services/notifications"
"github.com/netbirdio/netbird/client/ui/authsession"
"github.com/netbirdio/netbird/client/ui/i18n"
@@ -62,7 +63,7 @@ type registeredServices struct {
profiles *services.Profiles
update *services.Update
daemonFeed *services.DaemonFeed
notifier *Notifier
notifier *notifications.NotificationService
compat *services.Compat
profileSwitcher *services.ProfileSwitcher
bundle *i18n.Bundle
@@ -101,7 +102,7 @@ func main() {
updaterHolder := updater.NewHolder(app.Event)
update := services.NewUpdate(conn, updaterHolder)
daemonFeed := services.NewDaemonFeed(conn, app.Event, updaterHolder, debugLog)
notifier := newNotifier()
notifier := notifications.New()
compat := services.NewCompat(conn)
// macOS shows no toast until permission is requested. Run it after
// ApplicationStarted so the notifier's Startup has initialised the
@@ -209,7 +210,7 @@ func main() {
// requestNotificationAuthorization prompts for macOS notification permission.
// The request blocks until the user responds (up to 3 minutes), so callers run
// it in a goroutine. No-op on Linux/Windows.
func requestNotificationAuthorization(notifier *Notifier) {
func requestNotificationAuthorization(notifier *notifications.NotificationService) {
authorized, err := notifier.CheckNotificationAuthorization()
if err != nil {
logrus.Debugf("check notification authorization: %v", err)

View File

@@ -1,101 +0,0 @@
//go:build !android && !ios && !freebsd && !js
package main
import (
"context"
"errors"
"sync/atomic"
log "github.com/sirupsen/logrus"
"github.com/wailsapp/wails/v3/pkg/application"
"github.com/wailsapp/wails/v3/pkg/services/notifications"
)
var errNotificationsUnavailable = errors.New("notifications unavailable")
// Notifier wraps the Wails notification service so an unavailable backend
// disables notifications instead of aborting the app. Startup fails for
// environment reasons (a bare unbundled binary on macOS has no bundle
// identifier, a headless Linux session has no D-Bus session bus), and Wails
// treats a service startup error as fatal. After a failed startup every call
// is a no-op: on macOS, touching UNUserNotificationCenter without a bundle
// identifier raises an Objective-C exception that recover() cannot catch.
type Notifier struct {
inner *notifications.NotificationService
available atomic.Bool
}
func newNotifier() *Notifier {
return &Notifier{inner: notifications.New()}
}
// ServiceName implements the Wails service-name hook for startup logs.
func (n *Notifier) ServiceName() string {
return n.inner.ServiceName()
}
// ServiceStartup starts the platform notifier, downgrading failure to a
// warning so the app keeps running without notifications.
func (n *Notifier) ServiceStartup(ctx context.Context, options application.ServiceOptions) error {
if err := n.inner.ServiceStartup(ctx, options); err != nil {
log.Warnf("notifications disabled: %v", err)
return nil
}
n.available.Store(true)
return nil
}
func (n *Notifier) ServiceShutdown() error {
if !n.available.Load() {
return nil
}
return n.inner.ServiceShutdown()
}
func (n *Notifier) CheckNotificationAuthorization() (bool, error) {
if !n.available.Load() {
return false, errNotificationsUnavailable
}
return n.inner.CheckNotificationAuthorization()
}
func (n *Notifier) RequestNotificationAuthorization() (bool, error) {
if !n.available.Load() {
return false, errNotificationsUnavailable
}
return n.inner.RequestNotificationAuthorization()
}
// SendNotification delivers a notification, silently dropping it when the
// backend never started (notifications are best-effort everywhere).
func (n *Notifier) SendNotification(options notifications.NotificationOptions) error {
if !n.available.Load() {
log.Debugf("notifications disabled, dropping %q", options.ID)
return nil
}
return n.inner.SendNotification(options)
}
func (n *Notifier) SendNotificationWithActions(options notifications.NotificationOptions) error {
if !n.available.Load() {
log.Debugf("notifications disabled, dropping %q", options.ID)
return nil
}
return n.inner.SendNotificationWithActions(options)
}
func (n *Notifier) RegisterNotificationCategory(category notifications.NotificationCategory) error {
if !n.available.Load() {
return nil
}
return n.inner.RegisterNotificationCategory(category)
}
// OnNotificationResponse registers the response callback. Pure Go state, so
// it is safe (and simply inert) when the backend never started.
//
//wails:ignore
func (n *Notifier) OnNotificationResponse(callback func(result notifications.NotificationResult)) {
n.inner.OnNotificationResponse(callback)
}

View File

@@ -246,7 +246,6 @@ func (s *Store) ExistedAtLoad() bool {
func (s *Store) load() error {
if _, err := os.Stat(s.path); err != nil {
if errors.Is(err, os.ErrNotExist) {
log.Infof("no ui preferences file at %s; using defaults", s.path)
return nil
}
return fmt.Errorf("stat preferences: %w", err)

View File

@@ -33,21 +33,12 @@ type LoginResult struct {
UserCode string `json:"userCode"`
VerificationURI string `json:"verificationUri"`
VerificationURIComplete string `json:"verificationUriComplete"`
// ProfileID is the ID of the profile this login ran against, or "" when the
// caller named the profile itself and no ID was resolved. Pass it back in
// WaitSSOParams so the account email lands on this profile even if the
// active one changes during SSO.
ProfileID string `json:"profileId"`
}
// WaitSSOParams are the inputs to waitSSOLogin.
type WaitSSOParams struct {
UserCode string `json:"userCode"`
Hostname string `json:"hostname"`
// ProfileID is the profile the login was started for, used to file the
// account email against it rather than against whichever profile is active
// when the flow returns. Optional: empty falls back to the active profile.
ProfileID string `json:"profileId"`
}
// UpParams selects the profile to bring up.
@@ -86,16 +77,11 @@ func (s *Connection) Login(ctx context.Context, p LoginParams) (LoginResult, err
// Fall back to the daemon's active profile and the current OS user.
profileName := p.ProfileName
username := p.Username
// Only set when the daemon told us the ID. A caller-supplied ProfileName is
// a handle — a display name or an ID prefix resolve too — and the state file
// is named after the ID, so passing a handle on would name the wrong file.
profileID := ""
if profileName == "" {
if active, aerr := cli.GetActiveProfile(ctx, &proto.GetActiveProfileRequest{}); aerr == nil {
// Address the active profile by ID (the daemon resolves it as a
// handle); names can collide, the ID cannot.
profileName = active.GetId()
profileID = profileName
if username == "" {
username = active.GetUsername()
}
@@ -136,7 +122,6 @@ func (s *Connection) Login(ctx context.Context, p LoginParams) (LoginResult, err
UserCode: resp.GetUserCode(),
VerificationURI: resp.GetVerificationURI(),
VerificationURIComplete: resp.GetVerificationURIComplete(),
ProfileID: profileID,
}, nil
}
@@ -257,31 +242,6 @@ func (s *Connection) waitSSOLogin(ctx context.Context, p WaitSSOParams) (string,
return "", s.classifyDaemonError(err)
}
log.Infof("SSO login completed, daemon reported success")
// Persist the account email the same way the CLI does after its own
// WaitSSOLogin: the daemon returns it but cannot store it, since it runs as
// root and the per-profile state file is user-owned (see Logout below).
// Without this the profile has no email, so Profiles.List shows no account
// and later logins and session extends go out without a login_hint —
// leaving the IdP to guess which account was meant.
if email := resp.GetEmail(); email != "" {
state := &profilemanager.ProfileState{Email: email}
pm := profilemanager.NewProfileManager()
// Against the profile the login was started for: SSO spans seconds of
// user interaction, and a profile switch in that window would otherwise
// file the email under the wrong profile.
if p.ProfileID != "" {
err = pm.SetProfileState(profilemanager.ID(p.ProfileID), state)
} else {
err = pm.SetActiveProfileState(state)
}
if err != nil {
// Non-fatal: the login itself succeeded.
log.Warnf("failed to store account email: %v", err)
}
}
return resp.GetEmail(), nil
}

View File

@@ -59,13 +59,21 @@ type Emitter interface {
// SystemEvent is the frontend-facing shape of a daemon SystemEvent.
type SystemEvent struct {
ID string `json:"id"`
Severity string `json:"severity"`
Category string `json:"category"`
Message string `json:"message"`
UserMessage string `json:"userMessage"`
Timestamp int64 `json:"timestamp"`
Metadata map[string]string `json:"metadata"`
ID string `json:"id"`
Severity string `json:"severity"`
Category string `json:"category"`
Message string `json:"message"`
UserMessage string `json:"userMessage"`
// MessageKey names the localizable body for this event; empty on control
// events and on events from a daemon that predates the field. Resolve it
// against the UI bundle and fall back to UserMessage on a miss.
MessageKey string `json:"messageKey"`
MessageArgs map[string]string `json:"messageArgs"`
// TitleKey names the localizable notification title, empty when the event
// has none and the consumer should compose one from severity and category.
TitleKey string `json:"titleKey"`
Timestamp int64 `json:"timestamp"`
Metadata map[string]string `json:"metadata"`
}
// PeerStatus is the frontend-facing shape of a daemon PeerState.
@@ -563,6 +571,9 @@ func systemEventFromProto(e *proto.SystemEvent) SystemEvent {
Category: strings.ToLower(strings.TrimPrefix(e.GetCategory().String(), "SystemEvent_")),
Message: e.GetMessage(),
UserMessage: e.GetUserMessage(),
MessageKey: e.GetMessageKey(),
MessageArgs: e.GetMessageArgs(),
TitleKey: e.GetTitleKey(),
Metadata: map[string]string{},
}
if ts := e.GetTimestamp(); ts != nil {

View File

@@ -6,8 +6,6 @@ import (
"context"
"os/user"
log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/internal/profilemanager"
"github.com/netbirdio/netbird/client/proto"
)
@@ -153,31 +151,11 @@ func (s *Profiles) Remove(ctx context.Context, p ProfileRef) error {
if err != nil {
return err
}
resp, err := cli.RemoveProfile(ctx, &proto.RemoveProfileRequest{
_, err = cli.RemoveProfile(ctx, &proto.RemoveProfileRequest{
ProfileName: p.ProfileName,
Username: p.Username,
})
if err != nil {
return err
}
// The daemon deletes what it owns but runs as root, so it leaves the
// user-owned state file holding the account email behind (same split as
// Connection.Logout). Legacy profiles are keyed by name rather than by a
// generated ID, so a recreated profile of the same name would inherit the
// deleted one's email and offer it as the login_hint.
//
// Keyed on the ID the daemon resolved, not on the request handle: that may
// have been a display name or an ID prefix, which would name a different
// file (or none).
if id := resp.GetId(); id != "" {
if err := profilemanager.NewProfileManager().RemoveProfileState(id); err != nil {
// Non-fatal: the profile itself is gone.
log.Warnf("failed to remove profile state for %s: %v", id, err)
}
}
return nil
return err
}
// Rename changes a profile's display name. The on-disk ID is unaffected, so

View File

@@ -26,7 +26,6 @@ const (
notifyIDUpdatePrefix = "netbird-update-"
notifyIDEvent = "netbird-event-"
notifyIDTrayError = "netbird-tray-error"
notifyIDMDMPolicy = "netbird-mdm-policy"
statusError = "Error"
@@ -44,7 +43,7 @@ type TrayServices struct {
Profiles *services.Profiles
Networks *services.Networks
DaemonFeed *services.DaemonFeed
Notifier *Notifier
Notifier *notifications.NotificationService
Update *services.Update
ProfileSwitcher *services.ProfileSwitcher
WindowManager *services.WindowManager

View File

@@ -4,26 +4,17 @@ package main
// bindTrayClick wires the tray icon's left-click handler on Linux.
//
// Expected behaviour per tray host:
//
// Host Left click Right click
// KDE Plasma, Waybar main window (Activate) menu (host-rendered)
// GNOME Shell + AppIndicator menu only menu only
// Minimal WMs via XEmbed host main window (Activate) XEmbed GTK popup
//
// OnClick fires only on org.kde.StatusNotifierItem.Activate — a real left
// click. KDE/Waybar send it over D-Bus; the in-process XEmbed host
// (xembed_host_linux.go) maps a Button1 press to the same Activate call.
//
// GNOME Shell + AppIndicator never sends Activate: it renders the dbusmenu
// on ANY click and only reports the menu opening via dbusmenu
// Event("opened"). Upstream Wails treated that event as a click, so on GNOME
// both buttons raised the main window on top of the menu, and on KDE/Waybar
// a right click raised it over the freshly opened menu. The netbirdio/wails
// fork (go.mod replace) drops that heuristic: a menu open never fires
// OnClick. On GNOME the main window is reached via the "Open NetBird" menu
// entry; left-click-opens-window is not achievable there anyway, since the
// host always opens the menu itself.
// Both Linux click paths converge on Wails' linuxSystemTray.Activate, which
// fires the registered clickHandler:
// - Real SNI hosts (KDE Plasma, Waybar, GNOME Shell + AppIndicator) invoke
// org.kde.StatusNotifierItem.Activate over D-Bus on left-click.
// - The in-process StatusNotifierWatcher + XEmbed host used on minimal WMs
// (Fluxbox, i3, dwm, OpenBox) maps a Button1 press to that same Activate
// call itself (xembed_host_linux.go), so it routes through the same hook.
// Registering OnClick here therefore covers both paths with one handler — no
// changes to the watcher or XEmbed C code are needed. Left-click now opens the
// main window; right-click still opens the menu via Wails' default
// SecondaryActivate→OpenMenu handler (and the XEmbed GTK popup on minimal WMs).
//
// We do NOT register OnDoubleClick: Wails' Linux SNI backend never fires it
// (unlike Windows). And we deliberately skip AttachWindow — it plus Wails3's

View File

@@ -21,32 +21,14 @@ func (t *Tray) onSystemEvent(ev *application.CustomEvent) {
if !ok {
return
}
// config_changed carries no UserMessage, so handle it before the message gate below.
// config_changed carries no user-facing message, so handle it before the gate below.
if se.Category == "system" && se.Metadata[proto.MetadataTypeKey] == proto.MetadataTypeConfigChanged {
log.Infof("config_changed event received (source=%s); refreshing tray restrictions", se.Metadata[proto.MetadataSourceKey])
go t.refreshRestrictions()
go t.loadConfig()
// MDM gets a localised toast here; the daemon's English "policy_applied"
// event is suppressed in shouldSkipSystemEvent. Other sources stay silent.
if se.Metadata[proto.MetadataSourceKey] == proto.MetadataSourceMDM {
t.profileMu.Lock()
enabled := t.notificationsEnabled
t.profileMu.Unlock()
if enabled {
t.notify(
t.loc.T("notify.mdm.policyApplied.title"),
t.loc.T("notify.mdm.policyApplied.body"),
notifyIDMDMPolicy,
)
}
}
return
}
// Session-warning and deadline-rejected events build their body locally from
// metadata; every other event needs a UserMessage.
isSessionWarning := se.Metadata[authsession.MetaWarning] == "true"
isDeadlineRejected := se.Metadata[authsession.MetaDeadlineRejected] != ""
if !isSessionWarning && !isDeadlineRejected && se.UserMessage == "" {
if se.MessageKey == "" && se.UserMessage == "" {
return
}
if shouldSkipSystemEvent(se) {
@@ -61,56 +43,56 @@ func (t *Tray) onSystemEvent(ev *application.CustomEvent) {
return
}
// Session-warning events route via stable metadata flags rather than
// category/severity so a daemon-side reword still lands here. Final warning
// auto-opens the SessionExpiration dialog with no notification (the dialog is
// the last-chance reminder; doubling up would be noise).
if isDeadlineRejected {
t.notify(
t.loc.T("notify.sessionDeadlineRejected.title"),
t.loc.T("notify.sessionDeadlineRejected.body"),
notifyIDSessionExpired,
)
return
}
body := t.localizedEventMessage(se)
if se.Metadata != nil && se.Metadata[authsession.MetaWarning] == "true" {
// The final session warning auto-opens the SessionExpiration dialog instead of
// toasting: the dialog is the last-chance reminder and doubling up would be
// noise. This routes on metadata rather than the message key because it is a
// behavioural distinction, not a wording one.
if se.Metadata[authsession.MetaWarning] == "true" {
if se.Metadata[authsession.MetaFinal] == "true" {
t.openSessionExpiration()
return
}
t.notifySessionWarning(
t.loc.T("notify.sessionWarning.title"),
t.buildSessionWarningBody(se.Metadata),
)
t.notifySessionWarning(t.eventTitle(se), body)
return
}
body := se.UserMessage
if id := se.Metadata["id"]; id != "" {
body += fmt.Sprintf(" ID: %s", id)
}
t.notify(eventTitle(se), body, notifyIDEvent+se.ID)
t.notify(t.eventTitle(se), body, notifyIDEvent+se.ID)
}
// eventTitle composes a notification title, e.g. "Critical: DNS", "Warning: Authentication".
func eventTitle(e services.SystemEvent) string {
prefix := titleCase(e.Severity)
if prefix == "" {
prefix = "Info"
// localizedEventMessage resolves the daemon's message key against the active
// locale. A key this build does not ship — a daemon newer than the UI — falls
// back to the daemon's own English rendering rather than showing a bare key.
func (t *Tray) localizedEventMessage(se services.SystemEvent) string {
if body, ok := t.loc.Lookup(se.MessageKey, se.MessageArgs); ok {
return body
}
category := titleCase(e.Category)
if category == "" {
category = "System"
if se.MessageKey != "" {
log.Debugf("no translation for event message key %q, using the daemon's text", se.MessageKey)
}
return prefix + ": " + category
return se.UserMessage
}
func titleCase(s string) string {
if s == "" {
return ""
// eventTitle resolves the event's own title key, falling back to a title
// composed from severity and category, e.g. "Critical: DNS" in English. An enum
// value this build does not know falls back to the Info and System labels.
func (t *Tray) eventTitle(se services.SystemEvent) string {
if title, ok := t.loc.Lookup(se.TitleKey, nil); ok {
return title
}
return strings.ToUpper(s[:1]) + strings.ToLower(s[1:])
severity, ok := t.loc.Lookup("event.severity."+strings.ToLower(se.Severity), nil)
if !ok {
severity = t.loc.T("event.severity.info")
}
category, ok := t.loc.Lookup("event.category."+strings.ToLower(se.Category), nil)
if !ok {
category = t.loc.T("event.category.system")
}
return t.loc.T("event.title", "severity", severity, "category", category)
}
// shouldSkipSystemEvent reports whether a daemon SystemEvent must not surface as
@@ -119,11 +101,6 @@ func titleCase(s string) string {
// - install-progress signals (consumed by the install-progress window)
// - the ::/0 partner of an exit-node default route (0.0.0.0/0 already toasted)
func shouldSkipSystemEvent(se services.SystemEvent) bool {
// "policy_applied" carries a hardcoded English message; the localised toast
// fires on the paired config_changed (source=mdm) event instead.
if se.Metadata[proto.MetadataTypeKey] == proto.MetadataTypePolicyApplied {
return true
}
if _, isUpdate := se.Metadata["new_version_available"]; isUpdate {
return true
}

View File

@@ -0,0 +1,150 @@
//go:build !android && !ios && !freebsd && !js
package main
import (
"os"
"testing"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/netbirdio/netbird/client/proto"
"github.com/netbirdio/netbird/client/ui/i18n"
"github.com/netbirdio/netbird/client/ui/services"
)
// trayWithLocalizer builds the minimum Tray the message/title resolvers touch:
// they read t.loc and nothing else, so no app, window or daemon connection is
// needed. The shipped locale tree is used so the assertions below exercise the
// real bundles rather than a fixture.
func trayWithLocalizer(t *testing.T) *Tray {
t.Helper()
bundle, err := i18n.NewBundle(os.DirFS("i18n/locales"))
require.NoError(t, err, "the shipped locale tree must load")
return &Tray{loc: NewLocalizer(bundle, nil)}
}
func TestLocalizedEventMessageResolvesKey(t *testing.T) {
tray := trayWithLocalizer(t)
got := tray.localizedEventMessage(services.SystemEvent{
MessageKey: string(proto.UserMsgExitNodeConnected),
// A daemon always ships its English rendering too; the key must win.
UserMessage: "should not be used",
})
assert.Equal(t, "Exit node connected.", got)
}
func TestLocalizedEventMessageSubstitutesArgs(t *testing.T) {
tray := trayWithLocalizer(t)
got := tray.localizedEventMessage(services.SystemEvent{
MessageKey: string(proto.UserMsgUpdateCompleted),
MessageArgs: map[string]string{proto.ArgVersion: "0.60.1"},
})
assert.Equal(t, "Your NetBird client was auto-updated to version 0.60.1.", got)
}
// A daemon newer than the UI can publish a key this build has never heard of.
// Showing the raw key would be a visible regression, so the daemon's own English
// text has to win instead.
func TestLocalizedEventMessageFallsBackToDaemonText(t *testing.T) {
tray := trayWithLocalizer(t)
got := tray.localizedEventMessage(services.SystemEvent{
MessageKey: "event.somethingThisBuildNeverHeardOf",
UserMessage: "A message from a newer daemon.",
})
assert.Equal(t, "A message from a newer daemon.", got)
}
// An old daemon sends no key at all, only userMessage.
func TestLocalizedEventMessageWithoutKey(t *testing.T) {
tray := trayWithLocalizer(t)
got := tray.localizedEventMessage(services.SystemEvent{UserMessage: "Legacy English text."})
assert.Equal(t, "Legacy English text.", got)
}
func TestEventTitlePrefersTitleKey(t *testing.T) {
tray := trayWithLocalizer(t)
got := tray.eventTitle(services.SystemEvent{
Severity: "critical",
Category: "authentication",
TitleKey: string(proto.TitleSessionWarning),
})
assert.Equal(t, "Session expires soon", got, "a title key must beat the composed title")
}
func TestEventTitleComposesFromSeverityAndCategory(t *testing.T) {
tray := trayWithLocalizer(t)
tests := []struct {
name string
severity string
category string
want string
}{
{"warning dns", "warning", "dns", "Warning: DNS"},
{"critical system", "critical", "system", "Critical: System"},
{"info network", "info", "network", "Info: Network"},
{"error authentication", "error", "authentication", "Error: Authentication"},
// Enum values this build does not know, and the empty severity/category
// an event carries before the daemon fills them in.
{"unknown severity", "apocalyptic", "dns", "Info: DNS"},
{"unknown category", "warning", "quantum", "Warning: System"},
{"empty", "", "", "Info: System"},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
got := tray.eventTitle(services.SystemEvent{Severity: tc.severity, Category: tc.category})
assert.Equal(t, tc.want, got)
})
}
}
func TestShouldSkipSystemEvent(t *testing.T) {
tests := []struct {
name string
ev services.SystemEvent
want bool
}{
{
name: "update announcement handled by the tray updater",
ev: services.SystemEvent{Metadata: map[string]string{"new_version_available": "0.60.1"}},
want: true,
},
{
name: "install progress belongs to the progress window",
ev: services.SystemEvent{Metadata: map[string]string{"progress_window": "show"}},
want: true,
},
{
name: "the v6 half of a dual-stack default route is already toasted as v4",
ev: services.SystemEvent{Category: "network", Metadata: map[string]string{"network": "::/0"}},
want: true,
},
{
name: "the v4 default route is the one that toasts",
ev: services.SystemEvent{Category: "network", Metadata: map[string]string{"network": "0.0.0.0/0"}},
want: false,
},
{
// policy_applied used to be suppressed here while the tray toasted
// off the paired config_changed event; it now carries its own keys.
name: "mdm policy applied surfaces normally",
ev: services.SystemEvent{
MessageKey: string(proto.UserMsgMDMPolicyApplied),
Metadata: map[string]string{proto.MetadataTypeKey: proto.MetadataTypePolicyApplied},
},
want: false,
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
assert.Equal(t, tc.want, shouldSkipSystemEvent(tc.ev))
})
}
}

Some files were not shown because too many files have changed in this diff Show More