// Command harness is the MOPAC harness CLI. v0 surface: `harness once`, // `harness loop`, `harness events`. package main import ( "context" "errors" "flag" "fmt" "net/http" "os" "os/signal" "sort" "strings" "syscall" "time" "ukrrs.com/mopac/harness/internal/config" "ukrrs.com/mopac/harness/internal/events" "ukrrs.com/mopac/harness/internal/loop" "ukrrs.com/mopac/harness/internal/quota" "ukrrs.com/mopac/harness/internal/serve" ) const usage = `MOPAC harness (v0) Usage: harness once [-config PATH] [-dry-run] [-demo] [-task-id ID] harness loop [-config PATH] [-interval DUR] [-once] [-dry-run] harness events [-config PATH] [-listen ADDR] harness serve [-config PATH] [-listen ADDR] harness quota [-config PATH] once runs ONE conductor iteration and exits (chain by re-invoking; no daemon): intake (Redmine scope, or the [demo] issue) -> plan/model routing -> bounded turn via LiteLLM -> REPORT file writeback. loop is the self-hosting daemon: Redmine is the SoR, the loop is the worker. It polls the intake on an interval, runs ONE bounded turn per new/updated issue (sequentially), notes the REPORT back on the issue, transitions status per [redmine.status_map], optionally commits the REPORT to gitea ([gitea] commit_reports). State is append-only loop.jsonl, dedup by issue id + updated_on. --once = single scan (cron-able). SIGINT stops. events runs the webhook receiver until SIGINT/SIGTERM: Redmine/Discourse/ Gitea webhooks are verified, normalized, stored append-only (dedup by provider event id) and handed to the conductor (dispatch stub for now). Routes: POST /hooks/{redmine,discourse,gitea}, GET /healthz. serve runs the OpenAI-compatible front door until SIGINT/SIGTERM (the OpenWebUI connection): GET /v1/models lists the servable models (one per [models.classes] class, named mopac-), POST /v1/chat/completions runs ONE bounded stateless conductor turn over the sent history and returns the final text + usage. Bearer vkey auth ([serve] vkey_ref); non-streaming v0; tools off. Own port - coexists with events. Routes: POST /v1/chat/completions, GET /v1/models, GET /healthz. quota is the Redmine 490/491 gate surface: status one-shot: quota snapshot (polled or estimated), peak window, host resources, and the per-class token+credit usage table from the loop state (the Discourse usage-report feed) probe poll [quota] usage_url once and print the parsed buckets (or the raw error; bearer key never printed) gate evaluate the back-pressure gates NOW: the decision for every [models.classes] class + resource thresholds Flags: -config PATH config file (default $HARNESS_CONFIG or ./harness.toml) -dry-run once/loop: intake + plan only; no LLM call, no REPORT, no state writes -demo once: run the [demo] issue instead of Redmine intake -task-id ID once: run only the task/issue with this id -interval DUR loop: poll interval (overrides [loop] poll_interval_secs) -once loop: single scan then exit (cron-able) -listen ADDR events/serve: bind address (overrides [events]/[serve] listen) Exit codes: 0 ok (including "no tasks in scope"; loop: clean SIGINT stop) 1 usage / config / routing / writeback error 2 intake error 4 llm / turn error ` func main() { os.Exit(run(os.Args[1:])) } func run(args []string) int { if len(args) == 0 { fmt.Fprint(os.Stderr, usage) return 1 } switch args[0] { case "help", "-h", "--help": fmt.Print(usage) return 0 case "once": return runOnce(args[1:]) case "loop": return runLoop(args[1:]) case "events": return runEvents(args[1:]) case "serve": return runServe(args[1:]) case "quota": return runQuota(args[1:]) default: fmt.Fprintf(os.Stderr, "harness: unknown command %q\n\n%s", args[0], usage) return 1 } } func runOnce(args []string) int { fs := flag.NewFlagSet("once", flag.ContinueOnError) cfgPath := fs.String("config", "", "config file path") dryRun := fs.Bool("dry-run", false, "intake + plan only, no LLM call") demo := fs.Bool("demo", false, "use the [demo] issue instead of Redmine") taskID := fs.String("task-id", "", "run only this task id") if err := fs.Parse(args); err != nil { return 1 } if fs.NArg() > 0 { fmt.Fprintf(os.Stderr, "harness: unexpected argument %q\n", fs.Arg(0)) return 1 } if *cfgPath == "" { *cfgPath = os.Getenv("HARNESS_CONFIG") } if *cfgPath == "" { *cfgPath = "harness.toml" } cfg, err := config.Load(*cfgPath) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } conductor, err := loop.New(cfg, os.Stdout) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() _, err = conductor.Once(ctx, loop.OnceOpts{ DryRun: *dryRun, Demo: *demo, TaskID: *taskID, }) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) switch { case errors.Is(err, loop.ErrIntake): return 2 case errors.Is(err, loop.ErrLLM): return 4 default: return 1 } } return 0 } func runLoop(args []string) int { fs := flag.NewFlagSet("loop", flag.ContinueOnError) cfgPath := fs.String("config", "", "config file path") interval := fs.Duration("interval", 0, "poll interval (overrides [loop] poll_interval_secs)") once := fs.Bool("once", false, "single scan then exit (cron-able)") dryRun := fs.Bool("dry-run", false, "scan and print what would dispatch; no turns, no state writes") if err := fs.Parse(args); err != nil { return 1 } if fs.NArg() > 0 { fmt.Fprintf(os.Stderr, "harness: unexpected argument %q\n", fs.Arg(0)) return 1 } if *cfgPath == "" { *cfgPath = os.Getenv("HARNESS_CONFIG") } if *cfgPath == "" { *cfgPath = "harness.toml" } cfg, err := config.Load(*cfgPath) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } conductor, err := loop.New(cfg, os.Stdout) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() if err := conductor.RunLoop(ctx, loop.LoopOpts{ Interval: *interval, Once: *once, DryRun: *dryRun, }); err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) switch { case errors.Is(err, loop.ErrIntake): return 2 default: return 1 } } return 0 } func runQuota(args []string) int { sub := "status" var rest []string if len(args) > 0 && !strings.HasPrefix(args[0], "-") { sub = args[0] rest = args[1:] } switch sub { case "status", "probe", "gate": default: fmt.Fprintf(os.Stderr, "harness: quota: unknown subcommand %q (want status|probe|gate)\n", sub) return 1 } fs := flag.NewFlagSet("quota "+sub, flag.ContinueOnError) cfgPath := fs.String("config", "", "config file path") if err := fs.Parse(rest); err != nil { return 1 } if fs.NArg() > 0 { fmt.Fprintf(os.Stderr, "harness: unexpected argument %q\n", fs.Arg(0)) return 1 } if *cfgPath == "" { *cfgPath = os.Getenv("HARNESS_CONFIG") } if *cfgPath == "" { *cfgPath = "harness.toml" } cfg, err := config.Load(*cfgPath) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } conductor, err := loop.New(cfg, os.Stdout) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } gate := conductor.Gate() switch sub { case "probe": if gate == nil || cfg.Quota.UsageURL == "" { fmt.Println("harness: quota: no [quota] usage_url configured (running on estimates)") return 0 } snap := gate.Snapshot(context.Background()) fmt.Printf("harness: quota: account=%s source=%s fetched=%s\n", snap.Account, snap.Source, snap.FetchedAt.Format(time.RFC3339)) for _, b := range snap.Buckets { reset := "-" if !b.WindowReset.IsZero() { reset = b.WindowReset.Format(time.RFC3339) } fmt.Printf(" %-7s used=%.0f limit=%.0f (%.1f%%) resets=%s\n", b.ID, b.Used, b.Limit, b.UsedPct(), reset) } return 0 case "gate": if gate == nil { fmt.Println("harness: quota gate off ([quota] enabled = false)") return 0 } fmt.Printf("harness: quota gate: %s\n", conductor.GateStatusLine(context.Background())) for _, class := range sortedClasses(cfg) { d := gate.Decide(context.Background(), class) tag := "ALLOW" if d.Action == quota.ActionDefer { tag = "DEFER" } fmt.Printf(" %-6s class=%-12s %s\n", tag, class, d.Reason) } return 0 } // status if gate == nil && !cfg.Resources.Enabled { fmt.Println("harness: quota gate off ([quota]/[resources] enabled = false)") } else if gate != nil { fmt.Printf("harness: quota: %s\n", conductor.GateStatusLine(context.Background())) fmt.Printf(" peak window: %s-%s %s weekdays_only=%v currently=%v\n", cfg.Quota.PeakStart, cfg.Quota.PeakEnd, cfg.Quota.Timezone, cfg.Quota.PeakWeekdaysOnly, gate.InPeak()) } if st, ok := conductor.ResourceSample(); ok { fmt.Printf("harness: resources: load=%.2f memAvail=%.0fMB diskFree=%.0fMB ioDelay=%.1f%%(psi=%v)\n", st.LoadAvg1m, st.MemAvailable, st.DiskFree, st.IODelayPct, st.HasPSI) } report, err := loop.UsageReport(cfg.Loop.StateDir) if err != nil { fmt.Fprintf(os.Stderr, "harness: usage report: %v\n", err) return 0 } fmt.Printf("harness: usage by class (%s):\n%s", cfg.Loop.StateDir, report) return 0 } func sortedClasses(cfg *config.Config) []string { out := make([]string, 0, len(cfg.Models.Classes)) for class := range cfg.Models.Classes { out = append(out, class) } sort.Strings(out) return out } func runEvents(args []string) int { fs := flag.NewFlagSet("events", flag.ContinueOnError) cfgPath := fs.String("config", "", "config file path") listen := fs.String("listen", "", "bind address (overrides [events] listen)") if err := fs.Parse(args); err != nil { return 1 } if fs.NArg() > 0 { fmt.Fprintf(os.Stderr, "harness: unexpected argument %q\n", fs.Arg(0)) return 1 } if *cfgPath == "" { *cfgPath = os.Getenv("HARNESS_CONFIG") } if *cfgPath == "" { *cfgPath = "harness.toml" } cfg, err := config.Load(*cfgPath) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } // The conductor is the event dispatcher (stub until phase 3 wiring). conductor, err := loop.New(cfg, os.Stdout) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } srv, err := events.NewServer(cfg.Events, conductor, os.Stdout) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } defer srv.Close() if *listen == "" { *listen = cfg.Events.Listen } ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() httpSrv := &http.Server{ Addr: *listen, Handler: srv.Handler(), ReadHeaderTimeout: 10 * time.Second, } go func() { <-ctx.Done() shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() _ = httpSrv.Shutdown(shutdownCtx) }() fmt.Printf("harness: events receiver on %s (state: %s)\n", *listen, cfg.Events.StateDir) err = httpSrv.ListenAndServe() if err != nil && !errors.Is(err, http.ErrServerClosed) { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } fmt.Printf("harness: events receiver stopped\n") return 0 } func runServe(args []string) int { fs := flag.NewFlagSet("serve", flag.ContinueOnError) cfgPath := fs.String("config", "", "config file path") listen := fs.String("listen", "", "bind address (overrides [serve] listen)") if err := fs.Parse(args); err != nil { return 1 } if fs.NArg() > 0 { fmt.Fprintf(os.Stderr, "harness: unexpected argument %q\n", fs.Arg(0)) return 1 } if *cfgPath == "" { *cfgPath = os.Getenv("HARNESS_CONFIG") } if *cfgPath == "" { *cfgPath = "harness.toml" } cfg, err := config.Load(*cfgPath) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } // The conductor owns the turn machinery and the model router the // catalog is built from. conductor, err := loop.New(cfg, os.Stdout) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } srv, err := serve.NewServer(cfg.Serve, conductor.Router(), conductor, os.Stdout) if err != nil { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } if *listen == "" { *listen = cfg.Serve.Listen } ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) defer stop() httpSrv := &http.Server{ Addr: *listen, Handler: srv.Handler(), ReadHeaderTimeout: 10 * time.Second, } go func() { <-ctx.Done() shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() _ = httpSrv.Shutdown(shutdownCtx) }() fmt.Printf("harness: serve (OpenAI-compatible) on %s (models: %s)\n", *listen, strings.Join(srv.ModelNames(), ", ")) err = httpSrv.ListenAndServe() if err != nil && !errors.Is(err, http.ErrServerClosed) { fmt.Fprintf(os.Stderr, "harness: %v\n", err) return 1 } fmt.Printf("harness: serve stopped\n") return 0 }