From a3948e1d4062fb3db763537cc1689a0358723da2 Mon Sep 17 00:00:00 2001 From: Michael McQuade Date: Fri, 25 Sep 2026 22:35:35 -0500 Subject: [PATCH] feat(ui): hopperui web UI module (M8) --- .github/dependabot.yml | 2 +- .github/workflows/ci.yml | 3 + CHANGELOG.md | 4 + Makefile | 2 + README.md | 12 +- control.go | 5 +- docs/PLAN.md | 17 +- docs/operations.md | 5 + hopperui/cmd/hopperui/main.go | 83 +++++ hopperui/go.mod | 20 ++ hopperui/go.sum | 26 ++ hopperui/hopperui.go | 565 ++++++++++++++++++++++++++++++++++ hopperui/hopperui_test.go | 163 ++++++++++ hopperui/static/style.css | 57 ++++ hopperui/templates/ui.html | 233 ++++++++++++++ 15 files changed, 1185 insertions(+), 12 deletions(-) create mode 100644 hopperui/cmd/hopperui/main.go create mode 100644 hopperui/go.mod create mode 100644 hopperui/go.sum create mode 100644 hopperui/hopperui.go create mode 100644 hopperui/hopperui_test.go create mode 100644 hopperui/static/style.css create mode 100644 hopperui/templates/ui.html diff --git a/.github/dependabot.yml b/.github/dependabot.yml index 881c751..e92977d 100644 --- a/.github/dependabot.yml +++ b/.github/dependabot.yml @@ -1,7 +1,7 @@ version: 2 updates: - package-ecosystem: gomod - directories: ["/", "/tools"] + directories: ["/", "/tools", "/hopperotel", "/hopperui"] schedule: { interval: weekly } commit-message: { prefix: "build(deps)" } groups: diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 35c0e04..2b7d2b6 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -25,10 +25,12 @@ jobs: go mod tidy -diff (cd tools && go mod tidy -diff) (cd hopperotel && go mod tidy -diff) + (cd hopperui && go mod tidy -diff) - name: Lint run: | go tool -modfile=tools/go.mod golangci-lint run ./... (cd hopperotel && go tool -modfile=../tools/go.mod golangci-lint run ./...) + (cd hopperui && go tool -modfile=../tools/go.mod golangci-lint run ./...) test: runs-on: ubuntu-latest @@ -61,3 +63,4 @@ jobs: run: | go test -race -cover -timeout 10m ./... (cd hopperotel && go test -race -cover -timeout 10m ./...) + (cd hopperui && go test -race -cover -timeout 10m ./...) diff --git a/CHANGELOG.md b/CHANGELOG.md index 98e113f..cf63f93 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -41,6 +41,10 @@ versions may change the API. never skipped. `Config.StreamRetention`, `hopper streams consumers|seek` and `hopper subscriptions list`. Schema v5 adds `hopper_stream_events` and `hopper_stream_consumers`. +- `hopperui`, a separate module: an embeddable web UI for queues, jobs, + workflows (as a DAG), subscriptions, stream consumers and clients, with + actions behind an `Authorize` hook. +- `JobFilter.After`, to page job listings. ### Changed diff --git a/Makefile b/Makefile index edc5a82..fc134e9 100644 --- a/Makefile +++ b/Makefile @@ -39,6 +39,7 @@ check: lint test ## Run all linters and tests lint: ## Lint Go code $(GOTOOL) golangci-lint run ./... cd hopperotel && go tool -modfile=../tools/go.mod golangci-lint run ./... + cd hopperui && go tool -modfile=../tools/go.mod golangci-lint run ./... .PHONY: fmt fmt: ## Format Go code @@ -48,6 +49,7 @@ fmt: ## Format Go code test: ## Run tests (integration tests need `make pg` or HOPPER_TEST_DATABASE_URL) go test -race -cover -timeout 10m ./... cd hopperotel && go test -race -cover -timeout 10m ./... + cd hopperui && go test -race -cover -timeout 10m ./... .PHONY: bench bench: ## Run the benchmark harness against the test database diff --git a/README.md b/README.md index f907d72..70e5428 100644 --- a/README.md +++ b/README.md @@ -30,11 +30,13 @@ Postgres is the first engine, behind a driver interface designed for more. > LISTEN/NOTIFY wake-ups, unique jobs, cron and interval schedules, cancellation, > retry from the dead-letter queue, TTLs, pause and runtime queues, middleware, > awaitable results, listing, events and stats, plus the `hopper` CLI, -> `hoppertest`, `hopperotel` and the `drivertest` conformance suite. So is -> messaging: subscriptions with topic patterns, fan-out, dedup and ordering -> keys, and a SQL contract for producers in other languages. Flow control, -> batches and workflows follow, in the order laid out in -> [docs/PLAN.md](docs/PLAN.md), the plan of record. Start with +> `hoppertest`, `hopperotel`, the `hopperui` web UI and the `drivertest` +> conformance suite. So are messaging (subscriptions with topic patterns, +> fan-out, dedup and ordering keys, a SQL contract for producers in other +> languages, streams with consumers), flow control (global, rate and +> partitioned limits, priority aging), batches and workflows, and a +> `database/sql` driver, in the order laid out in [docs/PLAN.md](docs/PLAN.md), +> the plan of record. Start with > [docs/getting-started.md](docs/getting-started.md). Feedback is welcome > through issues and PRs. diff --git a/control.go b/control.go index 0cced20..b184f9b 100644 --- a/control.go +++ b/control.go @@ -58,6 +58,9 @@ type JobFilter struct { Queue string Kinds []string States []JobState + // After starts the listing after the given job, for paging: pass the + // last ID of one page to get the next. + After JobID } // jobsPageSize is how many jobs Jobs fetches at a time. @@ -78,7 +81,7 @@ func (c *Client[TTx]) JobsTx(ctx context.Context, tx TTx, filter JobFilter) iter func (c *Client[TTx]) jobs(ctx context.Context, exec driver.Executor, filter JobFilter) iter.Seq2[*JobRow, error] { return func(yield func(*JobRow, error) bool) { - var after JobID + after := filter.After for { page, err := exec.JobList(ctx, driver.JobListParams{ Queue: filter.Queue, Kinds: filter.Kinds, States: filter.States, After: after, Limit: jobsPageSize, diff --git a/docs/PLAN.md b/docs/PLAN.md index d3cec0d..9e2d772 100644 --- a/docs/PLAN.md +++ b/docs/PLAN.md @@ -9,7 +9,7 @@ in-process job framework and a separate message broker. - **Module:** `github.com/parallelworks/hopper` - **License:** Apache-2.0 - **Dependencies:** the Go standard library and `github.com/jackc/pgx/v5`. Nothing else in the core module. -- **Status:** M0 through M4, M6 and M7 are implemented (§16); the §8.2 targets still need a run on the reference hardware before v0.1.0 is tagged. This document is the plan of record, and changes to it go through PRs. +- **Status:** M0 through M4 and M6 through M8 are implemented (§16); the §8.2 targets still need a run on the reference hardware before v0.1.0 is tagged. This document is the plan of record, and changes to it go through PRs. The name refers to a feed hopper, which releases work into a machine one piece at a time, and to RADM Grace Hopper. It is also a fitting name for something @@ -1007,9 +1007,16 @@ observability without new machinery. `queues list|pause|resume|limit`, `clients list`, `workflows get`, `subscriptions list`, `streams consumers|seek` and `stats`, all with `-json`. Benchmarks are the separate `hopperbench` command. -- **Web UI (`hopperui` module):** an embeddable `http.Handler` for browsing queues, - jobs, history, subscriptions and workflows, with retry, cancel and pause actions - behind an application-supplied authorization hook. +- **Web UI (`hopperui` module):** `hopperui.New(client, cfg)` is an embeddable + `http.Handler`, mounted under a prefix with `http.StripPrefix`, for browsing queues + (depths, limits, paused state), jobs and history (filtered by queue, state and kind, + paged by ID), a job's args, metadata, output and errors, workflows as a DAG (SVG, laid + out by dependency depth), subscriptions, stream consumers and clients. Retry, cancel, + pause, resume and seek are offered only when the application supplies + `Config.Authorize`, which is asked before every action; without it the UI is + read-only. It is server-rendered HTML with no scripts and no external assets, so it + works behind any proxy and needs no build step. The client it is given need not be + started. ## 14. Testing strategy @@ -1050,7 +1057,7 @@ The estimates assume one engineer. Each milestone is one or more PRs. | M5 | First adoption | Move an internal service's `internal/jobs` package to hopper; drain and drop its old queue tables | 1d | | M6 | Messaging | Subscriptions, AMQP topic patterns, typed `Message[T]`, PublishTx fan-out, dedup, ordering keys, request/reply, SQL publish contract, `ReplayDiscarded`, the upgrade test. **Done**; **v0.2.0** follows v0.1.0. | 5d | | M7 | Flow control and batches | Global limits, rate limits, partitioned limits, priority aging, batches with callbacks, `hoppersql` driver, **v0.3.0**. **Done.** | 5d | -| M8 | Workflows, streams, UI | Job dependencies and DAG workflows (**done**), streams with consumers (**done**), `hopperui`. Each lands in its own PR. | 2–3w | +| M8 | Workflows, streams, UI | Job dependencies and DAG workflows, streams with consumers, `hopperui`. **Done**, one PR each. | 2–3w | | M9 | More engines (later) | `hoppersqlite`, then `hoppermongo`, each in its own module and passing `drivertest`. Not scheduled yet. | per engine | M0–M4 take roughly four weeks to a production-ready v0.1.0 that meets its performance diff --git a/docs/operations.md b/docs/operations.md index c2c61e3..4527bf8 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -119,6 +119,11 @@ one with `client.JobRetry` or `hopper jobs retry `. process. - The `hopperotel` module adds OpenTelemetry tracing (insert to work, through job metadata) and metrics. +- The `hopperui` module is a web UI: mount `hopperui.New(client, cfg)` under a + path of your application (`http.StripPrefix`) to browse queues, jobs, + workflows, subscriptions, consumers and clients. Pass `Config.Authorize` to + enable retry, cancel, pause, resume and seek; without it the UI is + read-only. - `pg_stat_activity` shows the listener connection as `hopper-listener:`. diff --git a/hopperui/cmd/hopperui/main.go b/hopperui/cmd/hopperui/main.go new file mode 100644 index 0000000..6d238a4 --- /dev/null +++ b/hopperui/cmd/hopperui/main.go @@ -0,0 +1,83 @@ +// Command hopperui serves the hopper web UI on its own, for operators who +// do not embed it in an application. +// +// hopperui -database-url postgres://... [-listen :8080] [-prefix /hopper] [-allow-actions] +// +// Without -allow-actions the UI is read-only. With it, every action is +// allowed for every request, so put the server behind your own +// authentication. +package main + +import ( + "context" + "errors" + "flag" + "fmt" + "log/slog" + "net/http" + "os" + "os/signal" + "syscall" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/parallelworks/hopper" + "github.com/parallelworks/hopper/driver/hopperpgx" + "github.com/parallelworks/hopper/hopperui" +) + +func main() { + var ( + url = flag.String("database-url", os.Getenv("HOPPER_DATABASE_URL"), "Postgres URL (or HOPPER_DATABASE_URL)") + listen = flag.String("listen", ":8080", "address to serve on") + prefix = flag.String("prefix", "", "path prefix to serve under, such as /hopper") + title = flag.String("title", "hopper", "title shown in the header") + actions = flag.Bool("allow-actions", false, "allow retry, cancel, pause, resume and seek for every request") + ) + flag.Parse() + if err := run(*url, *listen, *prefix, *title, *actions); err != nil { + fmt.Fprintln(os.Stderr, "hopperui:", err) + os.Exit(1) + } +} + +func run(url, listen, prefix, title string, actions bool) error { + if url == "" { + return errors.New("-database-url or HOPPER_DATABASE_URL is required") + } + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + pool, err := pgxpool.New(ctx, url) + if err != nil { + return err + } + defer pool.Close() + client, err := hopper.NewClient(hopperpgx.New(pool), &hopper.Config{Logger: slog.Default()}) + if err != nil { + return err + } + cfg := &hopperui.Config{Prefix: prefix, Title: title} + if actions { + cfg.Authorize = func(*http.Request, hopperui.Action) error { return nil } + } + mux := http.NewServeMux() + if prefix == "" { + mux.Handle("/", hopperui.New(client, cfg)) + } else { + mux.Handle(prefix+"/", http.StripPrefix(prefix, hopperui.New(client, cfg))) + mux.Handle("GET /{$}", http.RedirectHandler(prefix+"/", http.StatusFound)) + } + srv := &http.Server{Addr: listen, Handler: mux, ReadHeaderTimeout: 10 * time.Second} + go func() { + <-ctx.Done() + shutdown, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + _ = srv.Shutdown(shutdown) + }() + slog.Info("hopperui: serving", "listen", listen, "prefix", prefix, "actions", actions) + if err := srv.ListenAndServe(); !errors.Is(err, http.ErrServerClosed) { + return err + } + return nil +} diff --git a/hopperui/go.mod b/hopperui/go.mod new file mode 100644 index 0000000..b3df374 --- /dev/null +++ b/hopperui/go.mod @@ -0,0 +1,20 @@ +module github.com/parallelworks/hopper/hopperui + +go 1.27.0 + +godebug fips140=only + +require ( + github.com/jackc/pgx/v5 v5.11.0 + github.com/parallelworks/hopper v0.0.0 +) + +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + golang.org/x/sync v0.17.0 // indirect + golang.org/x/text v0.29.0 // indirect +) + +replace github.com/parallelworks/hopper => ../ diff --git a/hopperui/go.sum b/hopperui/go.sum new file mode 100644 index 0000000..3f716dd --- /dev/null +++ b/hopperui/go.sum @@ -0,0 +1,26 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.11.0 h1:IzBBtyK9AHqf98cctWFifYSci2hgQR/cd56wB4p+ogg= +github.com/jackc/pgx/v5 v5.11.0/go.mod h1:mal1tBGAFfLHvZzaYh77YS/eC6IX9OWbRV1QIIM0Jn4= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= +golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/text v0.29.0 h1:1neNs90w9YzJ9BocxfsQNHKuAT4pkghyXc4nhZ6sJvk= +golang.org/x/text v0.29.0/go.mod h1:7MhJOA9CD2qZyOKYazxdYMF85OwPdEr9jTtBpO7ydH4= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/hopperui/hopperui.go b/hopperui/hopperui.go new file mode 100644 index 0000000..8a61796 --- /dev/null +++ b/hopperui/hopperui.go @@ -0,0 +1,565 @@ +// Package hopperui is a web UI for hopper: an http.Handler that browses +// queues, jobs, history, workflows, subscriptions, stream consumers and +// clients, and can retry, cancel, pause, resume and seek when the +// application allows it. +// +// mux.Handle("/hopper/", http.StripPrefix("/hopper", hopperui.New(client, &hopperui.Config{ +// Prefix: "/hopper", +// Authorize: func(r *http.Request, action hopperui.Action) error { return requireAdmin(r) }, +// }))) +// +// The UI is server-rendered HTML with no scripts and no external assets. +// Without an Authorize hook it is read-only. +package hopperui + +import ( + "bytes" + "context" + "embed" + "encoding/json" + "errors" + "fmt" + "html/template" + "net/http" + "net/url" + "sort" + "strings" + "time" + + "github.com/parallelworks/hopper" +) + +//go:embed templates/ui.html +var templateFS embed.FS + +//go:embed static/style.css +var styleCSS []byte + +// Action is something the UI does to the cluster on the operator's behalf. +type Action string + +// The actions Config.Authorize is asked about. +const ( + ActionRetry Action = "retry" + ActionCancel Action = "cancel" + ActionPause Action = "pause" + ActionResume Action = "resume" + ActionSeek Action = "seek" +) + +// Config tunes the UI. A nil *Config means the defaults. +type Config struct { + // Prefix is the path the handler is mounted under ("/hopper"), used to + // build links. Mount with http.StripPrefix so the handler sees paths + // without it. + Prefix string + // Authorize is asked before every action; returning an error refuses + // it with 403. With no hook the UI is read-only. + Authorize func(r *http.Request, action Action) error + // Title is shown in the header. Defaults to "hopper". + Title string + // PageSize is how many jobs a page lists. Defaults to 50. + PageSize int +} + +// New returns the UI handler for a client. The client need not be started; +// an insert-only client works. +func New[TTx any](client *hopper.Client[TTx], cfg *Config) http.Handler { + u := &ui[TTx]{client: client} + if cfg != nil { + u.cfg = *cfg + } + u.cfg.Prefix = strings.TrimSuffix(u.cfg.Prefix, "/") + if u.cfg.Title == "" { + u.cfg.Title = "hopper" + } + if u.cfg.PageSize <= 0 { + u.cfg.PageSize = 50 + } + u.tmpl = template.Must(template.New("ui").Funcs(template.FuncMap{ + "since": since, + "ts": stamp, + "pretty": prettyJSON, + "dur": func(d time.Duration) string { return d.Round(time.Second).String() }, + "now": time.Now, + "short": shorten, + }).ParseFS(templateFS, "templates/ui.html")) + + mux := http.NewServeMux() + mux.HandleFunc("GET /{$}", u.index) + mux.HandleFunc("GET /queues", u.queues) + mux.HandleFunc("POST /queues/{name}/{action}", u.queueAction) + mux.HandleFunc("GET /jobs", u.jobs) + mux.HandleFunc("GET /jobs/{id}", u.job) + mux.HandleFunc("POST /jobs/{id}/{action}", u.jobAction) + mux.HandleFunc("GET /workflows/{id}", u.workflow) + mux.HandleFunc("GET /subscriptions", u.subscriptions) + mux.HandleFunc("GET /consumers", u.consumers) + mux.HandleFunc("POST /consumers/{name}/seek", u.seek) + mux.HandleFunc("GET /clients", u.clients) + mux.HandleFunc("GET /static/style.css", func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/css; charset=utf-8") + http.ServeContent(w, r, "style.css", time.Time{}, bytes.NewReader(styleCSS)) + }) + u.mux = mux + return u +} + +type ui[TTx any] struct { + client *hopper.Client[TTx] + cfg Config + tmpl *template.Template + mux *http.ServeMux +} + +func (u *ui[TTx]) ServeHTTP(w http.ResponseWriter, r *http.Request) { + u.mux.ServeHTTP(w, r) +} + +// page is what every template receives. +type page struct { + Prefix string + Title string + Page string + CanAct bool + Data any + Query url.Values + Message string +} + +func (u *ui[TTx]) render(w http.ResponseWriter, r *http.Request, name string, data any) { + p := page{Prefix: u.cfg.Prefix, Title: u.cfg.Title, Page: name, CanAct: u.cfg.Authorize != nil, Data: data, Query: r.URL.Query(), Message: r.URL.Query().Get("msg")} + var buf bytes.Buffer + if err := u.tmpl.ExecuteTemplate(&buf, "layout", p); err != nil { + http.Error(w, "hopperui: render "+name+": "+err.Error(), http.StatusInternalServerError) + return + } + w.Header().Set("Content-Type", "text/html; charset=utf-8") + _, _ = buf.WriteTo(w) +} + +func (u *ui[TTx]) fail(w http.ResponseWriter, err error) { + switch { + case errors.Is(err, hopper.ErrNotFound): + http.Error(w, "not found", http.StatusNotFound) + case errors.Is(err, context.Canceled): + http.Error(w, "cancelled", 499) + default: + http.Error(w, err.Error(), http.StatusInternalServerError) + } +} + +// authorize refuses an action unless the application allows it. +func (u *ui[TTx]) authorize(w http.ResponseWriter, r *http.Request, action Action) bool { + if u.cfg.Authorize == nil { + http.Error(w, "hopperui is read-only: no Authorize hook", http.StatusForbidden) + return false + } + if err := u.cfg.Authorize(r, action); err != nil { + http.Error(w, err.Error(), http.StatusForbidden) + return false + } + return true +} + +// back redirects to a page with a message. +func (u *ui[TTx]) back(w http.ResponseWriter, r *http.Request, path, msg string) { + // The target is the configured prefix plus a path built from route + // constants and parsed IDs, never a request value. + http.Redirect(w, r, u.cfg.Prefix+path+"?msg="+url.QueryEscape(msg), http.StatusSeeOther) //nolint:gosec // see above +} + +// Pages. + +type indexData struct { + Stats *hopper.Stats + Queues []queueView +} + +type queueView struct { + Name string + Stats *hopper.QueueStats + Row *hopper.QueueRow + Paused bool +} + +func (u *ui[TTx]) queueViews(ctx context.Context) (*hopper.Stats, []queueView, error) { + stats, err := u.client.Stats(ctx) + if err != nil { + return nil, nil, err + } + rows, err := u.client.Queues().List(ctx) + if err != nil { + return nil, nil, err + } + byName := map[string]*hopper.QueueRow{} + for _, q := range rows { + byName[q.Name] = q + } + names := map[string]struct{}{} + for n := range stats.Queues { + names[n] = struct{}{} + } + for n := range byName { + names[n] = struct{}{} + } + views := make([]queueView, 0, len(names)) + for n := range names { + v := queueView{Name: n, Stats: stats.Queues[n], Row: byName[n]} + if v.Stats == nil { + v.Stats = &hopper.QueueStats{} + } + v.Paused = v.Stats.Paused || (v.Row != nil && !v.Row.PausedAt.IsZero()) + views = append(views, v) + } + sort.Slice(views, func(i, j int) bool { return views[i].Name < views[j].Name }) + return stats, views, nil +} + +func (u *ui[TTx]) index(w http.ResponseWriter, r *http.Request) { + stats, queues, err := u.queueViews(r.Context()) + if err != nil { + u.fail(w, err) + return + } + u.render(w, r, "index", indexData{Stats: stats, Queues: queues}) +} + +func (u *ui[TTx]) queues(w http.ResponseWriter, r *http.Request) { + stats, queues, err := u.queueViews(r.Context()) + if err != nil { + u.fail(w, err) + return + } + u.render(w, r, "queues", indexData{Stats: stats, Queues: queues}) +} + +func (u *ui[TTx]) queueAction(w http.ResponseWriter, r *http.Request) { + name := r.PathValue("name") + var err error + switch r.PathValue("action") { + case "pause": + if !u.authorize(w, r, ActionPause) { + return + } + err = u.client.Queues().Pause(r.Context(), name) + case "resume": + if !u.authorize(w, r, ActionResume) { + return + } + err = u.client.Queues().Resume(r.Context(), name) + default: + http.NotFound(w, r) + return + } + if err != nil { + u.fail(w, err) + return + } + u.back(w, r, "/queues", "queue "+name+" "+r.PathValue("action")+"d") +} + +type jobsData struct { + Jobs []*hopper.JobRow + Queues []string + States []hopper.JobState + Filter hopper.JobFilter + Kind string + State string + Next string +} + +var allStates = []hopper.JobState{ + hopper.JobStatePending, hopper.JobStateAvailable, hopper.JobStateScheduled, hopper.JobStateRunning, + hopper.JobStateRetryable, hopper.JobStateCompleted, hopper.JobStateCancelled, hopper.JobStateDiscarded, +} + +func (u *ui[TTx]) jobs(w http.ResponseWriter, r *http.Request) { + ctx := r.Context() + q := r.URL.Query() + filter := hopper.JobFilter{Queue: q.Get("queue")} + if k := q.Get("kind"); k != "" { + filter.Kinds = []string{k} + } + if s := q.Get("state"); s != "" { + filter.States = []hopper.JobState{hopper.JobState(s)} + } + if a := q.Get("after"); a != "" { + id, err := hopper.ParseJobID(a) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + filter.After = id + } + data := jobsData{Filter: filter, Kind: q.Get("kind"), State: q.Get("state"), States: allStates} + for job, err := range u.client.Jobs(ctx, filter) { + if err != nil { + u.fail(w, err) + return + } + data.Jobs = append(data.Jobs, job) + if len(data.Jobs) == u.cfg.PageSize { + break + } + } + if len(data.Jobs) == u.cfg.PageSize { + q.Set("after", data.Jobs[len(data.Jobs)-1].ID.String()) + data.Next = u.cfg.Prefix + "/jobs?" + q.Encode() + } + if _, views, err := u.queueViews(ctx); err == nil { + for _, v := range views { + data.Queues = append(data.Queues, v.Name) + } + } + u.render(w, r, "jobs", data) +} + +type jobData struct { + Job *hopper.JobRow + Live bool + Errors []hopper.AttemptError +} + +func (u *ui[TTx]) job(w http.ResponseWriter, r *http.Request) { + id, err := hopper.ParseJobID(r.PathValue("id")) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + job, err := u.client.JobGet(r.Context(), id) + if err != nil { + u.fail(w, err) + return + } + u.render(w, r, "job", jobData{Job: job, Live: !job.State.Terminal(), Errors: job.Errors}) +} + +func (u *ui[TTx]) jobAction(w http.ResponseWriter, r *http.Request) { + id, err := hopper.ParseJobID(r.PathValue("id")) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + var job *hopper.JobRow + switch r.PathValue("action") { + case "retry": + if !u.authorize(w, r, ActionRetry) { + return + } + job, err = u.client.JobRetry(r.Context(), id) + case "cancel": + if !u.authorize(w, r, ActionCancel) { + return + } + job, err = u.client.JobCancel(r.Context(), id) + default: + http.NotFound(w, r) + return + } + if err != nil { + u.fail(w, err) + return + } + u.back(w, r, "/jobs/"+id.String(), fmt.Sprintf("job is now %s", job.State)) +} + +// workflowData lays a workflow out as a DAG: each job goes in the column +// after its furthest dependency, and edges are drawn between boxes. +type workflowData struct { + Workflow *hopper.WorkflowRow + Nodes []node + Edges []edge + Width int + Height int +} + +type node struct { + Job *hopper.JobRow + X, Y int +} + +type edge struct { + X1, Y1, X2, Y2 int +} + +const ( + nodeW, nodeH = 220, 44 + gapX, gapY = 60, 20 +) + +func (u *ui[TTx]) workflow(w http.ResponseWriter, r *http.Request) { + id, err := hopper.ParseJobID(r.PathValue("id")) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + wf, err := u.client.WorkflowGet(r.Context(), id) + if err != nil { + u.fail(w, err) + return + } + u.render(w, r, "workflow", layout(wf)) +} + +func layout(wf *hopper.WorkflowRow) workflowData { + index := map[hopper.JobID]int{} + for i, j := range wf.Jobs { + index[j.ID] = i + } + deps := map[hopper.JobID][]hopper.JobID{} + for _, e := range wf.Edges { + deps[e.Job] = append(deps[e.Job], e.DependsOn) + } + // Level = longest path from a root. Jobs are in insertion order and a + // step only depends on earlier steps, but iterate to a fixed point in + // case a graph came from elsewhere. + level := make([]int, len(wf.Jobs)) + for changed := true; changed; { + changed = false + for i, j := range wf.Jobs { + for _, d := range deps[j.ID] { + if k, ok := index[d]; ok && level[k]+1 > level[i] { + level[i] = level[k] + 1 + changed = true + } + } + } + } + rows := map[int]int{} + data := workflowData{Workflow: wf, Nodes: make([]node, len(wf.Jobs))} + for i, j := range wf.Jobs { + x := gapX + level[i]*(nodeW+gapX) + y := gapY + rows[level[i]]*(nodeH+gapY) + rows[level[i]]++ + data.Nodes[i] = node{Job: j, X: x, Y: y} + data.Width = max(data.Width, x+nodeW+gapX) + data.Height = max(data.Height, y+nodeH+gapY) + } + for _, e := range wf.Edges { + from, to := index[e.DependsOn], index[e.Job] + f, t := data.Nodes[from], data.Nodes[to] + data.Edges = append(data.Edges, edge{X1: f.X + nodeW, Y1: f.Y + nodeH/2, X2: t.X, Y2: t.Y + nodeH/2}) + } + return data +} + +func (u *ui[TTx]) subscriptions(w http.ResponseWriter, r *http.Request) { + subs, err := u.client.Subscriptions(r.Context()) + if err != nil { + u.fail(w, err) + return + } + u.render(w, r, "subscriptions", subs) +} + +func (u *ui[TTx]) consumers(w http.ResponseWriter, r *http.Request) { + consumers, err := u.client.Streams().Consumers(r.Context()) + if err != nil { + u.fail(w, err) + return + } + u.render(w, r, "consumers", consumers) +} + +func (u *ui[TTx]) seek(w http.ResponseWriter, r *http.Request) { + if !u.authorize(w, r, ActionSeek) { + return + } + if err := r.ParseForm(); err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + name := r.PathValue("name") + var opts hopper.SeekOpts + switch r.Form.Get("to") { + case "earliest": + opts.Earliest = true + case "latest": + opts.Latest = true + case "time": + t, err := time.ParseInLocation("2006-01-02T15:04", r.Form.Get("time"), time.Local) + if err != nil { + http.Error(w, "time: "+err.Error(), http.StatusBadRequest) + return + } + opts.Time = t + case "position": + p, err := hopper.ParseStreamPosition(r.Form.Get("position")) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + opts.Position = p + default: + http.Error(w, "to: one of earliest, latest, time or position", http.StatusBadRequest) + return + } + if err := u.client.Streams().Seek(r.Context(), name, opts); err != nil { + u.fail(w, err) + return + } + u.back(w, r, "/consumers", "consumer "+name+" moved") +} + +func (u *ui[TTx]) clients(w http.ResponseWriter, r *http.Request) { + clients, err := u.client.Clients(r.Context()) + if err != nil { + u.fail(w, err) + return + } + u.render(w, r, "clients", clients) +} + +// Template helpers. + +func since(t time.Time) string { + if t.IsZero() { + return "" + } + d := time.Since(t) + if d < 0 { + return "in " + humanDuration(-d) + } + return humanDuration(d) + " ago" +} + +func humanDuration(d time.Duration) string { + switch { + case d < time.Minute: + return fmt.Sprintf("%ds", int(d.Seconds())) + case d < time.Hour: + return fmt.Sprintf("%dm", int(d.Minutes())) + case d < 48*time.Hour: + return fmt.Sprintf("%dh", int(d.Hours())) + default: + return fmt.Sprintf("%dd", int(d.Hours()/24)) + } +} + +// shorten fits a JSON document on a node label. +func shorten(b json.RawMessage) string { + s := string(b) + if len(s) > 30 { + return s[:29] + "…" + } + return s +} + +func stamp(t time.Time) string { + if t.IsZero() { + return "" + } + return t.Local().Format("2006-01-02 15:04:05") +} + +func prettyJSON(b json.RawMessage) string { + if len(b) == 0 { + return "" + } + var buf bytes.Buffer + if err := json.Indent(&buf, b, "", " "); err != nil { + return string(b) + } + return buf.String() +} diff --git a/hopperui/hopperui_test.go b/hopperui/hopperui_test.go new file mode 100644 index 0000000..2fd7684 --- /dev/null +++ b/hopperui/hopperui_test.go @@ -0,0 +1,163 @@ +package hopperui_test + +import ( + "context" + "errors" + "io" + "net/http" + "net/http/httptest" + "net/url" + "strings" + "testing" + + "github.com/parallelworks/hopper" + "github.com/parallelworks/hopper/driver/hopperpgx" + "github.com/parallelworks/hopper/hopperui" + "github.com/parallelworks/hopper/internal/testdb" +) + +type ping struct { + N int `json:"n"` +} + +func (ping) Kind() string { return "ping" } + +type stepArgs struct { + Name string `json:"name"` +} + +func (stepArgs) Kind() string { return "step" } + +func TestUI(t *testing.T) { + ctx := context.Background() + pool := testdb.Pool(t) + client, err := hopper.NewClient(hopperpgx.New(pool), nil) + if err != nil { + t.Fatal(err) + } + job, err := client.Insert(ctx, ping{N: 1}, &hopper.InsertOpts{Queue: "emails"}) + if err != nil { + t.Fatal(err) + } + wf := hopper.NewWorkflow("ingest", nil) + first := wf.Add(stepArgs{Name: "fetch"}, nil) + second := wf.Add(stepArgs{Name: "parse"}, hopper.After(first)) + wf.Add(stepArgs{Name: "index"}, hopper.After(second)) + wf.Add(stepArgs{Name: "notify"}, hopper.After(first)) + wres, err := client.InsertWorkflow(ctx, wf) + if err != nil { + t.Fatal(err) + } + + readOnly := httptest.NewServer(http.StripPrefix("/hopper", hopperui.New(client, &hopperui.Config{Prefix: "/hopper"}))) + defer readOnly.Close() + denied := errors.New("not an admin") + admin := httptest.NewServer(http.StripPrefix("/hopper", hopperui.New(client, &hopperui.Config{ + Prefix: "/hopper", + Authorize: func(r *http.Request, action hopperui.Action) error { + if r.Header.Get("X-Admin") == "" { + return denied + } + return nil + }, + }))) + defer admin.Close() + + get := func(srv *httptest.Server, path string) (int, string) { + t.Helper() + res, err := http.Get(srv.URL + "/hopper" + path) + if err != nil { + t.Fatal(err) + } + defer res.Body.Close() + body, _ := io.ReadAll(res.Body) + return res.StatusCode, string(body) + } + post := func(srv *httptest.Server, path string, form url.Values, admin bool) (int, string) { + t.Helper() + req, _ := http.NewRequest(http.MethodPost, srv.URL+"/hopper"+path, strings.NewReader(form.Encode())) + req.Header.Set("Content-Type", "application/x-www-form-urlencoded") + if admin { + req.Header.Set("X-Admin", "1") + } + c := &http.Client{CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }} + res, err := c.Do(req) + if err != nil { + t.Fatal(err) + } + defer res.Body.Close() + body, _ := io.ReadAll(res.Body) + return res.StatusCode, string(body) + } + want := func(code int, body string, gotCode int, needles ...string) { + t.Helper() + if gotCode != code { + t.Errorf("status = %d, want %d: %s", gotCode, code, body) + } + for _, n := range needles { + if !strings.Contains(body, n) { + t.Errorf("body lacks %q:\n%s", n, body) + } + } + } + + code, body := get(readOnly, "/") + want(200, body, code, "emails", "read-only", `href="/hopper/jobs?queue=emails"`, `href="/hopper/static/style.css"`) + code, body = get(readOnly, "/static/style.css") + want(200, body, code, "svg.dag") + code, body = get(readOnly, "/queues") + want(200, body, code, "active") + if strings.Contains(body, "/queues/emails/pause") { + t.Error("read-only UI offers actions") + } + code, body = get(readOnly, "/jobs?queue=emails") + want(200, body, code, job.Job.ID.String(), "ping", `selected>emails`) + code, body = get(readOnly, "/jobs?state=pending") + want(200, body, code, "step", wres.Jobs[1].Job.ID.String(), wres.Jobs[3].Job.ID.String()) + if strings.Contains(body, wres.Jobs[0].Job.ID.String()) || strings.Contains(body, "ping") { + t.Error("state filter ignored") + } + code, body = get(readOnly, "/jobs/"+job.Job.ID.String()) + want(200, body, code, `"n": 1`, "available", "emails") + code, body = get(readOnly, "/jobs/not-an-id") + want(400, body, code) + code, body = get(readOnly, "/jobs/"+hopper.JobID{1}.String()) + want(404, body, code) + code, body = get(readOnly, "/workflows/"+wres.ID.String()) + want(200, body, code, "Workflow ingest", `step {"name": "fetch"}`, `step {"name": "notify"}`, `marker-end`, `unfinished 4`) + if n := strings.Count(body, ` + + + + +{{.Title}} · {{.Page}} + + + +
+ {{.Title}} + + {{if not .CanAct}}read-only{{end}} +
+
+{{if .Message}}

