Skip to content

hopper v0.1.0

Choose a tag to compare

@giraffesyo giraffesyo released this 29 Sep 16:26
· 29 commits to canary since this release
fc984e0

The first release: a job queue and message broker on PostgreSQL, with the
hopperotel and hopperui modules released alongside it at the same
version. The §8.2 performance targets in docs/PLAN.md were
run on the reference hardware before tagging; the results are recorded
there.

Jobs

  • Core engine: typed workers, Insert/InsertTx/InsertMany (one statement,
    or COPY for large batches), batched FOR UPDATE SKIP LOCKED claims,
    batched finalize into time-partitioned history, retries with backoff,
    Snooze and Cancel, timeouts, panic recovery, and a three-step graceful
    Run/Stop.
  • Reliability: per-client leases, rescue of jobs from crashed processes,
    fencing of paused processes, leader election, partition-drop retention,
    LISTEN/NOTIFY wake-ups with per-process coalescing and polling fallback,
    unique jobs with skip and replace.
  • Control: cron and interval periodic jobs with time zones, in-flight
    cancellation, retry from the dead-letter queue, TTLs, queue pause and
    resume, runtime queues, middleware, SetOutput/Await, the Jobs
    iterator (paged with JobFilter.After), events and Stats.
  • Flow control: cluster-wide GlobalLimit, RateLimit/RateBurst and
    PartitionLimit per queue (declared in QueueConfig or set at runtime
    with Queues().SetLimits and hopper queues limit), PriorityAging,
    and InsertOpts.PartitionKey.
  • Batches: NewBatch, Add, Insert/InsertTx, BatchGet, with
    OnSuccess, OnFailure and OnComplete callbacks inserted by the
    finalizing statement.
  • Workflows: NewWorkflow, Add with After, InsertWorkflow/
    InsertWorkflowTx, WorkflowGet and hopper workflows get. Steps with
    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.

Messaging

  • Subscriptions: Subscribe with AMQP topic patterns, typed Message[T],
    Publish/PublishTx fan-out in one statement, dedup keys, ordering keys
    (also on plain jobs through InsertOpts.OrderingKey), request/reply, and
    ReplayDiscarded.
  • 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.
  • The SQL contract functions hopper_insert and hopper_publish, for
    producers in other languages (docs/sql-contract.md).

Drivers and schema

  • hopperpgx: the pgx v5 driver, with COPY, LISTEN and pipelined statements.
  • hoppersql: a driver for database/sql (pgx's stdlib adapter or lib/pq)
    that passes the same conformance suite; it polls instead of listening and
    inserts without COPY. Both drivers share one implementation of the SQL.
  • hoppermigrate: embedded, versioned migrations under a cross-process lock.
    Schema versions 1 through 6: the job tables and history partitions,
    messaging, flow control and batches, workflows, streams, and the live
    table's autovacuum settings.
  • drivertest: the conformance, concurrency and chaos suite every driver
    must pass, including TestUpgradeUnderTraffic, which migrates to the
    latest schema while a client works jobs.

Operations

  • The leader maintains the live table (driver.Executor.JobsMaintain):
    it vacuums hopper_jobs once a hundred thousand dead rows have
    accumulated and re-analyzes it when the planner's row count is an order
    of magnitude off, every leader interval, without waiting for autovacuum's
    lock. Every claim walks the claim index past the
    entries of finished jobs until a vacuum removes them, and autovacuum
    looks only every minute by default: on the reference hardware pickup
    latency at 30,000 jobs/s climbed from 10 ms to seconds within a minute.
  • The claim's candidate subquery is a MATERIALIZED CTE. A plan cached
    while hopper_jobs was empty otherwise re-executed the locking subquery
    once per row of the table after a burst: a single claim ran for 20-30
    seconds holding locks, with every finalizer queued behind it.
  • cmd/hopper: migrations, jobs, queues, clients, workflows, subscriptions,
    stream consumers and stats from the shell.
  • cmd/hopperbench: the benchmark harness (including loaded, pickup
    latency under a paced insert load) and the CI performance gate.
  • hoppertest: test helpers for code that uses hopper.
  • hopperotel (separate module): OpenTelemetry tracing and metrics.
  • hopperui (separate module): an embeddable web UI for queues, jobs,
    workflows (as a DAG), subscriptions, stream consumers and clients, with
    actions behind an Authorize hook, and a standalone hopperui server.

Install

go get github.com/parallelworks/hopper@v0.1.0
go get github.com/parallelworks/hopper/hopperotel@v0.1.0   # OpenTelemetry
go get github.com/parallelworks/hopper/hopperui@v0.1.0     # web UI

Requires Go 1.27 and PostgreSQL 14 or later. Start with docs/getting-started.md.