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
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,11 @@ versions may change the API.
- `hoppersql`, a driver for `database/sql` (pgx's stdlib adapter or lib/pq)
that passes the same conformance suite as `hopperpgx`. It polls instead of
listening and inserts without COPY.
- 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`. Schema v4 adds `hopper_job_deps`.

### Changed

Expand Down
35 changes: 35 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"
"strings"
"text/tabwriter"
"time"

Expand Down Expand Up @@ -57,6 +58,7 @@ commands:
queues pause|resume <name>
queues limit <name> [-global N] [-rate R] [-burst B] [-partition P] [-aging D] (omitted limits are removed)
clients list
workflows get <id>
stats
bench (see the hopperbench command)

Expand Down Expand Up @@ -111,6 +113,8 @@ func run(ctx context.Context, args []string, out io.Writer) error {
return c.queues(ctx, rest[1:])
case "clients":
return c.clients(ctx, rest[1:])
case "workflows":
return c.workflows(ctx, rest[1:])
case "stats":
return c.stats(ctx)
default:
Expand Down Expand Up @@ -339,6 +343,37 @@ func (c *cli) clients(ctx context.Context, args []string) error {
})
}

func (c *cli) workflows(ctx context.Context, args []string) error {
if len(args) != 2 || args[0] != "get" {
return errors.New("workflows: get <id>")
}
id, err := hopper.ParseJobID(args[1])
if err != nil {
return err
}
wf, err := c.client.WorkflowGet(ctx, id)
if err != nil {
return err
}
return c.print(wf, func(w io.Writer) {
b := wf.Batch
fmt.Fprintf(w, "id: %s\nname: %s\nprogress: %d of %d finished, %d failed\n", b.ID, b.Name, b.Total-b.Pending, b.Total, b.Failed)
if !b.CompletedAt.IsZero() {
fmt.Fprintf(w, "completed: %s\n", b.CompletedAt.Local().Format(time.RFC3339))
}
deps := map[hopper.JobID][]string{}
for _, e := range wf.Edges {
deps[e.Job] = append(deps[e.Job], e.DependsOn.String())
}
tw := tabwriter.NewWriter(w, 0, 0, 2, ' ', 0)
fmt.Fprintln(tw, "ID\tKIND\tQUEUE\tSTATE\tATTEMPT\tAFTER")
for _, j := range wf.Jobs {
fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%d/%d\t%s\n", j.ID, j.Kind, j.Queue, j.State, j.Attempt, j.MaxAttempts, strings.Join(deps[j.ID], ","))
}
tw.Flush()
})
}

func (c *cli) stats(ctx context.Context) error {
stats, err := c.client.Stats(ctx)
if err != nil {
Expand Down
11 changes: 11 additions & 0 deletions cmd/hopper/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,17 @@ func TestCLI(t *testing.T) {
if got := hopperCmd("stats"); !strings.Contains(got, "q1") {
t.Errorf("stats = %q", got)
}
wf := hopper.NewWorkflow("ingest", nil)
first := wf.Add(ping{}, nil)
wf.Add(ping{}, hopper.After(first))
wres, err := client.InsertWorkflow(ctx, wf)
if err != nil {
t.Fatal(err)
}
if got := hopperCmd("workflows", "get", wres.ID.String()); !strings.Contains(got, "name: ingest") ||
!strings.Contains(got, "0 of 2 finished") || !strings.Contains(got, "pending") || !strings.Contains(got, wres.Jobs[0].Job.ID.String()) {
t.Errorf("workflows get = %q", got)
}
if got := hopperCmd("migrate", "down"); !strings.Contains(got, "schema version 0") {
t.Errorf("migrate down = %q", got)
}
Expand Down
38 changes: 26 additions & 12 deletions docs/PLAN.md
Original file line number Diff line number Diff line change
Expand Up @@ -231,12 +231,13 @@ b.Add(ProcessShard{Shard: 0}, nil)
b.Add(ProcessShard{Shard: 1}, nil)
err = b.InsertTx(ctx, tx)

wf := hopper.NewWorkflow("ingest-9")
fetch := wf.Add(Fetch{URL: u})
wf := hopper.NewWorkflow("ingest-9", &hopper.WorkflowOpts{OnFailure: AlertOps{RunID: 9}})
fetch := wf.Add(Fetch{URL: u}, nil)
parse := wf.Add(Parse{}, hopper.After(fetch))
wf.Add(Index{}, hopper.After(parse))
wf.Add(Notify{}, hopper.After(parse))
err = client.InsertWorkflowTx(ctx, tx, wf)
res, err := client.InsertWorkflowTx(ctx, tx, wf)
row, err := client.WorkflowGet(ctx, res.ID) // progress, jobs and edges
```

### 4.7 API conventions
Expand Down Expand Up @@ -444,7 +445,8 @@ CREATE TABLE hopper_schema (version int PRIMARY KEY, applied_at timestamptz NOT
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. M8 adds `hopper_job_deps` (§11).
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).

### 6.1 Job IDs
Job IDs are UUIDv7. They are safe to expose outside the application (in URLs, APIs
Expand Down Expand Up @@ -919,11 +921,23 @@ observability without new machinery.
failed, `on_failure` otherwise, `on_complete` either way, each with `batch_id` and
`batch_failed` in its metadata, and notifies their queues. `client.NewBatch(opts)`,
`Add`, `Insert`/`InsertTx` and `BatchGet` are the API.
- **Workflows (M8).** Jobs with dependencies are inserted as `pending`, with edges in
`hopper_job_deps`. When a job completes, the same finalize statement promotes
dependents whose dependencies are all complete to `available`. Failure policies are
`cancel dependents` (the default) and `ignore`. Workflows are inspectable as a DAG in
`hopperui`.
- **Workflows (M8).** A workflow is a batch whose jobs depend on each other, so it has
a name, the batch's progress and callbacks, and its graph (`hopper_batches.edges`).
`NewWorkflow` collects steps; a step can only depend on steps added before it, so the
graph is a DAG by construction. `InsertWorkflow` generates the IDs up front and
inserts the batch, the jobs (`pending` when they have dependencies, otherwise
available) and the edges in one transaction. `hopper_job_deps` is the working set of
unsatisfied edges, deleted as dependents finalize, and the statement that finalizes
a job is the one that acts on its dependents: it promotes those whose dependencies
have all left the live table (including the ones finalized in the same statement)
to `available` and notifies their queues, and it cancels, transitively, the
dependents of a cancelled or discarded job whose edge says `cancel` (the default;
`ignore` lets the dependent run once its dependencies have finished, whatever their
outcome). Cancelled dependents go to history with a "dependency failed" error and
count as batch failures. Every finalizing statement (the batched finalize,
cancelling a waiting job, expiring) ends with the same tail, so batches and
dependents behave the same however a job finishes. Workflows are inspectable with
`WorkflowGet`, `hopper workflows get` and, as a DAG, in `hopperui`.

## 12. Migrations

Expand Down Expand Up @@ -960,8 +974,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`, `clients list` and `stats`, all with `-json`. `queues limit`
and `subscriptions list` arrive with M7 and M6. Benchmarks are the separate
`queues list|pause|resume|limit`, `clients list`, `workflows get` and `stats`, all with
`-json`. `subscriptions list` arrives with M6. 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
Expand Down Expand Up @@ -1006,7 +1020,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, streams with consumer groups, `hopperui` | 2–3w |
| M8 | Workflows, streams, UI | Job dependencies and DAG workflows (**done**), streams with consumer groups, `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
Expand Down
23 changes: 23 additions & 0 deletions docs/getting-started.md
Original file line number Diff line number Diff line change
Expand Up @@ -165,6 +165,29 @@ Periodic: []hopper.PeriodicJob{

The leader inserts each slot exactly once, whichever replica is leader.

## Batches and workflows

```go
b := client.NewBatch(hopper.BatchOpts{OnSuccess: ReportDone{RunID: 9}, OnFailure: AlertOps{RunID: 9}})
b.Add(ProcessShard{Shard: 0}, nil)
b.Add(ProcessShard{Shard: 1}, nil)
res, err := b.InsertTx(ctx, tx)

wf := hopper.NewWorkflow("ingest-9", &hopper.WorkflowOpts{OnFailure: AlertOps{RunID: 9}})
fetch := wf.Add(Fetch{URL: u}, nil)
parse := wf.Add(Parse{}, hopper.After(fetch))
wf.Add(Index{}, hopper.After(parse))
wf.Add(Notify{}, hopper.After(parse))
res, err := client.InsertWorkflowTx(ctx, tx, wf)
```

A batch runs its callbacks when its last job finishes. A workflow's steps
run as their dependencies finish; a failed step cancels the steps that
depend on it (set `StepOpts.OnDependencyFailure` to `DependencyIgnore` to run
anyway), and the workflow is a batch, so it has the same callbacks.
`client.WorkflowGet(ctx, res.ID)` and `hopper workflows get <id>` show its
progress and graph.

## Results and waiting

```go
Expand Down
65 changes: 64 additions & 1 deletion driver/driver.go
Original file line number Diff line number Diff line change
Expand Up @@ -189,6 +189,13 @@ type Executor interface {
BatchInsert(ctx context.Context, params BatchInsertParams) (JobID, error)
// BatchGet returns a batch's progress.
BatchGet(ctx context.Context, id JobID) (*BatchRow, error)
// WorkflowInsert inserts a batch, its jobs and the dependencies between
// them atomically. Jobs with dependencies start pending. Results are in
// input order.
WorkflowInsert(ctx context.Context, params WorkflowInsertParams) (*WorkflowInsertResult, error)
// WorkflowGet returns a workflow's progress, its jobs (live and
// finished) and its edges.
WorkflowGet(ctx context.Context, id JobID) (*WorkflowRow, error)

// PeriodicInsert inserts the job for a periodic slot if that slot has
// not been inserted yet, atomically, and reports whether it did. A slot
Expand Down Expand Up @@ -419,7 +426,9 @@ type BatchInsertParams struct {

// BatchRow is a batch's progress.
type BatchRow struct {
ID JobID
ID JobID
// Name is set for workflows.
Name string
Pending int
Failed int
Total int
Expand All @@ -429,6 +438,60 @@ type BatchRow struct {
Metadata json.RawMessage
}

// DependencyFailure says what happens to a dependent job when the job it
// depends on is cancelled or discarded.
type DependencyFailure string

const (
// DependencyCancel cancels the dependent, and so on down the graph.
DependencyCancel DependencyFailure = "cancel"
// DependencyIgnore lets the dependent run once its dependencies have
// all finished, whatever their outcome.
DependencyIgnore DependencyFailure = "ignore"
)

// JobDependency is an edge of a workflow: Job waits for DependsOn. Both
// are indexes into WorkflowInsertParams.Jobs.
type JobDependency struct {
Job, DependsOn int
OnFailure DependencyFailure
}

// WorkflowInsertParams describes a workflow: a named batch whose jobs
// depend on each other.
type WorkflowInsertParams struct {
Name string
// Jobs must have no unique keys; a skipped job would leave its
// dependents waiting forever.
Jobs []JobInsertParams
Deps []JobDependency
// Callbacks and Metadata are the batch's; see BatchInsertParams.
OnSuccess, OnFailure, OnComplete *JobInsertParams
Metadata json.RawMessage
// Notify sends insert notifications for the jobs that start
// available from inside the statement.
Notify bool
}

// WorkflowInsertResult reports an inserted workflow.
type WorkflowInsertResult struct {
ID JobID
Jobs []JobInsertResult
}

// JobEdge is an edge of an inserted workflow.
type JobEdge struct {
Job, DependsOn JobID
}

// WorkflowRow is a workflow's progress, jobs and graph.
type WorkflowRow struct {
Batch BatchRow
// Jobs are in insertion order, live and finished alike.
Jobs []*JobRow
Edges []JobEdge
}

// PeriodicInsertParams identifies a periodic slot and the job to insert
// for it.
type PeriodicInsertParams struct {
Expand Down
27 changes: 0 additions & 27 deletions driver/internal/pgsql/flow.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,33 +12,6 @@ import (
"github.com/parallelworks/hopper/driver"
)

// batchAccountingSQLWith returns the CTEs that count finalized jobs off
// their batches and insert the callbacks of batches that reached zero, all
// inside the finalizing statement. They expect a CTE named done whose rows
// carry batch_id and a text column final_state, and the insert channel as
// the given parameter. Rows without a batch cost one filtered scan of the
// CTE and nothing else.
func batchAccountingSQLWith(channelParam string) string {
return `batches AS (
UPDATE hopper_batches b
SET pending = b.pending - c.n, failed = b.failed + c.f,
completed_at = CASE WHEN b.pending - c.n <= 0 THEN now() ELSE b.completed_at END
FROM (SELECT batch_id, count(*) AS n, count(*) FILTER (WHERE final_state IN ('cancelled', 'discarded')) AS f
FROM done WHERE batch_id IS NOT NULL GROUP BY batch_id) c
WHERE b.id = c.batch_id
RETURNING b.id, b.pending, b.failed, b.on_success, b.on_failure, b.on_complete
),
callbacks AS (
INSERT INTO hopper_jobs (kind, queue, priority, max_attempts, args, metadata)
SELECT cb->>'kind', coalesce(cb->>'queue', 'default'), coalesce((cb->>'priority')::smallint, 2),
coalesce((cb->>'max_attempts')::smallint, 25), coalesce(cb->'args', '{}'),
coalesce(cb->'metadata', '{}') || jsonb_build_object('batch_id', b.id::text, 'batch_failed', b.failed)
FROM batches b, LATERAL (VALUES (CASE WHEN b.failed = 0 THEN b.on_success ELSE b.on_failure END), (b.on_complete)) AS v(cb)
WHERE b.pending <= 0 AND cb IS NOT NULL
RETURNING pg_notify(` + channelParam + `, queue)
)`
}

// claimLimited claims inside a transaction that locks the queue row, so that
// every client's claims on a limited queue are serialized and the limits
// hold exactly: the budget is the smallest of the free slots, what the
Expand Down
Loading
Loading