Files
mrcharles c9e86eefa8 serve: OpenAI-compatible front door (harness serve)
OpenWebUI becomes the interactive surface by talking to MOPAC like any
OpenAI provider: GET /v1/models lists the servable catalog (one model per
[models.classes] class, named mopac-<class>, routed through the same
tier table `once` uses; unknown model = 400 naming the valid ones) and
POST /v1/chat/completions runs ONE bounded stateless conductor turn over
the sent conversation history — no session storage, tools hard-off,
non-streaming (stream:true gets an explicit 400; OWUI tolerates
non-streaming providers). Bearer vkey auth compares SHA-256 digests in
constant time; missing/wrong keys get one byte-identical 401 body, and
the vkey never reaches logs or responses. temperature/max_tokens are
forwarded upstream; usage is summed across rounds and returned in the
reply. Upstream failures surface as a terse 502. Own port (:8090
default) so it coexists with the events receiver; dev.sh gets a serve
runner publishing 8090 on the LAN. Tests drive a scripted fake OpenAI
upstream through the real HTTP server: auth matrix, catalog + subset,
history assembly (client system message preserved, harness identity
prepended only when missing), multi-round usage accounting, refused
tool-call feedback, knob forwarding, 400/502 paths.

💘 Generated with Crush

Assisted-by: Crush:glm-5.2
2026-08-29 01:19:59 -05:00

329 lines
9.2 KiB
Go

// 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"
"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/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]
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-<class>), 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.
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:])
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 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
}