322 lines
13 KiB
Go
322 lines
13 KiB
Go
package cluster
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"sort"
|
|
"strings"
|
|
"time"
|
|
|
|
"scm.bstein.dev/bstein/ananke/internal/state"
|
|
)
|
|
|
|
const credentialEventMaxAge = 10 * time.Minute
|
|
const credentialRepairCooldown = 30 * time.Minute
|
|
|
|
type imagePullCredentialDeploymentList struct {
|
|
Items []struct {
|
|
Metadata struct {
|
|
Name string `json:"name"`
|
|
Labels map[string]string `json:"labels"`
|
|
} `json:"metadata"`
|
|
} `json:"items"`
|
|
}
|
|
|
|
// imagePullCredentialBlockerReasons classifies image pull failures caused by
|
|
// registry auth or missing imagePullSecret material.
|
|
// Signature: (o *Orchestrator) imagePullCredentialBlockerReasons(ctx context.Context) (map[string]string, error).
|
|
// Why: deleting a pod with bad registry credentials does not refresh Vault/CSI
|
|
// secret material; the secret sync path needs a direct repair signal.
|
|
func (o *Orchestrator) imagePullCredentialBlockerReasons(ctx context.Context) (map[string]string, error) {
|
|
eventsOut, err := o.kubectl(ctx, 30*time.Second, "get", "events", "-A", "-o", "json")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("query events for image-pull credential scan: %w", err)
|
|
}
|
|
reasons := map[string]string{}
|
|
if strings.TrimSpace(eventsOut) == "" {
|
|
return reasons, nil
|
|
}
|
|
var events eventList
|
|
if err := json.Unmarshal([]byte(eventsOut), &events); err != nil {
|
|
return nil, fmt.Errorf("decode events for image-pull credential scan: %w", err)
|
|
}
|
|
if len(events.Items) == 0 {
|
|
return reasons, nil
|
|
}
|
|
podsOut, err := o.kubectl(ctx, 30*time.Second, "get", "pods", "-A", "-o", "json")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("query current pods for image-pull credential scan: %w", err)
|
|
}
|
|
var pods podList
|
|
if err := json.Unmarshal([]byte(podsOut), &pods); err != nil {
|
|
return nil, fmt.Errorf("decode current pods for image-pull credential scan: %w", err)
|
|
}
|
|
blocked := map[string]string{}
|
|
for _, pod := range pods.Items {
|
|
if podHasCurrentImagePullFailure(pod) {
|
|
blocked[pod.Metadata.Namespace+"/"+pod.Metadata.Name] = pod.Metadata.UID
|
|
}
|
|
}
|
|
for _, event := range events.Items {
|
|
if !strings.EqualFold(strings.TrimSpace(event.Type), "Warning") {
|
|
continue
|
|
}
|
|
if !strings.EqualFold(strings.TrimSpace(event.InvolvedObject.Kind), "Pod") {
|
|
continue
|
|
}
|
|
reason := strings.TrimSpace(event.Reason)
|
|
message := strings.TrimSpace(event.Message)
|
|
if !imagePullCredentialEvent(reason, message) {
|
|
continue
|
|
}
|
|
namespace := strings.TrimSpace(event.InvolvedObject.Namespace)
|
|
if namespace == "" {
|
|
namespace = strings.TrimSpace(event.Metadata.Namespace)
|
|
}
|
|
name := strings.TrimSpace(event.InvolvedObject.Name)
|
|
if namespace == "" || name == "" {
|
|
continue
|
|
}
|
|
// Retained events outlive pods and successful pulls. Require current state,
|
|
// exact pod identity, and recent evidence before changing a sync helper.
|
|
uid := blocked[namespace+"/"+name]
|
|
observed := eventLastObservedAt(event)
|
|
if uid == "" || uid != event.InvolvedObject.UID || observed.IsZero() ||
|
|
time.Since(observed) > credentialEventMaxAge || time.Until(observed) > time.Minute {
|
|
continue
|
|
}
|
|
reasons[namespace+"/"+name] = "ImagePullCredentialBlocker:" + imagePullCredentialFailureClass(reason, message)
|
|
}
|
|
return reasons, nil
|
|
}
|
|
|
|
// healImagePullCredentialSync restarts namespace-local Vault sync deployments
|
|
// when pods are blocked on registry credentials.
|
|
// Signature: (o *Orchestrator) healImagePullCredentialSync(ctx context.Context) ([]string, error).
|
|
// Why: Secrets Store CSI secretObjects are materialized by mounted workloads; a
|
|
// restarted vault-sync deployment is the bounded repair that refreshes pull
|
|
// secrets such as harbor-regcred without mutating application deployments.
|
|
func (o *Orchestrator) healImagePullCredentialSync(ctx context.Context) ([]string, error) {
|
|
if o.runner.DryRun {
|
|
return nil, nil
|
|
}
|
|
blockers, err := o.imagePullCredentialBlockerReasons(ctx)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if len(blockers) == 0 {
|
|
return nil, nil
|
|
}
|
|
|
|
namespaces := map[string]struct{}{}
|
|
for key := range blockers {
|
|
namespace, _, ok := strings.Cut(key, "/")
|
|
if ok && strings.TrimSpace(namespace) != "" {
|
|
namespaces[namespace] = struct{}{}
|
|
}
|
|
}
|
|
orderedNamespaces := make([]string, 0, len(namespaces))
|
|
for namespace := range namespaces {
|
|
orderedNamespaces = append(orderedNamespaces, namespace)
|
|
}
|
|
sort.Strings(orderedNamespaces)
|
|
|
|
repaired := []string{}
|
|
errs := []string{}
|
|
for _, namespace := range orderedNamespaces {
|
|
namespaceRepairs, restartErr := o.restartVaultSyncDeployments(ctx, namespace)
|
|
repaired = append(repaired, namespaceRepairs...)
|
|
if restartErr != nil {
|
|
errs = append(errs, restartErr.Error())
|
|
}
|
|
}
|
|
sort.Strings(repaired)
|
|
if len(errs) > 0 {
|
|
return repaired, errors.New(strings.Join(errs, "; "))
|
|
}
|
|
return repaired, nil
|
|
}
|
|
|
|
// restartVaultSyncDeployments runs one orchestration or CLI step.
|
|
// Signature: (o *Orchestrator) restartVaultSyncDeployments(ctx context.Context, namespace string) ([]string, error).
|
|
// Why: namespaces using Vault/CSI secret material have tiny sync deployments
|
|
// named or labeled vault-sync; restarting those nudges secretObject rotation.
|
|
func (o *Orchestrator) restartVaultSyncDeployments(ctx context.Context, namespace string) ([]string, error) {
|
|
o.credentialRepairMu.Lock()
|
|
defer o.credentialRepairMu.Unlock()
|
|
out, err := o.kubectl(ctx, 20*time.Second, "-n", namespace, "get", "deployment", "-o", "json")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("query deployments in %s for image-pull credential repair: %w", namespace, err)
|
|
}
|
|
var deployments imagePullCredentialDeploymentList
|
|
if err := json.Unmarshal([]byte(out), &deployments); err != nil {
|
|
return nil, fmt.Errorf("decode deployments in %s for image-pull credential repair: %w", namespace, err)
|
|
}
|
|
|
|
repaired := []string{}
|
|
found := false
|
|
for _, deployment := range deployments.Items {
|
|
name := strings.TrimSpace(deployment.Metadata.Name)
|
|
if name == "" || !vaultSyncDeployment(name, deployment.Metadata.Labels) {
|
|
continue
|
|
}
|
|
found = true
|
|
allowed, err := o.reserveCredentialRepair(namespace + "/deployment/" + name)
|
|
if err != nil {
|
|
return repaired, err
|
|
}
|
|
if !allowed {
|
|
continue
|
|
}
|
|
if _, err := o.kubectl(ctx, 25*time.Second, "-n", namespace, "rollout", "restart", "deployment", name); err != nil {
|
|
return repaired, fmt.Errorf("restart %s/deployment/%s for image-pull credential repair: %w", namespace, name, err)
|
|
}
|
|
if _, err := o.kubectl(ctx, 75*time.Second, "-n", namespace, "rollout", "status", "deployment/"+name, "--timeout=60s"); err != nil {
|
|
return repaired, fmt.Errorf("wait for %s/deployment/%s after image-pull credential repair: %w", namespace, name, err)
|
|
}
|
|
repaired = append(repaired, namespace+"/deployment/"+name)
|
|
}
|
|
if !found {
|
|
return nil, fmt.Errorf("image-pull credential blocker in namespace %s but no vault-sync deployment was found", namespace)
|
|
}
|
|
return repaired, nil
|
|
}
|
|
|
|
// imagePullCredentialEvent runs one orchestration or CLI step.
|
|
// Signature: imagePullCredentialEvent(reason string, message string) bool.
|
|
// Why: image-pull auth detection must include kubelet pull-secret events while
|
|
// staying separate from DNS and ordinary transient pull backoff.
|
|
func imagePullCredentialEvent(reason string, message string) bool {
|
|
normalizedReason := strings.ToLower(strings.TrimSpace(reason))
|
|
if normalizedReason == "failedtoretrieveimagepullsecret" {
|
|
return true
|
|
}
|
|
if normalizedReason != "failed" &&
|
|
normalizedReason != "failedpull" &&
|
|
normalizedReason != "errimagepull" &&
|
|
normalizedReason != "imagepullbackoff" {
|
|
return false
|
|
}
|
|
return imagePullMessageHasCredentialFailure(message)
|
|
}
|
|
|
|
// imagePullMessageHasCredentialFailure runs one orchestration or CLI step.
|
|
// Signature: imagePullMessageHasCredentialFailure(message string) bool.
|
|
// Why: credential blockers need a tight predicate so missing images and registry
|
|
// DNS outages keep their own repair paths.
|
|
func imagePullMessageHasCredentialFailure(message string) bool {
|
|
lower := strings.ToLower(strings.TrimSpace(message))
|
|
if lower == "" || imagePullMessageHasDNSFailure(message) {
|
|
return false
|
|
}
|
|
return strings.Contains(lower, "failed to retrieve image pull secret") ||
|
|
strings.Contains(lower, "failedtoretrieveimagepullsecret") ||
|
|
strings.Contains(lower, "unable to retrieve some image pull secrets") ||
|
|
(strings.Contains(lower, "image pull secret") && strings.Contains(lower, "not found")) ||
|
|
(strings.Contains(lower, "pull secret") && strings.Contains(lower, "not found")) ||
|
|
(strings.Contains(lower, "secret ") && strings.Contains(lower, " not found") && strings.Contains(lower, "pull")) ||
|
|
strings.Contains(lower, "no basic auth credentials") ||
|
|
strings.Contains(lower, "unauthorized") ||
|
|
strings.Contains(lower, "authentication required") ||
|
|
strings.Contains(lower, "authorization failed") ||
|
|
strings.Contains(lower, "failed to authorize") ||
|
|
strings.Contains(lower, "invalid username/password") ||
|
|
strings.Contains(lower, "401 unauthorized") ||
|
|
strings.Contains(lower, "403 forbidden") ||
|
|
strings.Contains(lower, "pull access denied") ||
|
|
strings.Contains(lower, "requested access to the resource is denied") ||
|
|
strings.Contains(lower, "failed to fetch anonymous token")
|
|
}
|
|
|
|
// imagePullCredentialFailureClass runs one orchestration or CLI step.
|
|
// Signature: imagePullCredentialFailureClass(reason string, message string) string.
|
|
// Why: compact blocker classes keep logs and status useful without copying long
|
|
// kubelet event messages.
|
|
func imagePullCredentialFailureClass(reason string, message string) string {
|
|
lowerReason := strings.ToLower(strings.TrimSpace(reason))
|
|
lower := strings.ToLower(strings.TrimSpace(message))
|
|
switch {
|
|
case lowerReason == "failedtoretrieveimagepullsecret",
|
|
strings.Contains(lower, "failed to retrieve image pull secret"),
|
|
strings.Contains(lower, "unable to retrieve some image pull secrets"),
|
|
strings.Contains(lower, "image pull secret") && strings.Contains(lower, "not found"),
|
|
strings.Contains(lower, "pull secret") && strings.Contains(lower, "not found"),
|
|
strings.Contains(lower, "secret ") && strings.Contains(lower, " not found") && strings.Contains(lower, "pull"):
|
|
return "missing-pull-secret"
|
|
case strings.Contains(lower, "no basic auth credentials"):
|
|
return "no-basic-auth"
|
|
case strings.Contains(lower, "invalid username/password"):
|
|
return "invalid-credentials"
|
|
case strings.Contains(lower, "unauthorized") || strings.Contains(lower, "authentication required"):
|
|
return "unauthorized"
|
|
case strings.Contains(lower, "forbidden") || strings.Contains(lower, "authorization failed") || strings.Contains(lower, "failed to authorize") || strings.Contains(lower, "failed to fetch anonymous token"):
|
|
return "authorization-failed"
|
|
case strings.Contains(lower, "pull access denied") || strings.Contains(lower, "requested access to the resource is denied"):
|
|
return "pull-access-denied"
|
|
default:
|
|
return "registry-credential-error"
|
|
}
|
|
}
|
|
|
|
// vaultSyncDeployment runs one orchestration or CLI step.
|
|
// Signature: vaultSyncDeployment(name string, labels map[string]string) bool.
|
|
// Why: different apps may prefix their sync deployment names, but they consistently
|
|
// identify the narrow Vault sync helper with a vault-sync name or label.
|
|
func vaultSyncDeployment(name string, labels map[string]string) bool {
|
|
if strings.Contains(strings.ToLower(strings.TrimSpace(name)), "vault-sync") {
|
|
return true
|
|
}
|
|
for key, value := range labels {
|
|
if strings.Contains(strings.ToLower(strings.TrimSpace(key)), "vault-sync") ||
|
|
strings.Contains(strings.ToLower(strings.TrimSpace(value)), "vault-sync") {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// podHasCurrentImagePullFailure returns whether a live pod still needs a pull.
|
|
// Completed, deleting, and recovered pods cannot justify a credential repair.
|
|
// Signature: podHasCurrentImagePullFailure(pod podResource) bool.
|
|
// Why: historical warnings must not trigger repairs for recovered pods.
|
|
func podHasCurrentImagePullFailure(pod podResource) bool {
|
|
if pod.Metadata.DeletionTimestamp != nil || pod.Status.Phase == "Succeeded" || pod.Status.Phase == "Failed" {
|
|
return false
|
|
}
|
|
for _, statuses := range [][]podContainerStatus{pod.Status.InitContainerStatuses, pod.Status.ContainerStatuses} {
|
|
for _, status := range statuses {
|
|
if wait := status.State.Waiting; wait != nil && (wait.Reason == "ErrImagePull" || wait.Reason == "ImagePullBackOff") {
|
|
return true
|
|
}
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
// reserveCredentialRepair records an attempt before mutation and returns false
|
|
// during the cooldown. The existing run history preserves it over daemon restarts.
|
|
// A failed or timed-out rollout also consumes the cooldown to prevent churn.
|
|
// Signature: (o *Orchestrator) reserveCredentialRepair(target string) (bool, error).
|
|
// Why: a cooldown must survive daemon restarts and failed rollouts.
|
|
func (o *Orchestrator) reserveCredentialRepair(target string) (bool, error) {
|
|
if o.store == nil {
|
|
return false, errors.New("credential repair requires persistent run history")
|
|
}
|
|
records, err := o.store.Load()
|
|
if err != nil {
|
|
return false, fmt.Errorf("read credential repair history: %w", err)
|
|
}
|
|
for _, record := range records {
|
|
if record.Action == "image-pull-credential-repair" && record.Reason == target && time.Since(record.StartedAt) < credentialRepairCooldown {
|
|
return false, nil
|
|
}
|
|
}
|
|
now := time.Now()
|
|
if err := o.store.Append(state.RunRecord{ID: now.UTC().Format(time.RFC3339Nano), Action: "image-pull-credential-repair", Reason: target, StartedAt: now, EndedAt: now}); err != nil {
|
|
return false, fmt.Errorf("record credential repair attempt: %w", err)
|
|
}
|
|
return true, nil
|
|
}
|