From 984afba6d4e908a42d17d39cffb80ea672690a7f Mon Sep 17 00:00:00 2001
From: Michael McQuade
Date: Fri, 25 Sep 2026 22:22:35 -0500
Subject: [PATCH 1/2] feat(core): streams with consumers (M8)
---
CHANGELOG.md | 8 +
client.go | 9 +-
cmd/hopper/main.go | 94 +++++
cmd/hopper/main_test.go | 9 +
config.go | 11 +
docs/PLAN.md | 52 ++-
docs/getting-started.md | 19 +
docs/operations.md | 4 +
driver/driver.go | 143 +++++++
driver/hopperpgx/postgres_test.go | 5 +-
driver/internal/pgsql/history.go | 23 +-
driver/internal/pgsql/stream.go | 350 ++++++++++++++++++
drivertest/drivertest.go | 1 +
drivertest/stream.go | 234 ++++++++++++
export_test.go | 1 +
hoppermigrate/migrations/005_streams.down.sql | 2 +
hoppermigrate/migrations/005_streams.up.sql | 39 ++
leader.go | 15 +-
listener.go | 4 +-
messaging.go | 16 +-
notifier.go | 36 +-
stream.go | 309 ++++++++++++++++
stream_test.go | 114 ++++++
worker.go | 1 +
24 files changed, 1463 insertions(+), 36 deletions(-)
create mode 100644 driver/internal/pgsql/stream.go
create mode 100644 drivertest/stream.go
create mode 100644 hoppermigrate/migrations/005_streams.down.sql
create mode 100644 hoppermigrate/migrations/005_streams.up.sql
create mode 100644 stream.go
create mode 100644 stream_test.go
diff --git a/CHANGELOG.md b/CHANGELOG.md
index 391e7dc..98e113f 100644
--- a/CHANGELOG.md
+++ b/CHANGELOG.md
@@ -33,6 +33,14 @@ versions may change the API.
dependencies wait pending and are promoted by the statement that finalizes
the last of them; a failed step cancels its dependents unless they opt to
`DependencyIgnore`. Schema v4 adds `hopper_job_deps`.
+- Streams: `Streams().Append`/`AppendTx` write to a retained, time-partitioned
+ log; `hopper.Consume` registers a consumer that delivers matching events as
+ jobs from a position of its own, starting at the earliest or latest event
+ and movable with `Seek`; `Read` pages the log. Consumers read by snapshot
+ deltas, so a late-committing transaction is delivered when it commits and
+ never skipped. `Config.StreamRetention`, `hopper streams consumers|seek`
+ and `hopper subscriptions list`. Schema v5 adds `hopper_stream_events` and
+ `hopper_stream_consumers`.
### Changed
diff --git a/client.go b/client.go
index 0e4c036..95d104e 100644
--- a/client.go
+++ b/client.go
@@ -59,6 +59,7 @@ type Client[TTx any] struct {
listening atomic.Bool
isLeader atomic.Bool
leaderPoke chan struct{}
+ streamPoke chan struct{}
events *eventBus
@@ -98,6 +99,7 @@ func NewClient[TTx any](d driver.Driver[TTx], cfg *Config) (*Client[TTx], error)
logger: resolved.Logger,
workers: resolved.Workers,
leaderPoke: make(chan struct{}, 1),
+ streamPoke: make(chan struct{}, 1),
events: newEventBus(),
running: map[JobID]context.CancelCauseFunc{},
doneWaiters: map[JobID][]chan struct{}{},
@@ -246,8 +248,8 @@ func (c *Client[TTx]) finalized(job *driver.JobRow, result driver.JobFinalize) {
}
}
-// declare records what this client works: its queues, subscriptions and
-// declared limits, which take effect cluster-wide.
+// declare records what this client works: its queues, subscriptions,
+// stream consumers and declared limits, which take effect cluster-wide.
func (c *Client[TTx]) declare(ctx context.Context) error {
if err := c.exec.QueueEnsure(ctx, slices.Sorted(maps.Keys(c.cfg.Queues))); err != nil {
c.logger.WarnContext(ctx, "hopper: record queues", "error", err)
@@ -257,6 +259,9 @@ func (c *Client[TTx]) declare(ctx context.Context) error {
if err := c.exec.SubscriptionUpsert(ctx, c.workers.subscriptionRows()); err != nil {
return err
}
+ if err := c.exec.StreamConsumerUpsert(ctx, c.workers.consumerRows()); err != nil {
+ return err
+ }
for name, q := range c.cfg.Queues {
if l := q.limits(name); l.Limited() || l.Aging > 0 {
if err := c.exec.QueueSetLimits(ctx, l); err != nil {
diff --git a/cmd/hopper/main.go b/cmd/hopper/main.go
index 397a8b5..86b575b 100644
--- a/cmd/hopper/main.go
+++ b/cmd/hopper/main.go
@@ -23,6 +23,7 @@ import (
"os"
"os/signal"
"slices"
+ "strconv"
"strings"
"text/tabwriter"
"time"
@@ -59,6 +60,9 @@ commands:
queues limit [-global N] [-rate R] [-burst B] [-partition P] [-aging D] (omitted limits are removed)
clients list
workflows get
+ subscriptions list
+ streams consumers
+ streams seek -earliest|-latest|-time RFC3339|-position xid:seq
stats
bench (see the hopperbench command)
@@ -115,6 +119,10 @@ func run(ctx context.Context, args []string, out io.Writer) error {
return c.clients(ctx, rest[1:])
case "workflows":
return c.workflows(ctx, rest[1:])
+ case "subscriptions":
+ return c.subscriptions(ctx, rest[1:])
+ case "streams":
+ return c.streams(ctx, rest[1:])
case "stats":
return c.stats(ctx)
default:
@@ -374,6 +382,92 @@ func (c *cli) workflows(ctx context.Context, args []string) error {
})
}
+func (c *cli) subscriptions(ctx context.Context, args []string) error {
+ if len(args) != 1 || args[0] != "list" {
+ return errors.New("subscriptions: list")
+ }
+ subs, err := c.client.Subscriptions(ctx)
+ if err != nil {
+ return err
+ }
+ return c.print(subs, func(w io.Writer) {
+ tw := tabwriter.NewWriter(w, 0, 0, 2, ' ', 0)
+ fmt.Fprintln(tw, "NAME\tPATTERN\tQUEUE\tMAX_ATTEMPTS\tCREATED")
+ for _, s := range subs {
+ attempts := "default"
+ if s.MaxAttempts > 0 {
+ attempts = strconv.Itoa(s.MaxAttempts)
+ }
+ fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%s\n", s.Name, s.Pattern, s.Queue, attempts, s.CreatedAt.Local().Format(time.RFC3339))
+ }
+ tw.Flush()
+ })
+}
+
+func (c *cli) streams(ctx context.Context, args []string) error {
+ if len(args) == 0 {
+ return errors.New("streams: consumers or seek")
+ }
+ switch args[0] {
+ case "consumers":
+ consumers, err := c.client.Streams().Consumers(ctx)
+ if err != nil {
+ return err
+ }
+ return c.print(consumers, func(w io.Writer) {
+ tw := tabwriter.NewWriter(w, 0, 0, 2, ' ', 0)
+ fmt.Fprintln(tw, "NAME\tPATTERN\tQUEUE\tSNAPSHOT\tPOSITION\tDELIVERED")
+ for _, cs := range consumers {
+ delivered := "never"
+ if !cs.DeliveredAt.IsZero() {
+ delivered = cs.DeliveredAt.Local().Format(time.RFC3339)
+ }
+ fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%s\t%s\n", cs.Name, cs.Pattern, cs.Queue, cs.Snapshot, cs.Position, delivered)
+ }
+ tw.Flush()
+ })
+ case "seek":
+ fs := flag.NewFlagSet("hopper streams seek", flag.ContinueOnError)
+ fs.SetOutput(c.out)
+ var opts hopper.SeekOpts
+ var at, position string
+ fs.BoolVar(&opts.Earliest, "earliest", false, "replay every retained event")
+ fs.BoolVar(&opts.Latest, "latest", false, "skip to events appended from now on")
+ fs.StringVar(&at, "time", "", "replay from the first event at or after this RFC3339 time")
+ fs.StringVar(&position, "position", "", "resume after this position (xid:seq)")
+ if len(args) < 2 {
+ return errors.New("streams seek: a consumer name is required")
+ }
+ if err := fs.Parse(args[2:]); err != nil {
+ return err
+ }
+ if at != "" {
+ t, err := time.Parse(time.RFC3339, at)
+ if err != nil {
+ return fmt.Errorf("streams seek: -time: %w", err)
+ }
+ opts.Time = t
+ }
+ if position != "" {
+ p, err := hopper.ParseStreamPosition(position)
+ if err != nil {
+ return err
+ }
+ opts.Position = p
+ }
+ if !opts.Earliest && !opts.Latest && opts.Time.IsZero() && opts.Position.IsZero() {
+ return errors.New("streams seek: one of -earliest, -latest, -time or -position is required")
+ }
+ if err := c.client.Streams().Seek(ctx, args[1], opts); err != nil {
+ return err
+ }
+ fmt.Fprintf(c.out, "consumer %s moved\n", args[1])
+ return nil
+ default:
+ return fmt.Errorf("streams: unknown subcommand %q", args[0])
+ }
+}
+
func (c *cli) stats(ctx context.Context) error {
stats, err := c.client.Stats(ctx)
if err != nil {
diff --git a/cmd/hopper/main_test.go b/cmd/hopper/main_test.go
index 8adb5a2..d06f412 100644
--- a/cmd/hopper/main_test.go
+++ b/cmd/hopper/main_test.go
@@ -108,6 +108,15 @@ func TestCLI(t *testing.T) {
if got := hopperCmd("stats"); !strings.Contains(got, "q1") {
t.Errorf("stats = %q", got)
}
+ if got := hopperCmd("subscriptions", "list"); !strings.Contains(got, "NAME") {
+ t.Errorf("subscriptions list = %q", got)
+ }
+ if got := hopperCmd("streams", "consumers"); !strings.Contains(got, "NAME") {
+ t.Errorf("streams consumers = %q", got)
+ }
+ if err := run(ctx, []string{"streams", "seek", "nobody", "-earliest"}, new(bytes.Buffer)); err == nil {
+ t.Error("seek of an unknown consumer succeeded")
+ }
wf := hopper.NewWorkflow("ingest", nil)
first := wf.Add(ping{}, nil)
wf.Add(ping{}, hopper.After(first))
diff --git a/config.go b/config.go
index 5dc8227..b8d4374 100644
--- a/config.go
+++ b/config.go
@@ -56,6 +56,10 @@ type Config struct {
// history, which is the dead-letter queue. Defaults to 7 days; negative
// keeps them forever.
FailedRetention time.Duration
+ // StreamRetention is how long stream events stay in the log, for
+ // consumers that start from the earliest event or are moved back.
+ // Defaults to 7 days; negative keeps them forever.
+ StreamRetention time.Duration
// RescueStuckAfter, if set, rescues jobs that have been running longer
// than this even though their client's lease is live, for workers that
// ignore their context. Off by default; timeouts cover most cases.
@@ -109,6 +113,7 @@ const (
DefaultPollInterval = time.Second
DefaultCompletedRetention = 24 * time.Hour
DefaultFailedRetention = 7 * 24 * time.Hour
+ DefaultStreamRetention = 7 * 24 * time.Hour
)
// withDefaults validates cfg and fills in defaults. It does not modify cfg.
@@ -141,6 +146,9 @@ func (cfg *Config) withDefaults() (Config, error) {
if out.FailedRetention == 0 {
out.FailedRetention = DefaultFailedRetention
}
+ if out.StreamRetention == 0 {
+ out.StreamRetention = DefaultStreamRetention
+ }
if out.Hostname == "" {
host, err := os.Hostname()
if err != nil || host == "" {
@@ -216,6 +224,8 @@ type tuning struct {
// partitions and retention.
maintenanceInterval time.Duration
rescueBatch int
+ // streamBatch is how many events one pump delivers per consumer.
+ streamBatch int
}
var defaultTuning = tuning{
@@ -231,4 +241,5 @@ var defaultTuning = tuning{
leaderInterval: 5 * time.Second,
maintenanceInterval: 5 * time.Minute,
rescueBatch: 1000,
+ streamBatch: 500,
}
diff --git a/docs/PLAN.md b/docs/PLAN.md
index 8e86f0d..d3cec0d 100644
--- a/docs/PLAN.md
+++ b/docs/PLAN.md
@@ -221,6 +221,12 @@ res, err := client.PublishTx(ctx, tx, AllocationCreated{ID: 42}, &hopper.Publish
TTL: 10 * time.Minute, // expire undelivered
})
res.MessageID; res.Deliveries // one delivery per matching subscription
+
+// Streams: the retained log. A consumer reads from a position of its own.
+hopper.Consume(workers, hopper.Consumer{Name: "audit", Pattern: "#", Start: hopper.StreamStartEarliest},
+ func(ctx context.Context, msg *hopper.Message[hopper.Raw]) error { return audit(ctx, msg) })
+ev, err := client.Streams().AppendTx(ctx, tx, AllocationCreated{ID: 42}, &hopper.AppendOpts{Key: "allocation:42"})
+err = client.Streams().Seek(ctx, "audit", hopper.SeekOpts{Time: yesterday}) // replay
```
### 4.6 Batches and workflows
@@ -446,7 +452,9 @@ Schema v2 (M6) adds `max_attempts` and `metadata` to `hopper_subscriptions`, the
ordering-key indexes (§10) and the SQL contract functions. Schema v3 (M7) adds
`partition_limit` and `aging_seconds` to `hopper_queues`, the partition-key running
index, `hopper_batches` and the batch index. Schema v4 (M8) adds `hopper_job_deps`, the
-name and edges of a batch, and a `batch_id` index on history (§11).
+name and edges of a batch, and a `batch_id` index on history (§11). Schema v5 (M8) adds
+the stream log `hopper_stream_events` (partitioned by day) and `hopper_stream_consumers`
+(§10).
### 6.1 Job IDs
Job IDs are UUIDv7. They are safe to expose outside the application (in URLs, APIs
@@ -735,9 +743,9 @@ acknowledged.
| `completed` | archived for 24h | deleted on ack (`DeleteCompleted: true`) |
| `cancelled`, `discarded` | archived for 7d (the dead-letter queue) | archived for 7d (the dead-letter queue) |
-`DeleteCompleted` is configurable per queue; the two retention periods
-(`CompletedRetention`, `FailedRetention`) are per cluster, because partitions are cut
-by time, not by queue. Setting `DeleteCompleted: true` on a job queue gives maximum
+`DeleteCompleted` is configurable per queue; the retention periods
+(`CompletedRetention`, `FailedRetention`, and `StreamRetention` for the stream log,
+7 days by default) are per cluster, because partitions are cut by time, not by queue. Setting `DeleteCompleted: true` on a job queue gives maximum
throughput by skipping the history insert.
### 7.11 Notifications
@@ -905,10 +913,32 @@ observability without new machinery.
and `hopper_publish(topic, payload, opts)` SQL functions ship with schema v2, so
Python, shell or HPC batch scripts can enqueue work with plain SQL, inside their own
transactions. This is the supported path for producers not written in Go (§2).
-- **Streams (M8).** An append-only, time-partitioned topic log with consumer groups
- that track offsets, for replay and late subscribers. Readers only see events below
- the current snapshot's `xmin` (`pg_snapshot_xmin(pg_current_snapshot())`), so
- transactions that commit out of order can never make a reader skip an event.
+- **Streams (M8).** `hopper_stream_events` is an append-only, time-partitioned log:
+ `Streams().Append`/`AppendTx` write one event (topic, optional key, payload, headers,
+ message ID) and consumers read it. An event is identified by `(xid, seq)`, the
+ transaction that wrote it and the event within it. A **consumer** (`hopper.Consume`)
+ is a named position on the log plus a topic pattern; the leader turns the events
+ after its position into delivery jobs of kind `stream:` (with the event's key
+ as ordering key), so a consumer inherits everything a subscription has: typed
+ handlers, retries, dead-lettering, competing workers. Unlike a subscription, which is
+ fanned out to at publish time, a consumer can start at the earliest retained event
+ (`StreamStartEarliest`) and be moved (`Seek` to a position, a time, the start or
+ now), so late subscribers and replays see past events.
+ - **Reads never skip a late commit.** A consumer's row records the `pg_snapshot` it
+ has read to. A pump takes a fresh snapshot and delivers the events of the
+ transactions visible in the new one but not in the old (`pg_visible_in_snapshot`),
+ in `(xid, seq)` order, then stores the new snapshot; a pump that stops part way
+ through keeps both snapshots and the last event delivered, and continues the same
+ delta next time. An event whose transaction commits late is delivered when it
+ commits, whatever its xid, and a long-running transaction holds nothing back. Events
+ of one transaction are delivered together; across transactions the order is commit
+ order per delta, which is the only order two overlapping transactions have.
+ - Pool appends notify `hopper_stream` through the coalescer and poke the local
+ leader; transactional ones notify from the transaction. The leader also pumps every
+ `PollInterval`, for clients without a listener. Two leaders pumping the same consumer
+ serialize on its row.
+ - `Streams().Read` pages the committed log by position, for browsing and the UI.
+ - Retention drops daily partitions after `StreamRetention` (§7.10).
## 11. Batches and workflows
@@ -974,8 +1004,8 @@ observability without new machinery.
gauges for queue depth by state, the oldest claimable job's age and live clients from
`Stats`. Prometheus users export through the OTel exporter.
- **CLI (`cmd/hopper`, core module):** `migrate up|down|version`, `jobs list|get|retry|cancel`,
- `queues list|pause|resume|limit`, `clients list`, `workflows get` and `stats`, all with
- `-json`. `subscriptions list` arrives with M6. Benchmarks are the separate
+ `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
@@ -1020,7 +1050,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 consumer groups, `hopperui`. Each lands in its own PR. | 2–3w |
+| M8 | Workflows, streams, UI | Job dependencies and DAG workflows (**done**), streams with consumers (**done**), `hopperui`. Each lands in its own PR. | 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/getting-started.md b/docs/getting-started.md
index 0cc6114..44dafa8 100644
--- a/docs/getting-started.md
+++ b/docs/getting-started.md
@@ -188,6 +188,25 @@ anyway), and the workflow is a batch, so it has the same callbacks.
`client.WorkflowGet(ctx, res.ID)` and `hopper workflows get ` show its
progress and graph.
+## Streams
+
+Messaging fans out at publish time; a stream keeps the events, so a consumer
+can start from the beginning or be moved back.
+
+```go
+hopper.Consume(workers, hopper.Consumer{Name: "audit", Pattern: "#", Start: hopper.StreamStartEarliest},
+ func(ctx context.Context, msg *hopper.Message[hopper.Raw]) error {
+ return record(ctx, msg.Topic, msg.Payload, msg.Position)
+ })
+
+ev, err := client.Streams().AppendTx(ctx, tx, AllocationCreated{ID: 42}, &hopper.AppendOpts{Key: "allocation:42"})
+err = client.Streams().Seek(ctx, "audit", hopper.SeekOpts{Time: time.Now().Add(-time.Hour)})
+```
+
+Events with the same key are handled one at a time, in order. The log keeps
+`Config.StreamRetention` (7 days) of events; `hopper streams consumers` shows
+where each consumer is.
+
## Results and waiting
```go
diff --git a/docs/operations.md b/docs/operations.md
index c82a981..c2c61e3 100644
--- a/docs/operations.md
+++ b/docs/operations.md
@@ -43,6 +43,10 @@
session connection: point the listener at a direct address with
`hopperpgx.Config.ListenConnConfig`, or accept polling (the client logs
once and polls every `PollInterval`, one second by default).
+- **Streams.** The stream log is partitioned by day and dropped after
+ `StreamRetention`, like history. Consumers are delivered by the leader,
+ in transaction-commit order per batch, and a transaction that commits
+ late is delivered when it commits; nothing waits on it.
- **`database/sql`.** `hoppersql.New(db)` runs hopper on a `*sql.DB` opened
with pgx's `stdlib` adapter or lib/pq, for applications that already use
one. It has no `LISTEN`, so clients poll every `PollInterval`, and no COPY,
diff --git a/driver/driver.go b/driver/driver.go
index 3aad550..054737a 100644
--- a/driver/driver.go
+++ b/driver/driver.go
@@ -13,6 +13,7 @@ import (
"context"
"encoding/json"
"errors"
+ "fmt"
"time"
)
@@ -72,6 +73,9 @@ const (
// ChannelDone carries the ID of a finalized job that was inserted with
// Await, so waiters return as soon as the finalize commits.
ChannelDone = "hopper_done"
+ // ChannelStream is signalled when events are appended to the stream
+ // log, with the topic as payload, so the leader pumps consumers at once.
+ ChannelStream = "hopper_stream"
)
// Listener receives notifications. One listener multiplexes every channel.
@@ -217,6 +221,26 @@ type Executor interface {
// whose dedup key matches a live one is reported as a duplicate.
MessagePublish(ctx context.Context, params MessagePublishParams) ([]JobInsertResult, error)
+ // StreamAppend appends one event to the stream log and returns it with
+ // its position.
+ StreamAppend(ctx context.Context, params StreamAppendParams) (*StreamEvent, error)
+ // StreamRead pages through the committed events after a position whose
+ // topic matches a pattern, in position order. It is for browsing; a
+ // consumer is what reads the log without missing late commits.
+ StreamRead(ctx context.Context, params StreamReadParams) ([]*StreamEvent, error)
+ // StreamConsumerUpsert records consumers, updating pattern, queue and
+ // settings of existing ones by name; a new consumer starts at Start.
+ StreamConsumerUpsert(ctx context.Context, consumers []StreamConsumerRow) error
+ // StreamConsumerList returns every consumer with its position.
+ StreamConsumerList(ctx context.Context) ([]*StreamConsumerRow, error)
+ // StreamConsumerSeek moves a consumer's position: events after it are
+ // delivered (again).
+ StreamConsumerSeek(ctx context.Context, name string, params StreamSeekParams) error
+ // StreamPump turns up to limit events after a consumer's position into
+ // delivery jobs and advances the position, atomically. Two leaders
+ // pumping the same consumer serialize on its row.
+ StreamPump(ctx context.Context, name string, limit int) (StreamPumpResult, error)
+
// HistoryMaintain enforces retention on finalized jobs and prepares
// storage for the near future, in whatever way suits the engine (dropping
// time partitions on Postgres). It is idempotent and may run on two
@@ -593,6 +617,8 @@ type HistoryMaintainParams struct {
CompletedRetention time.Duration
// FailedRetention is how long cancelled and discarded jobs are kept.
FailedRetention time.Duration
+ // StreamRetention is how long stream events are kept.
+ StreamRetention time.Duration
}
// HistoryMaintainResult reports what maintenance did.
@@ -640,6 +666,123 @@ type MessagePublishParams struct {
Notify bool
}
+// StreamPosition orders events in the stream log: by the transaction that
+// wrote them, then by sequence within it. The zero position is before every
+// event.
+type StreamPosition struct {
+ Xid uint64
+ Seq int64
+}
+
+// Less reports whether p comes before q.
+func (p StreamPosition) Less(q StreamPosition) bool {
+ return p.Xid < q.Xid || (p.Xid == q.Xid && p.Seq < q.Seq)
+}
+
+// IsZero reports the position before every event.
+func (p StreamPosition) IsZero() bool { return p == StreamPosition{} }
+
+// String formats the position as "xid:seq".
+func (p StreamPosition) String() string { return fmt.Sprintf("%d:%d", p.Xid, p.Seq) }
+
+// ParseStreamPosition parses the String form.
+func ParseStreamPosition(s string) (StreamPosition, error) {
+ var p StreamPosition
+ if _, err := fmt.Sscanf(s, "%d:%d", &p.Xid, &p.Seq); err != nil {
+ return StreamPosition{}, fmt.Errorf("hopper: stream position %q is not xid:seq", s)
+ }
+ return p, nil
+}
+
+// StreamEvent is one event of the stream log.
+type StreamEvent struct {
+ Position StreamPosition
+ Topic string
+ // Key, if set, becomes the ordering key of the event's deliveries.
+ Key string
+ Payload json.RawMessage
+ Headers map[string]string
+ MessageID string
+ CreatedAt time.Time
+}
+
+// StreamAppendParams describes an event to append.
+type StreamAppendParams struct {
+ Topic string
+ Key string
+ Payload json.RawMessage
+ Headers map[string]string
+ // Notify signals ChannelStream from inside the statement.
+ Notify bool
+}
+
+// StreamReadParams selects events.
+type StreamReadParams struct {
+ // Pattern is an AMQP topic pattern; empty matches every topic.
+ Pattern string
+ After StreamPosition
+ Limit int
+}
+
+// StreamStart says where a new consumer starts.
+type StreamStart string
+
+const (
+ // StreamStartLatest delivers events appended from now on.
+ StreamStartLatest StreamStart = "latest"
+ // StreamStartEarliest delivers every retained event first.
+ StreamStartEarliest StreamStart = "earliest"
+)
+
+// StreamConsumerRow is a consumer: a named position on the log whose
+// events, matched by pattern, are delivered as jobs of Kind on Queue.
+type StreamConsumerRow struct {
+ Name string
+ Pattern string
+ Kind string
+ Queue string
+ // MaxAttempts is the retry budget of deliveries, or 0 for the default
+ // (25, or 10 for events with a key).
+ MaxAttempts int
+ // Metadata is merged into every delivery's metadata.
+ Metadata json.RawMessage
+ // Start applies when the consumer is created; it defaults to latest.
+ Start StreamStart
+ // Snapshot is the engine's record of what the consumer has read to:
+ // on Postgres, the pg_snapshot every delivered transaction is visible
+ // in.
+ Snapshot string
+ // Position is the last event delivered of the delta being read, or
+ // zero between deltas.
+ Position StreamPosition
+ DeliveredAt time.Time
+ CreatedAt time.Time
+}
+
+// StreamSeekParams says where to move a consumer. Exactly one of Position,
+// Time, Earliest and Latest applies, in that order of precedence.
+type StreamSeekParams struct {
+ // Position, if set, becomes the consumer's position: the next event
+ // delivered is the one after it.
+ Position StreamPosition
+ // Time moves the consumer to just before the first event appended at
+ // or after it.
+ Time time.Time
+ // Earliest moves the consumer before every retained event.
+ Earliest bool
+ // Latest moves the consumer past every event appended so far.
+ Latest bool
+}
+
+// StreamPumpResult reports one pump.
+type StreamPumpResult struct {
+ Delivered int
+ // Position is the consumer's position afterwards.
+ Position StreamPosition
+ // Queues received deliveries.
+ Queues []string
+}
+
// ClientRegisterParams describes a client process taking out a lease.
type ClientRegisterParams struct {
Hostname string
diff --git a/driver/hopperpgx/postgres_test.go b/driver/hopperpgx/postgres_test.go
index cd6493a..08ca94d 100644
--- a/driver/hopperpgx/postgres_test.go
+++ b/driver/hopperpgx/postgres_test.go
@@ -95,8 +95,9 @@ func TestHistoryMaintain(t *testing.T) {
if err != nil {
t.Fatal(err)
}
- // 3 hourly and 2 daily periods ahead of the current ones.
- if len(res.Created) != 3+2 {
+ // 3 hourly and 2 daily periods ahead of the current ones, and 2 daily
+ // stream partitions.
+ if len(res.Created) != 3+2+2 {
t.Errorf("created = %v", res.Created)
}
next := "hopper_job_history_completed_" + time.Now().UTC().Add(time.Hour).Format("2006010215")
diff --git a/driver/internal/pgsql/history.go b/driver/internal/pgsql/history.go
index 45dbfb0..2ddaf24 100644
--- a/driver/internal/pgsql/history.go
+++ b/driver/internal/pgsql/history.go
@@ -10,11 +10,16 @@ import (
"github.com/parallelworks/hopper/driver"
)
-// historyGroup describes one outcome partition of hopper_job_history and how
-// it is split by time.
+// historyGroup describes one time-partitioned table: an outcome partition
+// of hopper_job_history, or the stream log.
type historyGroup struct {
- parent string
- period time.Duration
+ // root is the table whose shape new partitions copy; parent is the
+ // partitioned table they attach to (the same, except for history's
+ // outcome partitions).
+ root, parent string
+ // timeColumn is the range key.
+ timeColumn string
+ period time.Duration
// ahead is how many future periods to keep ready, beyond the current one.
ahead int
// format names a partition by its period start (in UTC).
@@ -27,8 +32,9 @@ type historyGroup struct {
const HistoryLockKey int32 = 0x686f7068 // "hoph"
var historyGroups = []historyGroup{
- {parent: "hopper_job_history_completed", period: time.Hour, ahead: 3, format: "2006010215"},
- {parent: "hopper_job_history_failed", period: 24 * time.Hour, ahead: 2, format: "20060102"},
+ {root: "hopper_job_history", parent: "hopper_job_history_completed", timeColumn: "finalized_at", period: time.Hour, ahead: 3, format: "2006010215"},
+ {root: "hopper_job_history", parent: "hopper_job_history_failed", timeColumn: "finalized_at", period: 24 * time.Hour, ahead: 2, format: "20060102"},
+ {root: "hopper_stream_events", parent: "hopper_stream_events", timeColumn: "created_at", period: 24 * time.Hour, ahead: 2, format: "20060102"},
}
func (g historyGroup) name(start time.Time) string {
@@ -67,6 +73,7 @@ func (e *Executor) HistoryMaintain(ctx context.Context, params driver.HistoryMai
retention := map[string]time.Duration{
"hopper_job_history_completed": params.CompletedRetention,
"hopper_job_history_failed": params.FailedRetention,
+ "hopper_stream_events": params.StreamRetention,
}
for _, g := range historyGroups {
existing, err := e.partitionsOf(ctx, g.parent)
@@ -111,7 +118,7 @@ func (e *Executor) HistoryMaintain(ctx context.Context, params driver.HistoryMai
}
// Rows in the DEFAULT partition cannot be dropped by period.
n, err := e.Conn.Exec(ctx,
- fmt.Sprintf("DELETE FROM %s_default WHERE finalized_at < now() - make_interval(secs => $1::float8)", g.parent),
+ fmt.Sprintf("DELETE FROM %s_default WHERE %s < now() - make_interval(secs => $1::float8)", g.parent, g.timeColumn),
keep.Seconds())
if err != nil {
return res, fmt.Errorf("hopper: prune %s_default: %w", g.parent, err)
@@ -172,7 +179,7 @@ func (e *Executor) createPartition(ctx context.Context, g historyGroup, name str
// SHARE UPDATE EXCLUSIVE on it. ATTACH does scan the DEFAULT partition
// for rows belonging to the new period, under an exclusive lock on that
// small table alone; the period is in the future, so it finds none.
- if _, err = tx.Exec(ctx, "CREATE TABLE "+ident+" (LIKE hopper_job_history INCLUDING DEFAULTS INCLUDING CONSTRAINTS)"); err != nil {
+ if _, err = tx.Exec(ctx, "CREATE TABLE "+ident+" (LIKE "+quoteIdent(g.root)+" INCLUDING DEFAULTS INCLUDING CONSTRAINTS)"); err != nil {
return false, fmt.Errorf("hopper: create partition %s: %w", name, err)
}
if _, err = tx.Exec(ctx, fmt.Sprintf("ALTER TABLE %s ATTACH PARTITION %s FOR VALUES FROM (%s) TO (%s)", parent, ident, fromLit, toLit)); err != nil {
diff --git a/driver/internal/pgsql/stream.go b/driver/internal/pgsql/stream.go
new file mode 100644
index 0000000..9e37224
--- /dev/null
+++ b/driver/internal/pgsql/stream.go
@@ -0,0 +1,350 @@
+package pgsql
+
+import (
+ "context"
+ "database/sql"
+ "encoding/json"
+ "errors"
+ "fmt"
+ "strconv"
+
+ "github.com/parallelworks/hopper/driver"
+)
+
+// Events are identified by (xid, seq): the transaction that wrote them,
+// then the event within it. A consumer reads the log by snapshot deltas.
+// Its row holds the snapshot it has read to (seen): every transaction
+// visible in it has been delivered. A pump takes a fresh snapshot and
+// delivers the events of the transactions visible in the new one but not
+// in seen, in (xid, seq) order; when a pump stops part way through a
+// delta, the row keeps the new snapshot (reading) and the last event
+// delivered, and the next pump continues the same delta. A transaction
+// that commits late is therefore delivered when it commits, whatever its
+// xid, and a long-running transaction holds nothing back. The cursor is
+// reset when a delta is exhausted, since the next delta's events may sort
+// before it; a seek sets it to skip part of a transaction.
+
+const streamEventColumns = `xid::text, seq, topic, key, payload, headers, message_id::text, created_at`
+
+func scanEvent(r Row) (*driver.StreamEvent, error) {
+ var (
+ e driver.StreamEvent
+ xid string
+ key sql.NullString
+ payload jsonText
+ headers jsonText
+ )
+ if err := r.Scan(&xid, &e.Position.Seq, &e.Topic, &key, &payload, &headers, &e.MessageID, &e.CreatedAt); err != nil {
+ return nil, err
+ }
+ var err error
+ if e.Position.Xid, err = strconv.ParseUint(xid, 10, 64); err != nil {
+ return nil, fmt.Errorf("xid %q: %w", xid, err)
+ }
+ e.Key, e.Payload = key.String, payload.raw()
+ if headers != "" && headers != "{}" {
+ if err := json.Unmarshal([]byte(headers), &e.Headers); err != nil {
+ return nil, fmt.Errorf("headers: %w", err)
+ }
+ }
+ return &e, nil
+}
+
+var streamEvents = collect(scanEvent)
+
+func xidParam(p driver.StreamPosition) string { return strconv.FormatUint(p.Xid, 10) }
+
+// StreamAppend implements driver.Executor.
+func (e *Executor) StreamAppend(ctx context.Context, params driver.StreamAppendParams) (*driver.StreamEvent, error) {
+ headers := "{}"
+ if params.Headers != nil {
+ b, err := json.Marshal(params.Headers)
+ if err != nil {
+ return nil, fmt.Errorf("hopper: append: encode headers: %w", err)
+ }
+ headers = string(b)
+ }
+ payload := "null"
+ if len(params.Payload) > 0 {
+ payload = string(params.Payload)
+ }
+ query := `INSERT INTO hopper_stream_events (topic, key, payload, headers) VALUES ($1, $2, $3::jsonb, $4::jsonb) RETURNING ` + streamEventColumns
+ args := []any{params.Topic, nullable(params.Key), payload, headers}
+ var (
+ events []*driver.StreamEvent
+ err error
+ )
+ if params.Notify {
+ events, err = streamEvents(e.Conn.QueryExec(ctx, query, args, NotifySQL, []any{driver.ChannelStream, textArray([]string{params.Topic})}))
+ } else {
+ events, err = streamEvents(e.Conn.Query(ctx, query, args...))
+ }
+ if err != nil {
+ return nil, fmt.Errorf("hopper: append: %w", err)
+ }
+ return events[0], nil
+}
+
+// streamReadSQL pages through committed events by position. The row
+// comparison walks the (xid, seq) index in order.
+const streamReadSQL = `
+SELECT ` + streamEventColumns + ` FROM hopper_stream_events
+WHERE (xid, seq) > ($1::xid8, $2::bigint) AND ($3::text = '' OR topic ~ hopper_topic_regex($3::text))
+ORDER BY xid, seq LIMIT $4`
+
+// StreamRead implements driver.Executor.
+func (e *Executor) StreamRead(ctx context.Context, params driver.StreamReadParams) ([]*driver.StreamEvent, error) {
+ events, err := streamEvents(e.Conn.Query(ctx, streamReadSQL, xidParam(params.After), params.After.Seq, params.Pattern, max(params.Limit, 1)))
+ if err != nil {
+ return nil, fmt.Errorf("hopper: read stream: %w", err)
+ }
+ return events, nil
+}
+
+// earliestSnapshot sees no transaction at all: 3 is the first transaction
+// ID Postgres assigns.
+const earliestSnapshot = "3:3:"
+
+// StreamConsumerUpsert implements driver.Executor.
+func (e *Executor) StreamConsumerUpsert(ctx context.Context, consumers []driver.StreamConsumerRow) error {
+ if len(consumers) == 0 {
+ return nil
+ }
+ n := len(consumers)
+ var (
+ names, patterns, kinds, queues, starts = make([]string, n), make([]string, n), make([]string, n), make([]string, n), make([]string, n)
+ maxAttempts = make([]*int64, n)
+ metadata = make([][]byte, n)
+ )
+ for i, c := range consumers {
+ names[i], patterns[i], kinds[i], queues[i] = c.Name, c.Pattern, c.Kind, c.Queue
+ starts[i] = string(c.Start)
+ if c.MaxAttempts > 0 {
+ v, err := smallint(c.MaxAttempts, "max_attempts")
+ if err != nil {
+ return err
+ }
+ n := int64(v)
+ maxAttempts[i] = &n
+ }
+ metadata[i] = []byte(jsonOrEmptyObject(c.Metadata))
+ }
+ // A consumer starting at latest has seen everything committed so far;
+ // transactions still running are delivered when they commit.
+ _, err := e.Conn.Exec(ctx, `
+ INSERT INTO hopper_stream_consumers (name, pattern, kind, queue, max_attempts, metadata, seen)
+ SELECT p.name, p.pattern, p.kind, p.queue, p.max_attempts, p.metadata,
+ CASE WHEN p.start = 'earliest' THEN $8::pg_snapshot ELSE pg_current_snapshot() END
+ FROM unnest($1::text[], $2::text[], $3::text[], $4::text[], $5::smallint[], $6::jsonb[], $7::text[])
+ AS p(name, pattern, kind, queue, max_attempts, metadata, start)
+ ON CONFLICT (name) DO UPDATE SET pattern = EXCLUDED.pattern, kind = EXCLUDED.kind, queue = EXCLUDED.queue,
+ max_attempts = EXCLUDED.max_attempts, metadata = EXCLUDED.metadata`,
+ textArray(names), textArray(patterns), textArray(kinds), textArray(queues), maxAttempts, jsonArray(metadata), textArray(starts), earliestSnapshot)
+ if err != nil {
+ return fmt.Errorf("hopper: upsert stream consumers: %w", err)
+ }
+ return nil
+}
+
+const streamConsumerColumns = `name, pattern, kind, queue, max_attempts, metadata, seen::text, xid::text, seq, delivered_at, created_at`
+
+func scanConsumer(r Row) (*driver.StreamConsumerRow, error) {
+ var (
+ c driver.StreamConsumerRow
+ maxAttempts sql.NullInt64
+ metadata jsonText
+ xid string
+ delivered sql.NullTime
+ )
+ if err := r.Scan(&c.Name, &c.Pattern, &c.Kind, &c.Queue, &maxAttempts, &metadata, &c.Snapshot, &xid, &c.Position.Seq, &delivered, &c.CreatedAt); err != nil {
+ return nil, err
+ }
+ var err error
+ if c.Position.Xid, err = strconv.ParseUint(xid, 10, 64); err != nil {
+ return nil, fmt.Errorf("xid %q: %w", xid, err)
+ }
+ c.MaxAttempts, c.Metadata, c.DeliveredAt = int(maxAttempts.Int64), metadata.raw(), delivered.Time
+ return &c, nil
+}
+
+// StreamConsumerList implements driver.Executor.
+func (e *Executor) StreamConsumerList(ctx context.Context) ([]*driver.StreamConsumerRow, error) {
+ consumers, err := collect(scanConsumer)(e.Conn.Query(ctx, `SELECT `+streamConsumerColumns+` FROM hopper_stream_consumers ORDER BY name`))
+ if err != nil {
+ return nil, fmt.Errorf("hopper: list stream consumers: %w", err)
+ }
+ return consumers, nil
+}
+
+// StreamConsumerSeek implements driver.Executor. Seeking to a position or
+// a time resumes the log from a transaction ID onwards: the snapshot
+// "x:x:" sees every transaction below x, and the cursor skips the events
+// of x itself up to the position.
+func (e *Executor) StreamConsumerSeek(ctx context.Context, name string, params driver.StreamSeekParams) error {
+ err := e.withTx(ctx, func(tx Tx) error {
+ var (
+ pos driver.StreamPosition
+ n int64
+ err error
+ )
+ switch {
+ case !params.Position.IsZero():
+ pos = params.Position
+ case !params.Time.IsZero():
+ // Just before the first event appended at or after the time.
+ // Events of one transaction share created_at, so seq - 1 is
+ // before all of them.
+ var xid string
+ err := tx.QueryRow(ctx, `SELECT xid::text, seq FROM hopper_stream_events WHERE created_at >= $1 ORDER BY xid, seq LIMIT 1`, params.Time).Scan(&xid, &pos.Seq)
+ if errors.Is(err, ErrNoRows) {
+ params.Latest = true
+ break
+ }
+ if err != nil {
+ return err
+ }
+ if pos.Xid, err = strconv.ParseUint(xid, 10, 64); err != nil {
+ return err
+ }
+ pos.Seq--
+ case params.Earliest:
+ case params.Latest:
+ default:
+ return errors.New("nothing to seek to")
+ }
+ switch {
+ case params.Latest:
+ n, err = tx.Exec(ctx, `UPDATE hopper_stream_consumers SET seen = pg_current_snapshot(), reading = NULL, xid = '0', seq = 0 WHERE name = $1`, name)
+ case pos.IsZero():
+ n, err = tx.Exec(ctx, `UPDATE hopper_stream_consumers SET seen = $2::pg_snapshot, reading = NULL, xid = '0', seq = 0 WHERE name = $1`, name, earliestSnapshot)
+ default:
+ snapshot := fmt.Sprintf("%d:%d:", pos.Xid, pos.Xid)
+ n, err = tx.Exec(ctx, `UPDATE hopper_stream_consumers SET seen = $2::pg_snapshot, reading = NULL, xid = $3::xid8, seq = $4 WHERE name = $1`, name, snapshot, xidParam(pos), pos.Seq)
+ }
+ if err != nil {
+ return err
+ }
+ if n == 0 {
+ return driver.ErrNotFound
+ }
+ return nil
+ })
+ if err != nil {
+ if errors.Is(err, driver.ErrNotFound) {
+ return err
+ }
+ return fmt.Errorf("hopper: seek stream consumer %s: %w", name, err)
+ }
+ return nil
+}
+
+// streamPumpSQL delivers the events of the transactions visible in one
+// snapshot but not in another, after a cursor, as jobs. The xid range
+// bounds the index scan: a transaction not visible in the old snapshot has
+// xid >= its xmin, and one visible in the new snapshot has xid < its xmax.
+// Deliveries carry the topic, message ID, headers, consumer name and
+// position in their metadata, and the event's key as their ordering key.
+// It returns how many, the last position, and the queues delivered to.
+const streamPumpSQL = `
+WITH ev AS (
+ SELECT xid, seq, topic, key, payload, headers, message_id FROM hopper_stream_events
+ WHERE xid >= pg_snapshot_xmin($1::pg_snapshot) AND xid < pg_snapshot_xmax($2::pg_snapshot)
+ AND (xid, seq) > ($3::xid8, $4::bigint)
+ AND pg_visible_in_snapshot(xid, $2::pg_snapshot) AND NOT pg_visible_in_snapshot(xid, $1::pg_snapshot)
+ AND topic ~ hopper_topic_regex($5::text)
+ ORDER BY xid, seq LIMIT $6
+),
+ins AS (
+ INSERT INTO hopper_jobs (kind, queue, priority, max_attempts, args, metadata, ordering_key)
+ SELECT $7::text, $8::text, 2, coalesce($9::smallint, CASE WHEN ev.key IS NOT NULL THEN 10 ELSE 25 END), ev.payload,
+ $10::jsonb || jsonb_build_object('topic', ev.topic, 'message_id', ev.message_id::text, 'headers', ev.headers,
+ 'stream', $11::text, 'position', ev.xid::text || ':' || ev.seq::text),
+ ev.key
+ FROM ev
+ RETURNING queue
+),
+last AS (SELECT xid::text AS xid, seq FROM ev ORDER BY xid DESC, seq DESC LIMIT 1)
+SELECT (SELECT count(*) FROM ins), (SELECT xid FROM last), (SELECT seq FROM last), (SELECT string_agg(DISTINCT queue, $12::text) FROM ins)`
+
+// StreamPump implements driver.Executor.
+func (e *Executor) StreamPump(ctx context.Context, name string, limit int) (driver.StreamPumpResult, error) {
+ var res driver.StreamPumpResult
+ limit = max(limit, 1)
+ err := e.withTx(ctx, func(tx Tx) error {
+ // Lock first, then read: a statement that waited for the lock keeps
+ // its snapshot, and would not see what the previous pump delivered.
+ var (
+ c driver.StreamConsumerRow
+ maxAttempts sql.NullInt64
+ metadata jsonText
+ seen, reading sql.NullString
+ xid string
+ )
+ err := tx.QueryRow(ctx, `SELECT pattern, kind, queue, max_attempts, metadata, seen::text, reading::text, xid::text, seq FROM hopper_stream_consumers WHERE name = $1 FOR UPDATE`, name).
+ Scan(&c.Pattern, &c.Kind, &c.Queue, &maxAttempts, &metadata, &seen, &reading, &xid, &c.Position.Seq)
+ if errors.Is(err, ErrNoRows) {
+ return driver.ErrNotFound
+ }
+ if err != nil {
+ return err
+ }
+ if c.Position.Xid, err = strconv.ParseUint(xid, 10, 64); err != nil {
+ return err
+ }
+ res.Position = c.Position
+ // Continue a delta left part way, or start a new one from a fresh
+ // snapshot: everything committed by now.
+ if !reading.Valid {
+ if err := tx.QueryRow(ctx, `SELECT pg_current_snapshot()::text`).Scan(&reading.String); err != nil {
+ return err
+ }
+ }
+ var attempts *int64
+ if maxAttempts.Valid {
+ attempts = &maxAttempts.Int64
+ }
+ var (
+ delivered int64
+ lastXid sql.NullString
+ lastSeq sql.NullInt64
+ queues joined
+ )
+ if err := tx.QueryRow(ctx, streamPumpSQL, seen.String, reading.String, xidParam(c.Position), c.Position.Seq, c.Pattern, limit,
+ c.Kind, c.Queue, attempts, jsonOrEmptyObject(metadata.raw()), name, joinSep).Scan(&delivered, &lastXid, &lastSeq, &queues); err != nil {
+ return err
+ }
+ res.Delivered = int(delivered)
+ res.Queues = queues
+ if lastXid.Valid {
+ if res.Position.Xid, err = strconv.ParseUint(lastXid.String, 10, 64); err != nil {
+ return err
+ }
+ res.Position.Seq = lastSeq.Int64
+ }
+ if res.Delivered < limit {
+ // The delta is exhausted: the new snapshot is what has been seen.
+ _, err = tx.Exec(ctx, `UPDATE hopper_stream_consumers SET seen = $2::pg_snapshot, reading = NULL, xid = '0', seq = 0,
+ delivered_at = CASE WHEN $3 > 0 THEN now() ELSE delivered_at END WHERE name = $1`,
+ name, reading.String, res.Delivered)
+ } else {
+ _, err = tx.Exec(ctx, `UPDATE hopper_stream_consumers SET reading = $2::pg_snapshot, xid = $3::xid8, seq = $4, delivered_at = now() WHERE name = $1`,
+ name, reading.String, xidParam(res.Position), res.Position.Seq)
+ }
+ if err != nil {
+ return err
+ }
+ if len(queues) > 0 {
+ if _, err := tx.Exec(ctx, NotifySQL, driver.ChannelInsert, textArray(queues)); err != nil {
+ return err
+ }
+ }
+ return nil
+ })
+ if err != nil {
+ if errors.Is(err, driver.ErrNotFound) {
+ return res, err
+ }
+ return res, fmt.Errorf("hopper: pump stream consumer %s: %w", name, err)
+ }
+ return res, nil
+}
diff --git a/drivertest/drivertest.go b/drivertest/drivertest.go
index 06e3356..a6fdb12 100644
--- a/drivertest/drivertest.go
+++ b/drivertest/drivertest.go
@@ -89,6 +89,7 @@ func Run[TTx any](t *testing.T, f Fixture[TTx]) {
{"Aging", testAging[TTx]},
{"Batches", testBatches[TTx]},
{"Workflows", testWorkflows[TTx]},
+ {"Streams", testStreams[TTx]},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
diff --git a/drivertest/stream.go b/drivertest/stream.go
new file mode 100644
index 0000000..5ad29c3
--- /dev/null
+++ b/drivertest/stream.go
@@ -0,0 +1,234 @@
+package drivertest
+
+import (
+ "context"
+ "encoding/json"
+ "errors"
+ "fmt"
+ "strings"
+ "testing"
+ "time"
+
+ "github.com/parallelworks/hopper/driver"
+)
+
+func testStreams[TTx any](t *testing.T, f Fixture[TTx]) {
+ ctx := context.Background()
+ d := f.NewDriver(t)
+ exec := d.Executor()
+ appendEvent := func(topic, key string, i int) *driver.StreamEvent {
+ t.Helper()
+ ev, err := exec.StreamAppend(ctx, driver.StreamAppendParams{
+ Topic: topic, Key: key, Payload: json.RawMessage(fmt.Sprintf(`{"i":%d}`, i)),
+ Headers: map[string]string{"n": fmt.Sprint(i)},
+ })
+ if err != nil {
+ t.Fatal(err)
+ }
+ return ev
+ }
+ read := func(pattern string, after driver.StreamPosition) []*driver.StreamEvent {
+ t.Helper()
+ events, err := exec.StreamRead(ctx, driver.StreamReadParams{Pattern: pattern, After: after, Limit: 100})
+ if err != nil {
+ t.Fatal(err)
+ }
+ return events
+ }
+ pump := func(name string) driver.StreamPumpResult {
+ t.Helper()
+ res, err := exec.StreamPump(ctx, name, 100)
+ if err != nil {
+ t.Fatal(err)
+ }
+ return res
+ }
+ deliveries := func(queue string) []*driver.JobRow {
+ t.Helper()
+ jobs, err := exec.JobList(ctx, driver.JobListParams{Queue: queue, Limit: 100})
+ if err != nil {
+ t.Fatal(err)
+ }
+ return jobs
+ }
+
+ // Appends are ordered and readable by pattern.
+ e1 := appendEvent("order.created", "o1", 1)
+ e2 := appendEvent("order.paid", "o1", 2)
+ e3 := appendEvent("user.created", "", 3)
+ if e1.Position.IsZero() || !e1.Position.Less(e2.Position) || !e2.Position.Less(e3.Position) {
+ t.Errorf("positions not ascending: %v %v %v", e1.Position, e2.Position, e3.Position)
+ }
+ if e1.Key != "o1" || e1.Headers["n"] != "1" || e1.MessageID == "" || e1.CreatedAt.IsZero() || string(e1.Payload) != `{"i": 1}` {
+ t.Errorf("event = %+v", e1)
+ }
+ if got := read("", driver.StreamPosition{}); len(got) != 3 || got[0].Position != e1.Position || got[2].Position != e3.Position {
+ t.Errorf("read all = %v", got)
+ }
+ if got := read("order.*", driver.StreamPosition{}); len(got) != 2 || got[1].Topic != "order.paid" {
+ t.Errorf("read order.* = %v", got)
+ }
+ if got := read("user.created", e1.Position); len(got) != 1 || got[0].Position != e3.Position {
+ t.Errorf("read after = %v", got)
+ }
+ if got := read("", e3.Position); len(got) != 0 {
+ t.Errorf("read past the end = %v", got)
+ }
+
+ // A consumer created at the earliest position delivers everything so
+ // far, once, as jobs carrying the event; one created at latest starts
+ // with what comes next.
+ err := exec.StreamConsumerUpsert(ctx, []driver.StreamConsumerRow{
+ {Name: "orders", Pattern: "order.#", Kind: "stream:orders", Queue: "orders", Start: driver.StreamStartEarliest, Metadata: json.RawMessage(`{"c":1}`)},
+ {Name: "audit", Pattern: "#", Kind: "stream:audit", Queue: "audit", Start: driver.StreamStartLatest, MaxAttempts: 3},
+ })
+ if err != nil {
+ t.Fatal(err)
+ }
+ consumers, err := exec.StreamConsumerList(ctx)
+ if err != nil || len(consumers) != 2 || consumers[0].Name != "audit" || consumers[1].Name != "orders" {
+ t.Fatalf("consumers = %v, %v", consumers, err)
+ }
+ if c := consumers[1]; !c.Position.IsZero() || c.Queue != "orders" || string(c.Metadata) != `{"c": 1}` || c.CreatedAt.IsZero() || c.Snapshot == "" {
+ t.Errorf("orders consumer = %+v", c)
+ }
+ if c := consumers[0]; c.MaxAttempts != 3 || c.Snapshot == consumers[1].Snapshot {
+ t.Errorf("audit consumer = %+v", c)
+ }
+ res := pump("orders")
+ if res.Delivered != 2 || res.Position != e2.Position || len(res.Queues) != 1 || res.Queues[0] != "orders" {
+ t.Errorf("first pump = %+v", res)
+ }
+ jobs := deliveries("orders")
+ if len(jobs) != 2 {
+ t.Fatalf("deliveries = %v", jobs)
+ }
+ j := jobs[0]
+ if argsN(j) != 1 {
+ j = jobs[1]
+ }
+ if j.Kind != "stream:orders" || j.OrderingKey != "o1" || j.MaxAttempts != 10 || string(j.Args) != `{"i": 1}` {
+ t.Errorf("delivery = %+v", j)
+ }
+ var meta struct {
+ Topic string `json:"topic"`
+ MessageID string `json:"message_id"`
+ Stream string `json:"stream"`
+ Position string `json:"position"`
+ Headers map[string]string `json:"headers"`
+ C int `json:"c"`
+ }
+ if err := json.Unmarshal(j.Metadata, &meta); err != nil || meta.Topic != "order.created" || meta.MessageID != e1.MessageID ||
+ meta.Stream != "orders" || meta.Position != e1.Position.String() || meta.Headers["n"] != "1" || meta.C != 1 {
+ t.Errorf("delivery metadata = %s (%v)", j.Metadata, err)
+ }
+ if res := pump("orders"); res.Delivered != 0 {
+ t.Errorf("second pump = %+v", res)
+ }
+ if res := pump("audit"); res.Delivered != 0 {
+ t.Errorf("latest consumer delivered old events: %+v", res)
+ }
+ if c, _ := exec.StreamConsumerList(ctx); c[1].DeliveredAt.IsZero() || !c[0].DeliveredAt.IsZero() {
+ t.Errorf("delivered_at: orders %v audit %v", c[1].DeliveredAt, c[0].DeliveredAt)
+ }
+
+ // New events reach both, and a pump of a small batch continues where
+ // it stopped.
+ e4 := appendEvent("user.deleted", "", 4)
+ e5 := appendEvent("order.shipped", "o1", 5)
+ if res := pump("orders"); res.Delivered != 1 || res.Position != e5.Position {
+ t.Errorf("pump after appends = %+v", res)
+ }
+ if res, err := exec.StreamPump(ctx, "audit", 1); err != nil || res.Delivered != 1 || res.Position != e4.Position {
+ t.Errorf("pump of one = %+v, %v", res, err)
+ }
+ if res, err := exec.StreamPump(ctx, "audit", 1); err != nil || res.Delivered != 1 || res.Position != e5.Position {
+ t.Errorf("pump of the next one = %+v, %v", res, err)
+ }
+ if res := pump("audit"); res.Delivered != 0 {
+ t.Errorf("audit pump after the delta = %+v", res)
+ }
+ if jobs := deliveries("audit"); len(jobs) != 2 || jobs[0].MaxAttempts != 3 || jobs[0].OrderingKey != "" {
+ t.Errorf("audit deliveries = %v", jobs)
+ }
+
+ // Seeking replays: to a position, to a time, to the start, to now.
+ seek := func(params driver.StreamSeekParams) {
+ t.Helper()
+ if err := exec.StreamConsumerSeek(ctx, "orders", params); err != nil {
+ t.Fatal(err)
+ }
+ }
+ seek(driver.StreamSeekParams{Position: e2.Position})
+ if res := pump("orders"); res.Delivered != 1 || res.Position != e5.Position {
+ t.Errorf("pump after seek to position = %+v", res)
+ }
+ seek(driver.StreamSeekParams{Earliest: true})
+ if res := pump("orders"); res.Delivered != 3 {
+ t.Errorf("pump after seek to earliest = %+v", res)
+ }
+ seek(driver.StreamSeekParams{Time: e4.CreatedAt})
+ if res := pump("orders"); res.Delivered != 1 || res.Position != e5.Position {
+ t.Errorf("pump after seek to time = %+v", res)
+ }
+ seek(driver.StreamSeekParams{Earliest: true})
+ seek(driver.StreamSeekParams{Latest: true})
+ if res := pump("orders"); res.Delivered != 0 {
+ t.Errorf("pump after seek to latest = %+v", res)
+ }
+ if err := exec.StreamConsumerSeek(ctx, "nobody", driver.StreamSeekParams{Earliest: true}); !errors.Is(err, driver.ErrNotFound) {
+ t.Errorf("seek unknown = %v", err)
+ }
+ if _, err := exec.StreamPump(ctx, "nobody", 10); !errors.Is(err, driver.ErrNotFound) {
+ t.Errorf("pump unknown = %v", err)
+ }
+
+ // A transaction that commits late is delivered when it commits, and
+ // does not hold up the events of younger transactions meanwhile.
+ tx, commit, rollback := f.Begin(ctx, t, d)
+ defer rollback() //nolint:errcheck // already committed on the happy path
+ old, err := d.UnwrapTx(tx).StreamAppend(ctx, driver.StreamAppendParams{Topic: "order.held", Payload: json.RawMessage(`{}`)})
+ if err != nil {
+ t.Fatal(err)
+ }
+ young := appendEvent("order.after", "", 6)
+ if !old.Position.Less(young.Position) {
+ t.Fatalf("older transaction sorts after the younger one: %v %v", old.Position, young.Position)
+ }
+ if res := pump("orders"); res.Delivered != 1 || res.Position != young.Position {
+ t.Errorf("pump with a transaction open = %+v", res)
+ }
+ if err := commit(); err != nil {
+ t.Fatal(err)
+ }
+ if res := pump("orders"); res.Delivered != 1 || res.Position != old.Position {
+ t.Errorf("pump after the late commit = %+v", res)
+ }
+ if got := read("", e5.Position); len(got) != 2 || got[0].Position != old.Position || got[1].Position != young.Position {
+ t.Errorf("read after commit = %v", got)
+ }
+ if res := pump("orders"); res.Delivered != 0 {
+ t.Errorf("late commit delivered twice: %+v", res)
+ }
+
+ // Retention keeps a partition ahead and drops nothing young.
+ m, err := exec.HistoryMaintain(ctx, driver.HistoryMaintainParams{CompletedRetention: time.Hour, FailedRetention: time.Hour, StreamRetention: time.Hour})
+ if err != nil {
+ t.Fatal(err)
+ }
+ if !containsPrefix(m.Created, "hopper_stream_events_") {
+ t.Errorf("no stream partition created: %v", m.Created)
+ }
+ if got := read("", driver.StreamPosition{}); len(got) != 7 {
+ t.Errorf("events after maintenance = %d, want 7", len(got))
+ }
+}
+
+func containsPrefix(names []string, prefix string) bool {
+ for _, n := range names {
+ if strings.HasPrefix(n, prefix) {
+ return true
+ }
+ }
+ return false
+}
diff --git a/export_test.go b/export_test.go
index dbf1618..c712e67 100644
--- a/export_test.go
+++ b/export_test.go
@@ -33,6 +33,7 @@ func (c *Client[TTx]) SetTuning(t Tuning) {
leaderInterval: t.LeaderInterval,
maintenanceInterval: t.Maintenance,
rescueBatch: t.RescueBatch,
+ streamBatch: defaultTuning.streamBatch,
}
c.notifier.interval = c.tuning.notifyInterval
}
diff --git a/hoppermigrate/migrations/005_streams.down.sql b/hoppermigrate/migrations/005_streams.down.sql
new file mode 100644
index 0000000..c8d47fa
--- /dev/null
+++ b/hoppermigrate/migrations/005_streams.down.sql
@@ -0,0 +1,2 @@
+DROP TABLE hopper_stream_consumers;
+DROP TABLE hopper_stream_events;
diff --git a/hoppermigrate/migrations/005_streams.up.sql b/hoppermigrate/migrations/005_streams.up.sql
new file mode 100644
index 0000000..0010487
--- /dev/null
+++ b/hoppermigrate/migrations/005_streams.up.sql
@@ -0,0 +1,39 @@
+-- hopper schema v5: streams.
+--
+-- A stream is an append-only log of events, partitioned by day and dropped
+-- by retention like history. Events are identified by (xid, seq): the
+-- transaction that wrote them, then the event within it. A consumer reads
+-- the log by snapshot deltas: the events of transactions committed since
+-- the snapshot it last read to, whatever their xid, so a transaction that
+-- commits late is delivered when it commits and never skipped, and a long
+-- transaction never holds the others back. The leader turns the events of
+-- each delta into delivery jobs.
+
+CREATE TABLE hopper_stream_events (
+ xid xid8 NOT NULL DEFAULT pg_current_xact_id(),
+ seq bigint GENERATED ALWAYS AS IDENTITY,
+ topic text NOT NULL,
+ key text, -- becomes the delivery's ordering key
+ payload jsonb NOT NULL DEFAULT 'null',
+ headers jsonb NOT NULL DEFAULT '{}',
+ message_id uuid NOT NULL DEFAULT hopper_uuidv7(),
+ created_at timestamptz NOT NULL DEFAULT now()
+) PARTITION BY RANGE (created_at);
+
+CREATE INDEX hopper_stream_events_position ON hopper_stream_events (xid, seq);
+CREATE TABLE hopper_stream_events_default PARTITION OF hopper_stream_events DEFAULT;
+
+CREATE TABLE hopper_stream_consumers (
+ name text PRIMARY KEY,
+ pattern text NOT NULL, -- AMQP topic syntax
+ kind text NOT NULL, -- the delivery job kind
+ queue text NOT NULL,
+ max_attempts smallint,
+ metadata jsonb NOT NULL DEFAULT '{}',
+ seen pg_snapshot NOT NULL DEFAULT '3:3:', -- every transaction visible here has been delivered
+ reading pg_snapshot, -- the delta being delivered, when a pump stopped part way
+ xid xid8 NOT NULL DEFAULT '0', -- the last event delivered of that delta
+ seq bigint NOT NULL DEFAULT 0,
+ delivered_at timestamptz,
+ created_at timestamptz NOT NULL DEFAULT now()
+);
diff --git a/leader.go b/leader.go
index d84534e..5ed7498 100644
--- a/leader.go
+++ b/leader.go
@@ -3,6 +3,7 @@ package hopper
import (
"context"
"slices"
+ "sync"
"time"
"github.com/parallelworks/hopper/driver"
@@ -65,8 +66,8 @@ func (c *Client[TTx]) leaderLoop(ctx context.Context) {
}
}
-// leaderTerm is one stretch of leadership. The periodic loop runs for its
-// duration.
+// leaderTerm is one stretch of leadership. The periodic and stream loops
+// run for its duration.
type leaderTerm struct {
cancel context.CancelFunc
done chan struct{}
@@ -75,11 +76,16 @@ type leaderTerm struct {
func (c *Client[TTx]) startTerm(ctx context.Context) *leaderTerm {
termCtx, cancel := context.WithCancel(ctx)
t := &leaderTerm{cancel: cancel, done: make(chan struct{})}
- go func() {
- defer close(t.done)
+ var wg sync.WaitGroup
+ wg.Go(func() {
if len(c.cfg.Periodic) > 0 {
c.periodicLoop(termCtx)
}
+ })
+ wg.Go(func() { c.streamLoop(termCtx) })
+ go func() {
+ defer close(t.done)
+ wg.Wait()
}()
return t
}
@@ -274,6 +280,7 @@ func (c *Client[TTx]) maintainHistory(ctx context.Context) {
res, err := c.exec.HistoryMaintain(ctx, driver.HistoryMaintainParams{
CompletedRetention: c.cfg.CompletedRetention,
FailedRetention: c.cfg.FailedRetention,
+ StreamRetention: c.cfg.StreamRetention,
})
if err != nil {
if ctx.Err() == nil {
diff --git a/listener.go b/listener.go
index b067fdd..b5bd7ab 100644
--- a/listener.go
+++ b/listener.go
@@ -49,7 +49,7 @@ func (c *Client[TTx]) listenOnce(ctx context.Context) error {
return err
}
defer l.Close(context.WithoutCancel(ctx)) //nolint:errcheck // best effort on a failed connection
- if err := l.Listen(ctx, driver.ChannelInsert, driver.ChannelLeader, driver.ChannelControl, driver.ChannelDone); err != nil {
+ if err := l.Listen(ctx, driver.ChannelInsert, driver.ChannelLeader, driver.ChannelControl, driver.ChannelDone, driver.ChannelStream); err != nil {
return err
}
c.listening.Store(true)
@@ -70,6 +70,8 @@ func (c *Client[TTx]) listenOnce(ctx context.Context) error {
if id, err := ParseJobID(n.Payload); err == nil {
c.signalDone(id)
}
+ case driver.ChannelStream:
+ c.pokeStreams()
}
}
}
diff --git a/messaging.go b/messaging.go
index bd34cbe..53e40df 100644
--- a/messaging.go
+++ b/messaging.go
@@ -32,6 +32,10 @@ type Message[T any] struct {
ID string
Headers map[string]string
OrderingKey string
+ // Stream and Position are set for a consumer's deliveries: the
+ // consumer's name and the event's place in the log.
+ Stream string
+ Position StreamPosition
}
// Subscription routes messages whose topic matches Pattern to a queue, as
@@ -118,14 +122,22 @@ type messageMetadata struct {
Topic string `json:"topic"`
MessageID string `json:"message_id"`
Headers map[string]string `json:"headers"`
+ Stream string `json:"stream,omitempty"`
+ Position string `json:"position,omitempty"`
}
func decodeMessage[T any](row *JobRow, codec Codec) (*Message[T], error) {
var meta messageMetadata
- if err := json.Unmarshal(row.Metadata, &meta); err != nil {
+ err := json.Unmarshal(row.Metadata, &meta)
+ if err != nil {
return nil, fmt.Errorf("hopper: decode message metadata: %w", err)
}
- msg := &Message[T]{JobRow: row, Topic: meta.Topic, ID: meta.MessageID, Headers: meta.Headers, OrderingKey: row.OrderingKey}
+ msg := &Message[T]{JobRow: row, Topic: meta.Topic, ID: meta.MessageID, Headers: meta.Headers, OrderingKey: row.OrderingKey, Stream: meta.Stream}
+ if meta.Position != "" {
+ if msg.Position, err = ParseStreamPosition(meta.Position); err != nil {
+ return nil, err
+ }
+ }
if err := codec.Unmarshal(row.Args, &msg.Payload); err != nil {
return nil, fmt.Errorf("hopper: decode message payload: %w", err)
}
diff --git a/notifier.go b/notifier.go
index d82422b..a58f468 100644
--- a/notifier.go
+++ b/notifier.go
@@ -13,7 +13,8 @@ import (
// transaction, coalesced so that a client sends at most one notification per
// queue per interval however many jobs it inserts. It fires on the trailing
// edge, after the inserts of the window have committed, so a notified worker
-// always finds the jobs.
+// always finds the jobs. Stream appends are coalesced the same way, per
+// topic, on the stream channel.
//
// Inserts inside a caller's transaction notify from within the transaction
// instead, since only Postgres knows when it commits.
@@ -24,11 +25,12 @@ type notifier struct {
mu sync.Mutex
pending map[string]struct{}
+ streams map[string]struct{}
timer *time.Timer
}
func newNotifier(exec driver.Executor, logger *slog.Logger, interval time.Duration) *notifier {
- return ¬ifier{exec: exec, logger: logger, interval: interval, pending: map[string]struct{}{}}
+ return ¬ifier{exec: exec, logger: logger, interval: interval, pending: map[string]struct{}{}, streams: map[string]struct{}{}}
}
// mark schedules a notification for queue.
@@ -41,6 +43,16 @@ func (n *notifier) mark(queue string) {
}
}
+// markStream schedules a stream notification for topic.
+func (n *notifier) markStream(topic string) {
+ n.mu.Lock()
+ defer n.mu.Unlock()
+ n.streams[topic] = struct{}{}
+ if n.timer == nil {
+ n.timer = time.AfterFunc(n.interval, n.flush)
+ }
+}
+
// flushNow sends anything pending without waiting for the timer.
func (n *notifier) flushNow() {
n.mu.Lock()
@@ -58,15 +70,27 @@ func (n *notifier) flush() {
queues = append(queues, q)
}
clear(n.pending)
+ topics := make([]string, 0, len(n.streams))
+ for t := range n.streams {
+ topics = append(topics, t)
+ }
+ clear(n.streams)
n.timer = nil
n.mu.Unlock()
- if len(queues) == 0 {
+ if len(queues) == 0 && len(topics) == 0 {
return
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
- if err := n.exec.Notify(ctx, driver.ChannelInsert, queues); err != nil {
- // Workers still poll, so a lost notification costs latency, not work.
- n.logger.WarnContext(ctx, "hopper: notify inserts", "queues", queues, "error", err)
+ if len(queues) > 0 {
+ if err := n.exec.Notify(ctx, driver.ChannelInsert, queues); err != nil {
+ // Workers still poll, so a lost notification costs latency, not work.
+ n.logger.WarnContext(ctx, "hopper: notify inserts", "queues", queues, "error", err)
+ }
+ }
+ if len(topics) > 0 {
+ if err := n.exec.Notify(ctx, driver.ChannelStream, topics); err != nil {
+ n.logger.WarnContext(ctx, "hopper: notify stream appends", "topics", topics, "error", err)
+ }
}
}
diff --git a/stream.go b/stream.go
new file mode 100644
index 0000000..8895796
--- /dev/null
+++ b/stream.go
@@ -0,0 +1,309 @@
+package hopper
+
+import (
+ "context"
+ "encoding/json"
+ "errors"
+ "fmt"
+ "iter"
+ "strings"
+ "time"
+ "unicode"
+
+ "github.com/parallelworks/hopper/driver"
+)
+
+// StreamEvent is one event of the stream log.
+type StreamEvent = driver.StreamEvent
+
+// StreamPosition is a place in the stream log.
+type StreamPosition = driver.StreamPosition
+
+// ParseStreamPosition parses a position's String form ("xid:seq").
+func ParseStreamPosition(s string) (StreamPosition, error) { return driver.ParseStreamPosition(s) }
+
+// StreamStart says where a new consumer starts.
+type StreamStart = driver.StreamStart
+
+const (
+ // StreamStartLatest delivers events appended from now on. It is the
+ // default.
+ StreamStartLatest = driver.StreamStartLatest
+ // StreamStartEarliest delivers every retained event first.
+ StreamStartEarliest = driver.StreamStartEarliest
+)
+
+// Consumer reads the stream log from a position of its own and delivers
+// the events whose topic matches Pattern to a queue, as jobs of kind
+// "stream:". Unlike a Subscription, which is fanned out to at publish
+// time, a consumer can start anywhere in the retained log and be moved
+// (Streams.Seek), so late subscribers and replays see past events.
+type Consumer struct {
+ // Name identifies the consumer; the delivery job kind is derived from
+ // it.
+ Name string
+ Pattern string
+ // Queue is where deliveries go. Defaults to QueueDefault.
+ Queue string
+ // Start applies when the consumer is created. Defaults to
+ // StreamStartLatest.
+ Start StreamStart
+ // MaxAttempts is the retry budget of deliveries. Zero means 25, or 10
+ // for events with a key, whose retries block their key.
+ MaxAttempts int
+ // Metadata is merged into every delivery's metadata.
+ Metadata json.RawMessage
+}
+
+// Kind returns the job kind of the consumer's deliveries.
+func (c Consumer) Kind() string { return "stream:" + c.Name }
+
+func (c *Consumer) validate() error {
+ if c.Name == "" || strings.ContainsFunc(c.Name, unicode.IsSpace) {
+ return fmt.Errorf("hopper: consumer name %q must be a non-empty word", c.Name)
+ }
+ if c.Pattern == "" {
+ return fmt.Errorf("hopper: consumer %q has no pattern", c.Name)
+ }
+ for word := range strings.SplitSeq(c.Pattern, ".") {
+ if word == "" {
+ return fmt.Errorf("hopper: consumer %q: pattern %q has an empty word", c.Name, c.Pattern)
+ }
+ }
+ if c.Queue == "" {
+ c.Queue = QueueDefault
+ }
+ switch c.Start {
+ case "":
+ c.Start = StreamStartLatest
+ case StreamStartLatest, StreamStartEarliest:
+ default:
+ return fmt.Errorf("hopper: consumer %q: unknown start %q", c.Name, c.Start)
+ }
+ if c.MaxAttempts < 0 {
+ return fmt.Errorf("hopper: consumer %q: MaxAttempts must not be negative", c.Name)
+ }
+ if len(c.Metadata) > 0 && !json.Valid(c.Metadata) {
+ return fmt.Errorf("hopper: consumer %q: Metadata is not valid JSON", c.Name)
+ }
+ return nil
+}
+
+// Consume registers a handler for a consumer's deliveries. The consumer is
+// created (at its Start) or updated when a client using workers starts,
+// and the leader delivers its events in order as jobs; fn receives them as
+// Subscribe's handler does, with Message.Stream and Message.Position set.
+// An event's key becomes the delivery's ordering key, so events with the
+// same key are handled one at a time, in order.
+//
+// It panics on an invalid consumer or a name registered twice, like
+// AddWorker.
+func Consume[T any](workers *Workers, consumer Consumer, fn func(ctx context.Context, msg *Message[T]) error) {
+ if fn == nil {
+ panic("hopper: Consume called with a nil function")
+ }
+ if err := consumer.validate(); err != nil {
+ panic(err)
+ }
+ workers.add(&workerInfo{
+ kind: consumer.Kind(),
+ newUnit: func(row *JobRow, codec Codec) (workUnit, error) {
+ msg, err := decodeMessage[T](row, codec)
+ if err != nil {
+ return nil, err
+ }
+ return &messageUnit[T]{fn: fn, msg: msg}, nil
+ },
+ })
+ workers.mu.Lock()
+ workers.consumers = append(workers.consumers, consumer)
+ workers.mu.Unlock()
+}
+
+// consumerRows returns the registered consumers as rows.
+func (w *Workers) consumerRows() []driver.StreamConsumerRow {
+ w.mu.RLock()
+ defer w.mu.RUnlock()
+ rows := make([]driver.StreamConsumerRow, len(w.consumers))
+ for i, c := range w.consumers {
+ rows[i] = driver.StreamConsumerRow{Name: c.Name, Pattern: c.Pattern, Kind: c.Kind(), Queue: c.Queue, Start: c.Start, MaxAttempts: c.MaxAttempts, Metadata: c.Metadata}
+ }
+ return rows
+}
+
+// AppendOpts customizes an append. A nil *AppendOpts means the defaults.
+type AppendOpts struct {
+ // Key orders the event's deliveries with others of the same key: a
+ // consumer handles them one at a time, in order.
+ Key string
+ // Headers travel with the event and are available on Message.
+ Headers map[string]string
+}
+
+// Streams is the stream log: an append-only, retained record of events
+// that consumers read from positions of their own.
+type Streams[TTx any] struct {
+ c *Client[TTx]
+}
+
+// Streams returns the client's stream log.
+func (c *Client[TTx]) Streams() *Streams[TTx] {
+ return &Streams[TTx]{c: c}
+}
+
+// Append appends an event to the log. Consumers whose pattern matches its
+// topic deliver it, in order, after everything appended before it.
+func (s *Streams[TTx]) Append(ctx context.Context, event Topic, opts *AppendOpts) (*StreamEvent, error) {
+ return s.append(ctx, s.c.exec, event, opts, false)
+}
+
+// AppendTx is Append in the caller's transaction: the event exists only if
+// tx commits, and is delivered after events of every transaction that
+// started before tx, however they commit.
+func (s *Streams[TTx]) AppendTx(ctx context.Context, tx TTx, event Topic, opts *AppendOpts) (*StreamEvent, error) {
+ return s.append(ctx, s.c.driver.UnwrapTx(tx), event, opts, true)
+}
+
+func (s *Streams[TTx]) append(ctx context.Context, exec driver.Executor, event Topic, opts *AppendOpts, inTx bool) (*StreamEvent, error) {
+ if event == nil {
+ return nil, errors.New("hopper: event must not be nil")
+ }
+ topic := event.Topic()
+ if topic == "" || strings.ContainsFunc(topic, unicode.IsSpace) {
+ return nil, fmt.Errorf("hopper: topic %q must be a non-empty word", topic)
+ }
+ if opts == nil {
+ opts = &AppendOpts{}
+ }
+ payload, err := s.c.cfg.Codec.Marshal(event)
+ if err != nil {
+ return nil, fmt.Errorf("hopper: encode event for %q: %w", topic, err)
+ }
+ ev, err := exec.StreamAppend(ctx, driver.StreamAppendParams{Topic: topic, Key: opts.Key, Payload: payload, Headers: opts.Headers, Notify: inTx})
+ if err != nil {
+ return nil, err
+ }
+ if !inTx {
+ s.c.pokeStreams()
+ s.c.notifier.markStream(topic)
+ }
+ return ev, nil
+}
+
+// Read returns the events after a position whose topic matches pattern
+// (empty for every topic), oldest first, in pages of limit. It sees only
+// events that every earlier transaction has finished around, as consumers
+// do, so paging by the last position never skips an event.
+func (s *Streams[TTx]) Read(ctx context.Context, pattern string, after StreamPosition, limit int) iter.Seq2[*StreamEvent, error] {
+ if limit <= 0 {
+ limit = 100
+ }
+ return func(yield func(*StreamEvent, error) bool) {
+ pos := after
+ for {
+ events, err := s.c.exec.StreamRead(ctx, driver.StreamReadParams{Pattern: pattern, After: pos, Limit: limit})
+ if err != nil {
+ yield(nil, err)
+ return
+ }
+ for _, ev := range events {
+ if !yield(ev, nil) {
+ return
+ }
+ pos = ev.Position
+ }
+ if len(events) < limit {
+ return
+ }
+ }
+ }
+}
+
+// ConsumerRow is a consumer as recorded in the database, with its
+// position.
+type ConsumerRow = driver.StreamConsumerRow
+
+// Consumers lists every consumer known to the cluster.
+func (s *Streams[TTx]) Consumers(ctx context.Context) ([]*ConsumerRow, error) {
+ return s.c.exec.StreamConsumerList(ctx)
+}
+
+// SeekOpts says where to move a consumer. Exactly one field applies, in
+// this order of precedence.
+type SeekOpts struct {
+ // Position makes the event after it the next delivered.
+ Position StreamPosition
+ // Time replays from the first event appended at or after it.
+ Time time.Time
+ // Earliest replays every retained event.
+ Earliest bool
+ // Latest skips to the events appended from now on.
+ Latest bool
+}
+
+// Seek moves a consumer's position; the events after it are delivered
+// (again) from the next pump. Deliveries already inserted are unaffected.
+func (s *Streams[TTx]) Seek(ctx context.Context, consumer string, opts SeekOpts) error {
+ if err := s.c.exec.StreamConsumerSeek(ctx, consumer, driver.StreamSeekParams(opts)); err != nil {
+ return err
+ }
+ s.c.pokeStreams()
+ if err := s.c.exec.Notify(ctx, driver.ChannelStream, []string{""}); err != nil {
+ s.c.logger.WarnContext(ctx, "hopper: notify stream seek", "error", err)
+ }
+ return nil
+}
+
+// pokeStreams asks the leader's stream loop to pump now.
+func (c *Client[TTx]) pokeStreams() {
+ select {
+ case c.streamPoke <- struct{}{}:
+ default:
+ }
+}
+
+// streamLoop runs on the leader: it pumps every consumer when events are
+// appended (a notification, or a local append) and every poll interval
+// regardless, since a client without a listener cannot be told.
+func (c *Client[TTx]) streamLoop(ctx context.Context) {
+ ticker := time.NewTicker(c.cfg.PollInterval)
+ defer ticker.Stop()
+ for {
+ c.pumpStreams(ctx)
+ select {
+ case <-ctx.Done():
+ return
+ case <-ticker.C:
+ case <-c.streamPoke:
+ }
+ }
+}
+
+// pumpStreams delivers the pending events of every consumer.
+func (c *Client[TTx]) pumpStreams(ctx context.Context) {
+ consumers, err := c.exec.StreamConsumerList(ctx)
+ if err != nil {
+ if ctx.Err() == nil {
+ c.logger.WarnContext(ctx, "hopper: list stream consumers", "error", err)
+ }
+ return
+ }
+ limit := max(c.tuning.streamBatch, 1)
+ for _, consumer := range consumers {
+ for {
+ res, err := c.exec.StreamPump(ctx, consumer.Name, limit)
+ if err != nil {
+ if ctx.Err() == nil && !errors.Is(err, ErrNotFound) {
+ c.logger.ErrorContext(ctx, "hopper: pump stream consumer", "consumer", consumer.Name, "error", err)
+ }
+ break
+ }
+ for _, q := range res.Queues {
+ c.wakeQueue(q)
+ }
+ if res.Delivered < limit {
+ break
+ }
+ }
+ }
+}
diff --git a/stream_test.go b/stream_test.go
new file mode 100644
index 0000000..0e30f2d
--- /dev/null
+++ b/stream_test.go
@@ -0,0 +1,114 @@
+package hopper_test
+
+import (
+ "context"
+ "sync"
+ "testing"
+
+ "github.com/parallelworks/hopper"
+)
+
+type orderEvent struct {
+ ID int `json:"id"`
+ At string `json:"at"`
+}
+
+func (orderEvent) Topic() string { return "order.placed" }
+
+type auditEvent struct {
+ What string `json:"what"`
+}
+
+func (auditEvent) Topic() string { return "audit.note" }
+
+func TestStreams(t *testing.T) {
+ t.Parallel()
+ h := newHarness(t)
+ ctx := context.Background()
+ var (
+ mu sync.Mutex
+ received []string // the orders consumer's deliveries, in order
+ streams []string
+ audited int
+ )
+ workers := hopper.NewWorkers()
+ hopper.Consume(workers, hopper.Consumer{Name: "orders", Pattern: "order.*", Start: hopper.StreamStartEarliest},
+ func(_ context.Context, msg *hopper.Message[orderEvent]) error {
+ mu.Lock()
+ defer mu.Unlock()
+ received = append(received, msg.Payload.At)
+ streams = append(streams, msg.Stream+"@"+msg.Position.String()+"/"+msg.OrderingKey)
+ return nil
+ })
+ hopper.Consume(workers, hopper.Consumer{Name: "audit", Pattern: "#"},
+ func(_ context.Context, msg *hopper.Message[hopper.Raw]) error {
+ mu.Lock()
+ defer mu.Unlock()
+ audited++
+ return nil
+ })
+ c := h.started(workers, 4)
+ s := c.Streams()
+
+ // Every event reaches both consumers, in order, through the pool and
+ // through a transaction alike.
+ before, err := s.Append(ctx, orderEvent{ID: 1, At: "before"}, &hopper.AppendOpts{Key: "o1"})
+ if err != nil {
+ t.Fatal(err)
+ }
+ tx, err := h.pool.Begin(ctx)
+ if err != nil {
+ t.Fatal(err)
+ }
+ if _, err := s.AppendTx(ctx, tx, orderEvent{ID: 2, At: "in-tx"}, &hopper.AppendOpts{Key: "o1"}); err != nil {
+ t.Fatal(err)
+ }
+ if err := tx.Commit(ctx); err != nil {
+ t.Fatal(err)
+ }
+ if _, err := s.Append(ctx, auditEvent{What: "x"}, nil); err != nil {
+ t.Fatal(err)
+ }
+ waitFor(t, func() bool { mu.Lock(); defer mu.Unlock(); return len(received) == 2 && audited == 3 })
+ mu.Lock()
+ if received[0] != "before" || received[1] != "in-tx" {
+ t.Errorf("received = %v", received)
+ }
+ if streams[0] != "orders@"+before.Position.String()+"/o1" {
+ t.Errorf("stream metadata = %v", streams)
+ }
+ received = nil
+ mu.Unlock()
+
+ // Read pages the log; consumers report their positions; a seek
+ // replays.
+ var topics []string
+ for ev, err := range s.Read(ctx, "", hopper.StreamPosition{}, 2) {
+ if err != nil {
+ t.Fatal(err)
+ }
+ topics = append(topics, ev.Topic)
+ }
+ if len(topics) != 3 || topics[0] != "order.placed" || topics[2] != "audit.note" {
+ t.Errorf("read = %v", topics)
+ }
+ consumers, err := s.Consumers(ctx)
+ if err != nil || len(consumers) != 2 || consumers[1].Name != "orders" || consumers[1].DeliveredAt.IsZero() {
+ t.Errorf("consumers = %v, %v", consumers, err)
+ }
+ if err := s.Seek(ctx, "orders", hopper.SeekOpts{Earliest: true}); err != nil {
+ t.Fatal(err)
+ }
+ waitFor(t, func() bool { mu.Lock(); defer mu.Unlock(); return len(received) == 2 })
+ mu.Lock()
+ if received[0] != "before" || received[1] != "in-tx" {
+ t.Errorf("replayed = %v", received)
+ }
+ mu.Unlock()
+ if err := s.Seek(ctx, "nobody", hopper.SeekOpts{Latest: true}); err == nil {
+ t.Error("seek of an unknown consumer succeeded")
+ }
+ if _, err := s.Append(ctx, nil, nil); err == nil {
+ t.Error("nil event accepted")
+ }
+}
diff --git a/worker.go b/worker.go
index 9abdd9b..5763d24 100644
--- a/worker.go
+++ b/worker.go
@@ -40,6 +40,7 @@ type Workers struct {
mu sync.RWMutex
kinds map[string]*workerInfo
subscriptions []Subscription
+ consumers []Consumer
}
// NewWorkers returns an empty registry.
From 2e82dc5b9badcb178e1d61cbfe0cfd00ca93660c Mon Sep 17 00:00:00 2001
From: Michael McQuade
Date: Sat, 26 Sep 2026 09:53:43 -0500
Subject: [PATCH 2/2] feat(ui): hopperui web UI module (M8) (#11)
---
.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", `
{{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
+| attempt | at | error |
+{{range .Errors}}| {{.Attempt}} | {{ts .At}} | {{.Error}}{{if .Trace}}
+
+{{.Trace}}{{end}} |
{{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}}
+
+{{end}}
+
+| id | kind | queue | state | attempt | created | finalized |
+
+{{range .Data.Workflow.Jobs}}
+
+ {{.ID}} |
+ {{.Kind}} | {{.Queue}} |
+ {{.State}} |
+ {{.Attempt}}/{{.MaxAttempts}} |
+ {{since .CreatedAt}} |
+ {{since .FinalizedAt}} |
+
+{{end}}
+
+
+{{end}}
+
+{{define "subscriptions"}}
+Subscriptions
+
+| name | pattern | queue | kind | max attempts | created |
+
+{{range .Data}}
+| {{.Name}} | {{.Pattern}} | {{.Queue}} | {{.Kind}} | {{if .MaxAttempts}}{{.MaxAttempts}}{{else}}default{{end}} | {{ts .CreatedAt}} |
+{{else}}| no subscriptions |
{{end}}
+
+
+{{end}}
+
+{{define "consumers"}}
+Stream consumers
+
+| name | pattern | queue | snapshot | position | last delivery | {{if .CanAct}}seek | {{end}}
+
+{{range .Data}}
+
+ | {{.Name}} | {{.Pattern}} |
+ {{.Queue}} |
+ {{.Snapshot}} | {{.Position}} |
+ {{if .DeliveredAt.IsZero}}never{{else}}{{since .DeliveredAt}}{{end}} |
+ {{if $.CanAct}}
+
+
+
+ | {{end}}
+
+{{else}}| no consumers |
{{end}}
+
+
+{{end}}
+
+{{define "clients"}}
+Clients
+
+| id | host | started | lease | info |
+
+{{range .Data}}
+| {{.ID}} | {{.Hostname}} | {{since .StartedAt}} | {{if .ExpiresAt.After now}}live{{else}}expired{{end}} | {{printf "%s" .Info}} |
+{{else}}| no clients |
{{end}}
+
+
+{{end}}