mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-10 23:49:09 +02:00
* implement certificate posture check * log signal address * add keychain and cert store support * read the console user's keychain through a user session helper A root daemon cannot reach a login keychain: securityd is per session and a key ACL needs a session to prompt in, so dropping uid is not enough. The daemon now answers certificate challenges from the System keychain itself, where MDM installs device identities, and launches "netbird posture cert-proof" into the console user's desktop session with launchctl asuser for the login keychain. Only the signature and the chain cross back, never the private key. The console user comes from SCDynamicStoreCopyConsoleUser, bound with purego like the keychain calls. The login window reports no user, root, or "loginwindow", and all three are treated as no keychain to read, so a Mac at the lock screen sends device proofs alone. Adds info logging across the path: the keychain search list, per class query status and item counts, the chain built per candidate, and the verification error for every rejected candidate. A run that sends nothing now says why. README.md documents the trust model, the console user limitation and how to read the logs. * read the signed-in user's certificate store on Windows A service reads LocalMachine\MY, where AD and Intune enrol device certificates. CurrentUser\MY lives in the signed-in user's registry hive with keys protected against their profile, and a service that opens it does not fail: "current user" resolves to HKU\S-1-5-18, so it silently reads the service account's own empty store. The service therefore reads the machine store itself and launches "netbird posture cert-proof" with the session token for the rest, mirroring the macOS console user helper. Windows lets a privileged service assume a user identity, so the token goes straight into the child process and no external tooling is involved. CREATE_NO_WINDOW keeps a console window from flashing on the desktop every sync. In-process impersonation would also work but is per OS thread while goroutines migrate, so the child process avoids that class of bug. Session selection prefers the physical console and falls back to any active session, so remote desktop and VDI hosts are covered. WTSQueryUserToken needs SE_TCB_NAME, so a user-run client skips the helper and reads the machine store alone. SystemStore takes a store location, gaining NewUserStore alongside NewSystemStore and the per candidate logging macOS already had. The request building and proof merging move to helper_spawn.go, shared by both platforms, and helperStore picks what the helper reads per platform. * start TPM support * split goreleaser to support pkcs11 and exclude on docker * update goreleaser * go mod tidy * add tpm pin to netbird config * split cert and key location and allow key lookup on tpm * add unsupported flag for mobile devices * Isolate the cert proof helper from the service environment and cap its output * Read the PKCS#11 token PIN from NB_TPM_PIN instead of the profile config * Bound certificate proof collection so a stuck token or keychain cannot hold the sync loop * Stop retrying a PKCS#11 PIN the token rejected * Log certificate posture details at debug level * Sign only nonces and peer keys of the size management issues * Skip certificate files whose key belongs to another certificate * Bound PKCS#11 driver sizes, pin template values, and log out only a login the session owns * Never pass NULL to CFRelease and skip unreadable keychain identities * Keep the macOS keychain code out of iOS and the PKCS#11 driver out of Android * Find a chain to each challenge's CAs through every intermediate the store holds * Require a token label whenever a PKCS#11 PIN is set * Read user certificates only from the session of the active profile's owner * Collect certificate proofs again when the owner's session changes and report lost proofs * Test the PKCS#11 build against SoftHSM in CI and warn once where the build has no driver * Document where an inline PKCS#11 PIN is stored and how it is protected * Refuse PKCS#11 URIs that this client cannot honour instead of widening the match * Trust certificate and key files only when no other user can write or redirect them * Explain a Windows certificate whose key only a legacy CryptoAPI provider holds * Use platform absolute module paths in tests and add a real owner session test for Windows * Match the Windows profile owner by name instead of resolving it through the domain controller * Keep the certificate stores and TPM library out of the WebAssembly build * [client] Read TSS2 key files on go-tpm, checked against the library it replaces The TSS2 parser was the only reason this repository depended on a crypto suite whose own build tooling it inherits. The replacement sits on go-tpm, which was already a direct dependency and is in fact what that suite calls underneath, so this removes a wrapper rather than porting onto a different library: the load, the derived storage root key and the signing commands are the same calls. Swapping a parser on the one path a customer actually runs is not something to assert, so the two are held side by side for this commit. One test feeds the replacement bytes the old library wrote and requires the same key type, empty auth flag, parent handle, blobs and decoded public key; the other feeds both the fixtures the tests are built on, so those are the shape the format calls for and not merely the shape the new parser reads. The scaffolding goes away with the dependency in the commit that follows. The encoder behind the fixtures is written out separately from the parser under test, so an encoder bug and a decoder bug cannot cancel each other out. * [client] Drop go.step.sm/crypto and the repo-wide upgrades it imposed The TSS2 parser was the only thing in the repository that used this module, and it brought 302 modules into the graph to do it — 35 of them linters, along with Google Cloud KMS and IAM, the AWS SDK and a terminal styling library. Those are the module's own development dependencies, which minimal version selection turns into floors in ours, and they are the whole reason gRPC, protobuf, the AWS SDK, OpenTelemetry, logrus and five x/ packages had moved. Management, signal, relay and proxy inherited every one of them for a feature none of them runs. Removing the import is not enough, because tidy never downgrades: the raised floors stay written in go.mod. Each one is pinned back to the version main had, then tidy is left to raise again whatever something still genuinely needs. It raised nothing: all 43 are back where they were, and go-tpm was already in the graph at the same version, so the certificate feature now costs no new module at all. The differential tests go with it. They existed to check the swap against the library while both were present, and there is nothing left to compare against. * [client] Clear the lint findings only the macOS and Windows runners see golangci-lint analyses one build at a time, so running it on Linux says nothing about the two platforms CI also lints. Against those builds the feature's packages reported eight findings, and the structural one is Config.dir: it is dead on macOS and Windows because neither reads a directory at all, their collectors take the configuration and discard it. Moving the method beside its only callers makes that visible in the layout instead of in a linter, and leaves the gap itself — no file or token store on those platforms — where it belongs, as something to decide rather than something to silence. An absent key file beside a certificate was reported as a nil signer with a nil error, which the caller then had to recognise by its nilness. It is a sentinel now, so the meaning is in the error rather than in the absence of one. The rest follow the standard library: the elliptic coordinates and the private scalar come from the encoding helpers rather than the deprecated big.Int fields, and an error string loses its trailing colon. Lint is clean on linux, darwin and windows; the hardware TPM path was exercised separately against a real device and passes. * Accept the TSS2 emptyAuth boolean OpenSSL writes and persistent parents on 32-bit builds * Count the certificates field in the peer meta store test * Check the store directory before listing it, refuse group-writable files, and reject a URI with two PIN sources * Share a PKCS#11 login between sessions and send each PIN at most once at a time * Collect certificate proofs again when the meta sync carrying them failed * Use no Windows user store when a domainless owner matches accounts of several domains * Use no user certificate store when the active profile's owner cannot be read * Document the PIN sources on CertPKCS11URI and keep the README PIN example off the command line * Test that the PKCS#11 URI stays out of the debug bundle and run the wrong-PIN test only on a disposable token * Refuse a TPM PSS signature request for the maximum salt length * Add the certificate fields to the network map golden data * Retry posture checks whose meta sync timed out instead of dropping them * Start no system info gathering while a timed-out one is still running * Guard the applied posture checks across goroutines and keep refreshing proofs while a pending update times out * Log what a successful certificate proof helper wrote to stderr * Send recollected certificate proofs to management only when the proven chains changed * Explain a macOS keychain key whose access list does not allow netbird * Kill the whole macOS certificate helper process group when it times out * End sudo option parsing before the macOS certificate helper binary * Hold off system info gathering only while a timed-out one is still running * Collect certificate proofs on the posture watcher instead of under the sync lock * Read the certificate store directory and PKCS#11 URI from the daemon environment, not the profile config * Install the RPM sysconfig file readable by root only and show the certificate posture variables * Move the certificate posture README into the package doc and the docs site * Name NB_CERT_PKCS11_URI in the PIN-without-token error * Keep the file check results of the latest-started system info refresh * Give the full import command for a keychain key netbird may not use, and correct the package doc * Restrict the service environment file to root on every package install * Search only the System keychain in the macOS daemon and only the login keychain in the user helper * Let the certificate proof helper read the PKCS#11 token from the environment on Linux * Ask a macOS user's keychain again only after an hour when it proved nothing * Clear the lint findings in certificate posture * Hold off the keychain helper only after a completed or timed-out run, independent of CA order * Keep free functions out of the method lists of PKCS11Store, URI and Challenger * Name the post-install permission helper in snake case and shorten the sysconfig certificate block * Drop the certificate store directory from certproof.Config, which only NB_CERT_STORE_DIR sets * [management] Renew certificate challenge nonces on quiet accounts A certificate challenge nonce is accepted for its own window and the one before it, and it only reaches a peer attached to a network map. An account where nothing changes sends no map, so after a day the peer re-sends the nonce it still holds, verification rejects its whole proof set, and the certificates stored for it are dropped. It fails the certificate check and loses every policy gated on it until some unrelated change happens to push a map. The outage repairs itself in seconds, which is what makes it expensive: it is intermittent, it only hits stable networks, and it is not reproducible on demand. Push the account's peers an update often enough that the nonce they hold is never close to expiring. Only accounts whose posture checks actually ask for a certificate are tracked, so a deployment without the feature does no extra work. The refresh runs from one goroutine over a map of accounts rather than a timer per account: the period is hours, so one pass every few minutes costs nothing next to it, and there is no timer to re-arm when an account that falls due sooner appears. Each account's first run is offset by a hash of its ID, because the challenge window is global and an instance restart would otherwise arm every account in the same moment. The push carries no administrative change, so it is counted as a refresh rather than an update and stays out of the figures that track what was edited. (cherry picked from commit7ad4a0df37) * [management] Make the certificate challenge window one knob to turn Renewal was timed against the window in two different ways: the period derived from it, the sweep interval did not. Shortening the window to watch a renewal in an end-to-end run would have left the refresher still looking for due accounts every quarter of an hour, so nothing would have been renewed in time and the test would have reported the feature broken. Derive the sweep from the period, within bounds that keep a very short window from spinning and a normal one from checking less often than is useful, and allow the window itself to be set through the environment so a run can take seconds instead of half a day. A value that cannot be parsed or falls outside the bounds keeps the default, because a window nobody intended is a security property nobody chose, and an override is logged at warning level since it sets how long a device keeps passing the check after its key is gone. Every instance has to be given the same value: the window is part of the nonce, so instances that disagree reject each other's. (cherry picked from commit0e38fcf409) * [management] Pin the property that makes per-peer nonce state unnecessary A nonce carries the window it was minted in, not the instant, and is accepted for that window and the one before it. So a peer re-stamped at least once per window can never be left holding one outside the accepted pair, whenever it was last served and however much life its own nonce had left. That is the whole reason management tracks nothing per peer, and it was resting on an argument rather than a test. The phases are part of the property, not decoration: accounts are deliberately given a refresh phase of their own, so the guarantee has to hold off the window boundary too. The negative case shows why that matters — a cadence of exactly two windows lands inside the grace window when it is aligned to the boundary and leaves a gap when it is not. (cherry picked from commitdee68facfd) * [management] Renew challenges only for the peers that answer one The refresh pushed an update to every connected peer of the account, while only the peers a certificate check applies to carry a nonce. On an account where a handful of peers sit behind the check and the rest do not, everyone was woken several times a day to be handed a map that changed nothing for them. Push to the sources of the enabled policies whose posture checks include a certificate check, which is exactly the set that is sent a challenge. Resolving the set the other way round than the gRPC layer does is the risk here: a peer the refresh forgets stops being renewed and falls out of its policies silently, which is the failure this whole mechanism exists to prevent. So the selection is held against processPeerPostureChecks, the per-peer rule that decides who receives a challenge in the first place, by a test that asks both the same question and requires the same answer. (cherry picked from commitdc4d0e0274) * [management] Derive certificate challenge nonces from the stored encryption key The nonce secret came from the server's WireGuard key, which is generated afresh in every process and never persisted. A nonce carries no state, so the only thing that lets one instance verify what another issued is deriving the same secret — and that premise, written in the comment above the challenger, was not met: every instance had its own key. A peer reconnecting after a restart therefore presented a nonce minted under the previous secret, verification failed with a mismatch, its whole proof set was rejected and the certificates stored for it were dropped until it signed again. Reproduced three times on the lab, each one logging "nonce was not issued to this peer", which only a changed secret produces. On a single instance it costs seconds of lost policy access per restart; across instances it is not transient at all, because every reconnect that lands elsewhere is rejected the same way. Derive from the data store encryption key instead: it is generated once, written back to the configuration and read by every instance, so it survives restarts and is shared. Where none is configured the secret falls back to the WireGuard key with a warning — degraded but still unpredictable, which is the property that matters most: a peer able to guess it could mint the nonces of future windows, sign them while its key is present and keep passing after it is gone. The challenger is now built once and passed to the two places that need it, rather than re-derived per message. (cherry picked from commit278f2f3807) * [management] Register an account for renewal where its nonce is issued Renewal was armed when a peer connected or when a posture check was saved, both of which ask the store whether the account has a certificate check. That misses the case it most needs to catch: the check is created through one instance while the peers are connected to another, so the instance serving them never learns it has anything to renew and their nonce expires. It also charged a query to every peer connect in every account, including the ones that will never use the feature, which a fleet reconnecting after a restart pays all at once. Register where the nonce is actually stamped instead. A nonce is verified from a shared secret and so travels between instances, but the renewal that keeps it fresh cannot: only the instance holding a peer's stream can push to it. Issuing and renewing now line up by construction — an instance renews exactly the accounts it has issued nonces for — and an instance that never issues one has nothing to renew, so there is no case left to miss. The registration is a map insert with no store access, which is what lets it sit on a path taken by every login and every initial sync. Reported by Viktor Liu, who also proposed registering at the point of issue. (cherry picked from commit 2d16dd7d7cf54762f2e64c5630ea092f32ef63ab) * [management] Register for renewal on pushed updates, not only on connect Registering where the nonce is stamped only covered the login and the initial sync, which both happen when a peer opens a stream. That left out the path the mechanism exists for. On the cloud the network map controller is wrapped so that an update publishes to an event bus instead of pushing locally: an instance handling a REST change broadcasts, and every instance holding a peer of that account pushes to its own. Those pushes stamp a nonce through the update handler, and nothing there registered, so an instance learned about an account only when one of its peers happened to reconnect. For a quiet fleet that is the original bug: the check is created, the peers are told about it, and nobody renews what they were told. Registering on the pushed update closes it, and is the difference between stamping and marking a peer connected — one happens on every push, the other only when a stream opens. Reported by Viktor Liu; the broadcast that makes it work was pointed out by Pascal Fischer. (cherry picked from commit 59efe8d93e53bacdf57cb546f4ab2c19dc4eddab) * [management] Let the challenge refresh loop stop with the manager that owns it The loop was started on a context explicitly detached from the caller's, so nothing could ever stop it. Production is unaffected either way, since BuildManager is called with context.Background(), but a test that builds a manager leaked a sweeping goroutine for the rest of the run, and a shutdown path added later would have had no way to reach it. Take the manager's context as the request buffer built on the line above already does. The test pins the contract the loop offers, so a detached context cannot come back inside Start either. * [management] Bound one account's challenge refresh so it cannot starve the rest Resolving which peers answer a challenge reads the store three times, and the refresher sweeps accounts one after another on a single goroutine. A read that never returns held the sweep for the life of the process, so every other account on the instance stopped being renewed and its peers fell out of the policies gated on the check: one account's bad luck became an outage for all of them. Give each refresh the sweep interval it is allowed to occupy, capped at 30s so a 12-hour window does not grant minutes to a query that should take milliseconds. A refresh that runs out of time keeps its account tracked, since a deadline says nothing about whether that account still has a certificate check. * [management] Send challenge refreshes down the path the rest of management uses The refresh dispatched through UpdateAffectedPeers, the one variant that takes no reason, so it was missing from the update counters and coalesced with nothing. An administrator editing a policy while the sweep ran made the account's network map twice over, and UpdateOperationRefresh, added for exactly this caller, was never referenced. Buffer it with a posture_check/refresh reason instead. The periodic push is now visible in the metrics as what it is, distinct from an edit, and the send detaches from the sweep deadline on its own, so that deadline bounds the store reads it was meant for. * [management] Keep the certificate challenge comments to what the history does not say Four of these ran to three and four times the comment budget, the longest at 992 characters. Most of the excess argued against designs that were never written or explained a bug that no longer exists in the code, which is what the commit that fixed it is for. What is left is the part a reader cannot recover from the code: that the nonce secret has to be persisted and unpredictable, that stamping and renewing are decided together because only the serving instance can push, and that the target rule is the inverse of processPeerPostureChecks. * Keep the newest posture checks pending whatever made their meta sync fail * Report no lost certificate when the engine stops during a proof collection * Share the proof collection single-flight across engine restarts * Close a PKCS#11 module that loads but cannot be used * Fix the pending checks comments * Renew certificate challenges only for the peers streamed to this instance * Ignore a challenge stamp from an older sync stream of the same peer * Kill the Windows certificate proof helper with its whole process tree * Expect the challenge untrack in the session ownership test * Drop an invalid certificate proof without discarding the valid ones * Start a system info gathering beside one that has been stuck for ten timeouts --------- Co-authored-by: pascal <pascal@netbird.io> Co-authored-by: mlsmaycon <mlsmaycon@gmail.com> Co-authored-by: riccardom <riccardomanfrin@gmail.com>
1315 lines
48 KiB
Go
1315 lines
48 KiB
Go
package grpc
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
"net/netip"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
jwtv5 "github.com/golang-jwt/jwt/v5"
|
|
pb "github.com/golang/protobuf/proto" // nolint
|
|
"github.com/golang/protobuf/ptypes/timestamp"
|
|
"github.com/grpc-ecosystem/go-grpc-middleware/v2/interceptors/realip"
|
|
log "github.com/sirupsen/logrus"
|
|
"golang.zx2c4.com/wireguard/wgctrl/wgtypes"
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/peer"
|
|
"google.golang.org/grpc/status"
|
|
|
|
"github.com/netbirdio/netbird/shared/management/client/common"
|
|
"github.com/netbirdio/netbird/shared/management/grpc"
|
|
|
|
"github.com/netbirdio/netbird/management/internals/controllers/network_map"
|
|
rpservice "github.com/netbirdio/netbird/management/internals/modules/reverseproxy/service"
|
|
nbconfig "github.com/netbirdio/netbird/management/internals/server/config"
|
|
"github.com/netbirdio/netbird/management/server/idp"
|
|
"github.com/netbirdio/netbird/management/server/job"
|
|
|
|
"github.com/netbirdio/netbird/management/server/integrations/integrated_validator"
|
|
"github.com/netbirdio/netbird/management/server/store"
|
|
|
|
"github.com/netbirdio/netbird/encryption"
|
|
"github.com/netbirdio/netbird/management/server/account"
|
|
"github.com/netbirdio/netbird/management/server/activity"
|
|
"github.com/netbirdio/netbird/management/server/auth"
|
|
nbContext "github.com/netbirdio/netbird/management/server/context"
|
|
nbpeer "github.com/netbirdio/netbird/management/server/peer"
|
|
"github.com/netbirdio/netbird/management/server/settings"
|
|
"github.com/netbirdio/netbird/management/server/telemetry"
|
|
"github.com/netbirdio/netbird/management/server/types"
|
|
"github.com/netbirdio/netbird/shared/management/certposture"
|
|
"github.com/netbirdio/netbird/shared/management/networkmap/nmdata"
|
|
"github.com/netbirdio/netbird/shared/management/proto"
|
|
internalStatus "github.com/netbirdio/netbird/shared/management/status"
|
|
)
|
|
|
|
const (
|
|
envLogBlockedPeers = "NB_LOG_BLOCKED_PEERS"
|
|
envBlockPeers = "NB_BLOCK_SAME_PEERS"
|
|
envConcurrentSyncs = "NB_MAX_CONCURRENT_SYNCS"
|
|
|
|
defaultSyncLim = 1000
|
|
)
|
|
|
|
// Server an instance of a Management gRPC API server
|
|
type Server struct {
|
|
accountManager account.Manager
|
|
settingsManager settings.Manager
|
|
proto.UnimplementedManagementServiceServer
|
|
jobManager *job.Manager
|
|
config *nbconfig.Config
|
|
secretsManager SecretsManager
|
|
appMetrics telemetry.AppMetrics
|
|
peerLocks sync.Map
|
|
authManager auth.Manager
|
|
sessionStore *auth.SessionStore
|
|
challenger *certposture.Challenger
|
|
|
|
logBlockedPeers bool
|
|
blockPeersWithSameConfig bool
|
|
integratedPeerValidator integrated_validator.IntegratedValidator
|
|
|
|
loginFilter *loginFilter
|
|
|
|
networkMapController network_map.Controller
|
|
|
|
oAuthConfigProvider idp.OAuthConfigProvider
|
|
|
|
syncSem atomic.Int32
|
|
syncLimEnabled bool
|
|
syncLim int32
|
|
|
|
reverseProxyManager rpservice.Manager
|
|
reverseProxyMu sync.RWMutex
|
|
}
|
|
|
|
// NewServer creates a new Management server
|
|
func NewServer(
|
|
config *nbconfig.Config,
|
|
accountManager account.Manager,
|
|
settingsManager settings.Manager,
|
|
jobManager *job.Manager,
|
|
secretsManager SecretsManager,
|
|
appMetrics telemetry.AppMetrics,
|
|
authManager auth.Manager,
|
|
integratedPeerValidator integrated_validator.IntegratedValidator,
|
|
networkMapController network_map.Controller,
|
|
oAuthConfigProvider idp.OAuthConfigProvider,
|
|
sessionStore *auth.SessionStore,
|
|
) (*Server, error) {
|
|
if appMetrics != nil {
|
|
// update gauge based on number of connected peers which is equal to open gRPC streams
|
|
err := appMetrics.GRPCMetrics().RegisterConnectedStreams(func() int64 {
|
|
return int64(networkMapController.CountStreams())
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
logBlockedPeers := strings.ToLower(os.Getenv(envLogBlockedPeers)) == "true"
|
|
blockPeersWithSameConfig := strings.ToLower(os.Getenv(envBlockPeers)) == "true"
|
|
|
|
syncLim := int32(defaultSyncLim)
|
|
syncLimEnabled := true
|
|
if syncLimStr := os.Getenv(envConcurrentSyncs); syncLimStr != "" {
|
|
syncLimParsed, err := strconv.Atoi(syncLimStr)
|
|
if err != nil {
|
|
log.Errorf("invalid value for %s: %v using %d", envConcurrentSyncs, err, defaultSyncLim)
|
|
} else {
|
|
//nolint:gosec
|
|
syncLim = int32(syncLimParsed)
|
|
if syncLim < 0 {
|
|
syncLimEnabled = false
|
|
}
|
|
}
|
|
}
|
|
|
|
serverKey, err := secretsManager.GetWGKey()
|
|
if err != nil {
|
|
return nil, fmt.Errorf("get server WireGuard key: %w", err)
|
|
}
|
|
|
|
return &Server{
|
|
challenger: newCertChallenger(config.DataStoreEncryptionKey, serverKey),
|
|
jobManager: jobManager,
|
|
accountManager: accountManager,
|
|
settingsManager: settingsManager,
|
|
config: config,
|
|
secretsManager: secretsManager,
|
|
authManager: authManager,
|
|
appMetrics: appMetrics,
|
|
logBlockedPeers: logBlockedPeers,
|
|
blockPeersWithSameConfig: blockPeersWithSameConfig,
|
|
integratedPeerValidator: integratedPeerValidator,
|
|
networkMapController: networkMapController,
|
|
oAuthConfigProvider: oAuthConfigProvider,
|
|
sessionStore: sessionStore,
|
|
|
|
loginFilter: newLoginFilter(),
|
|
|
|
syncLim: syncLim,
|
|
syncLimEnabled: syncLimEnabled,
|
|
}, nil
|
|
}
|
|
|
|
func (s *Server) GetServerKey(ctx context.Context, req *proto.Empty) (*proto.ServerKeyResponse, error) {
|
|
ip := ""
|
|
p, ok := peer.FromContext(ctx)
|
|
if ok {
|
|
ip = p.Addr.String()
|
|
}
|
|
|
|
log.WithContext(ctx).Tracef("GetServerKey request from %s", ip)
|
|
|
|
// todo introduce something more meaningful with the key expiration/rotation
|
|
if s.appMetrics != nil {
|
|
s.appMetrics.GRPCMetrics().CountGetKeyRequest()
|
|
}
|
|
now := time.Now().Add(24 * time.Hour)
|
|
secs := int64(now.Second())
|
|
nanos := int32(now.Nanosecond())
|
|
expiresAt := ×tamp.Timestamp{Seconds: secs, Nanos: nanos}
|
|
|
|
key, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
log.WithContext(ctx).Errorf("failed to get wireguard key: %v", err)
|
|
return nil, errors.New("failed to get wireguard key")
|
|
}
|
|
|
|
return &proto.ServerKeyResponse{
|
|
Key: key.PublicKey().String(),
|
|
ExpiresAt: expiresAt,
|
|
}, nil
|
|
}
|
|
|
|
func getRealIP(ctx context.Context) net.IP {
|
|
if addr, ok := realip.FromContext(ctx); ok {
|
|
return net.IP(addr.AsSlice())
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (s *Server) Job(srv proto.ManagementService_JobServer) error {
|
|
reqStart := time.Now()
|
|
ctx := srv.Context()
|
|
|
|
peerKey, err := s.handleHandshake(ctx, srv)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
accountID, err := s.accountManager.GetAccountIDForPeerKey(ctx, peerKey.String())
|
|
if err != nil {
|
|
// nolint:staticcheck
|
|
ctx = context.WithValue(ctx, nbContext.AccountIDKey, "UNKNOWN")
|
|
log.WithContext(ctx).Tracef("peer %s is not registered", peerKey.String())
|
|
if errStatus, ok := internalStatus.FromError(err); ok && errStatus.Type() == internalStatus.NotFound {
|
|
return status.Errorf(codes.PermissionDenied, "peer is not registered")
|
|
}
|
|
return err
|
|
}
|
|
// nolint:staticcheck
|
|
ctx = context.WithValue(ctx, nbContext.AccountIDKey, accountID)
|
|
peer, err := s.accountManager.GetStore().GetPeerByPeerPubKey(ctx, store.LockingStrengthNone, peerKey.String())
|
|
if err != nil {
|
|
return status.Errorf(codes.Unauthenticated, "peer is not registered")
|
|
}
|
|
|
|
s.startResponseReceiver(ctx, srv)
|
|
|
|
updates := s.jobManager.CreateJobChannel(ctx, accountID, peer.ID)
|
|
log.WithContext(ctx).Debugf("Job: took %v", time.Since(reqStart))
|
|
|
|
return s.sendJobsLoop(ctx, accountID, peerKey, peer, updates, srv)
|
|
}
|
|
|
|
// Sync validates the existence of a connecting peer, sends an initial state (all available for the connecting peers) and
|
|
// notifies the connected peer of any updates (e.g. new peers under the same account)
|
|
func (s *Server) Sync(req *proto.EncryptedMessage, srv proto.ManagementService_SyncServer) error {
|
|
if s.syncLimEnabled && s.syncSem.Load() >= s.syncLim {
|
|
return status.Errorf(codes.ResourceExhausted, "too many concurrent sync requests, please try again later")
|
|
}
|
|
s.syncSem.Add(1)
|
|
|
|
reqStart := time.Now()
|
|
syncStart := reqStart.UTC()
|
|
|
|
ctx := srv.Context()
|
|
|
|
syncReq := &proto.SyncRequest{}
|
|
peerKey, err := s.parseRequest(ctx, req, syncReq)
|
|
if err != nil {
|
|
s.syncSem.Add(-1)
|
|
return err
|
|
}
|
|
realIP := getRealIP(ctx)
|
|
sRealIP := realIP.String()
|
|
peerMeta := extractPeerMeta(ctx, syncReq.GetMeta())
|
|
peerMeta.Certificates = s.verifiedCertificates(ctx, peerKey, syncReq.GetMeta().GetCertificateProofs())
|
|
|
|
metahashed := metaHash(peerMeta)
|
|
if !s.loginFilter.allowLogin(peerKey.String(), metahashed) {
|
|
if s.appMetrics != nil {
|
|
s.appMetrics.GRPCMetrics().CountSyncRequestBlocked()
|
|
}
|
|
if s.logBlockedPeers {
|
|
log.WithContext(ctx).Tracef("peer %s with meta hash %d is blocked from syncing", peerKey.String(), metahashed)
|
|
}
|
|
if s.blockPeersWithSameConfig {
|
|
s.syncSem.Add(-1)
|
|
return mapError(ctx, internalStatus.ErrPeerAlreadyLoggedIn)
|
|
}
|
|
}
|
|
|
|
if s.appMetrics != nil {
|
|
s.appMetrics.GRPCMetrics().CountSyncRequest()
|
|
}
|
|
|
|
// nolint:staticcheck
|
|
ctx = context.WithValue(ctx, nbContext.PeerIDKey, peerKey.String())
|
|
|
|
accountID, err := s.accountManager.GetAccountIDForPeerKey(ctx, peerKey.String())
|
|
if err != nil {
|
|
// nolint:staticcheck
|
|
ctx = context.WithValue(ctx, nbContext.AccountIDKey, "UNKNOWN")
|
|
log.WithContext(ctx).Tracef("peer %s is not registered", peerKey.String())
|
|
if errStatus, ok := internalStatus.FromError(err); ok && errStatus.Type() == internalStatus.NotFound {
|
|
s.syncSem.Add(-1)
|
|
return status.Errorf(codes.PermissionDenied, "peer is not registered")
|
|
}
|
|
s.syncSem.Add(-1)
|
|
return err
|
|
}
|
|
|
|
// nolint:staticcheck
|
|
ctx = context.WithValue(ctx, nbContext.AccountIDKey, accountID)
|
|
|
|
start := time.Now()
|
|
unlock := s.acquirePeerLockByUID(ctx, peerKey.String())
|
|
defer func() {
|
|
if unlock != nil {
|
|
unlock()
|
|
}
|
|
}()
|
|
log.WithContext(ctx).Tracef("acquired peer lock for peer %s took %v", peerKey.String(), time.Since(start))
|
|
|
|
log.WithContext(ctx).Debugf("Sync request from peer [%s] [%s]", req.WgPubKey, sRealIP)
|
|
|
|
if syncReq.GetMeta() == nil {
|
|
log.WithContext(ctx).Tracef("peer system meta has to be provided on sync. Peer %s, remote addr %s", peerKey.String(), realIP)
|
|
}
|
|
|
|
metahash := metaHash(peerMeta)
|
|
s.loginFilter.addLogin(peerKey.String(), metahash)
|
|
|
|
peer, netMap, postureChecks, dnsFwdPort, err := s.accountManager.SyncAndMarkPeer(ctx, accountID, peerKey.String(), peerMeta, realIP, syncStart)
|
|
if err != nil {
|
|
log.WithContext(ctx).Debugf("error while syncing peer %s: %v", peerKey.String(), err)
|
|
s.syncSem.Add(-1)
|
|
return mapError(ctx, err)
|
|
}
|
|
|
|
trackChallenge := func() { s.accountManager.TrackCertificateChallenges(ctx, accountID, peer.ID, syncStart) }
|
|
err = s.sendInitialSync(ctx, peerKey, peer, netMap, postureChecks, srv, dnsFwdPort, trackChallenge)
|
|
if err != nil {
|
|
log.WithContext(ctx).Debugf("error while sending initial sync for %s: %v", peerKey.String(), err)
|
|
s.syncSem.Add(-1)
|
|
s.cancelPeerRoutinesWithoutLock(ctx, accountID, peer, syncStart, nil)
|
|
return err
|
|
}
|
|
|
|
updates, err := s.networkMapController.OnPeerConnected(ctx, accountID, peer.ID)
|
|
if err != nil {
|
|
log.WithContext(ctx).Debugf("error while notify peer connected for %s: %v", peerKey.String(), err)
|
|
s.syncSem.Add(-1)
|
|
s.cancelPeerRoutinesWithoutLock(ctx, accountID, peer, syncStart, nil)
|
|
return err
|
|
}
|
|
|
|
s.secretsManager.SetupRefresh(ctx, accountID, peer.ID)
|
|
|
|
unlock()
|
|
unlock = nil
|
|
|
|
if s.appMetrics != nil {
|
|
s.appMetrics.GRPCMetrics().CountSyncRequestDuration(time.Since(reqStart), accountID)
|
|
}
|
|
log.WithContext(ctx).Debugf("Sync took %s", time.Since(reqStart))
|
|
|
|
s.syncSem.Add(-1)
|
|
|
|
return PeerUpdateHandlerFactory(peerKey, updates, s.secretsManager, s.challenger,
|
|
trackChallenge,
|
|
srv, func() { s.cancelPeerRoutines(ctx, accountID, peer, syncStart, updates) }).
|
|
WithMetrics(s.appMetrics).HandleUpdates(ctx)
|
|
}
|
|
|
|
func (s *Server) handleHandshake(ctx context.Context, srv proto.ManagementService_JobServer) (wgtypes.Key, error) {
|
|
hello, err := srv.Recv()
|
|
if err != nil {
|
|
return wgtypes.Key{}, status.Errorf(codes.InvalidArgument, "missing hello: %v", err)
|
|
}
|
|
|
|
jobReq := &proto.JobRequest{}
|
|
peerKey, err := s.parseRequest(ctx, hello, jobReq)
|
|
if err != nil {
|
|
return wgtypes.Key{}, err
|
|
}
|
|
|
|
return peerKey, nil
|
|
}
|
|
|
|
func (s *Server) startResponseReceiver(ctx context.Context, srv proto.ManagementService_JobServer) {
|
|
go func() {
|
|
for {
|
|
msg, err := srv.Recv()
|
|
if err != nil {
|
|
if errors.Is(err, io.EOF) || errors.Is(err, context.Canceled) {
|
|
return
|
|
}
|
|
log.WithContext(ctx).Warnf("recv job response error: %v", err)
|
|
return
|
|
}
|
|
|
|
jobResp := &proto.JobResponse{}
|
|
if _, err := s.parseRequest(ctx, msg, jobResp); err != nil {
|
|
log.WithContext(ctx).Warnf("invalid job response: %v", err)
|
|
continue
|
|
}
|
|
|
|
if err := s.jobManager.HandleResponse(ctx, jobResp, msg.WgPubKey); err != nil {
|
|
log.WithContext(ctx).Errorf("handle job response failed: %v", err)
|
|
}
|
|
}
|
|
}()
|
|
}
|
|
|
|
func (s *Server) sendJobsLoop(ctx context.Context, accountID string, peerKey wgtypes.Key, peer *nbpeer.Peer, updates *job.Channel, srv proto.ManagementService_JobServer) error {
|
|
// todo figure out better error handling strategy
|
|
defer s.jobManager.CloseChannel(ctx, accountID, peer.ID, updates)
|
|
|
|
for {
|
|
event, err := updates.Event(ctx)
|
|
if err != nil {
|
|
if errors.Is(err, job.ErrJobChannelClosed) {
|
|
log.WithContext(ctx).Debugf("jobs channel for peer %s was closed", peerKey.String())
|
|
return nil
|
|
}
|
|
|
|
// happens when connection drops, e.g. client disconnects
|
|
log.WithContext(ctx).Debugf("stream of peer %s has been closed", peerKey.String())
|
|
return ctx.Err()
|
|
}
|
|
|
|
if err := s.sendJob(ctx, peerKey, event, srv); err != nil {
|
|
log.WithContext(ctx).Warnf("send job failed: %v", err)
|
|
return nil
|
|
}
|
|
}
|
|
}
|
|
|
|
// sendJob encrypts the update message using the peer key and the server's wireguard key,
|
|
// then sends the encrypted message to the connected peer via the sync server.
|
|
func (s *Server) sendJob(ctx context.Context, peerKey wgtypes.Key, job *job.Event, srv proto.ManagementService_JobServer) error {
|
|
wgKey, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
log.WithContext(ctx).Errorf("failed to get wg key for peer %s: %v", peerKey.String(), err)
|
|
return status.Errorf(codes.Internal, "failed processing job message")
|
|
}
|
|
|
|
encryptedResp, err := encryption.EncryptMessage(peerKey, wgKey, job.Request)
|
|
if err != nil {
|
|
log.WithContext(ctx).Errorf("failed to encrypt job for peer %s: %v", peerKey.String(), err)
|
|
return status.Errorf(codes.Internal, "failed processing job message")
|
|
}
|
|
err = srv.Send(&proto.EncryptedMessage{
|
|
WgPubKey: wgKey.PublicKey().String(),
|
|
Body: encryptedResp,
|
|
})
|
|
if err != nil {
|
|
return status.Errorf(codes.Internal, "failed sending job message")
|
|
}
|
|
log.WithContext(ctx).Debugf("sent a job to peer: %s", peerKey.String())
|
|
return nil
|
|
}
|
|
|
|
func (s *Server) cancelPeerRoutines(ctx context.Context, accountID string, peer *nbpeer.Peer, streamStartTime time.Time, session chan *network_map.UpdateMessage) {
|
|
uncanceledCTX := context.WithoutCancel(ctx)
|
|
unlock := s.acquirePeerLockByUID(uncanceledCTX, peer.Key)
|
|
defer unlock()
|
|
|
|
s.cancelPeerRoutinesWithoutLock(uncanceledCTX, accountID, peer, streamStartTime, session)
|
|
}
|
|
|
|
// cancelPeerRoutinesWithoutLock tears down the stream of the session identified by streamStartTime
|
|
// and its updates channel. A nil session means the stream failed before it registered a channel;
|
|
// the controller then closes any channel still registered.
|
|
func (s *Server) cancelPeerRoutinesWithoutLock(ctx context.Context, accountID string, peer *nbpeer.Peer, streamStartTime time.Time, session chan *network_map.UpdateMessage) {
|
|
err := s.accountManager.OnPeerDisconnected(ctx, accountID, peer.Key, streamStartTime)
|
|
if err != nil {
|
|
log.WithContext(ctx).Errorf("failed to disconnect peer %s properly: %v", peer.Key, err)
|
|
}
|
|
if !s.networkMapController.OnPeerDisconnected(ctx, accountID, peer.ID, session) {
|
|
log.WithContext(ctx).Debugf("skipped peer routines teardown for %s: a newer session owns the peer", peer.Key)
|
|
return
|
|
}
|
|
s.secretsManager.CancelRefresh(peer.ID)
|
|
s.accountManager.UntrackCertificateChallenges(accountID, peer.ID, streamStartTime)
|
|
|
|
log.WithContext(ctx).Debugf("peer %s has been disconnected", peer.Key)
|
|
}
|
|
|
|
func (s *Server) validateToken(ctx context.Context, peerKey, jwtToken string) (string, error) {
|
|
if s.authManager == nil {
|
|
return "", status.Errorf(codes.Internal, "missing auth manager")
|
|
}
|
|
|
|
userAuth, token, err := s.authManager.ValidateAndParseToken(ctx, jwtToken)
|
|
if err != nil {
|
|
return "", status.Errorf(codes.InvalidArgument, "invalid jwt token, err: %v", err)
|
|
}
|
|
|
|
if err := s.claimLoginToken(ctx, peerKey, jwtToken, token); err != nil {
|
|
return "", err
|
|
}
|
|
|
|
// we need to call this method because if user is new, we will automatically add it to existing or create a new account
|
|
accountId, _, err := s.accountManager.GetAccountIDFromUserAuth(ctx, userAuth)
|
|
if err != nil {
|
|
return "", status.Errorf(codes.Internal, "unable to fetch account with claims, err: %v", err)
|
|
}
|
|
|
|
if userAuth.AccountId != accountId {
|
|
log.WithContext(ctx).Debugf("gRPC server sets accountId from ensure, before %s, now %s", userAuth.AccountId, accountId)
|
|
userAuth.AccountId = accountId
|
|
}
|
|
|
|
userAuth, err = s.authManager.EnsureUserAccessByJWTGroups(ctx, userAuth, token)
|
|
if err != nil {
|
|
return "", status.Error(codes.PermissionDenied, err.Error())
|
|
}
|
|
|
|
err = s.accountManager.SyncUserJWTGroups(ctx, userAuth)
|
|
if err != nil {
|
|
log.WithContext(ctx).Errorf("gRPC server failed to sync user JWT groups: %s", err)
|
|
}
|
|
|
|
return userAuth.UserId, nil
|
|
}
|
|
|
|
func (s *Server) acquirePeerLockByUID(ctx context.Context, uniqueID string) (unlock func()) {
|
|
log.WithContext(ctx).Tracef("acquiring peer lock for ID %s", uniqueID)
|
|
|
|
start := time.Now()
|
|
value, _ := s.peerLocks.LoadOrStore(uniqueID, &sync.RWMutex{})
|
|
mtx := value.(*sync.RWMutex)
|
|
mtx.Lock()
|
|
log.WithContext(ctx).Tracef("acquired peer lock for ID %s in %v", uniqueID, time.Since(start))
|
|
start = time.Now()
|
|
|
|
unlock = func() {
|
|
mtx.Unlock()
|
|
log.WithContext(ctx).Tracef("released peer lock for ID %s in %v", uniqueID, time.Since(start))
|
|
}
|
|
|
|
return unlock
|
|
}
|
|
|
|
// maps internal internalStatus.Error to gRPC status.Error
|
|
func mapError(ctx context.Context, err error) error {
|
|
if e, ok := internalStatus.FromError(err); ok {
|
|
switch e.Type() {
|
|
case internalStatus.PermissionDenied:
|
|
return status.Error(codes.PermissionDenied, e.Message)
|
|
case internalStatus.Unauthorized:
|
|
return status.Error(codes.PermissionDenied, e.Message)
|
|
case internalStatus.Unauthenticated:
|
|
return status.Error(codes.PermissionDenied, e.Message)
|
|
case internalStatus.PreconditionFailed:
|
|
return status.Error(codes.FailedPrecondition, e.Message)
|
|
case internalStatus.NotFound:
|
|
return status.Error(codes.NotFound, e.Message)
|
|
default:
|
|
}
|
|
}
|
|
if errors.Is(err, internalStatus.ErrPeerAlreadyLoggedIn) {
|
|
return status.Error(codes.PermissionDenied, internalStatus.ErrPeerAlreadyLoggedIn.Error())
|
|
}
|
|
log.WithContext(ctx).Errorf("got an unhandled error: %s", err)
|
|
return status.Errorf(codes.Internal, "failed handling request")
|
|
}
|
|
|
|
func extractPeerMeta(ctx context.Context, meta *proto.PeerSystemMeta) nbpeer.PeerSystemMeta {
|
|
if meta == nil {
|
|
return nbpeer.PeerSystemMeta{}
|
|
}
|
|
|
|
osVersion := meta.GetOSVersion()
|
|
if osVersion == "" {
|
|
osVersion = meta.GetCore()
|
|
}
|
|
|
|
networkAddresses := make([]nbpeer.NetworkAddress, 0, len(meta.GetNetworkAddresses()))
|
|
for _, addr := range meta.GetNetworkAddresses() {
|
|
netAddr, err := netip.ParsePrefix(addr.GetNetIP())
|
|
if err != nil {
|
|
log.WithContext(ctx).Warnf("failed to parse netip address, %s: %v", addr.GetNetIP(), err)
|
|
continue
|
|
}
|
|
networkAddresses = append(networkAddresses, nbpeer.NetworkAddress{
|
|
NetIP: netAddr,
|
|
Mac: addr.GetMac(),
|
|
})
|
|
}
|
|
|
|
files := make([]nbpeer.File, 0, len(meta.GetFiles()))
|
|
for _, file := range meta.GetFiles() {
|
|
files = append(files, nbpeer.File{
|
|
Path: file.GetPath(),
|
|
Exist: file.GetExist(),
|
|
ProcessIsRunning: file.GetProcessIsRunning(),
|
|
})
|
|
}
|
|
|
|
return nbpeer.PeerSystemMeta{
|
|
Hostname: meta.GetHostname(),
|
|
GoOS: meta.GetGoOS(),
|
|
Kernel: meta.GetKernel(),
|
|
Platform: meta.GetPlatform(),
|
|
OS: meta.GetOS(),
|
|
OSVersion: osVersion,
|
|
WtVersion: meta.GetNetbirdVersion(),
|
|
UIVersion: meta.GetUiVersion(),
|
|
KernelVersion: meta.GetKernelVersion(),
|
|
NetworkAddresses: networkAddresses,
|
|
SystemSerialNumber: meta.GetSysSerialNumber(),
|
|
SystemProductName: meta.GetSysProductName(),
|
|
SystemManufacturer: meta.GetSysManufacturer(),
|
|
Environment: nbpeer.Environment{
|
|
Cloud: meta.GetEnvironment().GetCloud(),
|
|
Platform: meta.GetEnvironment().GetPlatform(),
|
|
},
|
|
Flags: nbpeer.Flags{
|
|
RosenpassEnabled: meta.GetFlags().GetRosenpassEnabled(),
|
|
RosenpassPermissive: meta.GetFlags().GetRosenpassPermissive(),
|
|
ServerSSHAllowed: meta.GetFlags().GetServerSSHAllowed(),
|
|
RemoteJobsAllowed: meta.GetFlags().GetRemoteJobsAllowed(),
|
|
DisableClientRoutes: meta.GetFlags().GetDisableClientRoutes(),
|
|
DisableServerRoutes: meta.GetFlags().GetDisableServerRoutes(),
|
|
DisableDNS: meta.GetFlags().GetDisableDNS(),
|
|
DisableFirewall: meta.GetFlags().GetDisableFirewall(),
|
|
BlockLANAccess: meta.GetFlags().GetBlockLANAccess(),
|
|
BlockInbound: meta.GetFlags().GetBlockInbound(),
|
|
LazyConnectionEnabled: meta.GetFlags().GetLazyConnectionEnabled(),
|
|
DisableIPv6: meta.GetFlags().GetDisableIPv6(),
|
|
},
|
|
Files: files,
|
|
Capabilities: capabilitiesToInt32(meta.GetCapabilities()),
|
|
SyncMessageVersion: int(meta.GetSyncMessageVersion()),
|
|
}
|
|
}
|
|
|
|
func capabilitiesToInt32(caps []proto.PeerCapability) []int32 {
|
|
result := make([]int32, len(caps))
|
|
for i, c := range caps {
|
|
result[i] = int32(c)
|
|
}
|
|
return result
|
|
}
|
|
|
|
func (s *Server) parseRequest(ctx context.Context, req *proto.EncryptedMessage, parsed pb.Message) (wgtypes.Key, error) {
|
|
peerKey, err := wgtypes.ParseKey(req.GetWgPubKey())
|
|
if err != nil {
|
|
log.WithContext(ctx).Warnf("error while parsing peer's WireGuard public key %s.", req.WgPubKey)
|
|
return wgtypes.Key{}, status.Errorf(codes.InvalidArgument, "provided wgPubKey %s is invalid", req.WgPubKey)
|
|
}
|
|
|
|
key, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
return wgtypes.Key{}, status.Errorf(codes.Internal, "failed processing request")
|
|
}
|
|
|
|
err = encryption.DecryptMessage(peerKey, key, req.Body, parsed)
|
|
if err != nil {
|
|
return wgtypes.Key{}, status.Errorf(codes.InvalidArgument, "invalid request message")
|
|
}
|
|
|
|
return peerKey, nil
|
|
}
|
|
|
|
// Login endpoint first checks whether peer is registered under any account
|
|
// In case it is, the login is successful
|
|
// In case it isn't, the endpoint checks whether setup key is provided within the request and tries to register a peer.
|
|
// In case of the successful registration login is also successful
|
|
func (s *Server) Login(ctx context.Context, req *proto.EncryptedMessage) (*proto.EncryptedMessage, error) {
|
|
reqStart := time.Now()
|
|
realIP := getRealIP(ctx)
|
|
sRealIP := realIP.String()
|
|
|
|
loginReq := &proto.LoginRequest{}
|
|
peerKey, err := s.parseRequest(ctx, req, loginReq)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
peerMeta := extractPeerMeta(ctx, loginReq.GetMeta())
|
|
peerMeta.Certificates = s.verifiedCertificates(ctx, peerKey, loginReq.GetMeta().GetCertificateProofs())
|
|
metahashed := metaHash(peerMeta)
|
|
if !s.loginFilter.allowLogin(peerKey.String(), metahashed) {
|
|
if s.logBlockedPeers {
|
|
log.WithContext(ctx).Tracef("peer %s with meta hash %d is blocked from login", peerKey.String(), metahashed)
|
|
}
|
|
if s.appMetrics != nil {
|
|
s.appMetrics.GRPCMetrics().CountLoginRequestBlocked()
|
|
}
|
|
if s.blockPeersWithSameConfig {
|
|
return nil, internalStatus.ErrPeerAlreadyLoggedIn
|
|
}
|
|
}
|
|
|
|
if s.appMetrics != nil {
|
|
s.appMetrics.GRPCMetrics().CountLoginRequest()
|
|
}
|
|
|
|
//nolint
|
|
ctx = context.WithValue(ctx, nbContext.PeerIDKey, peerKey.String())
|
|
accountID, err := s.accountManager.GetAccountIDForPeerKey(ctx, peerKey.String())
|
|
if err != nil {
|
|
// this case should not happen and already indicates an issue but we don't want the system to fail due to being unable to log in detail
|
|
accountID = "UNKNOWN"
|
|
}
|
|
//nolint
|
|
ctx = context.WithValue(ctx, nbContext.AccountIDKey, accountID)
|
|
|
|
log.WithContext(ctx).Debugf("Login request from peer [%s] [%s]", req.WgPubKey, sRealIP)
|
|
|
|
if loginReq.GetMeta() == nil {
|
|
msg := status.Errorf(codes.FailedPrecondition,
|
|
"peer system meta has to be provided to log in. Peer %s, remote addr %s", peerKey.String(), realIP)
|
|
log.WithContext(ctx).Warn(msg)
|
|
return nil, msg
|
|
}
|
|
|
|
userID, err := s.processJwtToken(ctx, loginReq, peerKey)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
var sshKey []byte
|
|
if loginReq.GetPeerKeys() != nil {
|
|
sshKey = loginReq.GetPeerKeys().GetSshPubKey()
|
|
}
|
|
|
|
peer, network, postureChecks, enableSSH, err := s.accountManager.LoginPeer(ctx, types.PeerLogin{
|
|
WireGuardPubKey: peerKey.String(),
|
|
SSHKey: string(sshKey),
|
|
Meta: peerMeta,
|
|
UserID: userID,
|
|
SetupKey: loginReq.GetSetupKey(),
|
|
ConnectionIP: realIP,
|
|
ExtraDNSLabels: loginReq.GetDnsLabels(),
|
|
})
|
|
if err != nil {
|
|
if errors.Is(err, internalStatus.ErrNoAuthMethodProvided) {
|
|
log.WithContext(ctx).Tracef("failed logging in peer %s: %s", peerKey, err)
|
|
} else {
|
|
log.WithContext(ctx).Warnf("failed logging in peer %s: %s", peerKey, err)
|
|
}
|
|
return nil, mapError(ctx, err)
|
|
}
|
|
|
|
loginResp, err := s.prepareLoginResponse(ctx, peer, network, postureChecks, enableSSH)
|
|
if err != nil {
|
|
log.WithContext(ctx).Warnf("failed preparing login response for peer %s: %s", peerKey, err)
|
|
return nil, status.Errorf(codes.Internal, "failed logging in peer")
|
|
}
|
|
|
|
key, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
log.WithContext(ctx).Warnf("failed getting server's WireGuard private key: %s", err)
|
|
return nil, status.Errorf(codes.Internal, "failed logging in peer")
|
|
}
|
|
|
|
// Renewal is tracked by the sync stream that follows, on whichever instance it lands.
|
|
stampCertificateChallenges(loginResp.Checks, s.challenger, peerKey)
|
|
encryptedResp, err := encryption.EncryptMessage(peerKey, key, loginResp)
|
|
if err != nil {
|
|
log.WithContext(ctx).Warnf("failed encrypting peer %s message", peer.ID)
|
|
return nil, status.Errorf(codes.Internal, "failed logging in peer")
|
|
}
|
|
|
|
if s.appMetrics != nil {
|
|
s.appMetrics.GRPCMetrics().CountLoginRequestDuration(time.Since(reqStart), accountID)
|
|
}
|
|
log.WithContext(ctx).Debugf("Login took %s", time.Since(reqStart))
|
|
|
|
return &proto.EncryptedMessage{
|
|
WgPubKey: key.PublicKey().String(),
|
|
Body: encryptedResp,
|
|
}, nil
|
|
}
|
|
|
|
// ExtendAuthSession refreshes the peer's SSO session expiry deadline using a
|
|
// fresh JWT. The same JWT validation pipeline as Login is used. The tunnel
|
|
// stays up; no network map sync is performed. The new deadline is returned
|
|
// in ExtendAuthSessionResponse.SessionExpiresAt.
|
|
func (s *Server) ExtendAuthSession(ctx context.Context, req *proto.EncryptedMessage) (*proto.EncryptedMessage, error) {
|
|
extendReq := &proto.ExtendAuthSessionRequest{}
|
|
peerKey, err := s.parseRequest(ctx, req, extendReq)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
//nolint
|
|
ctx = context.WithValue(ctx, nbContext.PeerIDKey, peerKey.String())
|
|
if accountID, accErr := s.accountManager.GetAccountIDForPeerKey(ctx, peerKey.String()); accErr == nil {
|
|
//nolint
|
|
ctx = context.WithValue(ctx, nbContext.AccountIDKey, accountID)
|
|
}
|
|
|
|
jwt := extendReq.GetJwtToken()
|
|
if jwt == "" {
|
|
return nil, status.Errorf(codes.InvalidArgument, "jwt token is required")
|
|
}
|
|
|
|
var userID string
|
|
const attempts = 3
|
|
for i := 0; i < attempts; i++ {
|
|
userID, err = s.validateToken(ctx, peerKey.String(), jwt)
|
|
if err == nil {
|
|
break
|
|
}
|
|
if i == attempts-1 {
|
|
break
|
|
}
|
|
log.WithContext(ctx).Warnf("failed validating JWT token while extending session for peer %s: %v. Retrying (idP cache).", peerKey.String(), err)
|
|
select {
|
|
case <-time.After(200 * time.Millisecond):
|
|
case <-ctx.Done():
|
|
return nil, ctx.Err()
|
|
}
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if userID == "" {
|
|
return nil, status.Errorf(codes.Unauthenticated, "jwt token did not yield a user id")
|
|
}
|
|
|
|
deadline, err := s.accountManager.ExtendPeerSession(ctx, peerKey.String(), userID)
|
|
if err != nil {
|
|
log.WithContext(ctx).Warnf("failed extending session for peer %s: %v", peerKey.String(), err)
|
|
return nil, mapError(ctx, err)
|
|
}
|
|
|
|
// Success path normally returns a non-zero deadline. A defensive zero
|
|
// would still encode as the explicit "disabled" sentinel rather than nil,
|
|
// so the client clears any stale anchor instead of preserving it.
|
|
resp := &proto.ExtendAuthSessionResponse{
|
|
SessionExpiresAt: encodeSessionExpiresAt(deadline),
|
|
}
|
|
|
|
wgKey, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
return nil, status.Errorf(codes.Internal, "failed processing request")
|
|
}
|
|
encrypted, err := encryption.EncryptMessage(peerKey, wgKey, resp)
|
|
if err != nil {
|
|
return nil, status.Errorf(codes.Internal, "failed encrypting response")
|
|
}
|
|
return &proto.EncryptedMessage{
|
|
WgPubKey: wgKey.PublicKey().String(),
|
|
Body: encrypted,
|
|
}, nil
|
|
}
|
|
|
|
func (s *Server) prepareLoginResponse(ctx context.Context, peer *nbpeer.Peer, network *types.Network, postureChecks []*nmdata.PostureChecks, enableSSH bool) (*proto.LoginResponse, error) {
|
|
var relayToken *Token
|
|
var err error
|
|
if s.config.Relay != nil && len(s.config.Relay.Addresses) > 0 {
|
|
relayToken, err = s.secretsManager.GenerateRelayToken()
|
|
if err != nil {
|
|
log.Errorf("failed generating Relay token: %v", err)
|
|
}
|
|
}
|
|
|
|
settings, err := s.settingsManager.GetSettings(ctx, peer.AccountID, activity.SystemInitiator)
|
|
if err != nil {
|
|
log.WithContext(ctx).Warnf("failed getting settings for peer %s: %s", peer.Key, err)
|
|
return nil, status.Errorf(codes.Internal, "failed getting settings")
|
|
}
|
|
|
|
// if peer has reached this point then it has logged in
|
|
loginResp := &proto.LoginResponse{
|
|
NetbirdConfig: toNetbirdConfig(s.config, nil, relayToken, nil, types.TwinAccountSettings(settings)),
|
|
PeerConfig: toPeerConfig(types.TwinPeer(peer), types.TwinNetwork(network), s.networkMapController.GetDNSDomain(settings), types.TwinAccountSettings(settings), s.config.HttpConfig, s.config.DeviceAuthorizationFlow, enableSSH, false),
|
|
Checks: toProtocolChecks(ctx, postureChecks),
|
|
}
|
|
|
|
// settings is always non-nil here, so we never emit nil — encoder returns
|
|
// either a valid deadline or the explicit-zero "disabled" sentinel.
|
|
loginResp.SessionExpiresAt = encodeSessionExpiresAt(
|
|
peer.SessionExpiresAt(settings.PeerLoginExpirationEnabled, settings.PeerLoginExpiration),
|
|
)
|
|
|
|
return loginResp, nil
|
|
}
|
|
|
|
func (s *Server) claimLoginToken(ctx context.Context, peerKey, jwtToken string, token *jwtv5.Token) error {
|
|
if s.sessionStore == nil || token == nil {
|
|
return nil
|
|
}
|
|
|
|
exp, err := token.Claims.GetExpirationTime()
|
|
if err != nil || exp == nil {
|
|
log.WithContext(ctx).Warnf("JWT has no usable exp claim for peer %s", peerKey)
|
|
return status.Error(codes.Unauthenticated, "jwt token has no expiration")
|
|
}
|
|
|
|
err = s.sessionStore.RegisterToken(ctx, jwtToken, exp.Time)
|
|
if err == nil {
|
|
return nil
|
|
}
|
|
|
|
if errors.Is(err, auth.ErrTokenAlreadyUsed) || errors.Is(err, auth.ErrTokenExpired) {
|
|
log.WithContext(ctx).Warnf("%v for peer %s", err, peerKey)
|
|
return status.Error(codes.Unauthenticated, err.Error())
|
|
}
|
|
|
|
log.WithContext(ctx).Warnf("failed to claim JWT for peer %s: %v", peerKey, err)
|
|
return status.Error(codes.Unavailable, "failed to claim jwt token")
|
|
}
|
|
|
|
// processJwtToken validates the existence of a JWT token in the login request, and returns the corresponding user ID if
|
|
// the token is valid.
|
|
//
|
|
// The user ID can be empty if the token is not provided, which is acceptable if the peer is already
|
|
// registered or if it uses a setup key to register.
|
|
func (s *Server) processJwtToken(ctx context.Context, loginReq *proto.LoginRequest, peerKey wgtypes.Key) (string, error) {
|
|
userID := ""
|
|
if loginReq.GetJwtToken() != "" {
|
|
var err error
|
|
for i := 0; i < 3; i++ {
|
|
userID, err = s.validateToken(ctx, peerKey.String(), loginReq.GetJwtToken())
|
|
if err == nil {
|
|
break
|
|
}
|
|
log.WithContext(ctx).Warnf("failed validating JWT token sent from peer %s with error %v. "+
|
|
"Trying again as it may be due to the IdP cache issue", peerKey.String(), err)
|
|
time.Sleep(200 * time.Millisecond)
|
|
}
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
}
|
|
return userID, nil
|
|
}
|
|
|
|
// IsHealthy indicates whether the service is healthy
|
|
func (s *Server) IsHealthy(ctx context.Context, req *proto.Empty) (*proto.Empty, error) {
|
|
return &proto.Empty{}, nil
|
|
}
|
|
|
|
// sendInitialSync sends initial proto.SyncResponse to the peer requesting synchronization
|
|
func (s *Server) sendInitialSync(ctx context.Context, peerKey wgtypes.Key, peer *nbpeer.Peer, networkMap *types.NetworkMap, postureChecks []*nmdata.PostureChecks, srv proto.ManagementService_SyncServer, dnsFwdPort int64, onChallengeStamped func()) error {
|
|
var err error
|
|
var turnToken *Token
|
|
|
|
if s.config.TURNConfig != nil && s.config.TURNConfig.TimeBasedCredentials {
|
|
turnToken, err = s.secretsManager.GenerateTurnToken()
|
|
if err != nil {
|
|
log.Errorf("failed generating TURN token: %v", err)
|
|
}
|
|
}
|
|
|
|
var relayToken *Token
|
|
if s.config.Relay != nil && len(s.config.Relay.Addresses) > 0 {
|
|
relayToken, err = s.secretsManager.GenerateRelayToken()
|
|
if err != nil {
|
|
log.Errorf("failed generating Relay token: %v", err)
|
|
}
|
|
}
|
|
|
|
settings, err := s.settingsManager.GetSettings(ctx, peer.AccountID, activity.SystemInitiator)
|
|
if err != nil {
|
|
return status.Errorf(codes.Internal, "error handling request")
|
|
}
|
|
|
|
peerGroups, err := s.accountManager.GetStore().GetPeerGroupIDs(ctx, store.LockingStrengthNone, peer.AccountID, peer.ID)
|
|
if err != nil {
|
|
return status.Errorf(codes.Internal, "failed to get peer groups %s", err)
|
|
}
|
|
|
|
dnsName := s.networkMapController.GetDNSDomain(settings)
|
|
|
|
var plainResp *proto.SyncResponse
|
|
|
|
commonSyncMessageVersion := grpc.HighestCommonSyncMessageVersion(
|
|
s.perAccountOrGlobalSyncMessageVersions(peer.AccountID),
|
|
grpc.SyncMessageVersionFromConfig(&peer.Meta.SyncMessageVersion))
|
|
|
|
log.WithContext(ctx).
|
|
WithFields(log.Fields{
|
|
"sync_message_version": commonSyncMessageVersion,
|
|
"server_sync_message_version": s.perAccountOrGlobalSyncMessageVersions(peer.AccountID),
|
|
"peer_sync_message_version": grpc.SyncMessageVersionFromConfig(&peer.Meta.SyncMessageVersion),
|
|
}).Debug("common highest sync message version")
|
|
|
|
if commonSyncMessageVersion == grpc.ComponentNetworkMap {
|
|
// Capable peer: discard the legacy NetworkMap that SyncAndMarkPeer
|
|
// computed and recompute the raw components instead. This wastes one
|
|
// Calculate() call per initial-sync — the component-based wire
|
|
// format is what the peer actually consumes. The streaming path
|
|
// (network_map.Controller.UpdateAccountPeers) skips this duplication
|
|
// because it dispatches by capability before computing.
|
|
//
|
|
// TODO: refactor SyncPeer / SyncAndMarkPeer / their mocks + manager
|
|
// interfaces to return PeerNetworkMapResult so the initial-sync path
|
|
// stops doing duplicate work. Deferred until the client-side
|
|
// decoder lands and there's a real deployment of capability=3 peers
|
|
// worth optimizing for.
|
|
freshPeer, components, freshPostureChecks, freshDnsFwdPort, err := s.networkMapController.GetValidatedPeerWithComponents(ctx, false, peer.AccountID, peer)
|
|
if err != nil {
|
|
log.WithContext(ctx).Errorf("failed to build components for peer %s on initial sync: %v", peer.ID, err)
|
|
return status.Errorf(codes.Internal, "failed to build initial sync envelope")
|
|
}
|
|
plainResp = ToComponentSyncResponse(ctx, s.config, s.config.HttpConfig, s.config.DeviceAuthorizationFlow, types.TwinPeer(freshPeer), turnToken, relayToken, components, dnsName, freshPostureChecks, types.TwinAccountSettings(settings), settings.Extra, peerGroups, freshDnsFwdPort)
|
|
} else {
|
|
plainResp = ToSyncResponse(ctx, s.config, s.config.HttpConfig, s.config.DeviceAuthorizationFlow, types.TwinPeer(peer), turnToken, relayToken, networkMap, dnsName, postureChecks, nil, types.TwinAccountSettings(settings), settings.Extra, peerGroups, dnsFwdPort)
|
|
}
|
|
|
|
key, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
return status.Errorf(codes.Internal, "failed getting server key")
|
|
}
|
|
|
|
if stampCertificateChallenges(plainResp.Checks, s.challenger, peerKey) {
|
|
onChallengeStamped()
|
|
}
|
|
encryptedResp, err := encryption.EncryptMessage(peerKey, key, plainResp)
|
|
if err != nil {
|
|
return status.Errorf(codes.Internal, "error handling request")
|
|
}
|
|
|
|
err = srv.Send(&proto.EncryptedMessage{
|
|
WgPubKey: key.PublicKey().String(),
|
|
Body: encryptedResp,
|
|
})
|
|
|
|
if err != nil {
|
|
log.WithContext(ctx).Errorf("failed sending SyncResponse %v", err)
|
|
return status.Errorf(codes.Internal, "error handling request")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (s *Server) perAccountOrGlobalSyncMessageVersions(accountId string) grpc.SyncMessageVersion {
|
|
if version, ok := s.config.PerAccountHighestSupportedSyncMessageVersion[accountId]; ok {
|
|
return grpc.SyncMessageVersionFromConfig(&version)
|
|
}
|
|
return grpc.SyncMessageVersionFromConfig(s.config.HighestSupportedSyncMessageVersion)
|
|
}
|
|
|
|
// GetDeviceAuthorizationFlow returns a device authorization flow information
|
|
// This is used for initiating an Oauth 2 device authorization grant flow
|
|
// which will be used by our clients to Login
|
|
func (s *Server) GetDeviceAuthorizationFlow(ctx context.Context, req *proto.EncryptedMessage) (*proto.EncryptedMessage, error) {
|
|
log.WithContext(ctx).Tracef("GetDeviceAuthorizationFlow request for pubKey: %s", req.WgPubKey)
|
|
|
|
peerKey, err := wgtypes.ParseKey(req.GetWgPubKey())
|
|
if err != nil {
|
|
errMSG := fmt.Sprintf("error while parsing peer's Wireguard public key %s on GetDeviceAuthorizationFlow request.", req.WgPubKey)
|
|
log.WithContext(ctx).Warn(errMSG)
|
|
return nil, status.Error(codes.InvalidArgument, errMSG)
|
|
}
|
|
|
|
key, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
return nil, status.Errorf(codes.Internal, "failed to get server key")
|
|
}
|
|
|
|
err = encryption.DecryptMessage(peerKey, key, req.Body, &proto.DeviceAuthorizationFlowRequest{})
|
|
if err != nil {
|
|
errMSG := fmt.Sprintf("error while decrypting peer's message with Wireguard public key %s.", req.WgPubKey)
|
|
log.WithContext(ctx).Warn(errMSG)
|
|
return nil, status.Error(codes.InvalidArgument, errMSG)
|
|
}
|
|
|
|
var flowInfoResp *proto.DeviceAuthorizationFlow
|
|
|
|
// Use embedded IdP configuration if available
|
|
if s.oAuthConfigProvider != nil {
|
|
flowInfoResp = &proto.DeviceAuthorizationFlow{
|
|
Provider: proto.DeviceAuthorizationFlow_HOSTED,
|
|
ProviderConfig: &proto.ProviderConfig{
|
|
ClientID: s.oAuthConfigProvider.GetCLIClientID(),
|
|
Audience: s.oAuthConfigProvider.GetCLIClientID(),
|
|
DeviceAuthEndpoint: s.oAuthConfigProvider.GetDeviceAuthEndpoint(),
|
|
TokenEndpoint: s.oAuthConfigProvider.GetTokenEndpoint(),
|
|
Scope: s.oAuthConfigProvider.GetDefaultScopes(),
|
|
},
|
|
}
|
|
} else {
|
|
if s.config.DeviceAuthorizationFlow == nil || s.config.DeviceAuthorizationFlow.Provider == string(nbconfig.NONE) {
|
|
return nil, status.Error(codes.NotFound, "no device authorization flow information available")
|
|
}
|
|
|
|
provider, ok := proto.DeviceAuthorizationFlowProvider_value[strings.ToUpper(s.config.DeviceAuthorizationFlow.Provider)]
|
|
if !ok {
|
|
return nil, status.Errorf(codes.InvalidArgument, "no provider found in the protocol for %s", s.config.DeviceAuthorizationFlow.Provider)
|
|
}
|
|
|
|
flowInfoResp = &proto.DeviceAuthorizationFlow{
|
|
Provider: proto.DeviceAuthorizationFlowProvider(provider),
|
|
ProviderConfig: &proto.ProviderConfig{
|
|
ClientID: s.config.DeviceAuthorizationFlow.ProviderConfig.ClientID,
|
|
ClientSecret: s.config.DeviceAuthorizationFlow.ProviderConfig.ClientSecret, //nolint:staticcheck
|
|
Domain: s.config.DeviceAuthorizationFlow.ProviderConfig.Domain,
|
|
Audience: s.config.DeviceAuthorizationFlow.ProviderConfig.Audience,
|
|
DeviceAuthEndpoint: s.config.DeviceAuthorizationFlow.ProviderConfig.DeviceAuthEndpoint,
|
|
TokenEndpoint: s.config.DeviceAuthorizationFlow.ProviderConfig.TokenEndpoint,
|
|
Scope: s.config.DeviceAuthorizationFlow.ProviderConfig.Scope,
|
|
UseIDToken: s.config.DeviceAuthorizationFlow.ProviderConfig.UseIDToken,
|
|
},
|
|
}
|
|
}
|
|
|
|
encryptedResp, err := encryption.EncryptMessage(peerKey, key, flowInfoResp)
|
|
if err != nil {
|
|
return nil, status.Error(codes.Internal, "failed to encrypt device authorization flow information")
|
|
}
|
|
|
|
return &proto.EncryptedMessage{
|
|
WgPubKey: key.PublicKey().String(),
|
|
Body: encryptedResp,
|
|
}, nil
|
|
}
|
|
|
|
// GetPKCEAuthorizationFlow returns a pkce authorization flow information
|
|
// This is used for initiating an Oauth 2 pkce authorization grant flow
|
|
// which will be used by our clients to Login
|
|
func (s *Server) GetPKCEAuthorizationFlow(ctx context.Context, req *proto.EncryptedMessage) (*proto.EncryptedMessage, error) {
|
|
log.WithContext(ctx).Tracef("GetPKCEAuthorizationFlow request for pubKey: %s", req.WgPubKey)
|
|
|
|
peerKey, err := wgtypes.ParseKey(req.GetWgPubKey())
|
|
if err != nil {
|
|
errMSG := fmt.Sprintf("error while parsing peer's Wireguard public key %s on GetPKCEAuthorizationFlow request.", req.WgPubKey)
|
|
log.WithContext(ctx).Warn(errMSG)
|
|
return nil, status.Error(codes.InvalidArgument, errMSG)
|
|
}
|
|
|
|
key, err := s.secretsManager.GetWGKey()
|
|
if err != nil {
|
|
return nil, status.Errorf(codes.Internal, "failed to get server key")
|
|
}
|
|
|
|
flowReq := &proto.PKCEAuthorizationFlowRequest{}
|
|
err = encryption.DecryptMessage(peerKey, key, req.Body, flowReq)
|
|
if err != nil {
|
|
errMSG := fmt.Sprintf("error while decrypting peer's message with Wireguard public key %s.", req.WgPubKey)
|
|
log.WithContext(ctx).Warn(errMSG)
|
|
return nil, status.Error(codes.InvalidArgument, errMSG)
|
|
}
|
|
|
|
var initInfoFlow *proto.PKCEAuthorizationFlow
|
|
|
|
// Use embedded IdP configuration if available
|
|
if s.oAuthConfigProvider != nil {
|
|
initInfoFlow = &proto.PKCEAuthorizationFlow{
|
|
ProviderConfig: &proto.ProviderConfig{
|
|
Audience: s.oAuthConfigProvider.GetCLIClientID(),
|
|
ClientID: s.oAuthConfigProvider.GetCLIClientID(),
|
|
TokenEndpoint: s.oAuthConfigProvider.GetTokenEndpoint(),
|
|
AuthorizationEndpoint: s.oAuthConfigProvider.GetAuthorizationEndpoint(),
|
|
Scope: s.oAuthConfigProvider.GetDefaultScopes(),
|
|
RedirectURLs: s.oAuthConfigProvider.GetCLIRedirectURLs(),
|
|
LoginFlag: uint32(common.LoginFlagPromptLogin),
|
|
},
|
|
}
|
|
} else {
|
|
if s.config.PKCEAuthorizationFlow == nil {
|
|
return nil, status.Error(codes.NotFound, "no pkce authorization flow information available")
|
|
}
|
|
|
|
initInfoFlow = &proto.PKCEAuthorizationFlow{
|
|
ProviderConfig: &proto.ProviderConfig{
|
|
Audience: s.config.PKCEAuthorizationFlow.ProviderConfig.Audience,
|
|
ClientID: s.config.PKCEAuthorizationFlow.ProviderConfig.ClientID,
|
|
ClientSecret: s.config.PKCEAuthorizationFlow.ProviderConfig.ClientSecret, //nolint:staticcheck
|
|
TokenEndpoint: s.config.PKCEAuthorizationFlow.ProviderConfig.TokenEndpoint,
|
|
AuthorizationEndpoint: s.config.PKCEAuthorizationFlow.ProviderConfig.AuthorizationEndpoint,
|
|
Scope: s.config.PKCEAuthorizationFlow.ProviderConfig.Scope,
|
|
RedirectURLs: s.config.PKCEAuthorizationFlow.ProviderConfig.RedirectURLs,
|
|
UseIDToken: s.config.PKCEAuthorizationFlow.ProviderConfig.UseIDToken,
|
|
DisablePromptLogin: s.config.PKCEAuthorizationFlow.ProviderConfig.DisablePromptLogin,
|
|
LoginFlag: uint32(s.config.PKCEAuthorizationFlow.ProviderConfig.LoginFlag),
|
|
},
|
|
}
|
|
}
|
|
|
|
flowInfoResp := s.integratedPeerValidator.ValidateFlowResponse(ctx, peerKey.String(), initInfoFlow)
|
|
applySessionExtendFlowPolicy(flowInfoResp, flowReq.GetSessionExtend())
|
|
|
|
encryptedResp, err := encryption.EncryptMessage(peerKey, key, flowInfoResp)
|
|
if err != nil {
|
|
return nil, status.Error(codes.Internal, "failed to encrypt pkce authorization flow information")
|
|
}
|
|
|
|
return &proto.EncryptedMessage{
|
|
WgPubKey: key.PublicKey().String(),
|
|
Body: encryptedResp,
|
|
}, nil
|
|
}
|
|
|
|
// applySessionExtendFlowPolicy forces a prompt=login flow for a session extend.
|
|
//
|
|
// An extend renews the session of one specific peer, so its token has to come
|
|
// from the account that peer is registered under. A flow that does not prompt
|
|
// leaves the choice to the IdP, which answers a silent authorization from any
|
|
// session it already holds — not necessarily this peer's account when several
|
|
// are signed in, and login_hint is a suggestion the IdP may ignore. The token
|
|
// then fails the jwt.UserID == peer.UserID check in ExtendAuthSession, and the
|
|
// user is given no opportunity to pick a different account.
|
|
//
|
|
// LoginFlagPromptLogin rather than max_age=0: both re-authenticate, but with
|
|
// prompt=login the IdP honours login_hint and offers the peer's own account,
|
|
// whereas max_age=0 leaves the user to find it among every account signed in.
|
|
//
|
|
// DisablePromptLogin is left alone. It is set for IdPs that break on
|
|
// prompt=login — Authentik triggers a double authentication, and social logins
|
|
// fail outright — so overriding it would trade a recoverable session extend for
|
|
// a login that cannot complete at all. Those deployments keep the silent flow
|
|
// and, with several accounts signed in, an extend answered from the wrong one
|
|
// still fails the user match.
|
|
//
|
|
// Called after ValidateFlowResponse so that a per-peer override cannot reinstate
|
|
// the silent flow for an extend.
|
|
func applySessionExtendFlowPolicy(flow *proto.PKCEAuthorizationFlow, sessionExtend bool) {
|
|
if !sessionExtend {
|
|
return
|
|
}
|
|
cfg := flow.GetProviderConfig()
|
|
if cfg == nil || cfg.GetDisablePromptLogin() {
|
|
return
|
|
}
|
|
cfg.LoginFlag = uint32(common.LoginFlagPromptLogin)
|
|
}
|
|
|
|
// SyncMeta endpoint is used to synchronize peer's system metadata and notifies the connected,
|
|
// peer's under the same account of any updates.
|
|
func (s *Server) SyncMeta(ctx context.Context, req *proto.EncryptedMessage) (*proto.Empty, error) {
|
|
realIP := getRealIP(ctx)
|
|
log.WithContext(ctx).Debugf("Sync meta request from peer [%s] [%s]", req.WgPubKey, realIP.String())
|
|
|
|
syncMetaReq := &proto.SyncMetaRequest{}
|
|
peerKey, err := s.parseRequest(ctx, req, syncMetaReq)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if syncMetaReq.GetMeta() == nil {
|
|
msg := status.Errorf(codes.FailedPrecondition,
|
|
"peer system meta has to be provided on sync. Peer %s, remote addr %s", peerKey.String(), realIP)
|
|
log.WithContext(ctx).Warn(msg)
|
|
return nil, msg
|
|
}
|
|
|
|
peerMeta := extractPeerMeta(ctx, syncMetaReq.GetMeta())
|
|
peerMeta.Certificates = s.verifiedCertificates(ctx, peerKey, syncMetaReq.GetMeta().GetCertificateProofs())
|
|
err = s.accountManager.SyncPeerMeta(ctx, peerKey.String(), peerMeta, realIP)
|
|
if err != nil {
|
|
return nil, mapError(ctx, err)
|
|
}
|
|
|
|
return &proto.Empty{}, nil
|
|
}
|
|
|
|
func (s *Server) Logout(ctx context.Context, req *proto.EncryptedMessage) (*proto.Empty, error) {
|
|
log.WithContext(ctx).Debugf("Logout request from peer [%s]", req.WgPubKey)
|
|
start := time.Now()
|
|
|
|
empty := &proto.Empty{}
|
|
peerKey, err := s.parseRequest(ctx, req, empty)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
peer, err := s.accountManager.GetStore().GetPeerByPeerPubKey(ctx, store.LockingStrengthNone, peerKey.String())
|
|
if err != nil {
|
|
log.WithContext(ctx).Debugf("peer %s is not registered for logout", peerKey.String())
|
|
// TODO: consider idempotency
|
|
return nil, mapError(ctx, err)
|
|
}
|
|
|
|
// nolint:staticcheck
|
|
ctx = context.WithValue(ctx, nbContext.PeerIDKey, peer.ID)
|
|
// nolint:staticcheck
|
|
ctx = context.WithValue(ctx, nbContext.AccountIDKey, peer.AccountID)
|
|
|
|
userID := peer.UserID
|
|
if userID == "" {
|
|
userID = activity.SystemInitiator
|
|
}
|
|
|
|
if err = s.accountManager.DeletePeer(ctx, peer.AccountID, peer.ID, userID); err != nil {
|
|
log.WithContext(ctx).Errorf("failed to logout peer %s: %v", peerKey.String(), err)
|
|
return nil, mapError(ctx, err)
|
|
}
|
|
|
|
log.WithContext(ctx).Debugf("peer %s logged out successfully after %s", peerKey.String(), time.Since(start))
|
|
|
|
return &proto.Empty{}, nil
|
|
}
|
|
|
|
// toProtocolChecks converts posture checks to protocol checks.
|
|
func toProtocolChecks(ctx context.Context, postureChecks []*nmdata.PostureChecks) []*proto.Checks {
|
|
protoChecks := make([]*proto.Checks, 0, len(postureChecks))
|
|
for _, postureCheck := range postureChecks {
|
|
check := toProtocolCheck(postureCheck)
|
|
if check != nil {
|
|
protoChecks = append(protoChecks, check)
|
|
}
|
|
}
|
|
|
|
return protoChecks
|
|
}
|
|
|
|
// toProtocolCheck converts posture checks to a proto.Checks.
|
|
func toProtocolCheck(postureCheck *nmdata.PostureChecks) *proto.Checks {
|
|
protoCheck := &proto.Checks{}
|
|
|
|
if check := postureCheck.Checks.ProcessCheck; check != nil {
|
|
for _, process := range check.Processes {
|
|
if process.LinuxPath != "" {
|
|
protoCheck.Files = append(protoCheck.Files, process.LinuxPath)
|
|
}
|
|
if process.MacPath != "" {
|
|
protoCheck.Files = append(protoCheck.Files, process.MacPath)
|
|
}
|
|
if process.WindowsPath != "" {
|
|
protoCheck.Files = append(protoCheck.Files, process.WindowsPath)
|
|
}
|
|
}
|
|
}
|
|
|
|
if check := postureCheck.Checks.CertificateCheck; check != nil {
|
|
protoCheck.CertificateChallenge = &proto.CertificateChallenge{CaCertificates: check.CACertificates}
|
|
}
|
|
|
|
if len(protoCheck.Files) == 0 && protoCheck.CertificateChallenge == nil {
|
|
return nil
|
|
}
|
|
|
|
return protoCheck
|
|
}
|