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] }