mirror of
https://github.com/tailscale/tailscale.git
synced 2026-10-10 12:21:53 -04:00
Previously ProxyClass spec.staticEndpoints was only implemented for ProxyGroups and silently ignored when the ProxyClass was referenced by a Connector. Extract the NodePort Service provisioning, port allocation, and node ExternalIP discovery logic shared with ProxyGroup into a new k8s-operator/reconciler/staticendpoints package, and wire it into the StatefulSet reconciler so that each Connector replica gets a per-pod NodePort Service, the discovered endpoints are written to its tailscaled config as static endpoints, and tailscaled listens on the Services' target port via the PORT env var. Also reconcile Connectors on Node changes, clean up NodePort Services on scale down, Connector deletion, and when static endpoints are removed, and surface the endpoints in the Connector's device status. Add e2e tests for static endpoints on both Connectors and ProxyGroups. They verify the per-replica NodePort Services, the PORT env var, the endpoints reported in status and advertised to control, and the cleanup on scale down and on removal of the static endpoints configuration. The tests skip on clusters whose Nodes have no ExternalIP addresses, such as kind. Fixes #18819 Change-Id: Ia134916062471e6c0088692834257aafabd0e713 Signed-off-by: David Bond <davidsbond93@gmail.com>
1103 lines
42 KiB
Go
1103 lines
42 KiB
Go
// Copyright (c) Tailscale Inc & contributors
|
|
// SPDX-License-Identifier: BSD-3-Clause
|
|
|
|
//go:build !plan9
|
|
|
|
package main
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"net/netip"
|
|
"slices"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
|
|
dockerref "github.com/distribution/reference"
|
|
"go.uber.org/zap"
|
|
xslices "golang.org/x/exp/slices"
|
|
appsv1 "k8s.io/api/apps/v1"
|
|
corev1 "k8s.io/api/core/v1"
|
|
rbacv1 "k8s.io/api/rbac/v1"
|
|
apiequality "k8s.io/apimachinery/pkg/api/equality"
|
|
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
"k8s.io/apimachinery/pkg/types"
|
|
"k8s.io/client-go/tools/record"
|
|
"sigs.k8s.io/controller-runtime/pkg/client"
|
|
"sigs.k8s.io/controller-runtime/pkg/reconcile"
|
|
|
|
"tailscale.com/ipn"
|
|
tsoperator "tailscale.com/k8s-operator"
|
|
tsapi "tailscale.com/k8s-operator/apis/v1alpha1"
|
|
"tailscale.com/k8s-operator/reconciler/staticendpoints"
|
|
"tailscale.com/k8s-operator/reconciler/tailscaled"
|
|
"tailscale.com/k8s-operator/tsclient"
|
|
"tailscale.com/kube/egressservices"
|
|
"tailscale.com/kube/k8s-proxy/conf"
|
|
"tailscale.com/kube/kubetypes"
|
|
"tailscale.com/tailcfg"
|
|
"tailscale.com/tstime"
|
|
"tailscale.com/types/opt"
|
|
"tailscale.com/util/clientmetric"
|
|
"tailscale.com/util/mak"
|
|
"tailscale.com/util/set"
|
|
)
|
|
|
|
const (
|
|
reasonProxyGroupCreationFailed = "ProxyGroupCreationFailed"
|
|
reasonProxyGroupReady = "ProxyGroupReady"
|
|
reasonProxyGroupAvailable = "ProxyGroupAvailable"
|
|
reasonProxyGroupCreating = "ProxyGroupCreating"
|
|
reasonProxyGroupInvalid = "ProxyGroupInvalid"
|
|
reasonProxyGroupTailnetUnavailable = "ProxyGroupTailnetUnavailable"
|
|
reasonACMEAccountsPendingDeletion = "ACMEAccountsPendingDeletion"
|
|
|
|
// Copied from k8s.io/apiserver/pkg/registry/generic/registry/store.go@cccad306d649184bf2a0e319ba830c53f65c445c
|
|
optimisticLockErrorMsg = "the object has been modified; please apply your changes to the latest version and try again"
|
|
|
|
// The minimum tailcfg.CapabilityVersion that deployed clients are expected
|
|
// to support to be compatible with the current ProxyGroup controller.
|
|
// If the controller needs to depend on newer client behaviour, it should
|
|
// maintain backwards compatible logic for older capability versions for 3
|
|
// stable releases, as per documentation on supported version drift:
|
|
// https://tailscale.com/kb/1236/kubernetes-operator#supported-versions
|
|
//
|
|
// tailcfg.CurrentCapabilityVersion was 106 when the ProxyGroup controller was
|
|
// first introduced.
|
|
pgMinCapabilityVersion = 106
|
|
)
|
|
|
|
var (
|
|
gaugeEgressProxyGroupResources = clientmetric.NewGauge(kubetypes.MetricProxyGroupEgressCount)
|
|
gaugeIngressProxyGroupResources = clientmetric.NewGauge(kubetypes.MetricProxyGroupIngressCount)
|
|
gaugeAPIServerProxyGroupResources = clientmetric.NewGauge(kubetypes.MetricProxyGroupAPIServerCount)
|
|
)
|
|
|
|
// ProxyGroupReconciler ensures cluster resources for a ProxyGroup definition.
|
|
type ProxyGroupReconciler struct {
|
|
client.Client
|
|
log *zap.SugaredLogger
|
|
recorder record.EventRecorder
|
|
clock tstime.Clock
|
|
clients ClientProvider
|
|
|
|
// User-specified defaults from the helm installation.
|
|
tsNamespace string
|
|
tsProxyImage string
|
|
k8sProxyImage string
|
|
defaultTags []string
|
|
tsFirewallMode string
|
|
defaultProxyClass string
|
|
loginServer string
|
|
|
|
reissuer *tailscaled.Reissuer
|
|
|
|
mu sync.Mutex // protects following
|
|
egressProxyGroups set.Slice[types.UID] // for egress proxygroups gauge
|
|
ingressProxyGroups set.Slice[types.UID] // for ingress proxygroups gauge
|
|
apiServerProxyGroups set.Slice[types.UID] // for kube-apiserver proxygroups gauge
|
|
|
|
// sharedACMEAccountKey is the operator-wide default for the
|
|
// shared-ACME-account feature. When true, every ProxyGroup uses the
|
|
// shared per-tailnet account key unless the ProxyGroup explicitly
|
|
// opts out via tailscale.com/share-acme-account=false. When false,
|
|
// only ProxyGroups annotated with tailscale.com/share-acme-account=true
|
|
// use it.
|
|
sharedACMEAccountKey bool
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) logger(name string) *zap.SugaredLogger {
|
|
return r.log.With("ProxyGroup", name)
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) Reconcile(ctx context.Context, req reconcile.Request) (_ reconcile.Result, err error) {
|
|
logger := r.logger(req.Name)
|
|
logger.Debugf("starting reconcile")
|
|
defer logger.Debugf("reconcile finished")
|
|
|
|
pg := new(tsapi.ProxyGroup)
|
|
err = r.Get(ctx, req.NamespacedName, pg)
|
|
if apierrors.IsNotFound(err) {
|
|
logger.Debugf("ProxyGroup not found, assuming it was deleted")
|
|
return reconcile.Result{}, nil
|
|
} else if err != nil {
|
|
return reconcile.Result{}, fmt.Errorf("failed to get tailscale.com ProxyGroup: %w", err)
|
|
}
|
|
|
|
tsClient, err := r.clients.For(pg.Spec.Tailnet)
|
|
if err != nil {
|
|
oldPGStatus := pg.Status.DeepCopy()
|
|
nrr := ¬ReadyReason{
|
|
reason: reasonProxyGroupTailnetUnavailable,
|
|
message: fmt.Errorf("failed to get tailscale client and loginUrl: %w", err).Error(),
|
|
}
|
|
|
|
return reconcile.Result{}, errors.Join(err, r.maybeUpdateStatus(ctx, logger, pg, oldPGStatus, nrr, make(map[string][]netip.AddrPort)))
|
|
}
|
|
|
|
if markedForDeletion(pg) {
|
|
logger.Debugf("ProxyGroup is being deleted, cleaning up resources")
|
|
ix := xslices.Index(pg.Finalizers, FinalizerName)
|
|
if ix < 0 {
|
|
logger.Debugf("no finalizer, nothing to do")
|
|
return reconcile.Result{}, nil
|
|
}
|
|
|
|
if done, err := r.maybeCleanup(ctx, tsClient, pg); err != nil {
|
|
if strings.Contains(err.Error(), optimisticLockErrorMsg) {
|
|
logger.Infof("optimistic lock error, retrying: %s", err)
|
|
return reconcile.Result{}, nil
|
|
}
|
|
return reconcile.Result{}, err
|
|
} else if !done {
|
|
logger.Debugf("ProxyGroup resource cleanup not yet finished, will retry...")
|
|
return reconcile.Result{RequeueAfter: shortRequeue}, nil
|
|
}
|
|
|
|
pg.Finalizers = slices.Delete(pg.Finalizers, ix, ix+1)
|
|
if err := r.Update(ctx, pg); err != nil {
|
|
return reconcile.Result{}, err
|
|
}
|
|
return reconcile.Result{}, nil
|
|
}
|
|
|
|
oldPGStatus := pg.Status.DeepCopy()
|
|
staticEndpoints, nrr, err := r.reconcilePG(ctx, tsClient, pg, logger)
|
|
return reconcile.Result{}, errors.Join(err, r.maybeUpdateStatus(ctx, logger, pg, oldPGStatus, nrr, staticEndpoints))
|
|
}
|
|
|
|
// reconcilePG handles all reconciliation of a ProxyGroup that is not marked
|
|
// for deletion. It is separated out from Reconcile to make a clear separation
|
|
// between reconciling the ProxyGroup, and posting the status of its created
|
|
// resources onto the ProxyGroup status field.
|
|
func (r *ProxyGroupReconciler) reconcilePG(ctx context.Context, tsClient tsclient.Client, pg *tsapi.ProxyGroup, logger *zap.SugaredLogger) (map[string][]netip.AddrPort, *notReadyReason, error) {
|
|
if !slices.Contains(pg.Finalizers, FinalizerName) {
|
|
// This log line is printed exactly once during initial provisioning,
|
|
// because once the finalizer is in place this block gets skipped. So,
|
|
// this is a nice place to log that the high level, multi-reconcile
|
|
// operation is underway.
|
|
logger.Infof("ensuring ProxyGroup is set up")
|
|
pg.Finalizers = append(pg.Finalizers, FinalizerName)
|
|
if err := r.Update(ctx, pg); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error adding finalizer: %w", err)
|
|
}
|
|
}
|
|
|
|
proxyClassName := r.defaultProxyClass
|
|
if pg.Spec.ProxyClass != "" {
|
|
proxyClassName = pg.Spec.ProxyClass
|
|
}
|
|
|
|
var proxyClass *tsapi.ProxyClass
|
|
if proxyClassName != "" {
|
|
proxyClass = new(tsapi.ProxyClass)
|
|
err := r.Get(ctx, types.NamespacedName{Name: proxyClassName}, proxyClass)
|
|
if apierrors.IsNotFound(err) {
|
|
msg := fmt.Sprintf("the ProxyGroup's ProxyClass %q does not (yet) exist", proxyClassName)
|
|
logger.Info(msg)
|
|
return notReady(reasonProxyGroupCreating, msg)
|
|
}
|
|
if err != nil {
|
|
return r.notReadyErrf(pg, logger, "error getting ProxyGroup's ProxyClass %q: %w", proxyClassName, err)
|
|
}
|
|
if !tsoperator.ProxyClassIsReady(proxyClass) {
|
|
msg := fmt.Sprintf("the ProxyGroup's ProxyClass %q is not yet in a ready state, waiting...", proxyClassName)
|
|
logger.Info(msg)
|
|
return notReady(reasonProxyGroupCreating, msg)
|
|
}
|
|
}
|
|
|
|
if err := r.validate(ctx, pg, proxyClass, logger); err != nil {
|
|
return notReady(reasonProxyGroupInvalid, fmt.Sprintf("invalid ProxyGroup spec: %v", err))
|
|
}
|
|
|
|
staticEndpoints, nrr, err := r.maybeProvision(ctx, tsClient, pg, proxyClass)
|
|
if err != nil {
|
|
return nil, nrr, err
|
|
}
|
|
|
|
return staticEndpoints, nrr, nil
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) validate(ctx context.Context, pg *tsapi.ProxyGroup, pc *tsapi.ProxyClass, logger *zap.SugaredLogger) error {
|
|
// Our custom logic for ensuring minimum downtime ProxyGroup update rollouts relies on the local health check
|
|
// beig accessible on the replica Pod IP:9002. This address can also be modified by users, via
|
|
// TS_LOCAL_ADDR_PORT env var.
|
|
//
|
|
// Currently TS_LOCAL_ADDR_PORT controls Pod's health check and metrics address. _Probably_ there is no need for
|
|
// users to set this to a custom value. Users who want to consume metrics, should integrate with the metrics
|
|
// Service and/or ServiceMonitor, rather than Pods directly. The health check is likely not useful to integrate
|
|
// directly with for operator proxies (and we should aim for unified lifecycle logic in the operator, users
|
|
// shouldn't need to set their own).
|
|
//
|
|
// TODO(irbekrm): maybe disallow configuring this env var in future (in Tailscale 1.84 or later).
|
|
if pg.Spec.Type == tsapi.ProxyGroupTypeEgress && hasLocalAddrPortSet(pc) {
|
|
msg := fmt.Sprintf("ProxyClass %s applied to an egress ProxyGroup has TS_LOCAL_ADDR_PORT env var set to a custom value."+
|
|
"This will disable the ProxyGroup graceful failover mechanism, so you might experience downtime when ProxyGroup pods are restarted."+
|
|
"In future we will remove the ability to set custom TS_LOCAL_ADDR_PORT for egress ProxyGroups."+
|
|
"Please raise an issue if you expect that this will cause issues for your workflow.", pc.Name)
|
|
logger.Warn(msg)
|
|
}
|
|
|
|
// image is the value of pc.Spec.StatefulSet.Pod.TailscaleContainer.Image or ""
|
|
// imagePath is a slash-delimited path ending with the image name, e.g.
|
|
// "tailscale/tailscale" or maybe "k8s-proxy" if hosted at example.com/k8s-proxy.
|
|
var image, imagePath string
|
|
if pc != nil &&
|
|
pc.Spec.StatefulSet != nil &&
|
|
pc.Spec.StatefulSet.Pod != nil &&
|
|
pc.Spec.StatefulSet.Pod.TailscaleContainer != nil &&
|
|
pc.Spec.StatefulSet.Pod.TailscaleContainer.Image != "" {
|
|
image, err := dockerref.ParseNormalizedNamed(pc.Spec.StatefulSet.Pod.TailscaleContainer.Image)
|
|
if err != nil {
|
|
// Shouldn't be possible as the ProxyClass won't be marked ready
|
|
// without successfully parsing the image.
|
|
return fmt.Errorf("error parsing %q as a container image reference: %w", pc.Spec.StatefulSet.Pod.TailscaleContainer.Image, err)
|
|
}
|
|
imagePath = dockerref.Path(image)
|
|
}
|
|
|
|
var errs []error
|
|
if isAuthAPIServerProxy(pg) {
|
|
// Validate that the static ServiceAccount already exists.
|
|
sa := &corev1.ServiceAccount{}
|
|
if err := r.Get(ctx, types.NamespacedName{Namespace: r.tsNamespace, Name: authAPIServerProxySAName}, sa); err != nil {
|
|
if !apierrors.IsNotFound(err) {
|
|
return fmt.Errorf("error validating that ServiceAccount %q exists: %w", authAPIServerProxySAName, err)
|
|
}
|
|
|
|
errs = append(errs, fmt.Errorf("the ServiceAccount %q used for the API server proxy in auth mode does not exist but "+
|
|
"should have been created during operator installation; use apiServerProxyConfig.allowImpersonation=true "+
|
|
"in the helm chart, or authproxy-rbac.yaml from the static manifests", authAPIServerProxySAName))
|
|
}
|
|
} else {
|
|
// Validate that the ServiceAccount we create won't overwrite the static one.
|
|
// TODO(tomhjp): This doesn't cover other controllers that could create a
|
|
// ServiceAccount. Perhaps should have some guards to ensure that an update
|
|
// would never change the ownership of a resource we expect to already be owned.
|
|
if pgServiceAccountName(pg) == authAPIServerProxySAName {
|
|
errs = append(errs, fmt.Errorf("the name of the ProxyGroup %q conflicts with the static ServiceAccount used for the API server proxy in auth mode", pg.Name))
|
|
}
|
|
}
|
|
|
|
if pg.Spec.Type == tsapi.ProxyGroupTypeKubernetesAPIServer {
|
|
if strings.HasSuffix(imagePath, "tailscale") {
|
|
errs = append(errs, fmt.Errorf("the configured ProxyClass %q specifies to use image %q but expected a %q image for ProxyGroup of type %q", pc.Name, image, "k8s-proxy", pg.Spec.Type))
|
|
}
|
|
|
|
if pc != nil && pc.Spec.StatefulSet != nil && pc.Spec.StatefulSet.Pod != nil && pc.Spec.StatefulSet.Pod.TailscaleInitContainer != nil {
|
|
errs = append(errs, fmt.Errorf("the configured ProxyClass %q specifies Tailscale init container config, but ProxyGroups of type %q do not use init containers", pc.Name, pg.Spec.Type))
|
|
}
|
|
} else {
|
|
if strings.HasSuffix(imagePath, "k8s-proxy") {
|
|
errs = append(errs, fmt.Errorf("the configured ProxyClass %q specifies to use image %q but expected a %q image for ProxyGroup of type %q", pc.Name, image, "tailscale", pg.Spec.Type))
|
|
}
|
|
}
|
|
|
|
return errors.Join(errs...)
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) maybeProvision(ctx context.Context, tsClient tsclient.Client, pg *tsapi.ProxyGroup, proxyClass *tsapi.ProxyClass) (map[string][]netip.AddrPort, *notReadyReason, error) {
|
|
logger := r.logger(pg.Name)
|
|
r.mu.Lock()
|
|
r.ensureStateAddedForProxyGroup(pg)
|
|
r.mu.Unlock()
|
|
|
|
svcToNodePorts := make(map[string]uint16)
|
|
var tailscaledPort *uint16
|
|
if proxyClass != nil && proxyClass.Spec.StaticEndpoints != nil {
|
|
var err error
|
|
svcToNodePorts, tailscaledPort, err = r.ensureNodePortServiceCreated(ctx, pg, proxyClass)
|
|
if err != nil {
|
|
if _, ok := errors.AsType[*staticendpoints.AllocatePortsError](err); ok {
|
|
reason := reasonProxyGroupCreationFailed
|
|
msg := fmt.Sprintf("error provisioning NodePort Services for static endpoints: %v", err)
|
|
r.recorder.Event(pg, corev1.EventTypeWarning, reason, msg)
|
|
return notReady(reason, msg)
|
|
}
|
|
return r.notReadyErrf(pg, logger, "error provisioning NodePort Services for static endpoints: %w", err)
|
|
}
|
|
}
|
|
|
|
staticEndpoints, err := r.ensureConfigSecretsCreated(ctx, tsClient, pg, proxyClass, svcToNodePorts)
|
|
if err != nil {
|
|
if _, ok := errors.AsType[*staticendpoints.FindEndpointsError](err); ok {
|
|
reason := reasonProxyGroupCreationFailed
|
|
msg := fmt.Sprintf("error provisioning config Secrets: %v", err)
|
|
r.recorder.Event(pg, corev1.EventTypeWarning, reason, msg)
|
|
return notReady(reason, msg)
|
|
}
|
|
return r.notReadyErrf(pg, logger, "error provisioning config Secrets: %w", err)
|
|
}
|
|
|
|
// State secrets are precreated so we can use the ProxyGroup CR as their owner ref.
|
|
stateSecrets := pgStateSecrets(pg, r.tsNamespace)
|
|
for _, sec := range stateSecrets {
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, sec, func(s *corev1.Secret) {
|
|
s.ObjectMeta.Labels = sec.ObjectMeta.Labels
|
|
s.ObjectMeta.Annotations = sec.ObjectMeta.Annotations
|
|
s.ObjectMeta.OwnerReferences = sec.ObjectMeta.OwnerReferences
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning state Secrets: %w", err)
|
|
}
|
|
}
|
|
|
|
// auth mode kube-apiserver ProxyGroups use a statically created
|
|
// ServiceAccount to keep ClusterRole creation permissions limited to the
|
|
// helm chart installer.
|
|
if !isAuthAPIServerProxy(pg) {
|
|
sa := pgServiceAccount(pg, r.tsNamespace)
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, sa, func(s *corev1.ServiceAccount) {
|
|
s.ObjectMeta.Labels = sa.ObjectMeta.Labels
|
|
s.ObjectMeta.Annotations = sa.ObjectMeta.Annotations
|
|
s.ObjectMeta.OwnerReferences = sa.ObjectMeta.OwnerReferences
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning ServiceAccount: %w", err)
|
|
}
|
|
}
|
|
|
|
role := pgRole(pg, r.tsNamespace, r.sharedACMEAccountEnabledFor(pg))
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, role, func(r *rbacv1.Role) {
|
|
r.ObjectMeta.Labels = role.ObjectMeta.Labels
|
|
r.ObjectMeta.Annotations = role.ObjectMeta.Annotations
|
|
r.ObjectMeta.OwnerReferences = role.ObjectMeta.OwnerReferences
|
|
r.Rules = role.Rules
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning Role: %w", err)
|
|
}
|
|
|
|
roleBinding := pgRoleBinding(pg, r.tsNamespace)
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, roleBinding, func(r *rbacv1.RoleBinding) {
|
|
r.ObjectMeta.Labels = roleBinding.ObjectMeta.Labels
|
|
r.ObjectMeta.Annotations = roleBinding.ObjectMeta.Annotations
|
|
r.ObjectMeta.OwnerReferences = roleBinding.ObjectMeta.OwnerReferences
|
|
r.RoleRef = roleBinding.RoleRef
|
|
r.Subjects = roleBinding.Subjects
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning RoleBinding: %w", err)
|
|
}
|
|
|
|
if pg.Spec.Type == tsapi.ProxyGroupTypeEgress {
|
|
cm, hp := pgEgressCM(pg, r.tsNamespace)
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, cm, func(existing *corev1.ConfigMap) {
|
|
existing.ObjectMeta.Labels = cm.ObjectMeta.Labels
|
|
existing.ObjectMeta.OwnerReferences = cm.ObjectMeta.OwnerReferences
|
|
mak.Set(&existing.BinaryData, egressservices.KeyHEPPings, hp)
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning egress ConfigMap %q: %w", cm.Name, err)
|
|
}
|
|
}
|
|
|
|
if pg.Spec.Type == tsapi.ProxyGroupTypeIngress {
|
|
cm := pgIngressCM(pg, r.tsNamespace)
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, cm, func(existing *corev1.ConfigMap) {
|
|
existing.ObjectMeta.Labels = cm.ObjectMeta.Labels
|
|
existing.ObjectMeta.OwnerReferences = cm.ObjectMeta.OwnerReferences
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning ingress ConfigMap %q: %w", cm.Name, err)
|
|
}
|
|
|
|
// Ensure the shared ACME accounts Secret exists (with finalizer)
|
|
// when this ProxyGroup opts into the feature. Proxy pods
|
|
// populate its fields on first cert issuance. See #18251.
|
|
if r.sharedACMEAccountEnabledFor(pg) {
|
|
acmeSecret := pgACMEAccountSecret(r.tsNamespace)
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, acmeSecret, func(existing *corev1.Secret) {
|
|
if !existing.DeletionTimestamp.IsZero() {
|
|
// Deletion can't be undone; warn so the account keys
|
|
// get backed up before the finalizer is removed.
|
|
msg := fmt.Sprintf("shared ACME accounts Secret %q is marked for deletion but retained by the %q finalizer. Its data remains readable until the finalizer is removed - back it up first to preserve the ACME account keys.", existing.Name, kubetypes.ACMEAccountsFinalizer)
|
|
r.recorder.Event(existing, corev1.EventTypeWarning, reasonACMEAccountsPendingDeletion, msg)
|
|
logger.Warn(msg)
|
|
return
|
|
}
|
|
existing.Labels = acmeSecret.Labels
|
|
if !slices.Contains(existing.Finalizers, kubetypes.ACMEAccountsFinalizer) {
|
|
existing.Finalizers = append(existing.Finalizers, kubetypes.ACMEAccountsFinalizer)
|
|
}
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning shared ACME accounts Secret %q: %w", acmeSecret.Name, err)
|
|
}
|
|
}
|
|
}
|
|
|
|
defaultImage := r.tsProxyImage
|
|
if pg.Spec.Type == tsapi.ProxyGroupTypeKubernetesAPIServer {
|
|
defaultImage = r.k8sProxyImage
|
|
}
|
|
ss, err := pgStatefulSet(pg, r.tsNamespace, defaultImage, r.tsFirewallMode, tailscaledPort, proxyClass, r.sharedACMEAccountEnabledFor(pg))
|
|
if err != nil {
|
|
return r.notReadyErrf(pg, logger, "error generating StatefulSet spec: %w", err)
|
|
}
|
|
cfg := &tailscaleSTSConfig{
|
|
proxyType: string(pg.Spec.Type),
|
|
}
|
|
ss = applyProxyClassToStatefulSet(proxyClass, ss, cfg, logger)
|
|
|
|
if _, err := createOrUpdate(ctx, r.Client, r.tsNamespace, ss, func(s *appsv1.StatefulSet) {
|
|
s.Spec = ss.Spec
|
|
s.ObjectMeta.Labels = ss.ObjectMeta.Labels
|
|
s.ObjectMeta.Annotations = ss.ObjectMeta.Annotations
|
|
s.ObjectMeta.OwnerReferences = ss.ObjectMeta.OwnerReferences
|
|
}); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error provisioning StatefulSet: %w", err)
|
|
}
|
|
|
|
mo := &metricsOpts{
|
|
tsNamespace: r.tsNamespace,
|
|
proxyStsName: pg.Name,
|
|
proxyLabels: pgLabels(pg.Name, nil),
|
|
proxyType: "proxygroup",
|
|
}
|
|
if err := reconcileMetricsResources(ctx, logger, mo, proxyClass, r.Client); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error reconciling metrics resources: %w", err)
|
|
}
|
|
|
|
if err := r.cleanupDanglingResources(ctx, tsClient, pg, proxyClass); err != nil {
|
|
return r.notReadyErrf(pg, logger, "error cleaning up dangling resources: %w", err)
|
|
}
|
|
|
|
logger.Info("ProxyGroup resources synced")
|
|
|
|
return staticEndpoints, nil, nil
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) maybeUpdateStatus(ctx context.Context, logger *zap.SugaredLogger, pg *tsapi.ProxyGroup, oldPGStatus *tsapi.ProxyGroupStatus, nrr *notReadyReason, endpoints map[string][]netip.AddrPort) (err error) {
|
|
defer func() {
|
|
if !apiequality.Semantic.DeepEqual(*oldPGStatus, pg.Status) {
|
|
if updateErr := r.Client.Status().Update(ctx, pg); updateErr != nil {
|
|
if strings.Contains(updateErr.Error(), optimisticLockErrorMsg) {
|
|
logger.Infof("optimistic lock error updating status, retrying: %s", updateErr)
|
|
updateErr = nil
|
|
}
|
|
err = errors.Join(err, updateErr)
|
|
}
|
|
}
|
|
}()
|
|
|
|
devices, err := r.getRunningProxies(ctx, pg, endpoints)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to list running proxies: %w", err)
|
|
}
|
|
|
|
pg.Status.Devices = devices
|
|
|
|
desiredReplicas := int(pgReplicas(pg))
|
|
|
|
// Set ProxyGroupAvailable condition.
|
|
status := metav1.ConditionFalse
|
|
reason := reasonProxyGroupCreating
|
|
message := fmt.Sprintf("%d/%d ProxyGroup pods running", len(devices), desiredReplicas)
|
|
if len(devices) > 0 {
|
|
status = metav1.ConditionTrue
|
|
if len(devices) == desiredReplicas {
|
|
reason = reasonProxyGroupAvailable
|
|
}
|
|
}
|
|
tsoperator.SetProxyGroupCondition(pg, tsapi.ProxyGroupAvailable, status, reason, message, 0, r.clock, logger)
|
|
|
|
// Set ProxyGroupReady condition.
|
|
tsSvcValid, tsSvcSet := tsoperator.KubeAPIServerProxyValid(pg)
|
|
status = metav1.ConditionFalse
|
|
reason = reasonProxyGroupCreating
|
|
switch {
|
|
case nrr != nil:
|
|
// If we failed earlier, that reason takes precedence.
|
|
reason = nrr.reason
|
|
message = nrr.message
|
|
case pg.Spec.Type == tsapi.ProxyGroupTypeKubernetesAPIServer && tsSvcSet && !tsSvcValid:
|
|
reason = reasonProxyGroupInvalid
|
|
message = "waiting for config in spec.kubeAPIServer to be marked valid"
|
|
case len(devices) < desiredReplicas:
|
|
case len(devices) > desiredReplicas:
|
|
message = fmt.Sprintf("waiting for %d ProxyGroup pods to shut down", len(devices)-desiredReplicas)
|
|
case pg.Spec.Type == tsapi.ProxyGroupTypeKubernetesAPIServer && !tsoperator.KubeAPIServerProxyConfigured(pg):
|
|
reason = reasonProxyGroupCreating
|
|
message = "waiting for proxies to start advertising the kube-apiserver proxy's hostname"
|
|
default:
|
|
status = metav1.ConditionTrue
|
|
reason = reasonProxyGroupReady
|
|
message = reasonProxyGroupReady
|
|
}
|
|
tsoperator.SetProxyGroupCondition(pg, tsapi.ProxyGroupReady, status, reason, message, pg.Generation, r.clock, logger)
|
|
|
|
return nil
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) ensureNodePortServiceCreated(ctx context.Context, pg *tsapi.ProxyGroup, pc *tsapi.ProxyClass) (map[string]uint16, *uint16, error) {
|
|
return staticendpoints.EnsureNodePortServices(ctx, r.Client, staticendpoints.Config{
|
|
Namespace: r.tsNamespace,
|
|
ParentType: proxyTypeProxyGroup,
|
|
ParentName: pg.Name,
|
|
ProxyClassName: pc.Name,
|
|
Replicas: pgReplicas(pg),
|
|
PortRanges: pc.Spec.StaticEndpoints.NodePort.Ports,
|
|
MakeService: func(_ int32, name string) *corev1.Service {
|
|
return pgNodePortService(pg, name, r.tsNamespace)
|
|
},
|
|
})
|
|
}
|
|
|
|
// cleanupDanglingResources ensures we don't leak config secrets, state secrets, and
|
|
// tailnet devices when the number of replicas specified is reduced.
|
|
func (r *ProxyGroupReconciler) cleanupDanglingResources(ctx context.Context, tsClient tsclient.Client, pg *tsapi.ProxyGroup, pc *tsapi.ProxyClass) error {
|
|
logger := r.logger(pg.Name)
|
|
metadata, err := getNodeMetadata(ctx, pg, r.Client, r.tsNamespace)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
for _, m := range metadata {
|
|
if m.ordinal+1 <= pgReplicas(pg) {
|
|
continue
|
|
}
|
|
|
|
// Dangling resource, delete the config + state Secrets, as well as
|
|
// deleting the device from the tailnet.
|
|
if err := tailscaled.EnsureDeviceDeleted(ctx, tsClient, logger, m.tsID); err != nil {
|
|
return err
|
|
}
|
|
if err := r.Delete(ctx, m.stateSecret); err != nil && !apierrors.IsNotFound(err) {
|
|
return fmt.Errorf("error deleting state Secret %q: %w", m.stateSecret.Name, err)
|
|
}
|
|
configSecret := m.stateSecret.DeepCopy()
|
|
configSecret.Name += "-config"
|
|
if err := r.Delete(ctx, configSecret); err != nil && !apierrors.IsNotFound(err) {
|
|
return fmt.Errorf("error deleting config Secret %q: %w", configSecret.Name, err)
|
|
}
|
|
// NOTE(ChaosInTheCRD): we shouldn't need to get the service first, checking for a not found error should be enough
|
|
svc := &corev1.Service{
|
|
ObjectMeta: metav1.ObjectMeta{
|
|
Name: fmt.Sprintf("%s-nodeport", m.stateSecret.Name),
|
|
Namespace: m.stateSecret.Namespace,
|
|
},
|
|
}
|
|
if err := r.Delete(ctx, svc); err != nil {
|
|
if !apierrors.IsNotFound(err) {
|
|
return fmt.Errorf("error deleting static endpoints Kubernetes Service %q: %w", svc.Name, err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// If the ProxyClass has its StaticEndpoints config removed, we want to remove all of the NodePort Services
|
|
if pc != nil && pc.Spec.StaticEndpoints == nil {
|
|
labels := map[string]string{
|
|
kubetypes.LabelManaged: "true",
|
|
LabelParentType: proxyTypeProxyGroup,
|
|
LabelParentName: pg.Name,
|
|
}
|
|
if err := r.DeleteAllOf(ctx, &corev1.Service{}, client.InNamespace(r.tsNamespace), client.MatchingLabels(labels)); err != nil {
|
|
return fmt.Errorf("error deleting Kubernetes Services for static endpoints: %w", err)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// maybeCleanup just deletes the device from the tailnet. All the kubernetes
|
|
// resources linked to a ProxyGroup will get cleaned up via owner references
|
|
// (which we can use because they are all in the same namespace).
|
|
func (r *ProxyGroupReconciler) maybeCleanup(ctx context.Context, tsClient tsclient.Client, pg *tsapi.ProxyGroup) (bool, error) {
|
|
logger := r.logger(pg.Name)
|
|
|
|
metadata, err := getNodeMetadata(ctx, pg, r.Client, r.tsNamespace)
|
|
if err != nil {
|
|
return false, err
|
|
}
|
|
|
|
for _, m := range metadata {
|
|
if err := tailscaled.EnsureDeviceDeleted(ctx, tsClient, logger, m.tsID); err != nil {
|
|
return false, err
|
|
}
|
|
}
|
|
|
|
mo := &metricsOpts{
|
|
proxyLabels: pgLabels(pg.Name, nil),
|
|
tsNamespace: r.tsNamespace,
|
|
proxyType: "proxygroup",
|
|
}
|
|
if err := maybeCleanupMetricsResources(ctx, mo, r.Client); err != nil {
|
|
return false, fmt.Errorf("error cleaning up metrics resources: %w", err)
|
|
}
|
|
|
|
logger.Infof("cleaned up ProxyGroup resources")
|
|
r.mu.Lock()
|
|
r.ensureStateRemovedForProxyGroup(pg)
|
|
r.mu.Unlock()
|
|
return true, nil
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) ensureConfigSecretsCreated(
|
|
ctx context.Context,
|
|
tsClient tsclient.Client,
|
|
pg *tsapi.ProxyGroup,
|
|
proxyClass *tsapi.ProxyClass,
|
|
svcToNodePorts map[string]uint16,
|
|
) (endpoints map[string][]netip.AddrPort, err error) {
|
|
logger := r.logger(pg.Name)
|
|
endpoints = make(map[string][]netip.AddrPort, pgReplicas(pg)) // keyed by Service name.
|
|
for i := range pgReplicas(pg) {
|
|
logger = logger.With("Pod", fmt.Sprintf("%s-%d", pg.Name, i))
|
|
cfgSecret := &corev1.Secret{
|
|
ObjectMeta: metav1.ObjectMeta{
|
|
Name: pgConfigSecretName(pg.Name, i),
|
|
Namespace: r.tsNamespace,
|
|
Labels: pgSecretLabels(pg.Name, kubetypes.LabelSecretTypeConfig),
|
|
OwnerReferences: pgOwnerReference(pg),
|
|
},
|
|
}
|
|
|
|
var existingCfgSecret *corev1.Secret // unmodified copy of secret
|
|
if err = r.Get(ctx, client.ObjectKeyFromObject(cfgSecret), cfgSecret); err == nil {
|
|
logger.Debugf("Secret %s/%s already exists", cfgSecret.GetNamespace(), cfgSecret.GetName())
|
|
existingCfgSecret = cfgSecret.DeepCopy()
|
|
} else if !apierrors.IsNotFound(err) {
|
|
return nil, err
|
|
}
|
|
|
|
authKey, err := r.getAuthKey(ctx, tsClient, pg, existingCfgSecret, i, logger)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
nodePortSvcName := staticendpoints.NodePortServiceName(pg.Name, i)
|
|
if len(svcToNodePorts) > 0 {
|
|
replicaName := fmt.Sprintf("%s-%d", pg.Name, i)
|
|
port, ok := svcToNodePorts[nodePortSvcName]
|
|
if !ok {
|
|
return nil, fmt.Errorf("could not find configured NodePort for ProxyGroup replica %q", replicaName)
|
|
}
|
|
|
|
endpoints[nodePortSvcName], err = staticendpoints.FindEndpoints(ctx, r.Client, staticEndpointsFromConfigSecret(existingCfgSecret, logger), proxyClass, port, logger)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("could not find static endpoints for replica %q: %w", replicaName, err)
|
|
}
|
|
}
|
|
|
|
if pg.Spec.Type == tsapi.ProxyGroupTypeKubernetesAPIServer {
|
|
hostname := pgHostname(pg, i)
|
|
|
|
if authKey == nil && existingCfgSecret != nil {
|
|
deviceAuthed := false
|
|
for _, d := range pg.Status.Devices {
|
|
if d.Hostname == hostname {
|
|
deviceAuthed = true
|
|
break
|
|
}
|
|
}
|
|
if !deviceAuthed {
|
|
existingCfg := conf.ConfigV1Alpha1{}
|
|
if err := json.Unmarshal(existingCfgSecret.Data[kubetypes.KubeAPIServerConfigFile], &existingCfg); err != nil {
|
|
return nil, fmt.Errorf("error unmarshalling existing config: %w", err)
|
|
}
|
|
if existingCfg.AuthKey != nil {
|
|
authKey = existingCfg.AuthKey
|
|
}
|
|
}
|
|
}
|
|
|
|
mode := kubetypes.APIServerProxyModeAuth
|
|
if !isAuthAPIServerProxy(pg) {
|
|
mode = kubetypes.APIServerProxyModeNoAuth
|
|
}
|
|
cfg := conf.VersionedConfig{
|
|
Version: "v1alpha1",
|
|
ConfigV1Alpha1: &conf.ConfigV1Alpha1{
|
|
AuthKey: authKey,
|
|
State: new(fmt.Sprintf("kube:%s", pgPodName(pg.Name, i))),
|
|
App: new(kubetypes.AppProxyGroupKubeAPIServer),
|
|
LogLevel: new(logger.Level().String()),
|
|
|
|
// Reloadable fields.
|
|
Hostname: &hostname,
|
|
APIServerProxy: &conf.APIServerProxyConfig{
|
|
Enabled: opt.NewBool(true),
|
|
Mode: &mode,
|
|
// The first replica is elected as the cert issuer, same
|
|
// as containerboot does for ingress-pg-reconciler.
|
|
IssueCerts: opt.NewBool(i == 0),
|
|
},
|
|
LocalPort: new(uint16(9002)),
|
|
HealthCheckEnabled: opt.NewBool(true),
|
|
},
|
|
}
|
|
|
|
// Copy over config that the apiserver-proxy-service-reconciler sets.
|
|
if existingCfgSecret != nil {
|
|
if k8sProxyCfg, ok := cfgSecret.Data[kubetypes.KubeAPIServerConfigFile]; ok {
|
|
k8sCfg := &conf.ConfigV1Alpha1{}
|
|
if err := json.Unmarshal(k8sProxyCfg, k8sCfg); err != nil {
|
|
return nil, fmt.Errorf("failed to unmarshal kube-apiserver config: %w", err)
|
|
}
|
|
|
|
cfg.AdvertiseServices = k8sCfg.AdvertiseServices
|
|
if k8sCfg.APIServerProxy != nil {
|
|
cfg.APIServerProxy.ServiceName = k8sCfg.APIServerProxy.ServiceName
|
|
}
|
|
}
|
|
}
|
|
|
|
if tsClient.LoginURL() != "" {
|
|
cfg.ServerURL = new(tsClient.LoginURL())
|
|
}
|
|
|
|
if proxyClass != nil && proxyClass.Spec.TailscaleConfig != nil {
|
|
cfg.AcceptRoutes = opt.NewBool(proxyClass.Spec.TailscaleConfig.AcceptRoutes)
|
|
}
|
|
|
|
if proxyClass != nil && proxyClass.Spec.Metrics != nil {
|
|
cfg.MetricsEnabled = opt.NewBool(proxyClass.Spec.Metrics.Enable)
|
|
}
|
|
|
|
if len(endpoints[nodePortSvcName]) > 0 {
|
|
cfg.StaticEndpoints = endpoints[nodePortSvcName]
|
|
}
|
|
|
|
cfgB, err := json.Marshal(cfg)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("error marshalling k8s-proxy config: %w", err)
|
|
}
|
|
mak.Set(&cfgSecret.Data, kubetypes.KubeAPIServerConfigFile, cfgB)
|
|
} else {
|
|
// AdvertiseServices config is set by ingress-pg-reconciler, so make sure we
|
|
// don't overwrite it if already set.
|
|
existingAdvertiseServices, err := extractAdvertiseServicesConfig(existingCfgSecret)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
configs, err := pgTailscaledConfig(pg, tsClient.LoginURL(), proxyClass, i, authKey, endpoints[nodePortSvcName], existingAdvertiseServices)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("error creating tailscaled config: %w", err)
|
|
}
|
|
|
|
for cap, cfg := range configs {
|
|
cfgJSON, err := json.Marshal(cfg)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("error marshalling tailscaled config: %w", err)
|
|
}
|
|
mak.Set(&cfgSecret.Data, tsoperator.TailscaledConfigFileName(cap), cfgJSON)
|
|
}
|
|
}
|
|
|
|
if existingCfgSecret != nil {
|
|
if !apiequality.Semantic.DeepEqual(existingCfgSecret, cfgSecret) {
|
|
logger.Debugf("Updating the existing ProxyGroup config Secret %s", cfgSecret.Name)
|
|
if err := r.Update(ctx, cfgSecret); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
} else {
|
|
logger.Debugf("Creating a new config Secret %s for the ProxyGroup", cfgSecret.Name)
|
|
if err := r.Create(ctx, cfgSecret); err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
}
|
|
|
|
return endpoints, nil
|
|
}
|
|
|
|
// getAuthKey returns an auth key for the proxy, or nil if none is needed.
|
|
// A new key is created if the config Secret doesn't exist yet, or if the
|
|
// proxy has requested a reissue via its state Secret. An existing key is
|
|
// retained while the device hasn't authed or a reissue is in progress.
|
|
func (r *ProxyGroupReconciler) getAuthKey(ctx context.Context, tsClient tsclient.Client, pg *tsapi.ProxyGroup, existingCfgSecret *corev1.Secret, ordinal int32, logger *zap.SugaredLogger) (*string, error) {
|
|
// Get state Secret to check if it's already authed or has requested
|
|
// a fresh auth key.
|
|
stateSecret := &corev1.Secret{
|
|
ObjectMeta: metav1.ObjectMeta{
|
|
Name: pgStateSecretName(pg.Name, ordinal),
|
|
Namespace: r.tsNamespace,
|
|
},
|
|
}
|
|
if err := r.Get(ctx, client.ObjectKeyFromObject(stateSecret), stateSecret); err != nil && !apierrors.IsNotFound(err) {
|
|
return nil, err
|
|
}
|
|
|
|
var createAuthKey bool
|
|
var cfgAuthKey *string
|
|
if existingCfgSecret == nil {
|
|
createAuthKey = true
|
|
} else {
|
|
var err error
|
|
cfgAuthKey, err = authKeyFromSecret(existingCfgSecret)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("error retrieving auth key from existing config Secret: %w", err)
|
|
}
|
|
}
|
|
|
|
if !createAuthKey {
|
|
var err error
|
|
createAuthKey, err = r.reissuer.ShouldReissue(ctx, tsClient, r.log, tailscaled.ReissueInput{
|
|
ParentName: pg.Name,
|
|
ReplicaName: stateSecret.Name,
|
|
Kind: tailscaled.KindProxyGroup,
|
|
StateSecret: stateSecret,
|
|
CfgAuthKey: cfgAuthKey,
|
|
})
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
}
|
|
|
|
var authKey *string
|
|
if createAuthKey {
|
|
logger.Debugf("creating auth key for ProxyGroup proxy %q", stateSecret.Name)
|
|
|
|
tags := pg.Spec.Tags.Stringify()
|
|
if len(tags) == 0 {
|
|
tags = r.defaultTags
|
|
}
|
|
key, err := newAuthKey(ctx, tsClient, tags)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
authKey = &key
|
|
} else {
|
|
// Retain auth key if the device hasn't authed yet, or if a
|
|
// reissue is in progress (device_id is stale during reissue).
|
|
_, reissueRequested := stateSecret.Data[kubetypes.KeyReissueAuthkey]
|
|
if !deviceAuthed(stateSecret) || reissueRequested {
|
|
authKey = cfgAuthKey
|
|
}
|
|
}
|
|
|
|
return authKey, nil
|
|
}
|
|
|
|
// ensureStateAddedForProxyGroup ensures the gauge metric for the ProxyGroup resource is updated when the ProxyGroup
|
|
// is created, and initialises per-ProxyGroup auth key re-issuance state. r.mu must be held.
|
|
func (r *ProxyGroupReconciler) ensureStateAddedForProxyGroup(pg *tsapi.ProxyGroup) {
|
|
switch pg.Spec.Type {
|
|
case tsapi.ProxyGroupTypeEgress:
|
|
r.egressProxyGroups.Add(pg.UID)
|
|
case tsapi.ProxyGroupTypeIngress:
|
|
r.ingressProxyGroups.Add(pg.UID)
|
|
case tsapi.ProxyGroupTypeKubernetesAPIServer:
|
|
r.apiServerProxyGroups.Add(pg.UID)
|
|
}
|
|
gaugeEgressProxyGroupResources.Set(int64(r.egressProxyGroups.Len()))
|
|
gaugeIngressProxyGroupResources.Set(int64(r.ingressProxyGroups.Len()))
|
|
gaugeAPIServerProxyGroupResources.Set(int64(r.apiServerProxyGroups.Len()))
|
|
|
|
r.reissuer.EnsureState(pg.Name, int(pgReplicas(pg)))
|
|
}
|
|
|
|
// ensureStateRemovedForProxyGroup ensures the gauge metric for the ProxyGroup resource type is updated when the
|
|
// ProxyGroup is deleted, and drops the per-ProxyGroup auth key re-issuance state to free memory. r.mu must be held.
|
|
func (r *ProxyGroupReconciler) ensureStateRemovedForProxyGroup(pg *tsapi.ProxyGroup) {
|
|
switch pg.Spec.Type {
|
|
case tsapi.ProxyGroupTypeEgress:
|
|
r.egressProxyGroups.Remove(pg.UID)
|
|
case tsapi.ProxyGroupTypeIngress:
|
|
r.ingressProxyGroups.Remove(pg.UID)
|
|
case tsapi.ProxyGroupTypeKubernetesAPIServer:
|
|
r.apiServerProxyGroups.Remove(pg.UID)
|
|
}
|
|
gaugeEgressProxyGroupResources.Set(int64(r.egressProxyGroups.Len()))
|
|
gaugeIngressProxyGroupResources.Set(int64(r.ingressProxyGroups.Len()))
|
|
gaugeAPIServerProxyGroupResources.Set(int64(r.apiServerProxyGroups.Len()))
|
|
r.reissuer.RemoveState(pg.Name)
|
|
}
|
|
|
|
func pgTailscaledConfig(pg *tsapi.ProxyGroup, loginServer string, pc *tsapi.ProxyClass, idx int32, authKey *string, staticEndpoints []netip.AddrPort, oldAdvertiseServices []string) (tailscaledConfigs, error) {
|
|
conf := &ipn.ConfigVAlpha{
|
|
Version: "alpha0",
|
|
AcceptDNS: "false",
|
|
AcceptRoutes: "false", // AcceptRoutes defaults to true
|
|
Locked: "false",
|
|
Hostname: new(pgHostname(pg, idx)),
|
|
AdvertiseServices: oldAdvertiseServices,
|
|
AuthKey: authKey,
|
|
}
|
|
|
|
if loginServer != "" {
|
|
conf.ServerURL = &loginServer
|
|
}
|
|
|
|
if shouldAcceptRoutes(pc) {
|
|
conf.AcceptRoutes = "true"
|
|
}
|
|
|
|
if len(staticEndpoints) > 0 {
|
|
conf.StaticEndpoints = staticEndpoints
|
|
}
|
|
|
|
return map[tailcfg.CapabilityVersion]ipn.ConfigVAlpha{
|
|
pgMinCapabilityVersion: *conf,
|
|
}, nil
|
|
}
|
|
|
|
func extractAdvertiseServicesConfig(cfgSecret *corev1.Secret) ([]string, error) {
|
|
if cfgSecret == nil {
|
|
return nil, nil
|
|
}
|
|
|
|
cfg, err := latestConfigFromSecret(cfgSecret)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if cfg == nil {
|
|
return nil, nil
|
|
}
|
|
|
|
return cfg.AdvertiseServices, nil
|
|
}
|
|
|
|
// getNodeMetadata gets metadata for all the pods owned by this ProxyGroup by
|
|
// querying their state Secrets. It may not return the same number of items as
|
|
// specified in the ProxyGroup spec if e.g. it is getting scaled up or down, or
|
|
// some pods have failed to write state.
|
|
//
|
|
// The returned metadata will contain an entry for each state Secret that exists.
|
|
func getNodeMetadata(ctx context.Context, pg *tsapi.ProxyGroup, cl client.Client, tsNamespace string) (metadata []nodeMetadata, _ error) {
|
|
// List all state Secrets owned by this ProxyGroup.
|
|
secrets := &corev1.SecretList{}
|
|
if err := cl.List(ctx, secrets, client.InNamespace(tsNamespace), client.MatchingLabels(pgSecretLabels(pg.Name, kubetypes.LabelSecretTypeState))); err != nil {
|
|
return nil, fmt.Errorf("failed to list state Secrets: %w", err)
|
|
}
|
|
for _, secret := range secrets.Items {
|
|
var ordinal int32
|
|
if _, err := fmt.Sscanf(secret.Name, pg.Name+"-%d", &ordinal); err != nil {
|
|
return nil, fmt.Errorf("unexpected secret %s was labelled as owned by the ProxyGroup %s: %w", secret.Name, pg.Name, err)
|
|
}
|
|
|
|
nm := nodeMetadata{
|
|
ordinal: ordinal,
|
|
stateSecret: &secret,
|
|
}
|
|
|
|
prefs, ok, err := getDevicePrefs(&secret)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if ok {
|
|
nm.tsID = prefs.Config.NodeID
|
|
nm.dnsName = prefs.Config.UserProfile.LoginName
|
|
}
|
|
|
|
pod := &corev1.Pod{}
|
|
if err := cl.Get(ctx, client.ObjectKey{Namespace: tsNamespace, Name: fmt.Sprintf("%s-%d", pg.Name, ordinal)}, pod); err != nil && !apierrors.IsNotFound(err) {
|
|
return nil, err
|
|
} else if err == nil {
|
|
nm.podUID = string(pod.UID)
|
|
}
|
|
metadata = append(metadata, nm)
|
|
}
|
|
|
|
// Sort for predictable ordering and status.
|
|
sort.Slice(metadata, func(i, j int) bool {
|
|
return metadata[i].ordinal < metadata[j].ordinal
|
|
})
|
|
|
|
return metadata, nil
|
|
}
|
|
|
|
// getRunningProxies will return status for all proxy Pods whose state Secret
|
|
// has an up to date Pod UID and at least a hostname.
|
|
func (r *ProxyGroupReconciler) getRunningProxies(ctx context.Context, pg *tsapi.ProxyGroup, staticEndpoints map[string][]netip.AddrPort) (devices []tsapi.TailnetDevice, _ error) {
|
|
metadata, err := getNodeMetadata(ctx, pg, r.Client, r.tsNamespace)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
for i, m := range metadata {
|
|
if m.podUID == "" || !strings.EqualFold(string(m.stateSecret.Data[kubetypes.KeyPodUID]), m.podUID) {
|
|
// Current Pod has not yet written its UID to the state Secret, data may
|
|
// be stale.
|
|
continue
|
|
}
|
|
|
|
device := tsapi.TailnetDevice{}
|
|
if hostname, _, ok := strings.Cut(string(m.stateSecret.Data[kubetypes.KeyDeviceFQDN]), "."); ok {
|
|
device.Hostname = hostname
|
|
} else {
|
|
continue
|
|
}
|
|
|
|
if ipsB := m.stateSecret.Data[kubetypes.KeyDeviceIPs]; len(ipsB) > 0 {
|
|
ips := []string{}
|
|
if err := json.Unmarshal(ipsB, &ips); err != nil {
|
|
return nil, fmt.Errorf("failed to extract device IPs from state Secret %q: %w", m.stateSecret.Name, err)
|
|
}
|
|
device.TailnetIPs = ips
|
|
}
|
|
|
|
// TODO(tomhjp): This is our input to the proxy, but we should instead
|
|
// read this back from the proxy's state in some way to more accurately
|
|
// reflect its status.
|
|
if ep, ok := staticEndpoints[staticendpoints.NodePortServiceName(pg.Name, int32(i))]; ok && len(ep) > 0 {
|
|
eps := make([]string, 0, len(ep))
|
|
for _, e := range ep {
|
|
eps = append(eps, e.String())
|
|
}
|
|
device.StaticEndpoints = eps
|
|
}
|
|
|
|
devices = append(devices, device)
|
|
}
|
|
|
|
return devices, nil
|
|
}
|
|
|
|
type nodeMetadata struct {
|
|
ordinal int32
|
|
stateSecret *corev1.Secret
|
|
podUID string // or empty if the Pod no longer exists.
|
|
tsID tailcfg.StableNodeID
|
|
dnsName string
|
|
}
|
|
|
|
func notReady(reason, msg string) (map[string][]netip.AddrPort, *notReadyReason, error) {
|
|
return nil, ¬ReadyReason{
|
|
reason: reason,
|
|
message: msg,
|
|
}, nil
|
|
}
|
|
|
|
// sharedACMEAccountEnabledFor reports whether the shared-ACME-account
|
|
// feature should be applied to pg. The per-PG
|
|
// tailscale.com/share-acme-account annotation wins when set; otherwise
|
|
// the operator's OPERATOR_SHARED_ACME_ACCOUNT_KEY setting is the default
|
|
// for every ProxyGroup.
|
|
func (r *ProxyGroupReconciler) sharedACMEAccountEnabledFor(pg *tsapi.ProxyGroup) bool {
|
|
return sharedACMEAccountEnabled(pg, r.sharedACMEAccountKey)
|
|
}
|
|
|
|
// sharedACMEAccountEnabled reports whether pg should use the shared ACME
|
|
// account, with the tailscale.com/share-acme-account annotation overriding
|
|
// the operator-wide default.
|
|
func sharedACMEAccountEnabled(pg *tsapi.ProxyGroup, operatorDefault bool) bool {
|
|
if v, ok := pg.Annotations[AnnotationShareACMEAccount]; ok {
|
|
return v == "true"
|
|
}
|
|
return operatorDefault
|
|
}
|
|
|
|
func (r *ProxyGroupReconciler) notReadyErrf(pg *tsapi.ProxyGroup, logger *zap.SugaredLogger, format string, a ...any) (map[string][]netip.AddrPort, *notReadyReason, error) {
|
|
err := fmt.Errorf(format, a...)
|
|
if strings.Contains(err.Error(), optimisticLockErrorMsg) {
|
|
msg := fmt.Sprintf("optimistic lock error, retrying: %s", err.Error())
|
|
logger.Info(msg)
|
|
return notReady(reasonProxyGroupCreating, msg)
|
|
}
|
|
|
|
r.recorder.Event(pg, corev1.EventTypeWarning, reasonProxyGroupCreationFailed, err.Error())
|
|
return nil, ¬ReadyReason{
|
|
reason: reasonProxyGroupCreationFailed,
|
|
message: err.Error(),
|
|
}, err
|
|
}
|
|
|
|
type notReadyReason struct {
|
|
reason string
|
|
message string
|
|
}
|