453 lines
11 KiB
Go
453 lines
11 KiB
Go
package main
|
|
|
|
import (
|
|
"bufio"
|
|
"context"
|
|
"encoding/json"
|
|
"flag"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"os/exec"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
type reportRow struct {
|
|
Status string `json:"status"`
|
|
Area string `json:"area"`
|
|
Check string `json:"check"`
|
|
Summary string `json:"summary"`
|
|
Detail string `json:"detail"`
|
|
Fix string `json:"fix"`
|
|
Graph string `json:"graph"`
|
|
Query string `json:"query"`
|
|
}
|
|
|
|
func main() {
|
|
if len(os.Args) < 2 {
|
|
usage(os.Stderr)
|
|
os.Exit(2)
|
|
}
|
|
|
|
switch os.Args[1] {
|
|
case "report-render":
|
|
if err := runReportRender(os.Args[2:], os.Stdin, os.Stdout); err != nil {
|
|
fmt.Fprintln(os.Stderr, err)
|
|
os.Exit(1)
|
|
}
|
|
case "run-dag":
|
|
if err := runDAG(os.Args[2:], os.Stdout, os.Stderr); err != nil {
|
|
fmt.Fprintln(os.Stderr, err)
|
|
os.Exit(1)
|
|
}
|
|
case "help", "-h", "--help":
|
|
usage(os.Stdout)
|
|
default:
|
|
usage(os.Stderr)
|
|
os.Exit(2)
|
|
}
|
|
}
|
|
|
|
func usage(out io.Writer) {
|
|
fmt.Fprintln(out, "Usage: jeannie-core {report-render|run-dag}")
|
|
}
|
|
|
|
type dagSpec struct {
|
|
Tasks []dagTask `json:"tasks"`
|
|
Parallelism int `json:"parallelism"`
|
|
}
|
|
|
|
type dagTask struct {
|
|
ID string `json:"id"`
|
|
Title string `json:"title"`
|
|
Command []string `json:"command"`
|
|
Needs []string `json:"needs"`
|
|
}
|
|
|
|
type dagTaskResult struct {
|
|
id string
|
|
exitCode int
|
|
err error
|
|
}
|
|
|
|
func runDAG(args []string, out io.Writer, errOut io.Writer) error {
|
|
flags := flag.NewFlagSet("run-dag", flag.ContinueOnError)
|
|
flags.SetOutput(io.Discard)
|
|
specPath := flags.String("spec", "", "")
|
|
if err := flags.Parse(args); err != nil {
|
|
return err
|
|
}
|
|
if *specPath == "" {
|
|
return fmt.Errorf("run-dag requires --spec")
|
|
}
|
|
|
|
specData, err := os.ReadFile(*specPath)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var spec dagSpec
|
|
if err := json.Unmarshal(specData, &spec); err != nil {
|
|
return err
|
|
}
|
|
return executeDAG(context.Background(), spec, out, errOut)
|
|
}
|
|
|
|
func executeDAG(ctx context.Context, spec dagSpec, out io.Writer, errOut io.Writer) error {
|
|
if spec.Parallelism <= 0 {
|
|
spec.Parallelism = 4
|
|
}
|
|
|
|
taskByID := map[string]dagTask{}
|
|
dependents := map[string][]string{}
|
|
remainingNeeds := map[string]int{}
|
|
for _, task := range spec.Tasks {
|
|
if task.ID == "" {
|
|
return fmt.Errorf("dag task has empty id")
|
|
}
|
|
if _, exists := taskByID[task.ID]; exists {
|
|
return fmt.Errorf("duplicate dag task id %q", task.ID)
|
|
}
|
|
if len(task.Command) == 0 {
|
|
return fmt.Errorf("dag task %q has empty command", task.ID)
|
|
}
|
|
taskByID[task.ID] = task
|
|
remainingNeeds[task.ID] = len(task.Needs)
|
|
for _, need := range task.Needs {
|
|
dependents[need] = append(dependents[need], task.ID)
|
|
}
|
|
}
|
|
for _, task := range spec.Tasks {
|
|
for _, need := range task.Needs {
|
|
if _, exists := taskByID[need]; !exists {
|
|
return fmt.Errorf("dag task %q depends on unknown task %q", task.ID, need)
|
|
}
|
|
}
|
|
}
|
|
|
|
ctx, cancel := context.WithCancel(ctx)
|
|
defer cancel()
|
|
|
|
ready := make(chan string, len(spec.Tasks))
|
|
results := make(chan dagTaskResult, len(spec.Tasks))
|
|
sem := make(chan struct{}, spec.Parallelism)
|
|
var writersMu sync.Mutex
|
|
var wg sync.WaitGroup
|
|
started := map[string]bool{}
|
|
completed := map[string]bool{}
|
|
|
|
for id, count := range remainingNeeds {
|
|
if count == 0 {
|
|
ready <- id
|
|
}
|
|
}
|
|
|
|
failed := false
|
|
for len(completed) < len(spec.Tasks) {
|
|
select {
|
|
case id := <-ready:
|
|
if started[id] || failed {
|
|
continue
|
|
}
|
|
started[id] = true
|
|
task := taskByID[id]
|
|
wg.Add(1)
|
|
go func() {
|
|
defer wg.Done()
|
|
sem <- struct{}{}
|
|
defer func() { <-sem }()
|
|
results <- runDAGTask(ctx, task, out, errOut, &writersMu)
|
|
}()
|
|
case result := <-results:
|
|
completed[result.id] = true
|
|
if result.err != nil || result.exitCode != 0 {
|
|
failed = true
|
|
cancel()
|
|
wg.Wait()
|
|
if result.err != nil {
|
|
return result.err
|
|
}
|
|
return fmt.Errorf("dag task %q failed with exit code %d", result.id, result.exitCode)
|
|
}
|
|
for _, dependent := range dependents[result.id] {
|
|
remainingNeeds[dependent]--
|
|
if remainingNeeds[dependent] == 0 {
|
|
ready <- dependent
|
|
}
|
|
}
|
|
}
|
|
}
|
|
wg.Wait()
|
|
return nil
|
|
}
|
|
|
|
func runDAGTask(ctx context.Context, task dagTask, out io.Writer, errOut io.Writer, writersMu *sync.Mutex) dagTaskResult {
|
|
title := task.Title
|
|
if title == "" {
|
|
title = task.ID
|
|
}
|
|
started := time.Now()
|
|
writeDAGLine(out, writersMu, task.ID, "started "+title)
|
|
|
|
cmd := exec.CommandContext(ctx, task.Command[0], task.Command[1:]...)
|
|
cmd.Env = os.Environ()
|
|
stdout, err := cmd.StdoutPipe()
|
|
if err != nil {
|
|
return dagTaskResult{id: task.ID, exitCode: 1, err: err}
|
|
}
|
|
cmd.Stderr = cmd.Stdout
|
|
if err := cmd.Start(); err != nil {
|
|
return dagTaskResult{id: task.ID, exitCode: 1, err: err}
|
|
}
|
|
|
|
scanner := bufio.NewScanner(stdout)
|
|
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
|
|
for scanner.Scan() {
|
|
writeDAGLine(out, writersMu, task.ID, scanner.Text())
|
|
}
|
|
if scanErr := scanner.Err(); scanErr != nil {
|
|
writeDAGLine(errOut, writersMu, task.ID, "output read error: "+scanErr.Error())
|
|
}
|
|
|
|
err = cmd.Wait()
|
|
elapsed := time.Since(started).Round(time.Second)
|
|
exitCode := 0
|
|
if err != nil {
|
|
exitCode = 1
|
|
if exitErr, ok := err.(*exec.ExitError); ok {
|
|
exitCode = exitErr.ExitCode()
|
|
}
|
|
writeDAGLine(out, writersMu, task.ID, fmt.Sprintf("failed after %s", elapsed))
|
|
return dagTaskResult{id: task.ID, exitCode: exitCode, err: nil}
|
|
}
|
|
writeDAGLine(out, writersMu, task.ID, fmt.Sprintf("completed in %s", elapsed))
|
|
return dagTaskResult{id: task.ID, exitCode: 0, err: nil}
|
|
}
|
|
|
|
func writeDAGLine(out io.Writer, writersMu *sync.Mutex, id string, line string) {
|
|
writersMu.Lock()
|
|
defer writersMu.Unlock()
|
|
fmt.Fprintf(out, "[%s] %s\n", id, line)
|
|
}
|
|
|
|
func runReportRender(args []string, in io.Reader, out io.Writer) error {
|
|
flags := flag.NewFlagSet("report-render", flag.ContinueOnError)
|
|
flags.SetOutput(io.Discard)
|
|
title := flags.String("title", envDefault("REPORT_TITLE", "Jeannie Report"), "")
|
|
details := flags.Bool("details", false, "")
|
|
only := flags.String("only", "all", "")
|
|
asJSON := flags.Bool("json", false, "")
|
|
if err := flags.Parse(args); err != nil {
|
|
usage(os.Stderr)
|
|
return err
|
|
}
|
|
|
|
rows, err := parseReportRows(in)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
rows = filterReportRows(rows, *only)
|
|
|
|
if *asJSON {
|
|
encoder := json.NewEncoder(out)
|
|
encoder.SetIndent("", " ")
|
|
return encoder.Encode(map[string]any{"title": *title, "rows": rows})
|
|
}
|
|
|
|
renderReportText(out, *title, rows, *details, *only)
|
|
return nil
|
|
}
|
|
|
|
func envDefault(name, fallback string) string {
|
|
if value := os.Getenv(name); value != "" {
|
|
return value
|
|
}
|
|
return fallback
|
|
}
|
|
|
|
func parseReportRows(in io.Reader) ([]reportRow, error) {
|
|
data, err := io.ReadAll(in)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
var rows []reportRow
|
|
for _, raw := range strings.Split(string(data), "\n") {
|
|
raw = strings.TrimSuffix(raw, "\r")
|
|
if strings.TrimSpace(raw) == "" {
|
|
continue
|
|
}
|
|
parts := strings.Split(raw, "\t")
|
|
for len(parts) < 8 {
|
|
parts = append(parts, "")
|
|
}
|
|
|
|
row := reportRow{
|
|
Status: normalizeStatus(parts[0]),
|
|
Area: strings.TrimSpace(parts[1]),
|
|
Check: strings.TrimSpace(parts[2]),
|
|
Summary: strings.TrimSpace(parts[3]),
|
|
Detail: strings.TrimSpace(parts[4]),
|
|
Fix: strings.TrimSpace(parts[5]),
|
|
Graph: strings.TrimSpace(parts[6]),
|
|
Query: strings.TrimSpace(parts[7]),
|
|
}
|
|
if row.Area == "" {
|
|
row.Area = "General"
|
|
}
|
|
if row.Status == "ok" {
|
|
row.Graph = ""
|
|
row.Query = ""
|
|
}
|
|
rows = append(rows, row)
|
|
}
|
|
return rows, nil
|
|
}
|
|
|
|
func normalizeStatus(value string) string {
|
|
status := strings.ToLower(strings.TrimSpace(value))
|
|
if status == "" {
|
|
return "skip"
|
|
}
|
|
return status
|
|
}
|
|
|
|
func filterReportRows(rows []reportRow, only string) []reportRow {
|
|
var filtered []reportRow
|
|
for _, row := range rows {
|
|
if includeReportRow(row, only) {
|
|
filtered = append(filtered, row)
|
|
}
|
|
}
|
|
return filtered
|
|
}
|
|
|
|
func includeReportRow(row reportRow, only string) bool {
|
|
switch only {
|
|
case "all":
|
|
return true
|
|
case "failures":
|
|
return row.Status == "fail"
|
|
case "warnings":
|
|
return row.Status == "warn"
|
|
case "problems":
|
|
return row.Status == "fail" || row.Status == "warn"
|
|
default:
|
|
return true
|
|
}
|
|
}
|
|
|
|
func renderReportText(out io.Writer, title string, rows []reportRow, details bool, only string) {
|
|
counts := map[string]int{"fail": 0, "warn": 0, "ok": 0, "skip": 0}
|
|
for _, row := range rows {
|
|
counts[row.Status]++
|
|
}
|
|
|
|
fmt.Fprintln(out, title)
|
|
fmt.Fprintln(out, strings.Repeat("=", len(title)))
|
|
fmt.Fprintf(out, "fail=%d warn=%d ok=%d skip=%d\n", counts["fail"], counts["warn"], counts["ok"], counts["skip"])
|
|
|
|
if len(rows) == 0 {
|
|
switch only {
|
|
case "problems":
|
|
fmt.Fprintln(out, "\nNo failing or warning items.")
|
|
case "failures":
|
|
fmt.Fprintln(out, "\nNo failing items.")
|
|
case "warnings":
|
|
fmt.Fprintln(out, "\nNo warning items.")
|
|
default:
|
|
fmt.Fprintln(out, "\nNo report rows.")
|
|
}
|
|
return
|
|
}
|
|
|
|
for _, area := range orderedReportAreas(rows) {
|
|
fmt.Fprintf(out, "\n== %s ==\n", area)
|
|
areaRows := rowsForReportArea(rows, area)
|
|
sort.SliceStable(areaRows, func(i, j int) bool {
|
|
leftRank := reportStatusRank(areaRows[i].Status)
|
|
rightRank := reportStatusRank(areaRows[j].Status)
|
|
if leftRank != rightRank {
|
|
return leftRank < rightRank
|
|
}
|
|
return areaRows[i].Check < areaRows[j].Check
|
|
})
|
|
|
|
for _, row := range areaRows {
|
|
fmt.Fprintf(out, "%-5s %-30s %s\n", row.Status, truncate(row.Check, 30), row.Summary)
|
|
if details || row.Status == "fail" || row.Status == "warn" {
|
|
if row.Detail != "" {
|
|
fmt.Fprintf(out, " detail: %s\n", row.Detail)
|
|
}
|
|
if row.Fix != "" {
|
|
fmt.Fprintf(out, " fix: %s\n", row.Fix)
|
|
}
|
|
if row.Graph != "" {
|
|
fmt.Fprintf(out, " graph: %s\n", row.Graph)
|
|
}
|
|
if row.Query != "" {
|
|
fmt.Fprintf(out, " query: %s\n", row.Query)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
for _, row := range rows {
|
|
if row.Status == "fail" || row.Status == "warn" {
|
|
fmt.Fprintln(out, "\nNext Fix")
|
|
if row.Fix != "" {
|
|
fmt.Fprintln(out, row.Fix)
|
|
} else {
|
|
fmt.Fprintf(out, "Investigate %s / %s\n", row.Area, row.Check)
|
|
}
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func orderedReportAreas(rows []reportRow) []string {
|
|
seen := map[string]bool{}
|
|
var areas []string
|
|
for _, row := range rows {
|
|
if !seen[row.Area] {
|
|
seen[row.Area] = true
|
|
areas = append(areas, row.Area)
|
|
}
|
|
}
|
|
return areas
|
|
}
|
|
|
|
func rowsForReportArea(rows []reportRow, area string) []reportRow {
|
|
var areaRows []reportRow
|
|
for _, row := range rows {
|
|
if row.Area == area {
|
|
areaRows = append(areaRows, row)
|
|
}
|
|
}
|
|
return areaRows
|
|
}
|
|
|
|
func reportStatusRank(status string) int {
|
|
switch status {
|
|
case "fail":
|
|
return 0
|
|
case "warn":
|
|
return 1
|
|
case "skip":
|
|
return 2
|
|
case "ok":
|
|
return 3
|
|
default:
|
|
return 9
|
|
}
|
|
}
|
|
|
|
func truncate(value string, max int) string {
|
|
if len(value) <= max {
|
|
return value
|
|
}
|
|
return value[:max]
|
|
}
|