{{.Message}}

{{end}} +{{if eq .Page "index"}}{{template "index" .}} +{{else if eq .Page "queues"}}{{template "queues" .}} +{{else if eq .Page "jobs"}}{{template "jobs" .}} +{{else if eq .Page "job"}}{{template "job" .}} +{{else if eq .Page "workflow"}}{{template "workflow" .}} +{{else if eq .Page "subscriptions"}}{{template "subscriptions" .}} +{{else if eq .Page "consumers"}}{{template "consumers" .}} +{{else if eq .Page "clients"}}{{template "clients" .}} +{{end}} +
+ +{{end}} + +{{define "queuetable"}} + +{{if $.CanAct}}{{end}} + +{{range .Data.Queues}} + + + + + + + + {{if $.CanAct}}{{end}} + +{{else}}{{end}} + +
queueavailablescheduledretryablerunningoldest waitingcompleted/minlimitsstate
{{.Name}}{{.Stats.Available}}{{.Stats.Scheduled}}{{.Stats.Retryable}}{{.Stats.Running}}{{if .Stats.OldestAvailable}}{{dur .Stats.OldestAvailable}}{{else}}–{{end}}{{.Stats.CompletedLastMinute}}{{with .Row}}{{if .Limits.GlobalLimit}}global {{.Limits.GlobalLimit}} {{end}}{{if .Limits.RatePerSec}}rate {{.Limits.RatePerSec}}/s {{end}}{{if .Limits.PartitionLimit}}partition {{.Limits.PartitionLimit}} {{end}}{{if .Limits.Aging}}aging {{dur .Limits.Aging}}{{end}}{{end}}{{if .Paused}}paused{{else}}active{{end}} + {{if .Paused}}
+ {{else}}
{{end}} +
no queues yet
+{{end}} + +{{define "index"}} +

