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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/dependabot.yml
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
version: 2
updates:
- package-ecosystem: gomod
directories: ["/", "/tools"]
directories: ["/", "/tools", "/hopperotel", "/hopperui"]
schedule: { interval: weekly }
commit-message: { prefix: "build(deps)" }
groups:
Expand Down
3 changes: 3 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,12 @@ jobs:
go mod tidy -diff
(cd tools && go mod tidy -diff)
(cd hopperotel && go mod tidy -diff)
(cd hopperui && go mod tidy -diff)
- name: Lint
run: |
go tool -modfile=tools/go.mod golangci-lint run ./...
(cd hopperotel && go tool -modfile=../tools/go.mod golangci-lint run ./...)
(cd hopperui && go tool -modfile=../tools/go.mod golangci-lint run ./...)

test:
runs-on: ubuntu-latest
Expand Down Expand Up @@ -61,3 +63,4 @@ jobs:
run: |
go test -race -cover -timeout 10m ./...
(cd hopperotel && go test -race -cover -timeout 10m ./...)
(cd hopperui && go test -race -cover -timeout 10m ./...)
12 changes: 12 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,18 @@ 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`.
- `hopperui`, a separate module: an embeddable web UI for queues, jobs,
workflows (as a DAG), subscriptions, stream consumers and clients, with
actions behind an `Authorize` hook.
- `JobFilter.After`, to page job listings.

### Changed

Expand Down
2 changes: 2 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ check: lint test ## Run all linters and tests
lint: ## Lint Go code
$(GOTOOL) golangci-lint run ./...
cd hopperotel && go tool -modfile=../tools/go.mod golangci-lint run ./...
cd hopperui && go tool -modfile=../tools/go.mod golangci-lint run ./...

.PHONY: fmt
fmt: ## Format Go code
Expand All @@ -48,6 +49,7 @@ fmt: ## Format Go code
test: ## Run tests (integration tests need `make pg` or HOPPER_TEST_DATABASE_URL)
go test -race -cover -timeout 10m ./...
cd hopperotel && go test -race -cover -timeout 10m ./...
cd hopperui && go test -race -cover -timeout 10m ./...

.PHONY: bench
bench: ## Run the benchmark harness against the test database
Expand Down
12 changes: 7 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,11 +30,13 @@ Postgres is the first engine, behind a driver interface designed for more.
> LISTEN/NOTIFY wake-ups, unique jobs, cron and interval schedules, cancellation,
> retry from the dead-letter queue, TTLs, pause and runtime queues, middleware,
> awaitable results, listing, events and stats, plus the `hopper` CLI,
> `hoppertest`, `hopperotel` and the `drivertest` conformance suite. So is
> messaging: subscriptions with topic patterns, fan-out, dedup and ordering
> keys, and a SQL contract for producers in other languages. Flow control,
> batches and workflows follow, in the order laid out in
> [docs/PLAN.md](docs/PLAN.md), the plan of record. Start with
> `hoppertest`, `hopperotel`, the `hopperui` web UI and the `drivertest`
> conformance suite. So are messaging (subscriptions with topic patterns,
> fan-out, dedup and ordering keys, a SQL contract for producers in other
> languages, streams with consumers), flow control (global, rate and
> partitioned limits, priority aging), batches and workflows, and a
> `database/sql` driver, in the order laid out in [docs/PLAN.md](docs/PLAN.md),
> the plan of record. Start with
> [docs/getting-started.md](docs/getting-started.md). Feedback is welcome
> through issues and PRs.

Expand Down
9 changes: 7 additions & 2 deletions client.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ type Client[TTx any] struct {
listening atomic.Bool
isLeader atomic.Bool
leaderPoke chan struct{}
streamPoke chan struct{}

events *eventBus

Expand Down Expand Up @@ -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{}{},
Expand Down Expand Up @@ -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)
Expand All @@ -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 {
Expand Down
94 changes: 94 additions & 0 deletions cmd/hopper/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import (
"os"
"os/signal"
"slices"
"strconv"
"strings"
"text/tabwriter"
"time"
Expand Down Expand Up @@ -59,6 +60,9 @@ commands:
queues limit <name> [-global N] [-rate R] [-burst B] [-partition P] [-aging D] (omitted limits are removed)
clients list
workflows get <id>
subscriptions list
streams consumers
streams seek <consumer> -earliest|-latest|-time RFC3339|-position xid:seq
stats
bench (see the hopperbench command)

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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 {
Expand Down
9 changes: 9 additions & 0 deletions cmd/hopper/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
11 changes: 11 additions & 0 deletions config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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 == "" {
Expand Down Expand Up @@ -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{
Expand All @@ -231,4 +241,5 @@ var defaultTuning = tuning{
leaderInterval: 5 * time.Second,
maintenanceInterval: 5 * time.Minute,
rescueBatch: 1000,
streamBatch: 500,
}
5 changes: 4 additions & 1 deletion control.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,9 @@ type JobFilter struct {
Queue string
Kinds []string
States []JobState
// After starts the listing after the given job, for paging: pass the
// last ID of one page to get the next.
After JobID
}

// jobsPageSize is how many jobs Jobs fetches at a time.
Expand All @@ -78,7 +81,7 @@ func (c *Client[TTx]) JobsTx(ctx context.Context, tx TTx, filter JobFilter) iter

func (c *Client[TTx]) jobs(ctx context.Context, exec driver.Executor, filter JobFilter) iter.Seq2[*JobRow, error] {
return func(yield func(*JobRow, error) bool) {
var after JobID
after := filter.After
for {
page, err := exec.JobList(ctx, driver.JobListParams{
Queue: filter.Queue, Kinds: filter.Kinds, States: filter.States, After: after, Limit: jobsPageSize,
Expand Down
Loading
Loading