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 ", "cluster runbook list|show|run ", "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 [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 [name-or-id] --vm [--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 [name-or-id] --vm [--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 [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 [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 [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 [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 ") 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 ") 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 ") 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 ") 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 ") 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 ") 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 ") 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 ") 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 [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 ") 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 ") 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 ") 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 ") 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 ") 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 } }