Overview

+

+ leader {{if .Data.Stats.Leader}}client {{.Data.Stats.Leader}}{{else}}none{{end}} + live clients {{.Data.Stats.LiveClients}} +

+{{template "queuetable" .}} +{{end}} + +{{define "queues"}} +

Queues

+{{template "queuetable" .}} +{{end}} + +{{define "jobs"}} +

Jobs

+
+ + + + +
+ + + +{{range .Data.Jobs}} + + + + + + + + + +{{else}}{{end}} + +
idkindqueuestateattemptscheduledcreatedfinalized
{{.ID}}{{.Kind}}{{.Queue}}{{.State}}{{.Attempt}}/{{.MaxAttempts}}{{since .ScheduledAt}}{{since .CreatedAt}}{{since .FinalizedAt}}
no jobs match
+{{if .Data.Next}}

next page →

{{end}} +{{end}} + +{{define "job"}} +{{with .Data.Job}} +

Job {{.ID}}

+

+ kind {{.Kind}} + queue {{.Queue}} + state {{.State}} + priority {{.Priority}} + attempt {{.Attempt}}/{{.MaxAttempts}} +

+{{if $.CanAct}}

+ {{if $.Data.Live}}

{{end}} + {{if ne (print .State) "running"}}
{{end}} +

{{end}} +
+
created
{{ts .CreatedAt}}
+
scheduled
{{ts .ScheduledAt}}
+ {{if not .AttemptedAt.IsZero}}
attempted
{{ts .AttemptedAt}} by client {{.AttemptedBy}}
{{end}} + {{if not .FinalizedAt.IsZero}}
finalized
{{ts .FinalizedAt}}
{{end}} + {{if not .ExpiresAt.IsZero}}
expires
{{ts .ExpiresAt}}
{{end}} + {{if not .CancelRequestedAt.IsZero}}
cancel requested
{{ts .CancelRequestedAt}}
{{end}} + {{if .UniqueKey}}
unique key
{{.UniqueKey}}
{{end}} + {{if .OrderingKey}}
ordering key
{{.OrderingKey}}
{{end}} + {{if .PartitionKey}}
partition key
{{.PartitionKey}}
{{end}} + {{if not .BatchID.IsZero}}
batch
{{.BatchID}}
{{end}} +
+

