2069 lines
60 KiB
Go
2069 lines
60 KiB
Go
package cli
|
|
|
|
import (
|
|
"context"
|
|
"encoding/csv"
|
|
"encoding/json"
|
|
"flag"
|
|
"fmt"
|
|
"io"
|
|
"math"
|
|
"os"
|
|
"os/exec"
|
|
"path/filepath"
|
|
"sort"
|
|
"strconv"
|
|
"strings"
|
|
"text/tabwriter"
|
|
"time"
|
|
|
|
"pxmon/internal/cluster"
|
|
"pxmon/internal/history"
|
|
)
|
|
|
|
func (a *App) runExplain(args []string) int {
|
|
fs := flag.NewFlagSet("explain", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
find := fs.String("find", "", "Fast filter text")
|
|
if err := fs.Parse(args); err != nil {
|
|
return 2
|
|
}
|
|
query := strings.ToLower(strings.TrimSpace(*find))
|
|
if query == "" && fs.NArg() > 0 {
|
|
query = strings.ToLower(strings.TrimSpace(strings.Join(fs.Args(), " ")))
|
|
}
|
|
catalog := []string{
|
|
"cluster list",
|
|
"cluster show [name]",
|
|
"cluster usage [name] --range 1h|1d|1mo|all",
|
|
"cluster traffic [name] --range 1h|1d|1mo|all",
|
|
"cluster graph [name] --range 1d|1mo|all",
|
|
"cluster p95 [name] --iface eth0 --range 30d --graph",
|
|
"cluster tag add|rm|ls [name] --tags prod,billing",
|
|
"cluster kvm-tag add|rm|ls [name] --vm vm123 --tags critical",
|
|
"cluster alert show|set [name]",
|
|
"cluster alert-vm show|set|check [name]",
|
|
"cluster drift [name]",
|
|
"cluster report export --format json|csv --out ./report.json",
|
|
"cluster backup target add|ls|rm",
|
|
"cluster backup plan add|ls|rm",
|
|
"cluster backup run <plan>",
|
|
"cluster runbook list|show|run <id>",
|
|
"cluster schedule add|ls|rm|run-due",
|
|
"cluster change-history --tail 100",
|
|
"cluster agent status",
|
|
"cluster agent update [name] --restart-bot=true",
|
|
}
|
|
if query == "" {
|
|
for _, line := range catalog {
|
|
fmt.Fprintln(a.out, line)
|
|
}
|
|
return 0
|
|
}
|
|
matched := 0
|
|
for _, line := range catalog {
|
|
if strings.Contains(strings.ToLower(line), query) {
|
|
fmt.Fprintln(a.out, line)
|
|
matched++
|
|
}
|
|
}
|
|
if matched == 0 {
|
|
fmt.Fprintf(a.out, "No explain matches for %q\n", query)
|
|
return 1
|
|
}
|
|
return 0
|
|
}
|
|
|
|
func parseTagsCSV(v string) []string {
|
|
parts := strings.Split(v, ",")
|
|
out := make([]string, 0, len(parts))
|
|
for _, p := range parts {
|
|
p = strings.TrimSpace(p)
|
|
if p == "" {
|
|
continue
|
|
}
|
|
out = append(out, p)
|
|
}
|
|
return out
|
|
}
|
|
|
|
func parseSinceRange(raw string) (time.Time, string, error) {
|
|
r := strings.ToLower(strings.TrimSpace(raw))
|
|
now := time.Now().UTC()
|
|
switch r {
|
|
case "", "30d", "1mo", "month":
|
|
return now.Add(-30 * 24 * time.Hour), "30d", nil
|
|
case "1d", "24h", "day":
|
|
return now.Add(-24 * time.Hour), "1d", nil
|
|
case "7d", "week":
|
|
return now.Add(-7 * 24 * time.Hour), "7d", nil
|
|
case "all":
|
|
return time.Time{}, "all", nil
|
|
default:
|
|
d, err := time.ParseDuration(r)
|
|
if err != nil {
|
|
return time.Time{}, "", fmt.Errorf("invalid range %q (use 1d|7d|30d|all)", raw)
|
|
}
|
|
return now.Add(-d), r, nil
|
|
}
|
|
}
|
|
|
|
type multiStringFlag []string
|
|
|
|
func (m *multiStringFlag) String() string {
|
|
return strings.Join(*m, ",")
|
|
}
|
|
|
|
func (m *multiStringFlag) Set(v string) error {
|
|
*m = append(*m, strings.TrimSpace(v))
|
|
return nil
|
|
}
|
|
|
|
func (a *App) runClusterTag(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster tag <add|rm|ls> [name-or-id]")
|
|
return 2
|
|
}
|
|
sub := strings.ToLower(strings.TrimSpace(args[0]))
|
|
switch sub {
|
|
case "ls", "list", "show":
|
|
selector := ""
|
|
if len(args) > 1 {
|
|
selector = args[1]
|
|
}
|
|
if strings.TrimSpace(selector) != "" {
|
|
c, err := svc.Get(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster tag ls: %v\n", err)
|
|
return 1
|
|
}
|
|
tags := c.Tags
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"cluster": c.Name, "tags": tags})
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Cluster: %s\n", c.Name)
|
|
if len(tags) == 0 {
|
|
fmt.Fprintln(a.out, "Tags: (none)")
|
|
} else {
|
|
fmt.Fprintf(a.out, "Tags: %s\n", strings.Join(tags, ", "))
|
|
}
|
|
return 0
|
|
}
|
|
clusters, _, err := svc.List()
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster tag ls: %v\n", err)
|
|
return 1
|
|
}
|
|
type row struct {
|
|
Cluster string `json:"cluster"`
|
|
Tags []string `json:"tags"`
|
|
}
|
|
rows := make([]row, 0, len(clusters))
|
|
for _, c := range clusters {
|
|
rows = append(rows, row{Cluster: c.Name, Tags: c.Tags})
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, rows)
|
|
return 0
|
|
}
|
|
tw := tabwriter.NewWriter(a.out, 0, 2, 2, ' ', 0)
|
|
fmt.Fprintln(tw, "CLUSTER\tTAGS")
|
|
for _, r := range rows {
|
|
v := "(none)"
|
|
if len(r.Tags) > 0 {
|
|
v = strings.Join(r.Tags, ", ")
|
|
}
|
|
fmt.Fprintf(tw, "%s\t%s\n", r.Cluster, v)
|
|
}
|
|
_ = tw.Flush()
|
|
return 0
|
|
case "add", "rm", "remove", "del":
|
|
fs := flag.NewFlagSet("cluster tag "+sub, flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
tagsCSV := fs.String("tags", "", "Comma separated tags")
|
|
selector, parseArgs := splitLeadingSelector(args[1:])
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 {
|
|
if selector != "" {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster tag add|rm [name-or-id] --tags a,b")
|
|
return 2
|
|
}
|
|
selector = fs.Arg(0)
|
|
}
|
|
tags := parseTagsCSV(*tagsCSV)
|
|
if len(tags) == 0 {
|
|
tags = fs.Args()
|
|
}
|
|
if len(tags) == 0 {
|
|
fmt.Fprintln(a.err, "cluster tag: provide tags with --tags a,b")
|
|
return 2
|
|
}
|
|
var (
|
|
c cluster.Cluster
|
|
err error
|
|
)
|
|
if sub == "add" {
|
|
c, err = svc.AddClusterTags(selector, tags)
|
|
} else {
|
|
c, err = svc.RemoveClusterTags(selector, tags)
|
|
}
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster tag %s: %v\n", sub, err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"cluster": c.Name, "tags": c.Tags})
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "%s tags for %s: %s\n", strings.ToUpper(sub), c.Name, strings.Join(c.Tags, ", "))
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown tag subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func (a *App) runClusterKVMTag(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster kvm-tag <add|rm|ls> [name-or-id] --vm <name> [--tags a,b]")
|
|
return 2
|
|
}
|
|
sub := strings.ToLower(strings.TrimSpace(args[0]))
|
|
fs := flag.NewFlagSet("cluster kvm-tag "+sub, flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
vm := fs.String("vm", "", "KVM domain name")
|
|
tagsCSV := fs.String("tags", "", "Comma separated tags")
|
|
selector, parseArgs := splitLeadingSelector(args[1:])
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 {
|
|
if selector != "" {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster kvm-tag <add|rm|ls> [name-or-id] --vm <name> [--tags a,b]")
|
|
return 2
|
|
}
|
|
selector = fs.Arg(0)
|
|
}
|
|
vmName := strings.TrimSpace(*vm)
|
|
switch sub {
|
|
case "ls", "list", "show":
|
|
items, err := svc.ListKVMTags(selector, vmName)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "kvm-tag ls: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, items)
|
|
return 0
|
|
}
|
|
if len(items) == 0 {
|
|
fmt.Fprintln(a.out, "(no vm tags)")
|
|
return 0
|
|
}
|
|
keys := make([]string, 0, len(items))
|
|
for k := range items {
|
|
keys = append(keys, k)
|
|
}
|
|
sort.Strings(keys)
|
|
for _, k := range keys {
|
|
fmt.Fprintf(a.out, "%s: %s\n", k, strings.Join(items[k], ", "))
|
|
}
|
|
return 0
|
|
case "add", "rm", "remove", "del":
|
|
if vmName == "" {
|
|
fmt.Fprintln(a.err, "kvm-tag: --vm is required")
|
|
return 2
|
|
}
|
|
tags := parseTagsCSV(*tagsCSV)
|
|
if len(tags) == 0 {
|
|
tags = fs.Args()
|
|
}
|
|
if len(tags) == 0 {
|
|
fmt.Fprintln(a.err, "kvm-tag: provide tags with --tags a,b")
|
|
return 2
|
|
}
|
|
var (
|
|
c cluster.Cluster
|
|
err error
|
|
)
|
|
if sub == "add" {
|
|
c, err = svc.AddKVMTag(selector, vmName, tags)
|
|
} else {
|
|
c, err = svc.RemoveKVMTag(selector, vmName, tags)
|
|
}
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "kvm-tag %s: %v\n", sub, err)
|
|
return 1
|
|
}
|
|
vmTags := c.KVMTags[vmName]
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"cluster": c.Name, "vm": vmName, "tags": vmTags})
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "%s vm tags for %s/%s: %s\n", strings.ToUpper(sub), c.Name, vmName, strings.Join(vmTags, ", "))
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown kvm-tag subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func (a *App) runClusterVMAlert(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster alert-vm <show|set|check> [name-or-id]")
|
|
return 2
|
|
}
|
|
sub := strings.ToLower(strings.TrimSpace(args[0]))
|
|
switch sub {
|
|
case "show":
|
|
selector := ""
|
|
if len(args) > 1 {
|
|
selector = args[1]
|
|
}
|
|
p, err := svc.GetVMAlertPolicy(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "alert-vm show: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, p)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "enabled=%t warn_on_shutoff=%t min_running=%d\n", p.Enabled, p.WarnOnShutoff, p.MinRunning)
|
|
return 0
|
|
case "set":
|
|
fs := flag.NewFlagSet("cluster alert-vm set", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
enabled := fs.Bool("enabled", false, "Enable VM state alerts")
|
|
warnOnShutoff := fs.Bool("warn-on-shutoff", true, "Warn when any VM is shut off")
|
|
minRunning := fs.Int("min-running", 1, "Minimum expected running VM count")
|
|
selector, parseArgs := splitLeadingSelector(args[1:])
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 {
|
|
if selector != "" {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster alert-vm set [name-or-id] [--enabled --warn-on-shutoff --min-running 1]")
|
|
return 2
|
|
}
|
|
selector = fs.Arg(0)
|
|
}
|
|
current, err := svc.GetVMAlertPolicy(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "alert-vm set: %v\n", err)
|
|
return 1
|
|
}
|
|
changed := false
|
|
fs.Visit(func(f *flag.Flag) {
|
|
switch f.Name {
|
|
case "enabled":
|
|
current.Enabled = *enabled
|
|
changed = true
|
|
case "warn-on-shutoff":
|
|
current.WarnOnShutoff = *warnOnShutoff
|
|
changed = true
|
|
case "min-running":
|
|
current.MinRunning = *minRunning
|
|
changed = true
|
|
}
|
|
})
|
|
if !changed {
|
|
fmt.Fprintln(a.err, "alert-vm set: provide at least one flag to change")
|
|
return 2
|
|
}
|
|
c, err := svc.SetVMAlertPolicy(selector, current)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "alert-vm set: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, c.VMAlerts)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Updated vm alerts for %s: enabled=%t warn_on_shutoff=%t min_running=%d\n", c.Name, c.VMAlerts.Enabled, c.VMAlerts.WarnOnShutoff, c.VMAlerts.MinRunning)
|
|
return 0
|
|
case "check":
|
|
selector := ""
|
|
if len(args) > 1 {
|
|
selector = args[1]
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
|
defer cancel()
|
|
rep, err := svc.CheckVMAlerts(ctx, selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "alert-vm check: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, rep)
|
|
if len(rep.Warnings) > 0 {
|
|
return 1
|
|
}
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "VMs total=%d running=%d shut_off=%d paused=%d others=%d\n", rep.Total, rep.Running, rep.ShutOff, rep.Paused, rep.Others)
|
|
if len(rep.ShutOffNames) > 0 {
|
|
fmt.Fprintf(a.out, "shut off: %s\n", strings.Join(rep.ShutOffNames, ", "))
|
|
}
|
|
if len(rep.PausedNames) > 0 {
|
|
fmt.Fprintf(a.out, "paused: %s\n", strings.Join(rep.PausedNames, ", "))
|
|
}
|
|
if len(rep.OtherNames) > 0 {
|
|
fmt.Fprintf(a.out, "other: %s\n", strings.Join(rep.OtherNames, ", "))
|
|
}
|
|
if len(rep.Warnings) > 0 {
|
|
for _, w := range rep.Warnings {
|
|
fmt.Fprintf(a.out, "! %s\n", w)
|
|
}
|
|
return 1
|
|
}
|
|
fmt.Fprintln(a.out, "VM alerts OK")
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown alert-vm subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func (a *App) runClusterAlertRouting(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster alert-routing <show|set> [cluster]")
|
|
return 2
|
|
}
|
|
switch strings.ToLower(strings.TrimSpace(args[0])) {
|
|
case "show":
|
|
selector := ""
|
|
if len(args) > 1 {
|
|
selector = args[1]
|
|
}
|
|
p, err := svc.GetAlertRouting(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "alert-routing show: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, p)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "critical_immediate=%t warning_batch_mins=%d\n", p.CriticalImmediate, p.WarningBatchMins)
|
|
return 0
|
|
case "set":
|
|
fs := flag.NewFlagSet("cluster alert-routing set", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
critical := fs.Bool("critical-immediate", true, "Send critical alerts immediately")
|
|
batch := fs.Int("warning-batch-mins", 5, "Warning batching interval in minutes")
|
|
selector, parseArgs := splitLeadingSelector(args[1:])
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 && selector == "" {
|
|
selector = fs.Arg(0)
|
|
}
|
|
cur, err := svc.GetAlertRouting(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "alert-routing set: %v\n", err)
|
|
return 1
|
|
}
|
|
fs.Visit(func(f *flag.Flag) {
|
|
switch f.Name {
|
|
case "critical-immediate":
|
|
cur.CriticalImmediate = *critical
|
|
case "warning-batch-mins":
|
|
cur.WarningBatchMins = *batch
|
|
}
|
|
})
|
|
c, err := svc.SetAlertRouting(selector, cur)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "alert-routing set: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, c.AlertRouting)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Updated alert routing for %s: critical_immediate=%t warning_batch_mins=%d\n",
|
|
c.Name, c.AlertRouting.CriticalImmediate, c.AlertRouting.WarningBatchMins)
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown alert-routing subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func (a *App) runClusterRunbookTrigger(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster runbook-trigger <show|set> [cluster]")
|
|
return 2
|
|
}
|
|
switch strings.ToLower(strings.TrimSpace(args[0])) {
|
|
case "show":
|
|
selector := ""
|
|
if len(args) > 1 {
|
|
selector = args[1]
|
|
}
|
|
p, err := svc.GetRunbookTrigger(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook-trigger show: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, p)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "enabled=%t on_vm_shutoff=%t runbook_id=%s cooldown_mins=%d last_triggered=%s\n",
|
|
p.Enabled, p.OnVMShutoff, emptyFallback(p.RunbookID, "-"), p.CooldownMins,
|
|
emptyFallback(p.LastTriggered.Format(time.RFC3339), "-"))
|
|
return 0
|
|
case "set":
|
|
fs := flag.NewFlagSet("cluster runbook-trigger set", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
enabled := fs.Bool("enabled", false, "Enable auto-trigger")
|
|
onShutoff := fs.Bool("on-vm-shutoff", true, "Trigger on VM shutoff alert")
|
|
runbookID := fs.String("runbook-id", "", "Runbook ID to execute")
|
|
cooldown := fs.Int("cooldown-mins", 30, "Cooldown between triggers")
|
|
selector, parseArgs := splitLeadingSelector(args[1:])
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 && selector == "" {
|
|
selector = fs.Arg(0)
|
|
}
|
|
cur, err := svc.GetRunbookTrigger(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook-trigger set: %v\n", err)
|
|
return 1
|
|
}
|
|
fs.Visit(func(f *flag.Flag) {
|
|
switch f.Name {
|
|
case "enabled":
|
|
cur.Enabled = *enabled
|
|
case "on-vm-shutoff":
|
|
cur.OnVMShutoff = *onShutoff
|
|
case "runbook-id":
|
|
cur.RunbookID = strings.TrimSpace(*runbookID)
|
|
case "cooldown-mins":
|
|
cur.CooldownMins = *cooldown
|
|
}
|
|
})
|
|
c, err := svc.SetRunbookTrigger(selector, cur)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook-trigger set: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, c.RunbookTrigger)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Updated runbook trigger for %s: enabled=%t on_vm_shutoff=%t runbook_id=%s cooldown_mins=%d\n",
|
|
c.Name, c.RunbookTrigger.Enabled, c.RunbookTrigger.OnVMShutoff,
|
|
emptyFallback(c.RunbookTrigger.RunbookID, "-"), c.RunbookTrigger.CooldownMins)
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown runbook-trigger subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func (a *App) runClusterP95(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
fs := flag.NewFlagSet("cluster p95", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
iface := fs.String("iface", "", "Interface name (required)")
|
|
rangeStr := fs.String("range", "30d", "Time window: 1h|1d|30d|1mo|all")
|
|
graph := fs.Bool("graph", false, "Generate graph PNG")
|
|
outPath := fs.String("out", "", "Output PNG path for --graph")
|
|
selector, parseArgs := splitLeadingSelector(args)
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 {
|
|
if selector == "" {
|
|
selector = fs.Arg(0)
|
|
} else if strings.TrimSpace(*iface) == "" {
|
|
*iface = fs.Arg(0)
|
|
}
|
|
}
|
|
if fs.NArg() > 1 && strings.TrimSpace(*iface) == "" {
|
|
*iface = fs.Arg(1)
|
|
}
|
|
if strings.TrimSpace(*iface) == "" {
|
|
fmt.Fprintln(a.err, "cluster p95: --iface is required")
|
|
return 2
|
|
}
|
|
rng, ok := history.ParseRangeShortcut(strings.TrimSpace(*rangeStr))
|
|
if !ok {
|
|
fmt.Fprintf(a.err, "cluster p95: invalid --range %q\n", *rangeStr)
|
|
return 2
|
|
}
|
|
snap, err := svc.CollectInterfaceP95(selector, *iface, rng)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster p95: %v\n", err)
|
|
return 1
|
|
}
|
|
pngPath := ""
|
|
if *graph {
|
|
png, err := svc.RenderInterfaceP95GraphPNG(snap)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster p95: render graph: %v\n", err)
|
|
return 1
|
|
}
|
|
pngPath = strings.TrimSpace(*outPath)
|
|
if pngPath == "" {
|
|
pngPath = fmt.Sprintf("%s-%s-%s.png", sanitizeFilename(snap.ClusterName), sanitizeFilename(snap.Interface), rng)
|
|
}
|
|
if err := os.WriteFile(pngPath, png, 0o600); err != nil {
|
|
fmt.Fprintf(a.err, "cluster p95: write graph: %v\n", err)
|
|
return 1
|
|
}
|
|
}
|
|
if jsonOut {
|
|
resp := map[string]any{
|
|
"cluster": snap.ClusterName,
|
|
"interface": snap.Interface,
|
|
"range": rng,
|
|
"samples": snap.Samples,
|
|
"p95_mbps": snap.P95Mbps,
|
|
"avg_mbps": snap.AvgMbps,
|
|
"max_mbps": snap.MaxMbps,
|
|
}
|
|
if pngPath != "" {
|
|
resp["graph"] = pngPath
|
|
}
|
|
_ = writeJSON(a.out, resp)
|
|
return 0
|
|
}
|
|
if pngPath != "" {
|
|
fmt.Fprintf(a.out, "Saved: %s\n", pngPath)
|
|
}
|
|
reportPath := fmt.Sprintf("%s-%s-%s.txt",
|
|
sanitizeFilename(snap.ClusterName),
|
|
sanitizeFilename(snap.Interface),
|
|
time.Now().UTC().Format("20060102"))
|
|
reportBody := strings.Join([]string{
|
|
fmt.Sprintf("cluster=%s", snap.ClusterName),
|
|
fmt.Sprintf("iface=%s", snap.Interface),
|
|
fmt.Sprintf("range=%s", rng),
|
|
fmt.Sprintf("generated_at=%s", time.Now().UTC().Format(time.RFC3339)),
|
|
fmt.Sprintf("p95_mbps=%.3f", snap.P95Mbps),
|
|
fmt.Sprintf("avg_mbps=%.3f", snap.AvgMbps),
|
|
fmt.Sprintf("max_mbps=%.3f", snap.MaxMbps),
|
|
fmt.Sprintf("samples=%d", snap.Samples),
|
|
"",
|
|
}, "\n")
|
|
if err := os.WriteFile(reportPath, []byte(reportBody), 0o600); err == nil {
|
|
fmt.Fprintf(a.out, "Saved: %s\n", reportPath)
|
|
}
|
|
fmt.Fprintf(a.out, "Cluster: %s\n", snap.ClusterName)
|
|
fmt.Fprintf(a.out, "Iface: %s\n", snap.Interface)
|
|
fmt.Fprintf(a.out, "Range: %s\n", rng.Label())
|
|
fmt.Fprintf(a.out, "P95: %s over %d samples\n", formatMbpsHuman(snap.P95Mbps), snap.Samples)
|
|
fmt.Fprintf(a.out, "Max: %s\n", formatMbpsHuman(snap.MaxMbps))
|
|
fmt.Fprintf(a.out, "Avg: %s\n", formatMbpsHuman(snap.AvgMbps))
|
|
_ = svc.AppendChange("p95", snap.ClusterName, fmt.Sprintf("iface=%s range=%s p95=%.2f", snap.Interface, rng, snap.P95Mbps))
|
|
return 0
|
|
}
|
|
|
|
func (a *App) runClusterSLO(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
fs := flag.NewFlagSet("cluster slo", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
rangeStr := fs.String("range", "30d", "Time window: 1d|7d|30d|all")
|
|
vm := fs.String("vm", "", "Optional VM name filter")
|
|
selector, parseArgs := splitLeadingSelector(args)
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 && selector == "" {
|
|
selector = fs.Arg(0)
|
|
}
|
|
since, label, err := parseSinceRange(*rangeStr)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster slo: %v\n", err)
|
|
return 2
|
|
}
|
|
rep, err := svc.AvailabilityReport(selector, since, *vm)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster slo: %v\n", err)
|
|
return 1
|
|
}
|
|
rep.Range = label
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, rep)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Cluster: %s\n", rep.Cluster)
|
|
fmt.Fprintf(a.out, "Range: %s\n", rep.Range)
|
|
fmt.Fprintf(a.out, "Uptime: %.2f%% (%d/%d samples up)\n", rep.Availability, rep.UpSamples, rep.Samples)
|
|
if strings.TrimSpace(*vm) != "" {
|
|
if len(rep.VMs) == 0 {
|
|
fmt.Fprintf(a.out, "VM %s: no samples\n", *vm)
|
|
} else {
|
|
v := rep.VMs[0]
|
|
fmt.Fprintf(a.out, "VM %s uptime: %.2f%% (%d/%d running)\n", v.Name, v.Availability, v.Running, v.Samples)
|
|
}
|
|
} else if len(rep.VMs) > 0 {
|
|
maxRows := 10
|
|
if len(rep.VMs) < maxRows {
|
|
maxRows = len(rep.VMs)
|
|
}
|
|
tw := tabwriter.NewWriter(a.out, 0, 2, 2, ' ', 0)
|
|
fmt.Fprintln(tw, "VM\tUPTIME\tRUNNING/SAMPLES")
|
|
for i := 0; i < maxRows; i++ {
|
|
v := rep.VMs[i]
|
|
fmt.Fprintf(tw, "%s\t%.2f%%\t%d/%d\n", v.Name, v.Availability, v.Running, v.Samples)
|
|
}
|
|
_ = tw.Flush()
|
|
if len(rep.VMs) > maxRows {
|
|
fmt.Fprintf(a.out, "... +%d VM(s)\n", len(rep.VMs)-maxRows)
|
|
}
|
|
}
|
|
return 0
|
|
}
|
|
|
|
func (a *App) runClusterCapacity(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 || strings.EqualFold(strings.TrimSpace(args[0]), "forecast") {
|
|
rest := args
|
|
if len(rest) > 0 && strings.EqualFold(strings.TrimSpace(rest[0]), "forecast") {
|
|
rest = rest[1:]
|
|
}
|
|
fs := flag.NewFlagSet("cluster capacity forecast", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
rangeStr := fs.String("range", "30d", "Time window: 7d|30d|all")
|
|
selector, parseArgs := splitLeadingSelector(rest)
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 && selector == "" {
|
|
selector = fs.Arg(0)
|
|
}
|
|
since, label, err := parseSinceRange(*rangeStr)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster capacity: %v\n", err)
|
|
return 2
|
|
}
|
|
rep, err := svc.CapacityForecast(selector, since)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster capacity: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"range": label, "report": rep})
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Cluster: %s\n", rep.Cluster)
|
|
fmt.Fprintf(a.out, "Range: %s\n", label)
|
|
fmt.Fprintf(a.out, "Samples: %d\n", rep.Samples)
|
|
if len(rep.Items) == 0 {
|
|
fmt.Fprintln(a.out, "No capacity trend data yet")
|
|
return 0
|
|
}
|
|
tw := tabwriter.NewWriter(a.out, 0, 2, 2, ' ', 0)
|
|
fmt.Fprintln(tw, "MOUNT\tUSED\tDAYS_TO_90%\tDAYS_TO_95%\tSLOPE")
|
|
for _, it := range rep.Items {
|
|
d90 := "inf"
|
|
d95 := "inf"
|
|
if !math.IsInf(it.DaysTo90, 1) {
|
|
d90 = fmt.Sprintf("%.1f", it.DaysTo90)
|
|
}
|
|
if !math.IsInf(it.DaysTo95, 1) {
|
|
d95 = fmt.Sprintf("%.1f", it.DaysTo95)
|
|
}
|
|
fmt.Fprintf(tw, "%s\t%.1f%%\t%s\t%s\t%.2f MiB/day\n",
|
|
it.Mount, it.UsedPct, d90, d95, it.SlopeBytesSec*86400/1024/1024)
|
|
}
|
|
_ = tw.Flush()
|
|
return 0
|
|
}
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster capacity forecast [cluster] [--range 7d|30d|all]")
|
|
return 2
|
|
}
|
|
|
|
func (a *App) runClusterDrift(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) > 0 {
|
|
switch strings.ToLower(strings.TrimSpace(args[0])) {
|
|
case "baseline":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster drift baseline <set|show> [cluster]")
|
|
return 2
|
|
}
|
|
switch strings.ToLower(strings.TrimSpace(args[1])) {
|
|
case "set":
|
|
selector := ""
|
|
if len(args) > 2 {
|
|
selector = args[2]
|
|
}
|
|
c, err := svc.SetDriftBaseline(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "drift baseline set: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, c.Drift)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Baseline set for %s\n", c.Name)
|
|
return 0
|
|
case "show":
|
|
selector := ""
|
|
if len(args) > 2 {
|
|
selector = args[2]
|
|
}
|
|
d, err := svc.GetDriftControl(selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "drift baseline show: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, d.Baseline)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "enabled=%t set_at=%s agent=%s software=%s\n",
|
|
d.Baseline.Enabled,
|
|
emptyFallback(d.Baseline.SetAt.Format(time.RFC3339), "-"),
|
|
emptyFallback(d.Baseline.AgentVersion, "-"),
|
|
emptyFallback(d.Baseline.Software, "-"))
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.err, "unknown drift baseline subcommand %q\n", args[1])
|
|
return 2
|
|
case "ack":
|
|
fs := flag.NewFlagSet("cluster drift ack", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
kind := fs.String("kind", "", "Issue kind to ack (e.g. agent_version)")
|
|
forDur := fs.Duration("for", 24*time.Hour, "Ack duration (e.g. 24h)")
|
|
selector, parseArgs := splitLeadingSelector(args[1:])
|
|
if err := fs.Parse(parseArgs); err != nil {
|
|
return 2
|
|
}
|
|
if fs.NArg() > 0 && selector == "" {
|
|
selector = fs.Arg(0)
|
|
}
|
|
if strings.TrimSpace(*kind) == "" {
|
|
fmt.Fprintln(a.err, "drift ack: --kind is required")
|
|
return 2
|
|
}
|
|
until := time.Now().UTC().Add(*forDur)
|
|
_, err := svc.AckDriftIssue(selector, *kind, until)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "drift ack: %v\n", err)
|
|
return 1
|
|
}
|
|
fmt.Fprintf(a.out, "Acked %s until %s\n", *kind, until.Format(time.RFC3339))
|
|
return 0
|
|
}
|
|
}
|
|
|
|
selector := ""
|
|
if len(args) > 0 {
|
|
selector = args[0]
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
|
|
defer cancel()
|
|
rep, err := svc.DetectDrift(ctx, selector)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster drift: %v\n", err)
|
|
return 1
|
|
}
|
|
c, _ := svc.Get(selector)
|
|
filtered := make([]cluster.DriftIssue, 0, len(rep.Issues))
|
|
for _, it := range rep.Issues {
|
|
if svc.IsDriftIssueAcked(c, it.Kind, time.Now().UTC()) {
|
|
continue
|
|
}
|
|
filtered = append(filtered, it)
|
|
}
|
|
rep.Issues = filtered
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, rep)
|
|
if len(rep.Issues) > 0 {
|
|
return 1
|
|
}
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Drift report for %s (%s)\n", rep.Cluster, rep.Generated.Format(time.RFC3339))
|
|
if len(rep.Issues) == 0 {
|
|
fmt.Fprintln(a.out, "No drift detected")
|
|
return 0
|
|
}
|
|
for _, it := range rep.Issues {
|
|
fmt.Fprintf(a.out, "[%s] %s: %s\n", strings.ToUpper(it.Level), it.Kind, it.Message)
|
|
}
|
|
return 1
|
|
}
|
|
|
|
func (a *App) runClusterChangeHistory(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
fs := flag.NewFlagSet("cluster change-history", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
tail := fs.Int("tail", 100, "Last N entries")
|
|
if err := fs.Parse(args); err != nil {
|
|
return 2
|
|
}
|
|
items, err := svc.ListChanges(*tail)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster change-history: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, items)
|
|
return 0
|
|
}
|
|
if len(items) == 0 {
|
|
fmt.Fprintln(a.out, "(no history)")
|
|
return 0
|
|
}
|
|
tw := tabwriter.NewWriter(a.out, 0, 2, 2, ' ', 0)
|
|
fmt.Fprintln(tw, "TIME\tACTION\tTARGET\tDETAILS")
|
|
for _, it := range items {
|
|
fmt.Fprintf(tw, "%s\t%s\t%s\t%s\n", it.At.Format("2006-01-02 15:04:05"), it.Action, it.Target, it.Details)
|
|
}
|
|
_ = tw.Flush()
|
|
return 0
|
|
}
|
|
|
|
func (a *App) runClusterReport(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster report export --out ./report.json [--format json|csv]")
|
|
return 2
|
|
}
|
|
sub := strings.ToLower(strings.TrimSpace(args[0]))
|
|
if sub != "export" {
|
|
fmt.Fprintf(a.err, "unknown report subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
fs := flag.NewFlagSet("cluster report export", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
format := fs.String("format", "json", "Output format: json|csv")
|
|
out := fs.String("out", "", "Output file path")
|
|
probe := fs.Bool("probe", false, "Probe agent reachability for each cluster")
|
|
if err := fs.Parse(args[1:]); err != nil {
|
|
return 2
|
|
}
|
|
if strings.TrimSpace(*out) == "" {
|
|
fmt.Fprintln(a.err, "cluster report export: --out is required")
|
|
return 2
|
|
}
|
|
clusters, activeID, err := svc.List()
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster report export: %v\n", err)
|
|
return 1
|
|
}
|
|
type row struct {
|
|
Name string `json:"name"`
|
|
ID string `json:"id"`
|
|
Host string `json:"host"`
|
|
User string `json:"user"`
|
|
Active bool `json:"active"`
|
|
Agent string `json:"agent"`
|
|
AgentVersion string `json:"agent_version"`
|
|
Software string `json:"software"`
|
|
Tags []string `json:"tags,omitempty"`
|
|
UpdatedAt string `json:"updated_at"`
|
|
Reachable string `json:"reachable,omitempty"`
|
|
}
|
|
items := make([]row, 0, len(clusters))
|
|
for _, c := range clusters {
|
|
r := row{
|
|
Name: c.Name,
|
|
ID: c.ID,
|
|
Host: c.Host,
|
|
User: c.User,
|
|
Active: c.ID == activeID,
|
|
Agent: ternary(c.Agent.Installed, "installed", "none"),
|
|
AgentVersion: c.Agent.Version,
|
|
Software: c.Software.Summary(),
|
|
Tags: c.Tags,
|
|
UpdatedAt: c.UpdatedAt.Format(time.RFC3339),
|
|
}
|
|
if *probe && c.Agent.Installed {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 1800*time.Millisecond)
|
|
p, err := svc.PingAgent(ctx, c.ID)
|
|
cancel()
|
|
if err == nil && p.Reachable && p.StatusCode < 400 {
|
|
r.Reachable = "up"
|
|
} else {
|
|
r.Reachable = "down"
|
|
}
|
|
}
|
|
items = append(items, r)
|
|
}
|
|
if err := os.MkdirAll(filepath.Dir(*out), 0o700); err != nil {
|
|
fmt.Fprintf(a.err, "cluster report export: %v\n", err)
|
|
return 1
|
|
}
|
|
switch strings.ToLower(strings.TrimSpace(*format)) {
|
|
case "json":
|
|
b, _ := jsonMarshalIndent(map[string]any{"generated_at": time.Now().UTC().Format(time.RFC3339), "clusters": items})
|
|
if err := os.WriteFile(*out, append(b, '\n'), 0o600); err != nil {
|
|
fmt.Fprintf(a.err, "cluster report export: %v\n", err)
|
|
return 1
|
|
}
|
|
case "csv":
|
|
f, err := os.Create(*out)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "cluster report export: %v\n", err)
|
|
return 1
|
|
}
|
|
w := csv.NewWriter(f)
|
|
_ = w.Write([]string{"name", "id", "host", "user", "active", "agent", "agent_version", "software", "tags", "updated_at", "reachable"})
|
|
for _, r := range items {
|
|
_ = w.Write([]string{r.Name, r.ID, r.Host, r.User, fmt.Sprintf("%t", r.Active), r.Agent, r.AgentVersion, r.Software, strings.Join(r.Tags, ","), r.UpdatedAt, r.Reachable})
|
|
}
|
|
w.Flush()
|
|
_ = f.Close()
|
|
if err := w.Error(); err != nil {
|
|
fmt.Fprintf(a.err, "cluster report export: %v\n", err)
|
|
return 1
|
|
}
|
|
default:
|
|
fmt.Fprintf(a.err, "cluster report export: invalid --format %q\n", *format)
|
|
return 2
|
|
}
|
|
_ = svc.AppendChange("report.export", *out, *format)
|
|
if !jsonOut {
|
|
fmt.Fprintf(a.out, "Report exported: %s\n", *out)
|
|
}
|
|
return 0
|
|
}
|
|
|
|
func jsonMarshalIndent(v any) ([]byte, error) {
|
|
return json.MarshalIndent(v, "", " ")
|
|
}
|
|
|
|
func (a *App) runClusterBackup(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup <target|plan|run>")
|
|
return 2
|
|
}
|
|
sub := strings.ToLower(strings.TrimSpace(args[0]))
|
|
switch sub {
|
|
case "target", "targets":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup target <add|ls|rm|test>")
|
|
return 2
|
|
}
|
|
tSub := strings.ToLower(strings.TrimSpace(args[1]))
|
|
switch tSub {
|
|
case "ls", "list", "show":
|
|
items, err := svc.BackupListTargets()
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup target ls: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, items)
|
|
return 0
|
|
}
|
|
if len(items) == 0 {
|
|
fmt.Fprintln(a.out, "(no backup targets)")
|
|
return 0
|
|
}
|
|
tw := tabwriter.NewWriter(a.out, 0, 2, 2, ' ', 0)
|
|
fmt.Fprintln(tw, "ID\tNAME\tTYPE\tENABLED\tDEST")
|
|
for _, t := range items {
|
|
dest := "-"
|
|
if t.Type == "sftp" {
|
|
dest = fmt.Sprintf("%s@%s:%s", t.SFTPUser, t.SFTPHost, strings.TrimSpace(t.SFTPBasePath))
|
|
} else if t.Type == "s3" {
|
|
dest = fmt.Sprintf("%s/%s", t.S3Endpoint, t.S3Bucket)
|
|
}
|
|
fmt.Fprintf(tw, "%s\t%s\t%s\t%t\t%s\n", t.ID, t.Name, t.Type, t.Enabled, dest)
|
|
}
|
|
_ = tw.Flush()
|
|
return 0
|
|
case "add":
|
|
fs := flag.NewFlagSet("cluster backup target add", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
name := fs.String("name", "", "Target name")
|
|
tp := fs.String("type", "", "Target type: sftp|s3")
|
|
enabled := fs.Bool("enabled", true, "Enable target")
|
|
sftpHost := fs.String("sftp-host", "", "SFTP host")
|
|
sftpPort := fs.Int("sftp-port", 22, "SFTP port")
|
|
sftpUser := fs.String("sftp-user", "", "SFTP user")
|
|
sftpPass := fs.String("sftp-password", "", "SFTP password")
|
|
sftpKey := fs.String("sftp-key", "", "SFTP key path")
|
|
sftpBase := fs.String("sftp-base", "", "SFTP base directory")
|
|
s3Endpoint := fs.String("s3-endpoint", "", "S3 endpoint")
|
|
s3Region := fs.String("s3-region", "us-east-1", "S3 region")
|
|
s3Bucket := fs.String("s3-bucket", "", "S3 bucket")
|
|
s3Prefix := fs.String("s3-prefix", "", "S3 key prefix")
|
|
s3Access := fs.String("s3-access-key", "", "S3 access key")
|
|
s3Secret := fs.String("s3-secret-key", "", "S3 secret key")
|
|
s3SSL := fs.Bool("s3-ssl", true, "Use TLS for S3")
|
|
s3PathStyle := fs.Bool("s3-path-style", false, "Use S3 path-style URLs")
|
|
if err := fs.Parse(args[2:]); err != nil {
|
|
return 2
|
|
}
|
|
t, err := svc.BackupAddTarget(cluster.BackupTarget{
|
|
Name: strings.TrimSpace(*name),
|
|
Type: strings.ToLower(strings.TrimSpace(*tp)),
|
|
Enabled: *enabled,
|
|
SFTPHost: strings.TrimSpace(*sftpHost),
|
|
SFTPPort: *sftpPort,
|
|
SFTPUser: strings.TrimSpace(*sftpUser),
|
|
SFTPPassword: strings.TrimSpace(*sftpPass),
|
|
SFTPKeyPath: strings.TrimSpace(*sftpKey),
|
|
SFTPBasePath: strings.TrimSpace(*sftpBase),
|
|
S3Endpoint: strings.TrimSpace(*s3Endpoint),
|
|
S3Region: strings.TrimSpace(*s3Region),
|
|
S3Bucket: strings.TrimSpace(*s3Bucket),
|
|
S3Prefix: strings.TrimSpace(*s3Prefix),
|
|
S3AccessKey: strings.TrimSpace(*s3Access),
|
|
S3SecretKey: strings.TrimSpace(*s3Secret),
|
|
S3UseSSL: *s3SSL,
|
|
S3PathStyle: *s3PathStyle,
|
|
})
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup target add: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, t)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Backup target added: %s (%s)\n", t.Name, t.ID)
|
|
return 0
|
|
case "rm", "remove", "del":
|
|
if len(args) < 3 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup target rm <id-or-name>")
|
|
return 2
|
|
}
|
|
t, err := svc.BackupRemoveTarget(args[2])
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup target rm: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, t)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Backup target removed: %s (%s)\n", t.Name, t.ID)
|
|
return 0
|
|
case "test":
|
|
if len(args) < 3 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup target test <id-or-name>")
|
|
return 2
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 25*time.Second)
|
|
defer cancel()
|
|
msg, err := svc.BackupTestTarget(ctx, args[2])
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup target test: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"ok": true, "message": msg})
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Backup target OK: %s\n", msg)
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown backup target subcommand %q\n", args[1])
|
|
return 2
|
|
}
|
|
case "plan", "plans":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup plan <add|ls|rm|run>")
|
|
return 2
|
|
}
|
|
pSub := strings.ToLower(strings.TrimSpace(args[1]))
|
|
switch pSub {
|
|
case "ls", "list", "show":
|
|
items, err := svc.BackupListPlans()
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup plan ls: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, items)
|
|
return 0
|
|
}
|
|
if len(items) == 0 {
|
|
fmt.Fprintln(a.out, "(no backup plans)")
|
|
return 0
|
|
}
|
|
tw := tabwriter.NewWriter(a.out, 0, 2, 2, ' ', 0)
|
|
fmt.Fprintln(tw, "ID\tNAME\tCLUSTER\tTARGET\tPATHS\tEVERY\tENABLED\tLAST_STATUS\tLAST_RUN")
|
|
for _, p := range items {
|
|
lastRun := "-"
|
|
if !p.LastRunAt.IsZero() {
|
|
lastRun = p.LastRunAt.Format("2006-01-02 15:04:05")
|
|
}
|
|
clusterSel := p.Cluster
|
|
if strings.TrimSpace(clusterSel) == "" {
|
|
clusterSel = "(active)"
|
|
}
|
|
fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%s\t%s\t%t\t%s\t%s\n",
|
|
p.ID, p.Name, clusterSel, p.TargetID, strings.Join(p.Paths, ","), p.Every, p.Enabled, p.LastStatus, lastRun)
|
|
}
|
|
_ = tw.Flush()
|
|
return 0
|
|
case "add":
|
|
fs := flag.NewFlagSet("cluster backup plan add", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
name := fs.String("name", "", "Plan name")
|
|
clusterSel := fs.String("cluster", "", "Cluster selector (default active)")
|
|
target := fs.String("target", "", "Target id or name")
|
|
every := fs.String("every", "24h", "Backup interval, e.g. 6h")
|
|
retain := fs.Int("retain-days", 30, "Retention days metadata")
|
|
enabled := fs.Bool("enabled", true, "Enable plan")
|
|
schedule := fs.Bool("schedule", true, "Create scheduler task automatically")
|
|
pathList := multiStringFlag{}
|
|
fs.Var(&pathList, "path", "Remote path (repeatable)")
|
|
if err := fs.Parse(args[2:]); err != nil {
|
|
return 2
|
|
}
|
|
paths := normalizeCLIPathList(pathList)
|
|
p, err := svc.BackupAddPlan(cluster.BackupPlan{
|
|
Name: strings.TrimSpace(*name),
|
|
Cluster: strings.TrimSpace(*clusterSel),
|
|
TargetID: strings.TrimSpace(*target),
|
|
Paths: paths,
|
|
Every: strings.TrimSpace(*every),
|
|
RetainDays: *retain,
|
|
Enabled: *enabled,
|
|
})
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup plan add: %v\n", err)
|
|
return 1
|
|
}
|
|
if *schedule {
|
|
taskName := "backup-" + p.ID
|
|
cmd := "cluster backup run " + p.ID
|
|
_, sErr := svc.SchedulerAdd(taskName, "", cmd, p.Every, "observer", "30s", 5, 3, p.Enabled)
|
|
if sErr != nil && !strings.Contains(strings.ToLower(sErr.Error()), "already exists") {
|
|
fmt.Fprintf(a.err, "backup plan add: schedule create failed: %v\n", sErr)
|
|
return 1
|
|
}
|
|
if _, running, _, stErr := schedulerDaemonStatus(svc); stErr == nil && !running {
|
|
_, _, _ = startSchedulerDaemon(svc, 30*time.Second)
|
|
}
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, p)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Backup plan added: %s (%s)\n", p.Name, p.ID)
|
|
if *schedule {
|
|
fmt.Fprintf(a.out, "Scheduler task: backup-%s\n", p.ID)
|
|
}
|
|
return 0
|
|
case "rm", "remove", "del":
|
|
if len(args) < 3 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup plan rm <id-or-name>")
|
|
return 2
|
|
}
|
|
p, err := svc.BackupRemovePlan(args[2])
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup plan rm: %v\n", err)
|
|
return 1
|
|
}
|
|
_, _ = svc.SchedulerRemove("backup-" + p.ID)
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, p)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Backup plan removed: %s (%s)\n", p.Name, p.ID)
|
|
return 0
|
|
case "run":
|
|
if len(args) < 3 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup plan run <id-or-name>")
|
|
return 2
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
|
|
defer cancel()
|
|
res, err := svc.BackupRunPlan(ctx, args[2])
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup plan run: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, res)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Backup uploaded: %s\n", res.UploadedTo)
|
|
fmt.Fprintf(a.out, "Archive: %s (%d bytes)\n", res.ArchiveName, res.SizeBytes)
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown backup plan subcommand %q\n", args[1])
|
|
return 2
|
|
}
|
|
case "run":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster backup run <plan-id-or-name>")
|
|
return 2
|
|
}
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
|
|
defer cancel()
|
|
res, err := svc.BackupRunPlan(ctx, args[1])
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "backup run: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, res)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Backup uploaded: %s\n", res.UploadedTo)
|
|
fmt.Fprintf(a.out, "Archive: %s (%d bytes)\n", res.ArchiveName, res.SizeBytes)
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown backup subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func (a *App) runClusterRunbook(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster runbook <list|show|run> [id]")
|
|
return 2
|
|
}
|
|
sub := strings.ToLower(strings.TrimSpace(args[0]))
|
|
switch sub {
|
|
case "list", "ls":
|
|
items, err := svc.ListRunbooks()
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook list: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, items)
|
|
return 0
|
|
}
|
|
for _, rb := range items {
|
|
fmt.Fprintf(a.out, "%s\t%s\n", rb.ID, rb.Name)
|
|
}
|
|
return 0
|
|
case "add":
|
|
fs := flag.NewFlagSet("cluster runbook add", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
id := fs.String("id", "", "Runbook id")
|
|
name := fs.String("name", "", "Runbook name")
|
|
desc := fs.String("desc", "", "Description")
|
|
edit := fs.Bool("edit", false, "Open interactive template in $EDITOR")
|
|
from := fs.String("from", "", "Read runbook JSON from file")
|
|
step := multiStringFlag{}
|
|
fs.Var(&step, "step", "Step in format 'Title|command'")
|
|
if err := fs.Parse(args[1:]); err != nil {
|
|
return 2
|
|
}
|
|
var rbInput cluster.Runbook
|
|
fromPath := strings.TrimSpace(*from)
|
|
if fromPath != "" {
|
|
raw, err := os.ReadFile(fromPath)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook add: read --from: %v\n", err)
|
|
return 1
|
|
}
|
|
parsed, err := parseRunbookJSON(raw)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook add: invalid --from json: %v\n", err)
|
|
return 1
|
|
}
|
|
rbInput = parsed
|
|
} else if *edit {
|
|
input := cluster.Runbook{
|
|
ID: strings.TrimSpace(*id),
|
|
Name: strings.TrimSpace(*name),
|
|
Description: strings.TrimSpace(*desc),
|
|
Steps: []cluster.RunbookStep{
|
|
{Title: "Check agent", Command: "cluster agent status"},
|
|
{Title: "Check drift", Command: "cluster drift"},
|
|
},
|
|
}
|
|
raw, _ := jsonMarshalIndent(input)
|
|
if isEmbeddedConsoleMode() {
|
|
path, err := writeTemplateFile("pxmon-runbook-*.json", raw)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook add: %v\n", err)
|
|
return 1
|
|
}
|
|
fmt.Fprintf(a.out, "Template saved: %s\n", path)
|
|
fmt.Fprintf(a.out, "Edit it in external terminal and run:\n")
|
|
fmt.Fprintf(a.out, "pxmon cluster runbook add --from %s\n", shellQuoteArg(path))
|
|
return 0
|
|
}
|
|
edited, err := openEditorTemplate("pxmon-runbook-*.json", raw)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook add: %v\n", err)
|
|
return 1
|
|
}
|
|
parsed, err := parseRunbookJSON(edited)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook add: invalid template: %v\n", err)
|
|
return 1
|
|
}
|
|
rbInput = parsed
|
|
} else {
|
|
steps := make([]cluster.RunbookStep, 0, len(step))
|
|
for _, raw := range step {
|
|
parts := strings.SplitN(raw, "|", 2)
|
|
title := strings.TrimSpace(parts[0])
|
|
cmd := ""
|
|
if len(parts) > 1 {
|
|
cmd = strings.TrimSpace(parts[1])
|
|
}
|
|
if title == "" {
|
|
continue
|
|
}
|
|
steps = append(steps, cluster.RunbookStep{Title: title, Command: cmd})
|
|
}
|
|
rbInput = cluster.Runbook{
|
|
ID: strings.TrimSpace(*id),
|
|
Name: strings.TrimSpace(*name),
|
|
Description: strings.TrimSpace(*desc),
|
|
Steps: steps,
|
|
}
|
|
}
|
|
rb, err := svc.AddRunbook(rbInput)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook add: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, rb)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Runbook added: %s (%s)\n", rb.Name, rb.ID)
|
|
return 0
|
|
case "rm", "remove", "del":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster runbook rm <id>")
|
|
return 2
|
|
}
|
|
rb, err := svc.RemoveRunbook(args[1])
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "runbook rm: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, rb)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Runbook removed: %s (%s)\n", rb.Name, rb.ID)
|
|
return 0
|
|
case "show":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster runbook show <id>")
|
|
return 2
|
|
}
|
|
rb, ok := svc.GetRunbook(args[1])
|
|
if !ok {
|
|
fmt.Fprintf(a.err, "runbook %q not found\n", args[1])
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, rb)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "%s - %s\n", rb.ID, rb.Name)
|
|
if rb.Description != "" {
|
|
fmt.Fprintln(a.out, rb.Description)
|
|
}
|
|
for i, st := range rb.Steps {
|
|
fmt.Fprintf(a.out, "%d. %s\n", i+1, st.Title)
|
|
if st.Command != "" {
|
|
fmt.Fprintf(a.out, " %s\n", st.Command)
|
|
}
|
|
}
|
|
return 0
|
|
case "run":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster runbook run <id>")
|
|
return 2
|
|
}
|
|
rb, ok := svc.GetRunbook(args[1])
|
|
if !ok {
|
|
fmt.Fprintf(a.err, "runbook %q not found\n", args[1])
|
|
return 1
|
|
}
|
|
type stepResult struct {
|
|
Title string `json:"title"`
|
|
Command string `json:"command"`
|
|
ExitCode int `json:"exit_code"`
|
|
Output string `json:"output,omitempty"`
|
|
}
|
|
results := make([]stepResult, 0, len(rb.Steps))
|
|
for _, st := range rb.Steps {
|
|
if strings.TrimSpace(st.Command) == "" {
|
|
continue
|
|
}
|
|
out, code := runObserverScopedCommand(svc, svc.ConfigPath(), st.Command, observerCommandOptions{
|
|
AllowShellEscape: false,
|
|
StatsAutoOnce: true,
|
|
BlockBotRun: true,
|
|
StripANSI: true,
|
|
})
|
|
results = append(results, stepResult{Title: st.Title, Command: st.Command, ExitCode: code, Output: strings.TrimSpace(out)})
|
|
}
|
|
_ = svc.AppendChange("runbook.run", rb.ID, rb.Name)
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"runbook": rb.ID, "results": results})
|
|
for _, r := range results {
|
|
if r.ExitCode != 0 {
|
|
return 1
|
|
}
|
|
}
|
|
return 0
|
|
}
|
|
failed := false
|
|
for i, r := range results {
|
|
fmt.Fprintf(a.out, "%d. %s -> exit %d\n", i+1, r.Title, r.ExitCode)
|
|
if r.Output != "" {
|
|
fmt.Fprintln(a.out, r.Output)
|
|
}
|
|
if r.ExitCode != 0 {
|
|
failed = true
|
|
}
|
|
}
|
|
if failed {
|
|
return 1
|
|
}
|
|
return 0
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown runbook subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func (a *App) runClusterSchedule(svc *cluster.Service, args []string, jsonOut bool) int {
|
|
if len(args) == 0 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster schedule <add|ls|rm|run-due|start|stop|status|logs|worker>")
|
|
return 2
|
|
}
|
|
sub := strings.ToLower(strings.TrimSpace(args[0]))
|
|
switch sub {
|
|
case "ls", "list", "show":
|
|
items, err := svc.SchedulerList()
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule ls: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, items)
|
|
return 0
|
|
}
|
|
if len(items) == 0 {
|
|
fmt.Fprintln(a.out, "(no tasks)")
|
|
return 0
|
|
}
|
|
tw := tabwriter.NewWriter(a.out, 0, 2, 2, ' ', 0)
|
|
fmt.Fprintln(tw, "ID\tNAME\tCLUSTER\tMODE\tEVERY\tBACKOFF\tRETRY\tENABLED\tNEXT_RUN\tCOMMAND")
|
|
for _, t := range items {
|
|
next := "-"
|
|
if !t.NextRunAt.IsZero() {
|
|
next = t.NextRunAt.Format("2006-01-02 15:04:05")
|
|
}
|
|
clusterSel := t.Cluster
|
|
if strings.TrimSpace(clusterSel) == "" {
|
|
clusterSel = "(active)"
|
|
}
|
|
fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%s\t%s\t%d/%d\t%t\t%s\t%s\n",
|
|
t.ID, t.Name, clusterSel, t.Mode, t.Every, t.Backoff, t.RetryCur, t.RetryMax, t.Enabled, next, t.Command)
|
|
}
|
|
_ = tw.Flush()
|
|
return 0
|
|
case "add":
|
|
fs := flag.NewFlagSet("cluster schedule add", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
name := fs.String("name", "", "Task name")
|
|
cmd := fs.String("cmd", "", "Command to execute")
|
|
clusterSel := fs.String("cluster", "", "Cluster selector (default: active cluster)")
|
|
mode := fs.String("mode", "shell", "Execution mode: shell|observer")
|
|
every := fs.String("every", "5m", "Run interval duration, e.g. 5m, 1h")
|
|
backoff := fs.String("backoff", "30s", "Retry backoff base, e.g. 30s")
|
|
jitter := fs.Int("jitter-sec", 5, "Retry jitter in seconds")
|
|
retryMax := fs.Int("retry-max", 3, "Maximum retries before next regular interval")
|
|
enabled := fs.Bool("enabled", true, "Enable task immediately")
|
|
edit := fs.Bool("edit", false, "Open interactive template in $EDITOR")
|
|
from := fs.String("from", "", "Read task JSON from file")
|
|
if err := fs.Parse(args[1:]); err != nil {
|
|
return 2
|
|
}
|
|
addName := strings.TrimSpace(*name)
|
|
addCluster := strings.TrimSpace(*clusterSel)
|
|
addCmd := strings.TrimSpace(*cmd)
|
|
addEvery := strings.TrimSpace(*every)
|
|
addMode := strings.TrimSpace(*mode)
|
|
addBackoff := strings.TrimSpace(*backoff)
|
|
addJitter := *jitter
|
|
addRetryMax := *retryMax
|
|
addEnabled := *enabled
|
|
fromPath := strings.TrimSpace(*from)
|
|
if fromPath != "" {
|
|
raw, err := os.ReadFile(fromPath)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule add: read --from: %v\n", err)
|
|
return 1
|
|
}
|
|
parsed, err := parseScheduleJSON(raw)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule add: invalid --from json: %v\n", err)
|
|
return 1
|
|
}
|
|
addName = strings.TrimSpace(parsed.Name)
|
|
addCluster = strings.TrimSpace(parsed.Cluster)
|
|
addCmd = strings.TrimSpace(parsed.Command)
|
|
addEvery = strings.TrimSpace(parsed.Every)
|
|
addMode = strings.TrimSpace(parsed.Mode)
|
|
addBackoff = strings.TrimSpace(parsed.Backoff)
|
|
addJitter = parsed.JitterSec
|
|
addRetryMax = parsed.RetryMax
|
|
addEnabled = parsed.Enabled
|
|
} else if *edit {
|
|
input := cluster.ScheduledTask{
|
|
Name: addName,
|
|
Cluster: addCluster,
|
|
Command: addCmd,
|
|
Mode: addMode,
|
|
Every: addEvery,
|
|
Backoff: addBackoff,
|
|
JitterSec: addJitter,
|
|
RetryMax: addRetryMax,
|
|
Enabled: addEnabled,
|
|
}
|
|
if strings.TrimSpace(input.Name) == "" {
|
|
input.Name = "new-task"
|
|
}
|
|
if strings.TrimSpace(input.Command) == "" {
|
|
input.Command = "cluster drift"
|
|
}
|
|
if strings.TrimSpace(input.Mode) == "" {
|
|
input.Mode = "observer"
|
|
}
|
|
if strings.TrimSpace(input.Every) == "" {
|
|
input.Every = "30m"
|
|
}
|
|
if strings.TrimSpace(input.Backoff) == "" {
|
|
input.Backoff = "30s"
|
|
}
|
|
if input.JitterSec <= 0 {
|
|
input.JitterSec = 5
|
|
}
|
|
if input.RetryMax <= 0 {
|
|
input.RetryMax = 3
|
|
}
|
|
raw, _ := jsonMarshalIndent(input)
|
|
if isEmbeddedConsoleMode() {
|
|
path, err := writeTemplateFile("pxmon-schedule-*.json", raw)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule add: %v\n", err)
|
|
return 1
|
|
}
|
|
fmt.Fprintf(a.out, "Template saved: %s\n", path)
|
|
fmt.Fprintf(a.out, "Edit it in external terminal and run:\n")
|
|
fmt.Fprintf(a.out, "pxmon cluster schedule add --from %s\n", shellQuoteArg(path))
|
|
return 0
|
|
}
|
|
edited, err := openEditorTemplate("pxmon-schedule-*.json", raw)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule add: %v\n", err)
|
|
return 1
|
|
}
|
|
parsed, err := parseScheduleJSON(edited)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule add: invalid template: %v\n", err)
|
|
return 1
|
|
}
|
|
addName = strings.TrimSpace(parsed.Name)
|
|
addCluster = strings.TrimSpace(parsed.Cluster)
|
|
addCmd = strings.TrimSpace(parsed.Command)
|
|
addEvery = strings.TrimSpace(parsed.Every)
|
|
addMode = strings.TrimSpace(parsed.Mode)
|
|
addBackoff = strings.TrimSpace(parsed.Backoff)
|
|
addJitter = parsed.JitterSec
|
|
addRetryMax = parsed.RetryMax
|
|
addEnabled = parsed.Enabled
|
|
}
|
|
t, err := svc.SchedulerAdd(addName, addCluster, addCmd, addEvery, addMode, addBackoff, addJitter, addRetryMax, addEnabled)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule add: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, t)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Task added: %s (%s)\n", t.Name, t.ID)
|
|
return 0
|
|
case "rm", "remove", "del":
|
|
if len(args) < 2 {
|
|
fmt.Fprintln(a.err, "usage: pxmon cluster schedule rm <task-id-or-name>")
|
|
return 2
|
|
}
|
|
t, err := svc.SchedulerRemove(args[1])
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule rm: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, t)
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "Task removed: %s\n", t.Name)
|
|
return 0
|
|
case "run-due":
|
|
due, err := svc.SchedulerDue(time.Now().UTC())
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule run-due: %v\n", err)
|
|
return 1
|
|
}
|
|
type taskRun struct {
|
|
ID string `json:"id"`
|
|
Name string `json:"name"`
|
|
ExitCode int `json:"exit_code"`
|
|
Output string `json:"output,omitempty"`
|
|
}
|
|
results := make([]taskRun, 0, len(due))
|
|
failed := false
|
|
for _, t := range due {
|
|
out, code := runScheduledTaskNow(svc, t)
|
|
results = append(results, taskRun{ID: t.ID, Name: t.Name, ExitCode: code, Output: strings.TrimSpace(out)})
|
|
_ = svc.SchedulerMarkResult(t.ID, code == 0, time.Now().UTC())
|
|
if code != 0 {
|
|
failed = true
|
|
}
|
|
}
|
|
_ = svc.AppendChange("scheduler.run-due", "tasks", fmt.Sprintf("count=%d", len(results)))
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"ran": len(results), "results": results})
|
|
if failed {
|
|
return 1
|
|
}
|
|
return 0
|
|
}
|
|
if len(results) == 0 {
|
|
fmt.Fprintln(a.out, "No due tasks")
|
|
return 0
|
|
}
|
|
for _, r := range results {
|
|
fmt.Fprintf(a.out, "%s (%s) -> exit %d\n", r.Name, r.ID, r.ExitCode)
|
|
if r.Output != "" {
|
|
fmt.Fprintln(a.out, r.Output)
|
|
}
|
|
}
|
|
if failed {
|
|
return 1
|
|
}
|
|
return 0
|
|
case "start":
|
|
fs := flag.NewFlagSet("cluster schedule start", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
interval := fs.Duration("interval", 30*time.Second, "Worker polling interval")
|
|
if err := fs.Parse(args[1:]); err != nil {
|
|
return 2
|
|
}
|
|
pid, logPath, err := startSchedulerDaemon(svc, *interval)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule start: %v\n", err)
|
|
return 1
|
|
}
|
|
fmt.Fprintf(a.out, "Scheduler started (pid %d)\n", pid)
|
|
fmt.Fprintf(a.out, "Log: %s\n", logPath)
|
|
return 0
|
|
case "stop":
|
|
if err := stopSchedulerDaemon(svc); err != nil {
|
|
fmt.Fprintf(a.err, "schedule stop: %v\n", err)
|
|
return 1
|
|
}
|
|
fmt.Fprintln(a.out, "Scheduler stopped")
|
|
return 0
|
|
case "status":
|
|
pid, running, logPath, err := schedulerDaemonStatus(svc)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule status: %v\n", err)
|
|
return 1
|
|
}
|
|
if jsonOut {
|
|
_ = writeJSON(a.out, map[string]any{"pid": pid, "running": running, "log": logPath})
|
|
return 0
|
|
}
|
|
fmt.Fprintf(a.out, "running: %t\n", running)
|
|
if pid > 0 {
|
|
fmt.Fprintf(a.out, "pid: %d\n", pid)
|
|
}
|
|
fmt.Fprintf(a.out, "log: %s\n", logPath)
|
|
return 0
|
|
case "logs":
|
|
fs := flag.NewFlagSet("cluster schedule logs", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
tail := fs.Int("tail", 200, "Last N lines")
|
|
if err := fs.Parse(args[1:]); err != nil {
|
|
return 2
|
|
}
|
|
_, logPath := schedulerDaemonPaths(svc)
|
|
lines, err := readLastLines(logPath, *tail)
|
|
if err != nil {
|
|
fmt.Fprintf(a.err, "schedule logs: %v\n", err)
|
|
return 1
|
|
}
|
|
for _, ln := range lines {
|
|
fmt.Fprintln(a.out, ln)
|
|
}
|
|
return 0
|
|
case "worker":
|
|
fs := flag.NewFlagSet("cluster schedule worker", flag.ContinueOnError)
|
|
fs.SetOutput(a.err)
|
|
interval := fs.Duration("interval", 30*time.Second, "Polling interval")
|
|
if err := fs.Parse(args[1:]); err != nil {
|
|
return 2
|
|
}
|
|
return runSchedulerWorkerLoop(svc, *interval, a.err)
|
|
default:
|
|
fmt.Fprintf(a.err, "unknown schedule subcommand %q\n", args[0])
|
|
return 2
|
|
}
|
|
}
|
|
|
|
func normalizeCLIPathList(in []string) []string {
|
|
out := make([]string, 0, len(in))
|
|
seen := map[string]struct{}{}
|
|
for _, raw := range in {
|
|
for _, p := range strings.Split(raw, ",") {
|
|
v := strings.TrimSpace(p)
|
|
if v == "" {
|
|
continue
|
|
}
|
|
if _, ok := seen[v]; ok {
|
|
continue
|
|
}
|
|
seen[v] = struct{}{}
|
|
out = append(out, v)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func isEmbeddedConsoleMode() bool {
|
|
return strings.TrimSpace(os.Getenv("PXMON_EMBEDDED_CONSOLE")) == "1"
|
|
}
|
|
|
|
func writeTemplateFile(pattern string, initial []byte) (string, error) {
|
|
f, err := os.CreateTemp("", pattern)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
defer f.Close()
|
|
if len(initial) == 0 {
|
|
initial = []byte("{}\n")
|
|
}
|
|
if _, err := f.Write(initial); err != nil {
|
|
return "", err
|
|
}
|
|
return f.Name(), nil
|
|
}
|
|
|
|
func openEditorTemplate(pattern string, initial []byte) ([]byte, error) {
|
|
editor := strings.TrimSpace(os.Getenv("VISUAL"))
|
|
if editor == "" {
|
|
editor = strings.TrimSpace(os.Getenv("EDITOR"))
|
|
}
|
|
if editor == "" {
|
|
editor = "vi"
|
|
}
|
|
f, err := os.CreateTemp("", pattern)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
tmpPath := f.Name()
|
|
defer os.Remove(tmpPath)
|
|
if len(initial) == 0 {
|
|
initial = []byte("{}\n")
|
|
}
|
|
if _, err := f.Write(initial); err != nil {
|
|
_ = f.Close()
|
|
return nil, err
|
|
}
|
|
if err := f.Close(); err != nil {
|
|
return nil, err
|
|
}
|
|
cmd := exec.Command("sh", "-lc", shellQuoteArg(editor)+" "+shellQuoteArg(tmpPath))
|
|
cmd.Stdin = os.Stdin
|
|
cmd.Stdout = os.Stdout
|
|
cmd.Stderr = os.Stderr
|
|
if err := cmd.Run(); err != nil {
|
|
return nil, fmt.Errorf("open editor %q: %w", editor, err)
|
|
}
|
|
b, err := os.ReadFile(tmpPath)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if strings.TrimSpace(string(b)) == "" {
|
|
return nil, fmt.Errorf("template is empty")
|
|
}
|
|
return b, nil
|
|
}
|
|
|
|
func parseRunbookJSON(raw []byte) (cluster.Runbook, error) {
|
|
var rb cluster.Runbook
|
|
if err := json.Unmarshal(raw, &rb); err != nil {
|
|
return cluster.Runbook{}, err
|
|
}
|
|
rb.ID = strings.TrimSpace(rb.ID)
|
|
rb.Name = strings.TrimSpace(rb.Name)
|
|
rb.Description = strings.TrimSpace(rb.Description)
|
|
outSteps := make([]cluster.RunbookStep, 0, len(rb.Steps))
|
|
for _, st := range rb.Steps {
|
|
title := strings.TrimSpace(st.Title)
|
|
cmd := strings.TrimSpace(st.Command)
|
|
note := strings.TrimSpace(st.Note)
|
|
if title == "" {
|
|
continue
|
|
}
|
|
outSteps = append(outSteps, cluster.RunbookStep{Title: title, Command: cmd, Note: note})
|
|
}
|
|
rb.Steps = outSteps
|
|
return rb, nil
|
|
}
|
|
|
|
func parseScheduleJSON(raw []byte) (cluster.ScheduledTask, error) {
|
|
var t cluster.ScheduledTask
|
|
if err := json.Unmarshal(raw, &t); err != nil {
|
|
return cluster.ScheduledTask{}, err
|
|
}
|
|
t.Name = strings.TrimSpace(t.Name)
|
|
t.Cluster = strings.TrimSpace(t.Cluster)
|
|
t.Command = strings.TrimSpace(t.Command)
|
|
t.Mode = strings.ToLower(strings.TrimSpace(t.Mode))
|
|
t.Every = strings.TrimSpace(t.Every)
|
|
t.Backoff = strings.TrimSpace(t.Backoff)
|
|
if t.JitterSec < 0 {
|
|
t.JitterSec = 0
|
|
}
|
|
if t.RetryMax <= 0 {
|
|
t.RetryMax = 3
|
|
}
|
|
return t, nil
|
|
}
|
|
|
|
func shellQuoteArg(v string) string {
|
|
if v == "" {
|
|
return "''"
|
|
}
|
|
return "'" + strings.ReplaceAll(v, "'", `'\''`) + "'"
|
|
}
|
|
|
|
func runScheduledTaskNow(svc *cluster.Service, t cluster.ScheduledTask) (string, int) {
|
|
switch strings.ToLower(strings.TrimSpace(t.Mode)) {
|
|
case "observer":
|
|
return runObserverScopedCommand(svc, svc.ConfigPath(), t.Command, observerCommandOptions{
|
|
AllowShellEscape: false,
|
|
StatsAutoOnce: true,
|
|
BlockBotRun: true,
|
|
StripANSI: true,
|
|
})
|
|
default:
|
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
|
|
defer cancel()
|
|
out, err := svc.RunRemoteShell(ctx, t.Cluster, t.Command)
|
|
if err != nil {
|
|
return err.Error(), 1
|
|
}
|
|
return out, 0
|
|
}
|
|
}
|
|
|
|
func schedulerDaemonPaths(svc *cluster.Service) (string, string) {
|
|
base := filepath.Join(svc.DataDir(), "scheduler")
|
|
return filepath.Join(base, "daemon.pid"), filepath.Join(base, "daemon.log")
|
|
}
|
|
|
|
func startSchedulerDaemon(svc *cluster.Service, interval time.Duration) (int, string, error) {
|
|
if interval < 5*time.Second {
|
|
interval = 5 * time.Second
|
|
}
|
|
if err := stopSchedulerDaemon(svc); err != nil {
|
|
_ = err
|
|
}
|
|
pidPath, logPath := schedulerDaemonPaths(svc)
|
|
if err := os.MkdirAll(filepath.Dir(pidPath), 0o700); err != nil {
|
|
return 0, "", err
|
|
}
|
|
logf, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600)
|
|
if err != nil {
|
|
return 0, "", err
|
|
}
|
|
defer logf.Close()
|
|
exe, err := os.Executable()
|
|
if err != nil {
|
|
return 0, "", err
|
|
}
|
|
args := []string{"cluster", "schedule", "worker", "--interval", interval.String()}
|
|
cmd := exec.Command(exe, args...)
|
|
cmd.Stdout = logf
|
|
cmd.Stderr = logf
|
|
cmd.Stdin = nil
|
|
if err := cmd.Start(); err != nil {
|
|
return 0, "", err
|
|
}
|
|
pid := cmd.Process.Pid
|
|
_ = os.WriteFile(pidPath, []byte(fmt.Sprintf("%d\n", pid)), 0o600)
|
|
return pid, logPath, nil
|
|
}
|
|
|
|
func stopSchedulerDaemon(svc *cluster.Service) error {
|
|
pidPath, _ := schedulerDaemonPaths(svc)
|
|
raw, err := os.ReadFile(pidPath)
|
|
if err != nil {
|
|
if os.IsNotExist(err) {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
pid, _ := strconv.Atoi(strings.TrimSpace(string(raw)))
|
|
if pid > 0 {
|
|
proc, err := os.FindProcess(pid)
|
|
if err == nil {
|
|
_ = proc.Kill()
|
|
}
|
|
}
|
|
_ = os.Remove(pidPath)
|
|
return nil
|
|
}
|
|
|
|
func schedulerDaemonStatus(svc *cluster.Service) (int, bool, string, error) {
|
|
pidPath, logPath := schedulerDaemonPaths(svc)
|
|
raw, err := os.ReadFile(pidPath)
|
|
if err != nil {
|
|
if os.IsNotExist(err) {
|
|
return 0, false, logPath, nil
|
|
}
|
|
return 0, false, logPath, err
|
|
}
|
|
pid, _ := strconv.Atoi(strings.TrimSpace(string(raw)))
|
|
if pid <= 0 {
|
|
return 0, false, logPath, nil
|
|
}
|
|
if err := exec.Command("ps", "-p", strconv.Itoa(pid), "-o", "pid=").Run(); err != nil {
|
|
return pid, false, logPath, nil
|
|
}
|
|
return pid, true, logPath, nil
|
|
}
|
|
|
|
func runSchedulerWorkerLoop(svc *cluster.Service, interval time.Duration, w io.Writer) int {
|
|
if interval < 5*time.Second {
|
|
interval = 5 * time.Second
|
|
}
|
|
ticker := time.NewTicker(interval)
|
|
defer ticker.Stop()
|
|
for {
|
|
due, err := svc.SchedulerDue(time.Now().UTC())
|
|
if err != nil {
|
|
fmt.Fprintf(w, "scheduler worker: list due: %v\n", err)
|
|
} else {
|
|
for _, t := range due {
|
|
out, code := runScheduledTaskNow(svc, t)
|
|
if code != 0 {
|
|
fmt.Fprintf(w, "scheduler worker: task %s failed (%d): %s\n", t.Name, code, strings.TrimSpace(out))
|
|
} else if strings.TrimSpace(out) != "" {
|
|
fmt.Fprintf(w, "scheduler worker: task %s ok: %s\n", t.Name, strings.TrimSpace(out))
|
|
} else {
|
|
fmt.Fprintf(w, "scheduler worker: task %s ok\n", t.Name)
|
|
}
|
|
_ = svc.SchedulerMarkResult(t.ID, code == 0, time.Now().UTC())
|
|
}
|
|
}
|
|
<-ticker.C
|
|
}
|
|
}
|