Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/dependabot.yml
Original file line number Diff line number Diff line change
@@ -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:
Expand Down
3 changes: 3 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 ./...)
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 2 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
12 changes: 7 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down
5 changes: 4 additions & 1 deletion control.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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,
Expand Down
17 changes: 12 additions & 5 deletions docs/PLAN.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions docs/operations.md
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,11 @@ one with `client.JobRetry` or `hopper jobs retry <id>`.
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:<schema>`.

Expand Down
83 changes: 83 additions & 0 deletions hopperui/cmd/hopperui/main.go
Original file line number Diff line number Diff line change
@@ -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
}
20 changes: 20 additions & 0 deletions hopperui/go.mod
Original file line number Diff line number Diff line change
@@ -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 => ../
26 changes: 26 additions & 0 deletions hopperui/go.sum
Original file line number Diff line number Diff line change
@@ -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=
Loading
Loading