Args

{{pretty .Args}}
+

Metadata

{{pretty .Metadata}}
+{{if .Output}}

Output

{{pretty .Output}}
{{end}} +{{if .Errors}}

Errors

+ +{{range .Errors}}{{end}} +
attemptaterror
{{.Attempt}}{{ts .At}}
{{.Error}}{{if .Trace}}
+
+{{.Trace}}{{end}}
{{end}} +{{end}} +{{end}} + +{{define "workflow"}} +{{with .Data.Workflow.Batch}} +

{{if .Name}}Workflow {{.Name}}{{else}}Batch{{end}} {{.ID}}

+

+ jobs {{.Total}} + unfinished {{.Pending}} + failed {{.Failed}} + created {{ts .CreatedAt}} + {{if not .CompletedAt.IsZero}}completed {{ts .CompletedAt}}{{end}} +

+{{end}} +{{if .Data.Edges}} + + + {{range .Data.Edges}}{{end}} + {{range .Data.Nodes}} + + {{.Job.Kind}} {{printf "%s" .Job.Args}} + + {{.Job.Kind}} {{short .Job.Args}} + {{.Job.State}} + + {{end}} + +{{end}} + + + +{{range .Data.Workflow.Jobs}} + + + + + + + + +{{end}} + +
idkindqueuestateattemptcreatedfinalized
{{.ID}}{{.Kind}}{{.Queue}}{{.State}}{{.Attempt}}/{{.MaxAttempts}}{{since .CreatedAt}}{{since .FinalizedAt}}
+{{end}} + +{{define "subscriptions"}} +

Subscriptions

+ + + +{{range .Data}} + +{{else}}{{end}} + +
namepatternqueuekindmax attemptscreated
{{.Name}}{{.Pattern}}{{.Queue}}{{.Kind}}{{if .MaxAttempts}}{{.MaxAttempts}}{{else}}default{{end}}{{ts .CreatedAt}}
no subscriptions
+{{end}} + +{{define "consumers"}} +

Stream consumers

+ +{{if .CanAct}}{{end}} + +{{range .Data}} + + + + + + {{if $.CanAct}}{{end}} + +{{else}}{{end}} + +
namepatternqueuesnapshotpositionlast deliveryseek
{{.Name}}{{.Pattern}}{{.Queue}}{{.Snapshot}}{{.Position}}{{if .DeliveredAt.IsZero}}never{{else}}{{since .DeliveredAt}}{{end}} +
+
+
+
no consumers
+{{end}} + +{{define "clients"}} +

Clients

+ + + +{{range .Data}} + +{{else}}{{end}} + +
idhoststartedleaseinfo
{{.ID}}{{.Hostname}}{{since .StartedAt}}{{if .ExpiresAt.After now}}live{{else}}expired{{end}}{{printf "%s" .Info}}
no clients
+{{end}}