mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-09-10 10:35:51 +00:00
2001 lines
71 KiB
Go
2001 lines
71 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"flag"
|
|
"fmt"
|
|
"io"
|
|
"io/fs"
|
|
"math"
|
|
"net/http"
|
|
"os"
|
|
"os/exec"
|
|
"os/signal"
|
|
"path/filepath"
|
|
"reflect"
|
|
"runtime"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"syscall"
|
|
"time"
|
|
|
|
"github.com/prometheus/client_golang/prometheus"
|
|
"github.com/prometheus/client_golang/prometheus/promauto"
|
|
"github.com/prometheus/client_golang/prometheus/promhttp"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/agentexec"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/agenttarget"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/agentupdate"
|
|
pulseconfig "github.com/rcourtman/pulse-go-rewrite/internal/config"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/dockeragent"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/hostagent"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/kubernetesagent"
|
|
pulselogging "github.com/rcourtman/pulse-go-rewrite/internal/logging"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/remoteconfig"
|
|
"github.com/rcourtman/pulse-go-rewrite/internal/utils"
|
|
agentshost "github.com/rcourtman/pulse-go-rewrite/pkg/agents/host"
|
|
"github.com/rcourtman/pulse-go-rewrite/pkg/securityutil"
|
|
"github.com/rs/zerolog"
|
|
gohost "github.com/shirou/gopsutil/v4/host"
|
|
"golang.org/x/sync/errgroup"
|
|
)
|
|
|
|
var (
|
|
Version = "dev"
|
|
|
|
// Prometheus metrics
|
|
agentInfo = promauto.NewGaugeVec(prometheus.GaugeOpts{
|
|
Name: "pulse_agent_info",
|
|
Help: "Information about the Pulse agent",
|
|
}, []string{"version", "host_enabled", "docker_enabled", "kubernetes_enabled"})
|
|
|
|
agentUp = promauto.NewGauge(prometheus.GaugeOpts{
|
|
Name: "pulse_agent_up",
|
|
Help: "Whether the Pulse agent is running (1 = up, 0 = down)",
|
|
})
|
|
|
|
agentModuleEnabled = promauto.NewGaugeVec(prometheus.GaugeOpts{
|
|
Name: "pulse_agent_module_enabled",
|
|
Help: "Whether a Pulse Unified Agent module is configured to run (1 = enabled, 0 = disabled)",
|
|
}, []string{"module"})
|
|
|
|
agentModuleReady = promauto.NewGaugeVec(prometheus.GaugeOpts{
|
|
Name: "pulse_agent_module_ready",
|
|
Help: "Whether a configured Pulse Unified Agent module initialized successfully (1 = ready, 0 = not ready)",
|
|
}, []string{"module"})
|
|
)
|
|
|
|
const (
|
|
agentLogMaxSizeMB = 25
|
|
agentLogMaxAgeDays = 14
|
|
)
|
|
|
|
// Runnable is an interface for agents that can be run
|
|
type Runnable interface {
|
|
Run(ctx context.Context) error
|
|
}
|
|
|
|
type RemoteConfigApplier interface {
|
|
ApplyRemoteConfig(settings map[string]interface{}, commandsEnabled *bool)
|
|
}
|
|
|
|
// Runnable closer for the Docker / Podman collection module which needs cleanup.
|
|
type RunnableCloser interface {
|
|
Runnable
|
|
Close() error
|
|
}
|
|
|
|
var (
|
|
// For testing - wrappers to return interfaces
|
|
newDockerAgent func(dockeragent.Config) (RunnableCloser, error) = func(c dockeragent.Config) (RunnableCloser, error) {
|
|
return dockeragent.New(c)
|
|
}
|
|
newKubeAgent func(kubernetesagent.Config) (Runnable, error) = func(c kubernetesagent.Config) (Runnable, error) {
|
|
return kubernetesagent.New(c)
|
|
}
|
|
newHostAgent func(hostagent.Config) (Runnable, error) = func(c hostagent.Config) (Runnable, error) {
|
|
return hostagent.New(c)
|
|
}
|
|
newPrivilegeHelperTelemetry = hostagent.NewPrivilegeHelperTelemetry
|
|
newPrivilegeHelperUpdate = agentupdate.NewPrivilegeHelperUpdate
|
|
newUpdater func(agentupdate.Config) *agentupdate.Updater = agentupdate.New
|
|
lookPath = exec.LookPath
|
|
runAsWindowsServiceFunc = runAsWindowsService
|
|
|
|
// For testing
|
|
retryInitialDelay = 5 * time.Second
|
|
retryMaxDelay = 5 * time.Minute
|
|
remoteConfigRefreshInterval = 1 * time.Minute
|
|
)
|
|
|
|
// wireUpdaterHooks connects the self-updater to the host module's report loop:
|
|
// update status snapshots flow out on reports, and server versions carried on
|
|
// report acks nudge the updater so a server upgrade converges within one
|
|
// report cycle instead of the next hourly check.
|
|
func wireUpdaterHooks(hostCfg *hostagent.Config, updater *agentupdate.Updater) {
|
|
hostCfg.UpdateStatus = updater.Snapshot
|
|
hostCfg.OnServerVersion = updater.NudgeVersion
|
|
}
|
|
|
|
func pendingUpdatePreviousVersion(pending *agentupdate.PendingPrivilegedUpdate) string {
|
|
if pending == nil {
|
|
return ""
|
|
}
|
|
return pending.PreviousVersion
|
|
}
|
|
|
|
func currentCollectorExecutableSHA256() (string, error) {
|
|
file, err := os.Open("/proc/self/exe")
|
|
if err != nil {
|
|
return "", fmt.Errorf("open current collector executable: %w", err)
|
|
}
|
|
defer file.Close()
|
|
hasher := sha256.New()
|
|
const maximumCollectorBytes = 100 * 1024 * 1024
|
|
written, err := io.Copy(hasher, io.LimitReader(file, maximumCollectorBytes+1))
|
|
if err != nil {
|
|
return "", fmt.Errorf("hash current collector executable: %w", err)
|
|
}
|
|
if written > maximumCollectorBytes {
|
|
return "", errors.New("current collector executable exceeds the bounded agent size")
|
|
}
|
|
return hex.EncodeToString(hasher.Sum(nil)), nil
|
|
}
|
|
|
|
func supervisePendingPrivilegedUpdate(
|
|
ctx context.Context,
|
|
update agentupdate.PrivilegedUpdate,
|
|
pending *agentupdate.PendingPrivilegedUpdate,
|
|
stateDir string,
|
|
runningSHA256 string,
|
|
locallyReady func() bool,
|
|
reportAccepted <-chan struct{},
|
|
pollInterval time.Duration,
|
|
logger *zerolog.Logger,
|
|
) error {
|
|
if pending == nil {
|
|
return nil
|
|
}
|
|
if update == nil || locallyReady == nil || pollInterval <= 0 {
|
|
return errors.New("pending helper update supervisor is not configured")
|
|
}
|
|
activation := pending.Activation
|
|
rollback := func(reason error) error {
|
|
rollbackCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
defer cancel()
|
|
result, err := update.Rollback(rollbackCtx, activation)
|
|
if err != nil {
|
|
return errors.Join(reason, fmt.Errorf("typed helper rollback failed: %w", err))
|
|
}
|
|
if result.Action != "rolled_back" || !strings.EqualFold(result.ActiveSHA256, activation.RollbackSHA256) {
|
|
return errors.Join(reason, errors.New("typed helper returned an invalid rollback result"))
|
|
}
|
|
if clearErr := agentupdate.ClearPendingPrivilegedUpdate(stateDir); clearErr != nil {
|
|
return errors.Join(fmt.Errorf("%w; pending update rolled back", reason), fmt.Errorf("clear pending update handoff: %w", clearErr))
|
|
}
|
|
return fmt.Errorf("%w; pending update rolled back", reason)
|
|
}
|
|
runningSHA256 = strings.TrimSpace(runningSHA256)
|
|
switch {
|
|
case strings.EqualFold(runningSHA256, activation.RollbackSHA256):
|
|
rollbackCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
result, err := update.Rollback(rollbackCtx, activation)
|
|
cancel()
|
|
if err != nil {
|
|
return fmt.Errorf("confirm recovered helper update rollback: %w", err)
|
|
}
|
|
if result.Action != "rolled_back" || !strings.EqualFold(result.ActiveSHA256, activation.RollbackSHA256) {
|
|
return errors.New("typed helper returned an invalid recovered rollback result")
|
|
}
|
|
if err := agentupdate.ClearPendingPrivilegedUpdate(stateDir); err != nil {
|
|
return fmt.Errorf("clear recovered helper update handoff: %w", err)
|
|
}
|
|
if logger != nil {
|
|
logger.Info().Str("activation_id", activation.ActivationID).Msg("Cleared helper update handoff after durable rollback recovery")
|
|
}
|
|
return nil
|
|
case !strings.EqualFold(runningSHA256, activation.ActiveSHA256):
|
|
return errors.New("running collector executable does not match the pending update or its rollback identity")
|
|
}
|
|
|
|
deadlineDelay := time.Until(activation.RollbackDeadline)
|
|
if deadlineDelay <= 0 {
|
|
return rollback(errors.New("pending update health deadline expired before startup"))
|
|
}
|
|
deadline := time.NewTimer(deadlineDelay)
|
|
defer deadline.Stop()
|
|
ticker := time.NewTicker(pollInterval)
|
|
defer ticker.Stop()
|
|
reportOK := false
|
|
for {
|
|
if reportOK && locallyReady() {
|
|
commitCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
|
|
result, err := update.Commit(commitCtx, activation)
|
|
cancel()
|
|
if err == nil && result.Action == "committed" && strings.EqualFold(result.ActiveSHA256, activation.ActiveSHA256) {
|
|
if err := agentupdate.ClearPendingPrivilegedUpdate(stateDir); err != nil {
|
|
return fmt.Errorf("clear committed update handoff: %w", err)
|
|
}
|
|
if logger != nil {
|
|
logger.Info().Str("activation_id", activation.ActivationID).Msg("Committed helper-backed agent update after local readiness and server report")
|
|
}
|
|
return nil
|
|
}
|
|
if logger != nil && err != nil {
|
|
logger.Warn().Err(err).Str("activation_id", activation.ActivationID).Msg("Pending helper update commit failed; retrying until rollback deadline")
|
|
}
|
|
}
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
return rollback(fmt.Errorf("runtime stopped before pending update commit: %w", ctx.Err()))
|
|
case <-deadline.C:
|
|
return rollback(errors.New("pending update did not reach local readiness and accepted-report health floor before deadline"))
|
|
case <-reportAccepted:
|
|
reportOK = true
|
|
reportAccepted = nil
|
|
case <-ticker.C:
|
|
}
|
|
}
|
|
}
|
|
|
|
type multiValue []string
|
|
|
|
func (m *multiValue) String() string {
|
|
return strings.Join(*m, ",")
|
|
}
|
|
|
|
func (m *multiValue) Set(value string) error {
|
|
*m = append(*m, value)
|
|
return nil
|
|
}
|
|
|
|
func main() {
|
|
ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
|
defer cancel()
|
|
if len(os.Args) > 1 && isCollectorLifecycleCommand(os.Args[1]) {
|
|
err := runCollectorLifecycleCommand(ctx, os.Args[1], os.Args[2:], os.Stdout, os.Stderr)
|
|
if err != nil {
|
|
fmt.Fprintf(os.Stderr, "Error: %v\n", err)
|
|
}
|
|
if code := collectorLifecycleExitCode(err); code != 0 {
|
|
os.Exit(code)
|
|
}
|
|
return
|
|
}
|
|
|
|
if err := run(ctx, os.Args[1:], os.Getenv); err != nil {
|
|
if err == flag.ErrHelp {
|
|
os.Exit(0)
|
|
}
|
|
fmt.Fprintf(os.Stderr, "Error: %v\n", err)
|
|
os.Exit(1)
|
|
}
|
|
}
|
|
|
|
func run(ctx context.Context, args []string, getenv func(string) string) error {
|
|
// 1. Parse Configuration
|
|
cfg, err := loadConfig(args, getenv)
|
|
if err != nil {
|
|
if err == flag.ErrHelp {
|
|
return err
|
|
}
|
|
return fmt.Errorf("failed to load unified agent configuration: %w", err)
|
|
}
|
|
|
|
// 2. Setup Logging
|
|
logger, closeLogger, err := configureAgentLogger(cfg)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to configure unified agent logging: %w", err)
|
|
}
|
|
defer closeLogger()
|
|
cfg.Logger = &logger
|
|
|
|
if cfg.InsecureSkipVerify && cfg.ServerFingerprint == "" {
|
|
logger.Warn().
|
|
Str("component", "startup").
|
|
Str("action", "tls_skip_verify_enabled").
|
|
Msg("TLS verification disabled for agent connections (self-signed cert mode)")
|
|
} else if cfg.ServerFingerprint != "" {
|
|
logger.Info().
|
|
Str("component", "startup").
|
|
Str("action", "tls_server_fingerprint_enabled").
|
|
Msg("Using pinned server certificate fingerprint for agent connections")
|
|
}
|
|
|
|
// 2a. Handle Self-Test
|
|
if cfg.SelfTest {
|
|
logger.Info().Msg("Self-test passed: config loaded and logger initialized")
|
|
return nil
|
|
}
|
|
|
|
if err := secureAgentStateDir(cfg.StateDir); err != nil {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("path", cfg.StateDir).
|
|
Msg("Failed to enforce owner-only agent state directory permissions")
|
|
}
|
|
|
|
// 2b. Compute Agent ID if missing (needed for remote config)
|
|
// We replicate the logic from hostagent.New to ensure we get the same ID
|
|
lookupHostname := strings.TrimSpace(cfg.HostnameOverride)
|
|
if cfg.AgentID == "" && cfg.AgentIDFile != "" {
|
|
if persisted, err := readAgentIDFile(cfg.AgentIDFile); err == nil && persisted != "" {
|
|
cfg.AgentID = persisted
|
|
logger.Info().
|
|
Str("path", cfg.AgentIDFile).
|
|
Str("agentID", cfg.AgentID).
|
|
Msg("Loaded persisted agent ID")
|
|
} else if err != nil && !errors.Is(err, fs.ErrNotExist) {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("path", cfg.AgentIDFile).
|
|
Msg("Failed to read agent-id-file; will fall back to machine-id and rewrite the file")
|
|
}
|
|
}
|
|
if cfg.AgentID == "" {
|
|
// Use a short timeout for host info
|
|
hCtx, hCancel := context.WithTimeout(ctx, 5*time.Second)
|
|
info, err := gohost.InfoWithContext(hCtx)
|
|
hCancel()
|
|
if err == nil {
|
|
if lookupHostname == "" {
|
|
lookupHostname = strings.TrimSpace(info.Hostname)
|
|
}
|
|
collector := hostagent.NewDefaultCollector()
|
|
machineID := hostagent.GetReliableMachineID(collector, info.HostID, logger)
|
|
cfg.AgentID = machineID
|
|
if cfg.AgentID == "" {
|
|
// Fallback to hostname
|
|
cfg.AgentID = lookupHostname
|
|
}
|
|
} else {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "startup").
|
|
Str("action", "agent_id_host_info_failed").
|
|
Msg("Failed to fetch host info for Agent ID generation")
|
|
}
|
|
}
|
|
if cfg.AgentID != "" && cfg.AgentIDFile != "" {
|
|
if err := writeAgentIDFile(cfg.AgentIDFile, cfg.AgentID); err != nil {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("path", cfg.AgentIDFile).
|
|
Msg("Failed to persist agent ID; ID will be re-derived on next start")
|
|
}
|
|
}
|
|
if lookupHostname == "" {
|
|
lookupHostname = strings.TrimSpace(cfg.HostnameOverride)
|
|
if lookupHostname == "" {
|
|
if name, err := os.Hostname(); err == nil {
|
|
lookupHostname = strings.TrimSpace(name)
|
|
}
|
|
}
|
|
}
|
|
|
|
var remoteConfigClient *remoteconfig.Client
|
|
var remoteConfigAppliers []RemoteConfigApplier
|
|
|
|
// 2c. Fetch Remote Config
|
|
// Only if we have enough info to contact server
|
|
if cfg.PulseURL != "" && cfg.APIToken != "" && cfg.AgentID != "" {
|
|
logger.Debug().Msg("Fetching remote configuration...")
|
|
remoteConfigClient = remoteconfig.New(remoteconfig.Config{
|
|
PulseURL: cfg.PulseURL,
|
|
APIToken: cfg.APIToken,
|
|
AgentID: cfg.AgentID,
|
|
Hostname: lookupHostname,
|
|
InsecureSkipVerify: cfg.InsecureSkipVerify,
|
|
CACertPath: cfg.CACertPath,
|
|
ServerFingerprint: cfg.ServerFingerprint,
|
|
Logger: logger,
|
|
})
|
|
defer remoteConfigClient.Close()
|
|
|
|
// Use a short timeout for config fetch so we don't block startup too long
|
|
rcCtx, rcCancel := context.WithTimeout(ctx, 10*time.Second)
|
|
settings, commandsEnabled, err := remoteConfigClient.Fetch(rcCtx)
|
|
rcCancel()
|
|
|
|
if err != nil {
|
|
// Just log warning and proceed with local config
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "remote_config").
|
|
Str("action", "fetch_failed").
|
|
Msg("Failed to fetch remote config - using local (or previously cached) defaults")
|
|
} else {
|
|
logger.Info().
|
|
Str("component", "remote_config").
|
|
Str("action", "fetch_succeeded").
|
|
Msg("Successfully fetched remote configuration")
|
|
commandAuthorityApplied := true
|
|
if commandsEnabled != nil {
|
|
commandAuthorityApplied = applyInitialRemoteCommandAuthority(&cfg, commandsEnabled)
|
|
logger.Info().
|
|
Str("component", "remote_config").
|
|
Str("action", "apply_enable_commands").
|
|
Str("local_authority", string(cfg.CommandAuthorityProfile)).
|
|
Bool("accepted", commandAuthorityApplied).
|
|
Bool("enabled", cfg.EnableCommands).
|
|
Msg("Applied remote command execution setting")
|
|
if !commandAuthorityApplied {
|
|
logger.Warn().
|
|
Str("component", "remote_config").
|
|
Str("action", "reject_enable_commands").
|
|
Msg("Ignored remote command enablement because the local runtime is monitoring-only")
|
|
}
|
|
}
|
|
if len(settings) > 0 {
|
|
applyRemoteSettings(&cfg, settings, &logger)
|
|
}
|
|
if commandAuthorityApplied && remoteconfig.HasAppliedDesiredConfig(commandsEnabled, settings) {
|
|
metadata, metadataErr := remoteconfig.BuildDesiredConfigMetadata(commandsEnabled, settings)
|
|
if metadataErr != nil {
|
|
logger.Warn().Err(metadataErr).Msg("Failed to derive applied managed configuration fingerprint")
|
|
} else {
|
|
cfg.AppliedConfig = &agentshost.ConfigFingerprint{Version: metadata.Version, Hash: metadata.Hash}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// 3. Check if running as Windows service
|
|
ranAsService, err := runAsWindowsServiceFunc(cfg, logger)
|
|
if err != nil {
|
|
return fmt.Errorf("Windows service failed: %w", err)
|
|
}
|
|
if ranAsService {
|
|
return nil
|
|
}
|
|
|
|
g, ctx := errgroup.WithContext(ctx)
|
|
|
|
// Resolve automatic runtime selection before publishing startup and
|
|
// readiness state so the process never claims a module set it did not
|
|
// actually attempt to start.
|
|
if !cfg.EnableDocker && !cfg.DockerConfigured {
|
|
if _, err := lookPath("docker"); err == nil {
|
|
logger.Info().Msg("Auto-detected Docker binary, enabling Docker monitoring")
|
|
cfg.EnableDocker = true
|
|
} else if _, err := lookPath("podman"); err == nil {
|
|
logger.Info().Msg("Auto-detected Podman binary, enabling Docker monitoring")
|
|
cfg.EnableDocker = true
|
|
} else {
|
|
logger.Debug().Msg("Docker/Podman not found, skipping Docker monitoring")
|
|
}
|
|
}
|
|
|
|
logger.Info().
|
|
Str("version", Version).
|
|
Str("pulse_url", cfg.PulseURL).
|
|
Int("observer_destinations", len(cfg.Observers)).
|
|
Bool("host_enabled", cfg.EnableHost).
|
|
Bool("docker_enabled", cfg.EnableDocker).
|
|
Bool("kubernetes_enabled", cfg.EnableKubernetes).
|
|
Bool("proxmox_mode", cfg.EnableProxmox).
|
|
Bool("auto_update", !cfg.DisableAutoUpdate).
|
|
Msg("Starting Pulse Unified Agent")
|
|
|
|
if cfg.AllowPlaintextHTTP {
|
|
logger.Warn().
|
|
Str("pulse_url", cfg.PulseURL).
|
|
Msg("--allow-plaintext-http is set: the agent API token travels in cleartext to any non-loopback Pulse URL; only use this on a network you fully control")
|
|
}
|
|
for _, observer := range cfg.Observers {
|
|
if observer.AllowPlaintextHTTP {
|
|
logger.Warn().
|
|
Str("destination", observer.Name).
|
|
Msg("Observer destination explicitly permits plaintext HTTP; its API token may travel in cleartext")
|
|
}
|
|
}
|
|
|
|
// 5. Set prometheus info metric
|
|
agentInfo.WithLabelValues(
|
|
Version,
|
|
fmt.Sprintf("%t", cfg.EnableHost),
|
|
fmt.Sprintf("%t", cfg.EnableDocker),
|
|
fmt.Sprintf("%t", cfg.EnableKubernetes),
|
|
).Set(1)
|
|
agentUp.Set(1)
|
|
|
|
// 6. Start Health/Metrics Server
|
|
var ready atomic.Bool
|
|
runtimeStatus := newRuntimeHealth(&ready, map[string]bool{
|
|
"host": cfg.EnableHost,
|
|
"docker": cfg.EnableDocker,
|
|
"kubernetes": cfg.EnableKubernetes,
|
|
})
|
|
if cfg.HealthAddr != "" {
|
|
startHealthServer(ctx, cfg.HealthAddr, &ready, &logger, runtimeStatus)
|
|
}
|
|
|
|
// 7. Start Auto-Updater
|
|
var privilegedUpdate agentupdate.PrivilegedUpdate
|
|
helperSocket := strings.TrimSpace(os.Getenv("PULSE_AGENT_HELPER_SOCKET"))
|
|
var helperContainerInventory dockeragent.ContainerInventory
|
|
var privilegeHelperStatus *hostagent.PrivilegeHelperStatus
|
|
if helperSocket != "" {
|
|
privilegeHelperStatus = hostagent.NewPrivilegeHelperStatus()
|
|
privilegedUpdate, err = newPrivilegeHelperUpdate(helperSocket)
|
|
if err != nil {
|
|
return fmt.Errorf("configure typed privilege-helper updates: %w", err)
|
|
}
|
|
helperContainerInventory, err = dockeragent.NewPrivilegeHelperContainerInventory(helperSocket)
|
|
if err != nil {
|
|
return fmt.Errorf("configure typed privilege-helper container inventory: %w", err)
|
|
}
|
|
}
|
|
updater := newUpdater(agentupdate.Config{
|
|
PulseURL: cfg.PulseURL,
|
|
APIToken: cfg.APIToken,
|
|
AgentName: "pulse-agent",
|
|
CurrentVersion: Version,
|
|
StateDir: cfg.StateDir,
|
|
CheckInterval: 1 * time.Hour,
|
|
InsecureSkipVerify: cfg.InsecureSkipVerify,
|
|
CACertPath: cfg.CACertPath,
|
|
ServerFingerprint: cfg.ServerFingerprint,
|
|
Logger: &logger,
|
|
Disabled: cfg.DisableAutoUpdate,
|
|
PrivilegedUpdate: privilegedUpdate,
|
|
})
|
|
var pendingUpdate *agentupdate.PendingPrivilegedUpdate
|
|
if helperSocket != "" {
|
|
pendingUpdate, err = agentupdate.LoadPendingPrivilegedUpdate(cfg.StateDir)
|
|
if err != nil {
|
|
return fmt.Errorf("load pending helper update handoff: %w", err)
|
|
}
|
|
}
|
|
if pendingUpdate != nil && privilegedUpdate == nil {
|
|
return errors.New("pending helper update cannot be verified without the typed privilege helper")
|
|
}
|
|
pendingExecutableSHA256 := ""
|
|
if pendingUpdate != nil {
|
|
pendingExecutableSHA256, err = currentCollectorExecutableSHA256()
|
|
if err != nil {
|
|
return fmt.Errorf("bind pending helper update to the running collector: %w", err)
|
|
}
|
|
}
|
|
pendingReportAccepted := make(chan struct{})
|
|
var pendingReportOnce sync.Once
|
|
|
|
g.Go(func() error {
|
|
updater.RunLoop(ctx)
|
|
return nil
|
|
})
|
|
|
|
// The host module starts before the Docker module (which may only come up
|
|
// after daemon retries), so the typed container-update bridge late-binds.
|
|
dockerUpdaterBridge := &lateBoundDockerUpdater{}
|
|
|
|
// 8. Start Host Agent (if enabled)
|
|
if cfg.EnableHost {
|
|
privilegedTelemetry, err := newPrivilegeHelperTelemetry(os.Getenv("PULSE_AGENT_HELPER_SOCKET"))
|
|
if err != nil {
|
|
return fmt.Errorf("configure typed privilege helper: %w", err)
|
|
}
|
|
if privilegedTelemetry != nil {
|
|
healthCtx, cancelHealth := context.WithTimeout(ctx, 10*time.Second)
|
|
healthErr := privilegedTelemetry.Health(healthCtx)
|
|
cancelHealth()
|
|
if healthErr != nil {
|
|
return fmt.Errorf("verify typed privilege helper protocol: %w", healthErr)
|
|
}
|
|
}
|
|
hostCfg := hostagent.Config{
|
|
PulseURL: cfg.PulseURL,
|
|
APIToken: cfg.APIToken,
|
|
Interval: cfg.Interval,
|
|
HostnameOverride: cfg.HostnameOverride,
|
|
AgentID: cfg.AgentID,
|
|
AgentType: "unified",
|
|
AgentVersion: Version,
|
|
Tags: cfg.Tags,
|
|
InsecureSkipVerify: cfg.InsecureSkipVerify,
|
|
CACertPath: cfg.CACertPath,
|
|
ServerFingerprint: cfg.ServerFingerprint,
|
|
CustomSensorsFile: cfg.CustomSensorsFile,
|
|
DeploySSHUser: cfg.DeploySSHUser,
|
|
LogLevel: cfg.LogLevel,
|
|
Logger: &logger,
|
|
EnableProxmox: cfg.EnableProxmox,
|
|
ProxmoxType: cfg.ProxmoxType,
|
|
EnableCommands: cfg.EnableCommands,
|
|
CommandAuthorityProfile: cfg.CommandAuthorityProfile,
|
|
Enroll: cfg.Enroll,
|
|
DiskExclude: cfg.DiskExclude,
|
|
DiskInclude: cfg.DiskInclude,
|
|
StateDir: cfg.StateDir,
|
|
ReportIP: cfg.ReportIP,
|
|
DisableCeph: cfg.DisableCeph,
|
|
AvailabilityTargets: cfg.AvailabilityTargets,
|
|
AppliedConfig: cfg.AppliedConfig,
|
|
ModuleStatus: runtimeStatus.moduleStatuses,
|
|
PrivilegeHelperStatus: privilegeHelperStatus,
|
|
Observers: hostObserverTargets(cfg.Observers),
|
|
PrivilegedTelemetry: privilegedTelemetry,
|
|
UpdatedFromVersion: pendingUpdatePreviousVersion(pendingUpdate),
|
|
|
|
DockerContainerUpdater: dockerUpdaterBridge,
|
|
DockerContainerLifecycleOperator: dockerUpdaterBridge,
|
|
}
|
|
if pendingUpdate != nil {
|
|
hostCfg.OnPrimaryReportAccepted = func() {
|
|
pendingReportOnce.Do(func() { close(pendingReportAccepted) })
|
|
}
|
|
}
|
|
wireUpdaterHooks(&hostCfg, updater)
|
|
|
|
agent, err := newHostAgent(hostCfg)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to initialize host agent: %w", err)
|
|
}
|
|
if applier, ok := agent.(RemoteConfigApplier); ok {
|
|
remoteConfigAppliers = append(remoteConfigAppliers, applier)
|
|
}
|
|
runtimeStatus.setState("host", moduleStateRunning, nil)
|
|
|
|
g.Go(func() error {
|
|
logger.Info().Msg("Host agent module started")
|
|
return agent.Run(ctx)
|
|
})
|
|
}
|
|
|
|
// 9. Start Docker / Podman module (if enabled)
|
|
var dockerAgent RunnableCloser
|
|
if cfg.EnableDocker {
|
|
dockerCfg := dockeragent.Config{
|
|
PulseURL: cfg.PulseURL,
|
|
APIToken: cfg.APIToken,
|
|
Interval: cfg.Interval,
|
|
HostnameOverride: cfg.HostnameOverride,
|
|
AgentID: cfg.AgentID,
|
|
AgentType: "unified",
|
|
AgentVersion: Version,
|
|
InsecureSkipVerify: cfg.InsecureSkipVerify,
|
|
CACertPath: cfg.CACertPath,
|
|
ServerFingerprint: cfg.ServerFingerprint,
|
|
DisableAutoUpdate: cfg.DisableAutoUpdate,
|
|
DisableUpdateChecks: cfg.DisableDockerUpdateChecks,
|
|
DisableRegistryCredentials: cfg.DisableRegistryCredentials,
|
|
Runtime: cfg.DockerRuntime,
|
|
LogLevel: cfg.LogLevel,
|
|
Logger: &logger,
|
|
SwarmScope: "node",
|
|
IncludeContainers: true,
|
|
IncludeServices: true,
|
|
IncludeTasks: true,
|
|
CollectDiskMetrics: false,
|
|
DiskExclude: cfg.DiskExclude,
|
|
DiskInclude: cfg.DiskInclude,
|
|
Targets: dockerReportTargets(cfg),
|
|
HelperInventory: helperContainerInventory,
|
|
HelperOperationStatus: privilegeHelperStatus,
|
|
}
|
|
|
|
dockerAgent, err = newDockerAgent(dockerCfg)
|
|
if err == nil {
|
|
bindDockerActionBridge(dockerUpdaterBridge, dockerAgent)
|
|
}
|
|
if err != nil {
|
|
runtimeStatus.setState("docker", moduleStateRetrying, err)
|
|
// Docker isn't available yet - start retry loop in background
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "docker_agent").
|
|
Str("action", "initialization_failed_retry_scheduled").
|
|
Msg("Docker not available, will retry with exponential backoff")
|
|
|
|
g.Go(func() error {
|
|
agent := initDockerWithRetry(ctx, dockerCfg, &logger)
|
|
if agent != nil {
|
|
dockerAgent = agent
|
|
bindDockerActionBridge(dockerUpdaterBridge, agent)
|
|
runtimeStatus.setState("docker", moduleStateRunning, nil)
|
|
logger.Info().Msg("Docker / Podman module started (after retry)")
|
|
return agent.Run(ctx)
|
|
}
|
|
// Docker never became available, continue without it
|
|
return nil
|
|
})
|
|
} else {
|
|
runtimeStatus.setState("docker", moduleStateRunning, nil)
|
|
g.Go(func() error {
|
|
logger.Info().Msg("Docker / Podman module started")
|
|
return dockerAgent.Run(ctx)
|
|
})
|
|
}
|
|
}
|
|
|
|
// 10. Start Kubernetes Agent (if enabled)
|
|
if cfg.EnableKubernetes {
|
|
kubeCfg := kubernetesagent.Config{
|
|
PulseURL: cfg.PulseURL,
|
|
APIToken: cfg.APIToken,
|
|
Interval: cfg.Interval,
|
|
AgentID: cfg.AgentID,
|
|
AgentType: "unified",
|
|
AgentVersion: Version,
|
|
InsecureSkipVerify: cfg.InsecureSkipVerify,
|
|
CACertPath: cfg.CACertPath,
|
|
ServerFingerprint: cfg.ServerFingerprint,
|
|
LogLevel: cfg.LogLevel,
|
|
Logger: &logger,
|
|
KubeconfigPath: cfg.KubeconfigPath,
|
|
KubeContext: cfg.KubeContext,
|
|
IncludeNamespaces: cfg.KubeIncludeNamespaces,
|
|
ExcludeNamespaces: cfg.KubeExcludeNamespaces,
|
|
IncludeAllPods: cfg.KubeIncludeAllPods,
|
|
IncludeAllDeployments: cfg.KubeIncludeAllDeployments,
|
|
MaxPods: cfg.KubeMaxPods,
|
|
Targets: kubernetesReportTargets(cfg),
|
|
}
|
|
|
|
agent, err := newKubeAgent(kubeCfg)
|
|
if err != nil {
|
|
runtimeStatus.setState("kubernetes", moduleStateRetrying, err)
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "kubernetes_agent").
|
|
Str("action", "initialization_failed_retry_scheduled").
|
|
Msg("Kubernetes not available, will retry with exponential backoff")
|
|
|
|
g.Go(func() error {
|
|
retried := initKubernetesWithRetry(ctx, kubeCfg, &logger)
|
|
if retried != nil {
|
|
runtimeStatus.setState("kubernetes", moduleStateRunning, nil)
|
|
logger.Info().Msg("Kubernetes agent module started (after retry)")
|
|
return retried.Run(ctx)
|
|
}
|
|
return nil
|
|
})
|
|
} else {
|
|
runtimeStatus.setState("kubernetes", moduleStateRunning, nil)
|
|
g.Go(func() error {
|
|
logger.Info().Msg("Kubernetes agent module started")
|
|
return agent.Run(ctx)
|
|
})
|
|
}
|
|
}
|
|
|
|
if remoteConfigClient != nil && len(remoteConfigAppliers) > 0 {
|
|
g.Go(func() error {
|
|
runRemoteConfigLoop(ctx, remoteConfigClient, remoteConfigAppliers, &logger)
|
|
return nil
|
|
})
|
|
}
|
|
if pendingUpdate != nil {
|
|
g.Go(func() error {
|
|
return supervisePendingPrivilegedUpdate(
|
|
ctx,
|
|
privilegedUpdate,
|
|
pendingUpdate,
|
|
cfg.StateDir,
|
|
pendingExecutableSHA256,
|
|
ready.Load,
|
|
pendingReportAccepted,
|
|
250*time.Millisecond,
|
|
&logger,
|
|
)
|
|
})
|
|
}
|
|
|
|
// 11. Wait for all agents to exit
|
|
if err := g.Wait(); err != nil && err != context.Canceled {
|
|
logger.Error().Err(err).Msg("Agent terminated with error")
|
|
agentUp.Set(0)
|
|
cleanupDockerAgent(dockerAgent, &logger)
|
|
return fmt.Errorf("unified agent runtime failed: %w", err)
|
|
}
|
|
|
|
// 12. Cleanup
|
|
agentUp.Set(0)
|
|
cleanupDockerAgent(dockerAgent, &logger)
|
|
|
|
logger.Info().Msg("Pulse Unified Agent stopped")
|
|
return nil
|
|
}
|
|
|
|
type containerActionCapability interface {
|
|
ContainerActionsAvailable() bool
|
|
}
|
|
|
|
func bindDockerActionBridge(bridge *lateBoundDockerUpdater, agent RunnableCloser) {
|
|
if bridge == nil || agent == nil {
|
|
return
|
|
}
|
|
if capability, ok := agent.(containerActionCapability); ok && !capability.ContainerActionsAvailable() {
|
|
return
|
|
}
|
|
bridge.set(agent)
|
|
}
|
|
|
|
func configureAgentLogger(cfg Config) (zerolog.Logger, func(), error) {
|
|
zerolog.SetGlobalLevel(cfg.LogLevel)
|
|
if cfg.LogFile == "" {
|
|
logger := zerolog.New(os.Stdout).With().Timestamp().Logger()
|
|
return logger, func() {}, nil
|
|
}
|
|
|
|
logger, closer, err := pulselogging.NewStandaloneLogger(pulselogging.Config{
|
|
Format: "json",
|
|
Level: cfg.LogLevel.String(),
|
|
Component: "pulse-agent",
|
|
FilePath: cfg.LogFile,
|
|
MaxSizeMB: agentLogMaxSizeMB,
|
|
MaxAgeDays: agentLogMaxAgeDays,
|
|
Compress: true,
|
|
}, os.Stdout)
|
|
if err != nil {
|
|
return zerolog.Logger{}, func() {}, fmt.Errorf("initialize log file %q: %w", cfg.LogFile, err)
|
|
}
|
|
return logger, func() {
|
|
if closer != nil {
|
|
_ = closer.Close()
|
|
}
|
|
}, nil
|
|
}
|
|
|
|
func secureAgentStateDir(path string) error {
|
|
path = strings.TrimSpace(path)
|
|
if path == "" {
|
|
return nil
|
|
}
|
|
if err := os.MkdirAll(path, 0o700); err != nil {
|
|
return fmt.Errorf("create state directory: %w", err)
|
|
}
|
|
if err := os.Chmod(path, 0o700); err != nil {
|
|
return fmt.Errorf("chmod state directory: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// readAgentIDFile reads a persisted agent identifier from the given path.
|
|
func readAgentIDFile(path string) (string, error) {
|
|
if path == "" {
|
|
return "", nil
|
|
}
|
|
data, err := os.ReadFile(path)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
return strings.TrimSpace(string(data)), nil
|
|
}
|
|
|
|
// writeAgentIDFile persists the agent identifier atomically for future starts.
|
|
func writeAgentIDFile(path, id string) error {
|
|
if path == "" || id == "" {
|
|
return nil
|
|
}
|
|
|
|
dir := filepath.Dir(path)
|
|
if dir != "" && dir != "." {
|
|
if err := os.MkdirAll(dir, 0o700); err != nil {
|
|
return fmt.Errorf("create agent-id-file directory: %w", err)
|
|
}
|
|
}
|
|
|
|
tmp, err := os.CreateTemp(dir, ".agent-id-*")
|
|
if err != nil {
|
|
return fmt.Errorf("create agent-id-file temp: %w", err)
|
|
}
|
|
tmpPath := tmp.Name()
|
|
cleanup := func() { _ = os.Remove(tmpPath) }
|
|
|
|
if _, err := tmp.WriteString(id + "\n"); err != nil {
|
|
_ = tmp.Close()
|
|
cleanup()
|
|
return fmt.Errorf("write agent-id-file: %w", err)
|
|
}
|
|
if err := tmp.Chmod(0o600); err != nil {
|
|
_ = tmp.Close()
|
|
cleanup()
|
|
return fmt.Errorf("chmod agent-id-file: %w", err)
|
|
}
|
|
if err := tmp.Close(); err != nil {
|
|
cleanup()
|
|
return fmt.Errorf("close agent-id-file: %w", err)
|
|
}
|
|
if err := os.Rename(tmpPath, path); err != nil {
|
|
cleanup()
|
|
return fmt.Errorf("rename agent-id-file: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func applyInitialRemoteCommandAuthority(cfg *Config, desired *bool) bool {
|
|
if cfg == nil || desired == nil {
|
|
return true
|
|
}
|
|
effective, accepted := hostagent.ResolveCommandAuthority(cfg.CommandAuthorityProfile, *desired)
|
|
cfg.EnableCommands = effective
|
|
return accepted
|
|
}
|
|
|
|
func runRemoteConfigLoop(ctx context.Context, client *remoteconfig.Client, appliers []RemoteConfigApplier, logger *zerolog.Logger) {
|
|
if client == nil || len(appliers) == 0 {
|
|
return
|
|
}
|
|
ticker := time.NewTicker(remoteConfigRefreshInterval)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case <-ticker.C:
|
|
}
|
|
|
|
fetchCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
|
|
settings, commandsEnabled, err := client.Fetch(fetchCtx)
|
|
cancel()
|
|
if err != nil {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "remote_config").
|
|
Str("action", "refresh_failed").
|
|
Msg("Failed to refresh remote config")
|
|
continue
|
|
}
|
|
for _, applier := range appliers {
|
|
applier.ApplyRemoteConfig(settings, commandsEnabled)
|
|
}
|
|
logger.Debug().
|
|
Str("component", "remote_config").
|
|
Str("action", "refresh_applied").
|
|
Int("settings_count", len(settings)).
|
|
Bool("commands_enabled_present", commandsEnabled != nil).
|
|
Msg("Applied refreshed remote config")
|
|
}
|
|
}
|
|
|
|
func cleanupDockerAgent(agent RunnableCloser, logger *zerolog.Logger) {
|
|
if agent == nil || reflect.ValueOf(agent).IsNil() {
|
|
return
|
|
}
|
|
if err := agent.Close(); err != nil {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "docker_agent").
|
|
Str("action", "shutdown_failed").
|
|
Msg("Failed to close Docker / Podman module")
|
|
}
|
|
}
|
|
|
|
func healthHandler(ready *atomic.Bool, runtimes ...*runtimeHealth) http.Handler {
|
|
mux := http.NewServeMux()
|
|
var runtimeStatus *runtimeHealth
|
|
if len(runtimes) > 0 {
|
|
runtimeStatus = runtimes[0]
|
|
}
|
|
|
|
// Liveness probe - always returns 200 if server is running
|
|
mux.HandleFunc("/healthz", func(w http.ResponseWriter, r *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
w.Write([]byte("ok"))
|
|
})
|
|
|
|
// Readiness probe - returns 200 only when agents are initialized
|
|
mux.HandleFunc("/readyz", func(w http.ResponseWriter, r *http.Request) {
|
|
if runtimeStatus != nil {
|
|
snapshot := runtimeStatus.snapshot()
|
|
w.Header().Set("Content-Type", "application/json")
|
|
if !snapshot.Ready {
|
|
w.WriteHeader(http.StatusServiceUnavailable)
|
|
}
|
|
_ = json.NewEncoder(w).Encode(snapshot)
|
|
return
|
|
}
|
|
if ready.Load() {
|
|
w.WriteHeader(http.StatusOK)
|
|
w.Write([]byte("ok"))
|
|
} else {
|
|
w.WriteHeader(http.StatusServiceUnavailable)
|
|
w.Write([]byte("not ready"))
|
|
}
|
|
})
|
|
|
|
mux.HandleFunc("/status", func(w http.ResponseWriter, r *http.Request) {
|
|
w.Header().Set("Content-Type", "application/json")
|
|
if runtimeStatus == nil {
|
|
_ = json.NewEncoder(w).Encode(runtimeHealthSnapshot{Ready: ready.Load()})
|
|
return
|
|
}
|
|
_ = json.NewEncoder(w).Encode(runtimeStatus.snapshot())
|
|
})
|
|
|
|
// Prometheus metrics
|
|
mux.Handle("/metrics", promhttp.Handler())
|
|
return mux
|
|
}
|
|
|
|
func startHealthServer(ctx context.Context, addr string, ready *atomic.Bool, logger *zerolog.Logger, runtimes ...*runtimeHealth) {
|
|
srv := &http.Server{
|
|
Addr: addr,
|
|
Handler: healthHandler(ready, runtimes...),
|
|
ReadTimeout: 5 * time.Second,
|
|
WriteTimeout: 10 * time.Second,
|
|
IdleTimeout: 30 * time.Second,
|
|
}
|
|
|
|
go func() {
|
|
<-ctx.Done()
|
|
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
if err := srv.Shutdown(shutdownCtx); err != nil && err != http.ErrServerClosed {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "health_server").
|
|
Str("action", "shutdown_failed").
|
|
Str("addr", addr).
|
|
Msg("Failed to shut down health server")
|
|
}
|
|
}()
|
|
|
|
go func() {
|
|
logger.Info().
|
|
Str("component", "health_server").
|
|
Str("action", "listening").
|
|
Str("addr", addr).
|
|
Msg("Health/metrics server listening")
|
|
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
|
logger.Warn().
|
|
Err(err).
|
|
Str("component", "health_server").
|
|
Str("action", "stopped_unexpectedly").
|
|
Str("addr", addr).
|
|
Msg("Health server stopped unexpectedly")
|
|
}
|
|
}()
|
|
}
|
|
|
|
type Config struct {
|
|
PulseURL string
|
|
APIToken string
|
|
Interval time.Duration
|
|
HostnameOverride string
|
|
AgentID string
|
|
AgentIDFile string
|
|
Tags []string
|
|
InsecureSkipVerify bool
|
|
AllowPlaintextHTTP bool
|
|
CACertPath string
|
|
ServerFingerprint string
|
|
ObserversFile string
|
|
Observers []agenttarget.Observer
|
|
CustomSensorsFile string
|
|
DeploySSHUser string
|
|
LogLevel zerolog.Level
|
|
LogFile string
|
|
Logger *zerolog.Logger
|
|
|
|
// Module flags
|
|
EnableHost bool
|
|
EnableDocker bool
|
|
DockerConfigured bool
|
|
DockerExplicitlyDisabled bool
|
|
EnableKubernetes bool
|
|
EnableProxmox bool
|
|
ProxmoxType string // "pve", "pbs", or "" for auto-detect
|
|
|
|
// Auto-update
|
|
DisableAutoUpdate bool
|
|
DisableDockerUpdateChecks bool // Disable Docker image update detection
|
|
DisableRegistryCredentials bool // Do not read host Docker credentials for registry update checks
|
|
DockerRuntime string // Force Docker / Podman runtime: docker, podman, or auto
|
|
|
|
// Security
|
|
EnableCommands bool // Enable Pulse command execution for Patrol actions and Proxmox LXC Docker inventory (disabled by default)
|
|
CommandAuthorityProfile hostagent.CommandAuthorityProfile // Immutable local authority ceiling; empty/unmarked is legacy-compatible.
|
|
|
|
// Enrollment
|
|
Enroll bool // Exchange bootstrap token for runtime token on startup
|
|
|
|
// State directory
|
|
StateDir string // Persistent state directory for agent-id, proxmox registration, etc.
|
|
|
|
// Disk filtering
|
|
DiskExclude []string // Mount points or patterns to exclude from disk monitoring
|
|
DiskInclude []string // Devices or mount points to opt into monitoring despite automatic filtering
|
|
|
|
// Network configuration
|
|
ReportIP string // IP address to report (for multi-NIC systems)
|
|
DisableCeph bool // Disable local Ceph status polling
|
|
SelfTest bool // Perform self-test and exit
|
|
|
|
// AvailabilityTargets are externally probed availability checks assigned to
|
|
// this agent by the server. Remote config is the only source.
|
|
AvailabilityTargets []pulseconfig.AvailabilityTarget
|
|
|
|
// Health/metrics server
|
|
HealthAddr string
|
|
AppliedConfig *agentshost.ConfigFingerprint
|
|
|
|
// Kubernetes
|
|
KubeconfigPath string
|
|
KubeContext string
|
|
KubeIncludeNamespaces []string
|
|
KubeExcludeNamespaces []string
|
|
KubeIncludeAllPods bool
|
|
KubeIncludeAllDeployments bool
|
|
KubeMaxPods int
|
|
}
|
|
|
|
func hostObserverTargets(observers []agenttarget.Observer) []hostagent.ObserverTarget {
|
|
targets := make([]hostagent.ObserverTarget, 0, len(observers))
|
|
for _, observer := range observers {
|
|
targets = append(targets, hostagent.ObserverTarget{
|
|
Name: observer.Name,
|
|
ID: observer.ID,
|
|
PulseURL: observer.URL,
|
|
APIToken: observer.Token,
|
|
InsecureSkipVerify: observer.InsecureSkipVerify,
|
|
AllowPlaintextHTTP: observer.AllowPlaintextHTTP,
|
|
CACertPath: observer.CACertPath,
|
|
ServerFingerprint: observer.ServerFingerprint,
|
|
ProvisionProxmox: observer.ProvisionProxmox,
|
|
})
|
|
}
|
|
return targets
|
|
}
|
|
|
|
func dockerReportTargets(cfg Config) []dockeragent.TargetConfig {
|
|
targets := make([]dockeragent.TargetConfig, 0, len(cfg.Observers)+1)
|
|
targets = append(targets, dockeragent.TargetConfig{
|
|
Name: "primary",
|
|
URL: cfg.PulseURL,
|
|
Token: cfg.APIToken,
|
|
InsecureSkipVerify: cfg.InsecureSkipVerify,
|
|
AllowPlaintextHTTP: cfg.AllowPlaintextHTTP,
|
|
CACertPath: cfg.CACertPath,
|
|
ServerFingerprint: cfg.ServerFingerprint,
|
|
Authoritative: true,
|
|
})
|
|
for _, observer := range cfg.Observers {
|
|
targets = append(targets, dockeragent.TargetConfig{
|
|
Name: observer.Name,
|
|
URL: observer.URL,
|
|
Token: observer.Token,
|
|
InsecureSkipVerify: observer.InsecureSkipVerify,
|
|
AllowPlaintextHTTP: observer.AllowPlaintextHTTP,
|
|
CACertPath: observer.CACertPath,
|
|
ServerFingerprint: observer.ServerFingerprint,
|
|
})
|
|
}
|
|
return targets
|
|
}
|
|
|
|
func kubernetesReportTargets(cfg Config) []kubernetesagent.TargetConfig {
|
|
targets := make([]kubernetesagent.TargetConfig, 0, len(cfg.Observers)+1)
|
|
targets = append(targets, kubernetesagent.TargetConfig{
|
|
Name: "primary",
|
|
URL: cfg.PulseURL,
|
|
Token: cfg.APIToken,
|
|
InsecureSkipVerify: cfg.InsecureSkipVerify,
|
|
AllowPlaintextHTTP: cfg.AllowPlaintextHTTP,
|
|
CACertPath: cfg.CACertPath,
|
|
ServerFingerprint: cfg.ServerFingerprint,
|
|
Authoritative: true,
|
|
})
|
|
for _, observer := range cfg.Observers {
|
|
targets = append(targets, kubernetesagent.TargetConfig{
|
|
Name: observer.Name,
|
|
URL: observer.URL,
|
|
Token: observer.Token,
|
|
InsecureSkipVerify: observer.InsecureSkipVerify,
|
|
AllowPlaintextHTTP: observer.AllowPlaintextHTTP,
|
|
CACertPath: observer.CACertPath,
|
|
ServerFingerprint: observer.ServerFingerprint,
|
|
})
|
|
}
|
|
return targets
|
|
}
|
|
|
|
func loadConfig(args []string, getenv func(string) string) (Config, error) {
|
|
// Environment Variables
|
|
envURL := strings.TrimSpace(getenv("PULSE_URL"))
|
|
envToken := strings.TrimSpace(getenv("PULSE_TOKEN"))
|
|
envInterval := strings.TrimSpace(getenv("PULSE_INTERVAL"))
|
|
envHostname := strings.TrimSpace(getenv("PULSE_HOSTNAME"))
|
|
envAgentID := strings.TrimSpace(getenv("PULSE_AGENT_ID"))
|
|
envAgentIDFile := strings.TrimSpace(getenv("PULSE_AGENT_ID_FILE"))
|
|
envInsecure := strings.TrimSpace(getenv("PULSE_INSECURE_SKIP_VERIFY"))
|
|
envCACertPath := strings.TrimSpace(getenv("PULSE_CACERT"))
|
|
envServerFingerprint := strings.TrimSpace(getenv("PULSE_SERVER_FINGERPRINT"))
|
|
envDeploySSHUser := strings.TrimSpace(getenv("PULSE_DEPLOY_SSH_USER"))
|
|
envTags := strings.TrimSpace(getenv("PULSE_TAGS"))
|
|
envLogLevel := strings.TrimSpace(getenv("LOG_LEVEL"))
|
|
envLogFile := strings.TrimSpace(getenv("PULSE_LOG_FILE"))
|
|
envEnableHost := strings.TrimSpace(getenv("PULSE_ENABLE_HOST"))
|
|
envEnableDocker := strings.TrimSpace(getenv("PULSE_ENABLE_DOCKER"))
|
|
envEnableKubernetes := strings.TrimSpace(getenv("PULSE_ENABLE_KUBERNETES"))
|
|
envEnableProxmox := strings.TrimSpace(getenv("PULSE_ENABLE_PROXMOX"))
|
|
envProxmoxType := strings.TrimSpace(getenv("PULSE_PROXMOX_TYPE"))
|
|
envDisableAutoUpdate := strings.TrimSpace(getenv("PULSE_DISABLE_AUTO_UPDATE"))
|
|
envDisableDockerUpdateChecks := strings.TrimSpace(getenv("PULSE_DISABLE_DOCKER_UPDATE_CHECKS"))
|
|
envDisableRegistryCredentials := strings.TrimSpace(getenv("PULSE_DISABLE_REGISTRY_CREDENTIALS"))
|
|
envDockerRuntime := strings.TrimSpace(getenv("PULSE_DOCKER_RUNTIME"))
|
|
envEnableCommands := strings.TrimSpace(getenv("PULSE_ENABLE_COMMANDS"))
|
|
envCommandAuthority := strings.TrimSpace(getenv("PULSE_COMMAND_AUTHORITY"))
|
|
envDisableCommands := strings.TrimSpace(getenv("PULSE_DISABLE_COMMANDS")) // deprecated
|
|
envHealthAddr := strings.TrimSpace(getenv("PULSE_HEALTH_ADDR"))
|
|
envKubeconfig := strings.TrimSpace(getenv("PULSE_KUBECONFIG"))
|
|
envKubeContext := strings.TrimSpace(getenv("PULSE_KUBE_CONTEXT"))
|
|
envKubeIncludeNamespaces := strings.TrimSpace(getenv("PULSE_KUBE_INCLUDE_NAMESPACES"))
|
|
envKubeExcludeNamespaces := strings.TrimSpace(getenv("PULSE_KUBE_EXCLUDE_NAMESPACES"))
|
|
envKubeIncludeAllPods := strings.TrimSpace(getenv("PULSE_KUBE_INCLUDE_ALL_PODS"))
|
|
if envKubeIncludeAllPods == "" {
|
|
// Backwards compatibility for older env var name.
|
|
envKubeIncludeAllPods = strings.TrimSpace(getenv("PULSE_KUBE_INCLUDE_ALL_POD_FILES"))
|
|
}
|
|
envKubeIncludeAllDeployments := strings.TrimSpace(getenv("PULSE_KUBE_INCLUDE_ALL_DEPLOYMENTS"))
|
|
envKubeMaxPods := strings.TrimSpace(getenv("PULSE_KUBE_MAX_PODS"))
|
|
envStateDir := strings.TrimSpace(getenv("PULSE_STATE_DIR"))
|
|
envDiskExclude := strings.TrimSpace(getenv("PULSE_DISK_EXCLUDE"))
|
|
envDiskInclude := strings.TrimSpace(getenv("PULSE_DISK_INCLUDE"))
|
|
envReportIP := strings.TrimSpace(getenv("PULSE_REPORT_IP"))
|
|
envDisableCeph := strings.TrimSpace(getenv("PULSE_DISABLE_CEPH"))
|
|
envObserversFile := strings.TrimSpace(getenv("PULSE_OBSERVERS_FILE"))
|
|
envCustomSensorsFile := strings.TrimSpace(getenv("PULSE_CUSTOM_SENSORS_FILE"))
|
|
|
|
// Defaults
|
|
defaultInterval := 30 * time.Second
|
|
if envInterval != "" {
|
|
if parsed, err := time.ParseDuration(envInterval); err == nil {
|
|
defaultInterval = parsed
|
|
}
|
|
}
|
|
|
|
defaultEnableHost := true
|
|
if envEnableHost != "" {
|
|
defaultEnableHost = utils.ParseBool(envEnableHost)
|
|
}
|
|
|
|
defaultEnableDocker := false
|
|
if envEnableDocker != "" {
|
|
defaultEnableDocker = utils.ParseBool(envEnableDocker)
|
|
}
|
|
|
|
defaultEnableKubernetes := false
|
|
if envEnableKubernetes != "" {
|
|
defaultEnableKubernetes = utils.ParseBool(envEnableKubernetes)
|
|
}
|
|
|
|
defaultEnableProxmox := false
|
|
if envEnableProxmox != "" {
|
|
defaultEnableProxmox = utils.ParseBool(envEnableProxmox)
|
|
}
|
|
|
|
healthAddrDisabledByEnv := false
|
|
defaultHealthAddr := envHealthAddr
|
|
switch strings.ToLower(defaultHealthAddr) {
|
|
case "off", "none", "disabled":
|
|
defaultHealthAddr = ""
|
|
healthAddrDisabledByEnv = true
|
|
}
|
|
if defaultHealthAddr == "" && !healthAddrDisabledByEnv {
|
|
defaultHealthAddr = "127.0.0.1:9191"
|
|
}
|
|
|
|
// Flags
|
|
fs := flag.NewFlagSet("pulse-agent", flag.ContinueOnError)
|
|
urlFlag := fs.String("url", envURL, "Pulse server URL")
|
|
tokenFlag := fs.String("token", envToken, "Pulse API token (prefer --token-file for security)")
|
|
tokenFileFlag := fs.String("token-file", "", "Path to file containing Pulse API token (more secure than --token)")
|
|
intervalFlag := fs.Duration("interval", defaultInterval, "Reporting interval")
|
|
hostnameFlag := fs.String("hostname", envHostname, "Override hostname")
|
|
agentIDFlag := fs.String("agent-id", envAgentID, "Override agent identifier")
|
|
agentIDFileFlag := fs.String("agent-id-file", envAgentIDFile, "Path to a file storing the agent identifier (read on start, written on first start). Mount this file as a volume to keep the agent identity stable across container recreation.")
|
|
insecureFlag := fs.Bool("insecure", utils.ParseBool(envInsecure), "Skip TLS verification")
|
|
allowPlaintextHTTPFlag := fs.Bool("allow-plaintext-http", utils.ParseBool(strings.TrimSpace(getenv("PULSE_AGENT_ALLOW_PLAINTEXT_HTTP"))), "Allow plain HTTP to a Pulse server that does not look local (sends the API token in cleartext; only for networks you fully control)")
|
|
caCertFlag := fs.String("cacert", envCACertPath, "Path to custom CA bundle for agent HTTPS transport")
|
|
serverFingerprintFlag := fs.String("server-fingerprint", envServerFingerprint, "Expected Pulse server TLS certificate fingerprint (SHA256)")
|
|
observersFileFlag := fs.String("observers-file", envObserversFile, "Absolute path to a private JSON file defining report-only Pulse observer destinations")
|
|
customSensorsFileFlag := fs.String("custom-sensors-file", envCustomSensorsFile, "Absolute path to a private YAML file defining command or REST custom metrics")
|
|
deploySSHUserFlag := fs.String("deploy-ssh-user", envDeploySSHUser, "SSH user for peer deploy fan-out (default: root; non-root requires passwordless sudo)")
|
|
logLevelFlag := fs.String("log-level", defaultLogLevel(envLogLevel), "Log level")
|
|
logFileFlag := fs.String("log-file", envLogFile, "Write rotating JSON logs to this file")
|
|
|
|
enableHostFlag := fs.Bool("enable-host", defaultEnableHost, "Enable Host Agent module")
|
|
enableDockerFlag := fs.Bool("enable-docker", defaultEnableDocker, "Enable Docker / Podman Agent module")
|
|
enableKubernetesFlag := fs.Bool("enable-kubernetes", defaultEnableKubernetes, "Enable Kubernetes Agent module")
|
|
enableProxmoxFlag := fs.Bool("enable-proxmox", defaultEnableProxmox, "Enable Proxmox mode (creates API token, registers node)")
|
|
proxmoxTypeFlag := fs.String("proxmox-type", envProxmoxType, "Proxmox type: pve or pbs (auto-detected if not specified)")
|
|
disableAutoUpdateFlag := fs.Bool("disable-auto-update", utils.ParseBool(envDisableAutoUpdate), "Disable automatic updates")
|
|
disableDockerUpdateChecksFlag := fs.Bool("disable-docker-update-checks", utils.ParseBool(envDisableDockerUpdateChecks), "Disable Docker image update detection (avoids Docker Hub rate limits)")
|
|
disableRegistryCredentialsFlag := fs.Bool("disable-registry-credentials", utils.ParseBool(envDisableRegistryCredentials), "Do not read host Docker credentials (config.json / credential helpers) for registry update checks")
|
|
dockerRuntimeFlag := fs.String("docker-runtime", envDockerRuntime, "Docker / Podman runtime: auto, docker, or podman (default: auto)")
|
|
enableCommandsFlag := fs.Bool("enable-commands", utils.ParseBool(envEnableCommands), "Enable Pulse command execution for Patrol actions and Proxmox LXC Docker inventory (disabled by default)")
|
|
commandAuthorityFlag := fs.String("command-authority", envCommandAuthority, "Local command authority: monitoring-only, command-capable, or legacy")
|
|
disableCommandsFlag := fs.Bool("disable-commands", false, "[DEPRECATED] Commands are now disabled by default; use --enable-commands to enable")
|
|
enrollFlag := fs.Bool("enroll", false, "Exchange bootstrap token for runtime token (used by deploy wizard)")
|
|
healthAddrFlag := fs.String("health-addr", defaultHealthAddr, "Health/metrics server address (empty to disable)")
|
|
kubeconfigFlag := fs.String("kubeconfig", envKubeconfig, "Path to kubeconfig (optional; uses in-cluster config if available)")
|
|
kubeContextFlag := fs.String("kube-context", envKubeContext, "Kubeconfig context (optional)")
|
|
kubeIncludeAllPodsFlag := fs.Bool("kube-include-all-pods", utils.ParseBool(envKubeIncludeAllPods), "Include all non-succeeded pods (may be large)")
|
|
kubeIncludeAllDeploymentsFlag := fs.Bool("kube-include-all-deployments", utils.ParseBool(envKubeIncludeAllDeployments), "Include all deployments, not just problem ones")
|
|
kubeMaxPodsFlag := fs.Int("kube-max-pods", defaultInt(envKubeMaxPods, 200), "Max pods included in report")
|
|
stateDirFlag := fs.String("state-dir", envStateDir, "Persistent state directory (default: platform service state directory)")
|
|
reportIPFlag := fs.String("report-ip", envReportIP, "IP address to report (for multi-NIC systems)")
|
|
disableCephFlag := fs.Bool("disable-ceph", utils.ParseBool(envDisableCeph), "Disable local Ceph status polling")
|
|
showVersion := fs.Bool("version", false, "Print the agent version and exit")
|
|
selfTest := fs.Bool("self-test", false, "Perform self-test and exit (used during auto-update)")
|
|
|
|
var tagFlags multiValue
|
|
fs.Var(&tagFlags, "tag", "Tag to apply (repeatable)")
|
|
var kubeIncludeNamespaceFlags multiValue
|
|
fs.Var(&kubeIncludeNamespaceFlags, "kube-include-namespace", "Namespace to include (repeatable; default is all)")
|
|
var kubeExcludeNamespaceFlags multiValue
|
|
fs.Var(&kubeExcludeNamespaceFlags, "kube-exclude-namespace", "Namespace to exclude (repeatable)")
|
|
var diskExcludeFlags multiValue
|
|
fs.Var(&diskExcludeFlags, "disk-exclude", "Device name/path or mount point pattern to exclude from disk monitoring (repeatable)")
|
|
var diskIncludeFlags multiValue
|
|
fs.Var(&diskIncludeFlags, "disk-include", "Device name/path or mount point pattern to include despite automatic filesystem filtering (repeatable)")
|
|
|
|
if err := fs.Parse(args); err != nil {
|
|
return Config{}, err
|
|
}
|
|
|
|
if *showVersion {
|
|
fmt.Println(Version)
|
|
return Config{}, flag.ErrHelp
|
|
}
|
|
|
|
// Validation
|
|
pulseURL := strings.TrimSpace(*urlFlag)
|
|
if pulseURL == "" {
|
|
pulseURL = "http://localhost:7655"
|
|
}
|
|
|
|
// Record operator plaintext consent before any module validates the Pulse
|
|
// URL; the startup warning is emitted once the logger exists in run().
|
|
securityutil.SetOperatorPlaintextHTTPConsent(*allowPlaintextHTTPFlag)
|
|
|
|
// Resolve the state directory before any implicit token or identity path.
|
|
// A custom instance must never borrow the default instance's credentials.
|
|
stateDir := strings.TrimSpace(*stateDirFlag)
|
|
if stateDir == "" {
|
|
stateDir = defaultAgentStateDir()
|
|
}
|
|
agentIDFile := strings.TrimSpace(*agentIDFileFlag)
|
|
if agentIDFile == "" {
|
|
agentIDFile = filepath.Join(stateDir, "agent-id")
|
|
}
|
|
|
|
// Resolve token with priority: --token > --token-file > env > state-dir file.
|
|
token := resolveToken(*tokenFlag, *tokenFileFlag, envToken, stateDir)
|
|
observers, err := agenttarget.LoadObservers(strings.TrimSpace(*observersFileFlag), pulseURL)
|
|
if err != nil {
|
|
return Config{}, fmt.Errorf("load observer destinations: %w", err)
|
|
}
|
|
|
|
// When --enroll is set and a runtime token already exists from a previous
|
|
// enrollment, use it instead of the bootstrap token embedded in the service
|
|
// config. This ensures the agent survives process and server restarts.
|
|
if *enrollFlag {
|
|
runtimeTokenPath := filepath.Join(stateDir, "runtime.token")
|
|
if content, err := os.ReadFile(runtimeTokenPath); err == nil {
|
|
if t := strings.TrimSpace(string(content)); t != "" {
|
|
token = t
|
|
}
|
|
}
|
|
}
|
|
|
|
if token == "" && *enrollFlag && !*selfTest {
|
|
return Config{}, fmt.Errorf("Pulse API token is required for enrollment (use --token, --token-file, PULSE_TOKEN env, or %s)", defaultTokenFilePath())
|
|
}
|
|
|
|
logLevel, err := parseLogLevel(*logLevelFlag)
|
|
if err != nil {
|
|
return Config{}, fmt.Errorf("invalid log level %q: %w", strings.TrimSpace(*logLevelFlag), err)
|
|
}
|
|
interval := *intervalFlag
|
|
if interval <= 0 {
|
|
return Config{}, fmt.Errorf("interval must be greater than 0 (got %s)", interval)
|
|
}
|
|
kubeMaxPods := *kubeMaxPodsFlag
|
|
if kubeMaxPods <= 0 {
|
|
return Config{}, fmt.Errorf("kube-max-pods must be greater than 0 (got %d)", kubeMaxPods)
|
|
}
|
|
dockerRuntime, err := normalizeDockerRuntime(*dockerRuntimeFlag)
|
|
if err != nil {
|
|
return Config{}, err
|
|
}
|
|
deploySSHUser, err := hostagent.NormalizeDeploySSHUser(*deploySSHUserFlag)
|
|
if err != nil {
|
|
return Config{}, err
|
|
}
|
|
enableCommands := resolveEnableCommands(*enableCommandsFlag, *disableCommandsFlag, envEnableCommands, envDisableCommands)
|
|
commandAuthorityRaw := strings.TrimSpace(*commandAuthorityFlag)
|
|
if commandAuthorityRaw == "" && enableCommands {
|
|
commandAuthorityRaw = string(hostagent.CommandAuthorityCommandCapable)
|
|
}
|
|
commandAuthorityProfile, err := hostagent.NormalizeCommandAuthorityProfile(commandAuthorityRaw)
|
|
if err != nil {
|
|
return Config{}, err
|
|
}
|
|
if commandAuthorityProfile == hostagent.CommandAuthorityMonitoringOnly && enableCommands {
|
|
return Config{}, fmt.Errorf("--enable-commands conflicts with --command-authority monitoring-only")
|
|
}
|
|
|
|
tags := gatherTags(envTags, tagFlags)
|
|
kubeIncludeNamespaces := gatherCSV(envKubeIncludeNamespaces, kubeIncludeNamespaceFlags)
|
|
kubeExcludeNamespaces := gatherCSV(envKubeExcludeNamespaces, kubeExcludeNamespaceFlags)
|
|
diskExclude := gatherCSV(envDiskExclude, diskExcludeFlags)
|
|
diskInclude := gatherCSV(envDiskInclude, diskIncludeFlags)
|
|
|
|
// Check if Docker was explicitly configured via fs or env. An explicit
|
|
// local disable is a privacy boundary and cannot be reversed by remote
|
|
// profile config or auto-detection.
|
|
dockerConfigured := envEnableDocker != ""
|
|
dockerExplicitlyDisabled := envEnableDocker != "" && !utils.ParseBool(envEnableDocker)
|
|
if !dockerConfigured {
|
|
fs.Visit(func(f *flag.Flag) {
|
|
if f.Name == "enable-docker" {
|
|
dockerConfigured = true
|
|
dockerExplicitlyDisabled = !*enableDockerFlag
|
|
}
|
|
})
|
|
} else {
|
|
fs.Visit(func(f *flag.Flag) {
|
|
if f.Name == "enable-docker" {
|
|
dockerExplicitlyDisabled = !*enableDockerFlag
|
|
}
|
|
})
|
|
}
|
|
|
|
return Config{
|
|
PulseURL: pulseURL,
|
|
APIToken: token,
|
|
Interval: interval,
|
|
HostnameOverride: strings.TrimSpace(*hostnameFlag),
|
|
AgentID: strings.TrimSpace(*agentIDFlag),
|
|
AgentIDFile: agentIDFile,
|
|
Tags: tags,
|
|
InsecureSkipVerify: *insecureFlag,
|
|
AllowPlaintextHTTP: *allowPlaintextHTTPFlag,
|
|
CACertPath: strings.TrimSpace(*caCertFlag),
|
|
ServerFingerprint: strings.TrimSpace(*serverFingerprintFlag),
|
|
ObserversFile: strings.TrimSpace(*observersFileFlag),
|
|
Observers: observers,
|
|
CustomSensorsFile: strings.TrimSpace(*customSensorsFileFlag),
|
|
DeploySSHUser: deploySSHUser,
|
|
LogLevel: logLevel,
|
|
LogFile: strings.TrimSpace(*logFileFlag),
|
|
EnableHost: *enableHostFlag,
|
|
EnableDocker: *enableDockerFlag,
|
|
DockerConfigured: dockerConfigured,
|
|
DockerExplicitlyDisabled: dockerExplicitlyDisabled,
|
|
EnableKubernetes: *enableKubernetesFlag,
|
|
EnableProxmox: *enableProxmoxFlag,
|
|
ProxmoxType: strings.TrimSpace(*proxmoxTypeFlag),
|
|
DisableAutoUpdate: *disableAutoUpdateFlag,
|
|
DisableDockerUpdateChecks: *disableDockerUpdateChecksFlag,
|
|
DisableRegistryCredentials: *disableRegistryCredentialsFlag,
|
|
DockerRuntime: dockerRuntime,
|
|
EnableCommands: enableCommands,
|
|
CommandAuthorityProfile: commandAuthorityProfile,
|
|
Enroll: *enrollFlag,
|
|
HealthAddr: strings.TrimSpace(*healthAddrFlag),
|
|
KubeconfigPath: strings.TrimSpace(*kubeconfigFlag),
|
|
KubeContext: strings.TrimSpace(*kubeContextFlag),
|
|
KubeIncludeNamespaces: kubeIncludeNamespaces,
|
|
KubeExcludeNamespaces: kubeExcludeNamespaces,
|
|
KubeIncludeAllPods: *kubeIncludeAllPodsFlag,
|
|
KubeIncludeAllDeployments: *kubeIncludeAllDeploymentsFlag,
|
|
KubeMaxPods: kubeMaxPods,
|
|
StateDir: stateDir,
|
|
DiskExclude: diskExclude,
|
|
DiskInclude: diskInclude,
|
|
ReportIP: strings.TrimSpace(*reportIPFlag),
|
|
DisableCeph: *disableCephFlag,
|
|
SelfTest: *selfTest,
|
|
}, nil
|
|
}
|
|
|
|
func gatherTags(env string, flags []string) []string {
|
|
tags := make([]string, 0)
|
|
if env != "" {
|
|
for _, tag := range strings.Split(env, ",") {
|
|
tag = strings.TrimSpace(tag)
|
|
if tag != "" {
|
|
tags = append(tags, tag)
|
|
}
|
|
}
|
|
}
|
|
for _, tag := range flags {
|
|
tag = strings.TrimSpace(tag)
|
|
if tag != "" {
|
|
tags = append(tags, tag)
|
|
}
|
|
}
|
|
return tags
|
|
}
|
|
|
|
func gatherCSV(env string, flags []string) []string {
|
|
values := make([]string, 0)
|
|
if env != "" {
|
|
for _, value := range strings.Split(env, ",") {
|
|
value = strings.TrimSpace(value)
|
|
if value != "" {
|
|
values = append(values, value)
|
|
}
|
|
}
|
|
}
|
|
for _, value := range flags {
|
|
value = strings.TrimSpace(value)
|
|
if value != "" {
|
|
values = append(values, value)
|
|
}
|
|
}
|
|
return values
|
|
}
|
|
|
|
func defaultInt(value string, fallback int) int {
|
|
value = strings.TrimSpace(value)
|
|
if value == "" {
|
|
return fallback
|
|
}
|
|
parsed, err := strconv.Atoi(value)
|
|
if err != nil {
|
|
return fallback
|
|
}
|
|
return parsed
|
|
}
|
|
|
|
func normalizeDockerRuntime(value string) (string, error) {
|
|
runtime := strings.ToLower(strings.TrimSpace(value))
|
|
switch runtime {
|
|
case "", "auto", "default":
|
|
return "", nil
|
|
case "docker", "podman":
|
|
return runtime, nil
|
|
default:
|
|
return "", fmt.Errorf("invalid docker runtime %q: must be auto, docker, or podman", value)
|
|
}
|
|
}
|
|
|
|
func parseLogLevel(value string) (zerolog.Level, error) {
|
|
normalized := strings.ToLower(strings.TrimSpace(value))
|
|
if normalized == "" {
|
|
return zerolog.InfoLevel, nil
|
|
}
|
|
return zerolog.ParseLevel(normalized)
|
|
}
|
|
|
|
func defaultLogLevel(envValue string) string {
|
|
if strings.TrimSpace(envValue) == "" {
|
|
return "info"
|
|
}
|
|
return envValue
|
|
}
|
|
|
|
// resolveEnableCommands determines whether command execution should be enabled.
|
|
// Priority: --enable-commands > --disable-commands (deprecated) > PULSE_ENABLE_COMMANDS > PULSE_DISABLE_COMMANDS (deprecated)
|
|
// Default: disabled (false) for security
|
|
func resolveEnableCommands(enableFlag, disableFlag bool, envEnable, envDisable string) bool {
|
|
// If --enable-commands is explicitly set, use it
|
|
if enableFlag {
|
|
return true
|
|
}
|
|
|
|
// Backwards compat: if --disable-commands was used, log deprecation but respect it
|
|
// (disableFlag being true means commands should be disabled, which is already the default)
|
|
if disableFlag {
|
|
fmt.Fprintln(os.Stderr, "warning: --disable-commands is deprecated and no longer needed (commands are disabled by default). Use --enable-commands to enable.")
|
|
return false
|
|
}
|
|
|
|
// Check environment variables
|
|
if envEnable != "" {
|
|
return utils.ParseBool(envEnable)
|
|
}
|
|
|
|
// Backwards compat: PULSE_DISABLE_COMMANDS=true means commands disabled (already default)
|
|
// PULSE_DISABLE_COMMANDS=false means commands enabled (backwards compat)
|
|
if envDisable != "" {
|
|
fmt.Fprintln(os.Stderr, "warning: PULSE_DISABLE_COMMANDS is deprecated. Use PULSE_ENABLE_COMMANDS=true to enable commands.")
|
|
// Invert: DISABLE=false means enable
|
|
return !utils.ParseBool(envDisable)
|
|
}
|
|
|
|
// Default: commands disabled
|
|
return false
|
|
}
|
|
|
|
// resolveToken resolves the API token with priority:
|
|
// 1. --token flag (direct value)
|
|
// 2. --token-file flag (read from file)
|
|
// 3. PULSE_TOKEN environment variable
|
|
// 4. Token file under the resolved state directory
|
|
//
|
|
// Reading from a file is more secure than CLI args as tokens won't appear in `ps` output.
|
|
func resolveToken(tokenFlag, tokenFileFlag, envToken, stateDir string) string {
|
|
return resolveTokenInternal(tokenFlag, tokenFileFlag, envToken, stateDir, os.ReadFile)
|
|
}
|
|
|
|
// defaultAgentStateDir mirrors where each platform's installer keeps agent
|
|
// state: %ProgramData%\Pulse on Windows (see scripts/install.ps1), the
|
|
// systemd/launchd convention /var/lib/pulse-agent everywhere else.
|
|
func defaultAgentStateDir() string {
|
|
if runtime.GOOS == "windows" {
|
|
if pd := strings.TrimSpace(os.Getenv("ProgramData")); pd != "" {
|
|
return filepath.Join(pd, "Pulse")
|
|
}
|
|
return `C:\ProgramData\Pulse`
|
|
}
|
|
return "/var/lib/pulse-agent"
|
|
}
|
|
|
|
func defaultTokenFilePath() string {
|
|
return filepath.Join(defaultAgentStateDir(), "token")
|
|
}
|
|
|
|
func resolveTokenInternal(tokenFlag, tokenFileFlag, envToken, stateDir string, readFile func(string) ([]byte, error)) string {
|
|
// 1. Direct token from --token flag
|
|
if t := strings.TrimSpace(tokenFlag); t != "" {
|
|
return t
|
|
}
|
|
|
|
// 2. Token from --token-file flag
|
|
if tokenFileFlag != "" {
|
|
if content, err := readFile(tokenFileFlag); err == nil {
|
|
if t := strings.TrimSpace(string(content)); t != "" {
|
|
return t
|
|
}
|
|
}
|
|
}
|
|
|
|
// 3. PULSE_TOKEN environment variable
|
|
if t := strings.TrimSpace(envToken); t != "" {
|
|
return t
|
|
}
|
|
|
|
// 4. Token file in the resolved state directory. When stateDir is custom,
|
|
// do not fall through to the default instance's token.
|
|
tokenFile := filepath.Join(strings.TrimSpace(stateDir), "token")
|
|
if strings.TrimSpace(stateDir) == "" {
|
|
tokenFile = defaultTokenFilePath()
|
|
}
|
|
if content, err := readFile(tokenFile); err == nil {
|
|
if t := strings.TrimSpace(string(content)); t != "" {
|
|
return t
|
|
}
|
|
}
|
|
|
|
return ""
|
|
}
|
|
|
|
// initModuleWithRetry attempts to initialize an agent module with exponential
|
|
// backoff. It returns the module when its backend becomes available, or a zero
|
|
// value if the context is cancelled. component and displayName feed the
|
|
// structured log fields; unavailableMsg is the per-module retry log line.
|
|
// Retry intervals: 5s, 10s, 20s, 40s, 80s, 160s, then cap at 5 minutes.
|
|
func initModuleWithRetry[T any](ctx context.Context, logger *zerolog.Logger, component, displayName, unavailableMsg string, connect func() (T, error)) T {
|
|
const multiplier = 2.0
|
|
|
|
var zero T
|
|
delay := retryInitialDelay
|
|
attempt := 0
|
|
|
|
for {
|
|
if err := ctx.Err(); err != nil {
|
|
logger.Info().Msg(displayName + " retry cancelled, context done")
|
|
return zero
|
|
}
|
|
|
|
module, err := connect()
|
|
if err == nil {
|
|
logger.Info().
|
|
Str("component", component).
|
|
Str("action", "retry_connect_succeeded").
|
|
Int("attempts", attempt+1).
|
|
Msg("Successfully connected to " + displayName + " after retry")
|
|
return module
|
|
}
|
|
|
|
attempt++
|
|
retryLogEvent(logger, attempt).
|
|
Err(err).
|
|
Str("component", component).
|
|
Str("action", "retry_connect_failed").
|
|
Int("attempt", attempt).
|
|
Str("next_retry", delay.String()).
|
|
Msg(unavailableMsg)
|
|
|
|
if !waitForRetryDelay(ctx, delay) {
|
|
logger.Info().
|
|
Str("component", component).
|
|
Str("action", "retry_cancelled").
|
|
Msg(displayName + " retry cancelled, context done")
|
|
return zero
|
|
}
|
|
|
|
// Calculate next delay with exponential backoff, capped at retryMaxDelay
|
|
delay = time.Duration(float64(delay) * multiplier)
|
|
if delay > retryMaxDelay {
|
|
delay = retryMaxDelay
|
|
}
|
|
}
|
|
}
|
|
|
|
// initDockerWithRetry attempts to initialize the Docker / Podman collection module with exponential backoff.
|
|
// It returns the module when Docker / Podman becomes available, or nil if the context is cancelled.
|
|
// lateBoundDockerUpdater satisfies hostagent.DockerContainerUpdater while the
|
|
// Docker module comes up (or never does). The host command client holds this
|
|
// bridge for the process lifetime; set installs the module's implementation.
|
|
type lateBoundDockerUpdater struct {
|
|
mu sync.RWMutex
|
|
updater hostagent.DockerContainerUpdater
|
|
lifecycle hostagent.DockerContainerLifecycleOperator
|
|
}
|
|
|
|
func (b *lateBoundDockerUpdater) set(candidate any) {
|
|
updater, ok := candidate.(hostagent.DockerContainerUpdater)
|
|
if !ok {
|
|
fmt.Fprintf(os.Stderr, "docker update bridge: %T does not implement the typed container updater\n", candidate)
|
|
return
|
|
}
|
|
lifecycle, ok := candidate.(hostagent.DockerContainerLifecycleOperator)
|
|
if !ok {
|
|
fmt.Fprintf(os.Stderr, "docker lifecycle bridge: %T does not implement the typed container lifecycle operator\n", candidate)
|
|
return
|
|
}
|
|
b.mu.Lock()
|
|
b.updater = updater
|
|
b.lifecycle = lifecycle
|
|
b.mu.Unlock()
|
|
}
|
|
|
|
func (b *lateBoundDockerUpdater) InspectDockerContainerLifecycle(ctx context.Context, runtime, containerID string) (agentexec.DockerContainerLifecycleSnapshot, error) {
|
|
b.mu.RLock()
|
|
lifecycle := b.lifecycle
|
|
b.mu.RUnlock()
|
|
if lifecycle == nil {
|
|
return agentexec.DockerContainerLifecycleSnapshot{}, fmt.Errorf("docker module is not running on this agent")
|
|
}
|
|
return lifecycle.InspectDockerContainerLifecycle(ctx, runtime, containerID)
|
|
}
|
|
|
|
func (b *lateBoundDockerUpdater) MutateDockerContainerLifecycle(ctx context.Context, runtime, operation, containerID string) error {
|
|
b.mu.RLock()
|
|
lifecycle := b.lifecycle
|
|
b.mu.RUnlock()
|
|
if lifecycle == nil {
|
|
return fmt.Errorf("docker module is not running on this agent")
|
|
}
|
|
return lifecycle.MutateDockerContainerLifecycle(ctx, runtime, operation, containerID)
|
|
}
|
|
|
|
func (b *lateBoundDockerUpdater) TypedContainerUpdate(ctx context.Context, runtime, containerID, expectedImageDigest string, progress func(string)) (agentexec.DockerContainerUpdateOutcome, error) {
|
|
b.mu.RLock()
|
|
updater := b.updater
|
|
b.mu.RUnlock()
|
|
if updater == nil {
|
|
return agentexec.DockerContainerUpdateOutcome{}, fmt.Errorf("docker module is not running on this agent")
|
|
}
|
|
return updater.TypedContainerUpdate(ctx, runtime, containerID, expectedImageDigest, progress)
|
|
}
|
|
|
|
func (b *lateBoundDockerUpdater) TypedContainerUpdatePreflight(ctx context.Context, runtime, containerID, expectedImageDigest string) error {
|
|
b.mu.RLock()
|
|
updater := b.updater
|
|
b.mu.RUnlock()
|
|
if updater == nil {
|
|
return agentexec.NewActionPreflightError(agentexec.ActionRefusalCapabilityUnavailable, fmt.Errorf("docker module is not running on this agent"))
|
|
}
|
|
return updater.TypedContainerUpdatePreflight(ctx, runtime, containerID, expectedImageDigest)
|
|
}
|
|
|
|
func initDockerWithRetry(ctx context.Context, cfg dockeragent.Config, logger *zerolog.Logger) RunnableCloser {
|
|
return initModuleWithRetry(ctx, logger, "docker_agent", "Docker", "Docker not available, will retry", func() (RunnableCloser, error) {
|
|
return newDockerAgent(cfg)
|
|
})
|
|
}
|
|
|
|
// initKubernetesWithRetry attempts to initialize the Kubernetes agent with exponential backoff.
|
|
// It returns the agent when Kubernetes becomes available, or nil if the context is cancelled.
|
|
func initKubernetesWithRetry(ctx context.Context, cfg kubernetesagent.Config, logger *zerolog.Logger) Runnable {
|
|
return initModuleWithRetry(ctx, logger, "kubernetes_agent", "Kubernetes", "Kubernetes still not available, will retry", func() (Runnable, error) {
|
|
return newKubeAgent(cfg)
|
|
})
|
|
}
|
|
|
|
func waitForRetryDelay(ctx context.Context, delay time.Duration) bool {
|
|
timer := time.NewTimer(delay)
|
|
defer func() {
|
|
if !timer.Stop() {
|
|
select {
|
|
case <-timer.C:
|
|
default:
|
|
}
|
|
}
|
|
}()
|
|
|
|
select {
|
|
case <-ctx.Done():
|
|
return false
|
|
case <-timer.C:
|
|
return true
|
|
}
|
|
}
|
|
|
|
// retryLogEvent returns a zerolog event at a level that decreases with attempt count
|
|
// to avoid flooding logs on misconfigured systems with unbounded retries.
|
|
// - Attempts 1-10: Warn (initial visibility)
|
|
// - Attempts 11-50: Info (still visible, less noisy)
|
|
// - Attempts 51+: Debug (effectively silent unless debug logging enabled)
|
|
func retryLogEvent(logger *zerolog.Logger, attempt int) *zerolog.Event {
|
|
switch {
|
|
case attempt <= 10:
|
|
return logger.Warn()
|
|
case attempt <= 50:
|
|
return logger.Info()
|
|
default:
|
|
return logger.Debug()
|
|
}
|
|
}
|
|
|
|
// applyRemoteSettings merges remote settings into the local configuration.
|
|
// Supported keys:
|
|
// - enable_host (bool)
|
|
// - enable_docker (bool)
|
|
// - enable_kubernetes (bool)
|
|
// - enable_proxmox (bool)
|
|
// - proxmox_type (string)
|
|
// - docker_runtime (string)
|
|
// - disable_auto_update (bool)
|
|
// - disable_docker_update_checks (bool)
|
|
// - kube_include_all_pods (bool)
|
|
// - kube_include_all_deployments (bool)
|
|
// - log_level (string)
|
|
// - interval (string/duration)
|
|
// - report_ip (string)
|
|
// - disable_ceph (bool)
|
|
// - availabilityTargets (array of assigned availability checks)
|
|
func applyRemoteSettings(cfg *Config, settings map[string]interface{}, logger *zerolog.Logger) {
|
|
for k, v := range settings {
|
|
switch k {
|
|
case "enable_host":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.EnableHost = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: enable_host")
|
|
}
|
|
case "enable_docker":
|
|
if b, ok := v.(bool); ok {
|
|
if b && cfg.DockerExplicitlyDisabled {
|
|
cfg.DockerConfigured = true
|
|
logger.Info().Msg("Remote config: enable_docker ignored because Docker / Podman monitoring is locally disabled")
|
|
continue
|
|
}
|
|
cfg.EnableDocker = b
|
|
cfg.DockerConfigured = true
|
|
logger.Info().Bool("val", b).Msg("Remote config: enable_docker")
|
|
}
|
|
case "enable_kubernetes":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.EnableKubernetes = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: enable_kubernetes")
|
|
}
|
|
case "enable_proxmox":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.EnableProxmox = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: enable_proxmox")
|
|
}
|
|
case "proxmox_type":
|
|
if s, ok := v.(string); ok {
|
|
normalized := strings.TrimSpace(strings.ToLower(s))
|
|
if normalized == "auto" {
|
|
normalized = ""
|
|
}
|
|
cfg.ProxmoxType = normalized
|
|
logger.Info().Str("val", s).Msg("Remote config: proxmox_type")
|
|
}
|
|
case "docker_runtime":
|
|
if s, ok := v.(string); ok {
|
|
runtime, err := normalizeDockerRuntime(s)
|
|
if err != nil {
|
|
logger.Warn().Str("val", s).Msg("Remote config: ignoring invalid docker_runtime value")
|
|
continue
|
|
}
|
|
cfg.DockerRuntime = runtime
|
|
logger.Info().Str("val", s).Msg("Remote config: docker_runtime")
|
|
}
|
|
case "log_level":
|
|
if s, ok := v.(string); ok {
|
|
if l, err := zerolog.ParseLevel(s); err == nil {
|
|
cfg.LogLevel = l
|
|
zerolog.SetGlobalLevel(l)
|
|
logger.Info().Str("val", s).Msg("Remote config: log_level")
|
|
}
|
|
}
|
|
case "interval":
|
|
if s, ok := v.(string); ok {
|
|
if d, err := time.ParseDuration(s); err == nil && d > 0 {
|
|
cfg.Interval = d
|
|
logger.Info().Str("val", s).Msg("Remote config: interval")
|
|
} else {
|
|
logger.Warn().Str("val", s).Msg("Remote config: ignoring invalid interval value")
|
|
}
|
|
} else if f, ok := v.(float64); ok {
|
|
if math.IsNaN(f) || math.IsInf(f, 0) || f <= 0 {
|
|
logger.Warn().Float64("val", f).Msg("Remote config: ignoring invalid interval value")
|
|
continue
|
|
}
|
|
// JSON numbers are floats, assume seconds.
|
|
cfg.Interval = time.Duration(f * float64(time.Second))
|
|
logger.Info().Float64("val", f).Msg("Remote config: interval (s)")
|
|
}
|
|
case "disable_auto_update":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.DisableAutoUpdate = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: disable_auto_update")
|
|
}
|
|
case "disable_docker_update_checks":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.DisableDockerUpdateChecks = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: disable_docker_update_checks")
|
|
}
|
|
case "kube_include_all_pods":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.KubeIncludeAllPods = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: kube_include_all_pods")
|
|
}
|
|
case "kube_include_all_deployments":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.KubeIncludeAllDeployments = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: kube_include_all_deployments")
|
|
}
|
|
case "report_ip":
|
|
if s, ok := v.(string); ok {
|
|
cfg.ReportIP = s
|
|
logger.Info().Str("val", s).Msg("Remote config: report_ip")
|
|
}
|
|
case "disable_ceph":
|
|
if b, ok := v.(bool); ok {
|
|
cfg.DisableCeph = b
|
|
logger.Info().Bool("val", b).Msg("Remote config: disable_ceph")
|
|
}
|
|
case "availabilityTargets":
|
|
targets, err := hostagent.AvailabilityTargetsFromSetting(v)
|
|
if err != nil {
|
|
logger.Warn().Err(err).Msg("Remote config: ignoring unreadable availabilityTargets value")
|
|
continue
|
|
}
|
|
cfg.AvailabilityTargets = targets
|
|
logger.Info().Int("count", len(targets)).Msg("Remote config: availabilityTargets")
|
|
}
|
|
}
|
|
}
|
|
|
|
func remoteDurationSetting(settings map[string]interface{}, key string) (time.Duration, bool) {
|
|
value, ok := settings[key]
|
|
if !ok {
|
|
return 0, false
|
|
}
|
|
|
|
switch typed := value.(type) {
|
|
case string:
|
|
parsed, err := time.ParseDuration(typed)
|
|
if err != nil {
|
|
return 0, false
|
|
}
|
|
return parsed, true
|
|
case float64:
|
|
return time.Duration(typed * float64(time.Second)), true
|
|
case int:
|
|
return time.Duration(typed) * time.Second, true
|
|
case int64:
|
|
return time.Duration(typed) * time.Second, true
|
|
default:
|
|
return 0, false
|
|
}
|
|
}
|