A distributed task & workflow orchestration service in Go — built deliberately, in public, as a learning project.
This is a learning project, and the README is written to be honest about that. I am building one coherent distributed system instead of a pile of tutorials, specifically to close the gaps senior Go backend interviews probe: advanced concurrency, distributed coordination, Postgres depth, and production observability.
Phase 0 (Foundations) is complete and runs. Phases 1–4 are designed in detail but not yet built. Everything below is split accordingly: What works today links to real code you can read; The road ahead describes what is designed and explicitly does not link to code, because there isn't any yet.
If you are evaluating me: the fastest signal is probably Patterns & techniques → where they live, which maps each concept to the exact file and line that implements it.
- What is flowhand
- What works today
- Architecture
- Patterns & techniques → where they live
- Tech stack, and why each choice
- Engineering practices
- Run it in five minutes
- Repository map
- The road ahead
- What I have learned so far
A single Go service that fuses a task queue (priorities, retries, cron, idempotency) with a declarative workflow DAG engine (fan-out, fan-in, conditionals). Think "Celery + Airflow in one binary."
The domain is framed as the backend of a Stripe-style webhook and notification delivery platform: external services POST events, flowhand fans each event out to its subscriber set, retries on failure, and guarantees each subscriber receives each event effectively once.
That framing is not decoration — it is what makes each hard requirement necessary rather than bolted on:
| The domain demands | So the system needs |
|---|---|
| Producers retry webhooks | Idempotency keys with a durable floor |
| We promise delivery | Transactional outbox — no lost events on crash |
| One event → N subscribers → per-channel transforms | DAG workflows |
| Duplicate deliveries break customers | Effectively-once end to end |
| Downstream endpoints are fragile | Rate limiting, circuit breaking |
| One bad endpoint must not block the queue | Dead-letter queue |
Read task as "one delivery attempt to one subscriber", and workflow as "one event's fan-out to its subscriber set".
Phase 0 is a complete, runnable vertical slice: an HTTP request becomes a row in Postgres, and you can see the whole thing as a distributed trace in Grafana.
task up && task migrate:up && task run:server # terminal A
curl -X POST localhost:8080/v1/tasks \
-H 'content-type: application/json' \
-d '{"handler":"echo","payload":{"msg":"hello"}}'
# → 201 {"id":"019f7013-…","status":"pending","created_at":"…"}| Working today | Where |
|---|---|
POST /v1/tasks — OpenAPI-first, generated server, validated input |
handler.go#L46 |
GET /v1/tasks/{id} — 200 / 404 / 500, all typed |
handler.go#L81 |
| Idempotent submit — a repeated key replays the original task, never a duplicate | commander.go#L54 |
| Layered architecture, machine-enforced | .go-arch-lint.yml |
| Distributed tracing HTTP → handler → SQL | tracing.go#L15 |
Logs carrying trace_id, correlated to traces in Grafana |
log.go#L17 |
| Prometheus metrics with SLO-tuned histogram buckets | metrics.go#L32 |
| Graceful shutdown that actually flushes telemetry | shutdown.go#L21 |
| 12-service dev stack, one command | docker-compose.yml |
| Reference load producer that verifies replay on the wire | flowhand-demo |
Full walkthrough with expected output at every step:
docs/getting-started.md.
The core rule: dependencies point inward, and the domain depends on nothing.
flowchart TD
subgraph root["composition root"]
CLI["internal/cli<br/><i>wires everything, once</i>"]
end
subgraph transport["transport"]
SRV["internal/server<br/><i>http.Server, mux, pprof</i>"]
API["internal/api<br/><i>ogen Handler impl</i>"]
OAS["internal/api/oas<br/><i>generated — never hand-edited</i>"]
end
subgraph app["application"]
SVC["internal/service/task<br/><i>Commander · Querier · ports</i>"]
end
subgraph core["domain — the leaf"]
DOM["internal/domain/tasks<br/><i>Task aggregate, invariants, events</i>"]
end
subgraph infra["infrastructure"]
REPO["internal/repository/tasks<br/><i>row ↔ aggregate mapping</i>"]
TXM["internal/storage/txmgr<br/><i>transaction scope</i>"]
QRY["internal/storage/queries<br/><i>generated by sqlc</i>"]
OBS["internal/obs<br/><i>slog · OTel</i>"]
end
CLI --> SRV & API & SVC & REPO & TXM & OBS
SRV --> API --> OAS
API --> SVC
SVC --> DOM
SVC -.->|"ports it declares"| REPO
REPO --> DOM
REPO --> QRY
TXM --> QRY
classDef leaf fill:#1f6f43,stroke:#0d3d24,color:#fff
classDef gen fill:#4a4a4a,stroke:#2a2a2a,color:#fff
class DOM leaf
class OAS,QRY gen
internal/service/task never names pgx or sqlc — it declares the interfaces it
needs in ports.go and the persistence layer
satisfies them. That inversion is what makes the write path unit-testable with no
database at all, and it is checked by two tools on every commit:
.go-arch-lint.yml and depguard in
.golangci.yml. Rules and their traps:
docs/design/architecture/import-rules.md.
sequenceDiagram
autonumber
participant C as Client
participant OT as otelhttp<br/>(middleware)
participant OG as ogen<br/>(generated router)
participant H as api.Handler
participant CMD as task.Commander
participant TX as txmgr
participant R as repository/tasks
participant PG as Postgres
participant T as Tempo / Loki / Prometheus
C->>OT: POST /v1/tasks {handler, payload, idempotency_key}
OT->>OT: start root span, inject trace context
OT->>OG: ServeHTTP
OG->>OG: decode + validate against openapi.yaml
Note over OG: malformed body never reaches the handler —<br/>it goes to api.ErrorHandler as a JSON envelope
OG->>H: CreateTask(ctx, req)
H->>CMD: Submit(ctx, SubmitCommand)
CMD->>CMD: mint UUIDv7, build tasks.Task aggregate
CMD->>TX: WithinTx(...)
TX->>R: Insert(ctx, task, key)
R->>PG: INSERT INTO tasks (sqlc + pgx)
alt idempotency key already used
PG-->>R: 23505 unique_violation
R-->>CMD: domain ErrConflict (never a SQLSTATE)
CMD->>R: GetByIdempotencyKey — outside the aborted tx
R-->>CMD: the original task
Note over CMD: replay: caller cannot tell this<br/>from the first submit
else first submit
PG-->>R: row
CMD->>TX: outbox.Append(TaskSubmitted)
end
CMD-->>H: tasks.Task
H->>H: toOASTask — explicit status mapping
H-->>C: 201 {id, status, created_at}
OT-->>T: span + metrics — logs carry the same trace_id
Every arrow in that diagram exists in code today. The pieces worth reading are linked in the next section.
This is the map. Each row is a concept I set out to learn, why it is there, and the exact code that implements it.
| Pattern | Why it is here | Implementation |
|---|---|---|
| Layered architecture | Dependencies point inward; the domain is a leaf. Enforced, not merely intended. | .go-arch-lint.yml · rules |
| DDD-lite: aggregate root | Task has unexported fields, so no code outside the package can build one in an invalid state. |
task.go#L11 |
| Bounded contexts | tasks / workflows / schedules are separate packages that may not import each other — a Phase 2 mistake becomes a lint failure. |
domain/ |
| CQRS-lite | Writes (Commander) and reads (Querier) split at the boundary, so Phase 1 can grow transactional writes without touching the read path. |
commander.go · querier.go |
| Ports & adapters | The service declares the interfaces it consumes; persistence implements them. Dependency inversion in ~30 lines. | ports.go#L14 |
| Three-types pattern | One "task" wears three shapes — storage row, domain aggregate, transport DTO — so schema changes cannot leak into the API. | row · domain · DTO |
| Adapter (GoF) | Mappers translate at each boundary, so neither type system leaks into the other. | repo mapper · API mapper |
| Decorator (GoF) | A slog.Handler wrapping a slog.Handler to stamp trace_id — same interface in and out, so it composes. |
log.go#L17 |
| Chain of responsibility | otelhttp → ogen → handler. Phase 1 grows it to auth → rate-limit → idempotency. |
server.go#L30 |
| Composition root | One place constructs the whole object graph, bottom-up. Nothing else wires anything. | cli/server.go#L55 |
| Functional options | Injected clock and ID generator, so time and UUIDs are deterministic in tests. | commander.go#L28 |
| Mechanism | The idea | Implementation |
|---|---|---|
| Idempotency, durable floor | A partial unique index is the source of truth. Redis will be a fast path in front of it, never a replacement — the database is what survives a cache flush. | migration · replay |
| Replay vs. 409 | A repeated key returns the original task, so a retrying client cannot tell a retry from the first call. | commander.go#L76 |
| Error translation at the boundary | 23505 becomes ErrConflict, pgx.ErrNoRows becomes ErrNotFound. No layer above persistence knows SQLSTATE exists. |
postgres.go#L67 |
| Transaction scope as a port | The service says "these writes are atomic" without naming pgx. Phase 1 swaps the stub for a real pgx.Tx in the context — no service code changes. |
txmgr |
| Transactional outbox | State change and its event committed together. The seam exists now; Phase 1 fills it in. | outbox |
| UUIDv7 over UUIDv4 | Time-ordered keys give B-tree locality on insert; v4 scatters writes across the index. | commander.go#L55 |
| Graceful shutdown | signal.NotifyContext → drain in-flight requests → then flush telemetry, on a fresh context because the original is already cancelled. |
server.go#L60 · shutdown.go#L21 |
| Signal | What I did | Implementation |
|---|---|---|
| Traces | OTel → OTLP/gRPC → Tempo, with DB spans via otelpgx, so one trace spans HTTP and SQL. |
tracing.go#L15 · pool.go#L14 |
| Logs ↔ traces | Every *Context log record carries trace_id/span_id, so Grafana jumps from a span to its log lines. |
log.go#L21 |
| Metrics | OTel → Prometheus bridge, with runtime metrics and a shared resource. | metrics.go#L17 |
| SLO-tuned histograms | Default buckets top out at 10s — useless against a p99 < 50ms SLO. Custom buckets, pinned by unit, because otelhttp switched ms → s and a unit-blind View would silently make every SLO look met. |
metrics.go#L32 |
| Continuous profiling | pprof exposed and scraped by Pyroscope, so yesterday's p99 flame graph exists without having predicted you'd want it. |
server.go#L37 |
| Cardinality discipline | No high-cardinality label ever crosses a histogram; buckets are themselves a cardinality multiplier. | metrics.go#L32 |
| Technique | Implementation |
|---|---|
Custom slog.Handler with correct WithAttrs/WithGroup re-wrapping (the commonly-botched part) |
log.go#L33 |
Errors as values: sentinels, %w wrapping, errors.Is/errors.As |
errors.go · postgres.go#L74 |
errors.Join for multi-resource shutdown |
shutdown.go#L57 |
| Interface segregation — the consumer declares the narrow interface it needs | handler.go#L23 |
| Layered config: defaults → file → env → flags, last wins | config.go#L51 |
ReadHeaderTimeout against slowloris (Go's default is unlimited) |
server.go#L43 |
Stdlib-only HTTP — Go 1.22+ ServeMux, no gin/echo/chi |
server.go#L29 |
Every dependency below was a deliberate decision, not a default.
| Area | Choice | Why this and not the obvious alternative |
|---|---|---|
| HTTP API | ogen | OpenAPI-first: the spec is the contract and the server is generated from it, so the two cannot drift. Not gin/echo/chi — stdlib ServeMux is enough since Go 1.22. |
| Database | pgx/v5 native | Skipping database/sql buys real batching, COPY FROM, LISTEN/NOTIFY and proper Postgres types — all of which Phase 1 needs. |
| Queries | sqlc | Write SQL, get type-safe Go. No ORM magic, no interface{}, no string interpolation. Not GORM/ent — I want SQL depth, which is half the point of the project. |
| Migrations | goose | Plain .sql files, usable as a CLI and as a library — which is how the integration tests migrate a throwaway database. |
| Config | koanf | Composable and test-friendly. Not viper — it couples to cobra and carries global state. |
| CLI | cobra | One binary, many subcommands, matching the process topology. |
| Logging | stdlib log/slog |
Deliberately not zap/zerolog/logrus. slog loses some benchmarks; it wins on being the ecosystem default and on Handler being trivially decoratable. |
| Telemetry | OpenTelemetry | Vendor-neutral. One API, swap the backend. |
| RPC | gRPC + buf | Scaffolded now, streaming worker dispatch in Phase 1. buf gives lint and breaking-change detection. |
| Testing | testify, testcontainers, native fuzzing, goleak |
Real Postgres in integration tests beats a mock that agrees with your misconceptions. |
| Task runner | Task | Every tool version pinned into ./bin, and CI runs the same targets — so "works locally" and "works in CI" cannot diverge. |
Infrastructure in the dev stack: Postgres 18, Redis 8, etcd 3.6,
Kafka 4 (KRaft), SeaweedFS (S3-compatible), and the full
Grafana LGTM stack plus Pyroscope and Alertmanager — twelve services,
one task up.
The habits matter as much as the code.
Contract tests on the wire — Handler tests structurally cannot see decode
failures, because those never reach a handler method. So the error envelope is
asserted through a real httptest server across all six failure classes: malformed
body, validation failure, missing field, bad path param, and handler errors on both
routes. Deleting the error handler turns it red.
→ errors_test.go
Fuzzing — 30 hand-picked adversarial seeds: rune-vs-byte length boundaries
(128 CJK characters = 384 bytes), lone surrogates, RTL overrides, duplicate JSON
keys, 1e309, 64-deep nesting. The property under test is if the validator
accepted it, the spec's constraints must actually hold. Runs nightly in CI.
→ validator_test.go ·
fuzz.yml
Integration tests against real Postgres — testcontainers spins a throwaway
database and applies the real migrations through the goose library, proving what
unit tests cannot: that a genuine 23505 becomes ErrConflict, and that handler
actually reaches the column.
→ postgres_integration_test.go
Testing call order, not just calls — a shared recorder across all three fakes
asserts [tx:begin, insert, append, tx:commit], and on conflict
[tx:begin, insert, tx:rollback, get_by_key] — which pins the subtle decision that
the replay read happens outside the aborted transaction.
→ commander_test.go#L104
Architecture as a CI gate — two overlapping tools, and the gate is verified to
fail on a planted violation. A gate you have never watched fail is one you are
trusting on faith.
→ ci.yml ·
how to verify
Generated code is never hand-edited — task gen regenerates ogen, sqlc and
protobuf. Change the spec, not the output.
→ Taskfile.yml
Reproducible tooling — every tool version is pinned into ./bin, and CI runs
the same task targets a developer runs, so "green locally, red in CI" cannot
come from version drift.
→ Taskfile.yml · ci.yml
Security scanning — govulncheck on every CI run, gosec in the linter set.
→ ci.yml
Quality gates on every commit — all green:
go build · go vet · go vet -tags=integration · go test -race ·
golangci-lint (16 linters, 0 issues) · go-arch-lint · govulncheck.
Coverage where it counts: internal/service/task 100%, internal/api 98%.
task setup # install every pinned tool into ./bin
task up # 12-service dev stack
task migrate:up # apply migrations
task run:server # terminal A
task demo # terminal B — reference producerThen open Grafana at localhost:3000 → Explore → Tempo, and watch a request become a trace with its SQL span and its log lines attached.
Step-by-step with expected output at every stage, plus troubleshooting:
docs/getting-started.md.
All Taskfile targets
task setup # install every pinned tool into ./bin
task build # build the flowhand binary
task test # go test -race ./...
task test:int # integration tests (testcontainers; requires Docker)
task fuzz # run every fuzz target (FUZZTIME=30s by default)
task lint # golangci-lint
task lint:arch # go-arch-lint — enforce the layered architecture
task vuln # govulncheck
task gen # regenerate protobuf + ogen + sqlc
task run:server # build + run the control-plane server
task migrate:up # apply database migrations
task migrate:down # roll back one migration
task demo # launch the reference producer
task up / down / logsflowhand/
├── api/
│ ├── openapi.yaml # REST contract — source of truth for ogen
│ └── proto/ # gRPC scaffold (Phase 1)
├── cmd/
│ ├── flowhand/ # the binary: server|worker|relay|ingester|ctl
│ └── flowhand-demo/ # reference producer, verifies replay on the wire
├── internal/
│ ├── api/ # transport — ogen Handler, mappers, error envelope
│ │ └── oas/ # GENERATED by ogen — never hand-edited
│ ├── service/task/ # use cases — Commander, Querier, ports
│ ├── domain/ # the leaf: aggregates, events, invariants
│ │ ├── tasks/ # ← the only context with real code today
│ │ ├── workflows/ # ← Phase 2 placeholder
│ │ └── schedules/ # ← Phase 1 placeholder
│ ├── repository/ # persistence — row ↔ aggregate, error translation
│ ├── storage/ # pgx pool, txmgr, sqlc output
│ ├── obs/ # slog + OTel + shutdown coordination
│ ├── server/ # http.Server wiring, pprof, /metrics
│ ├── config/ # koanf loader
│ └── cli/ # cobra commands — the composition root
├── migrations/ # goose
├── deploy/compose/ # the 12-service dev stack
└── docs/
├── getting-started.md # five-minute walkthrough
└── design/architecture/ # import rules and how they are enforced
Every package carries a doc.go stating its layer and its import rules —
e.g. internal/domain/tasks/doc.go.
Designed in detail, not yet built. No code links here, because there is no code yet — that is the point of keeping this section separate.
flowchart LR
P0["<b>Phase 0</b><br/>Foundations<br/>✅ complete"]
P1["<b>Phase 1</b><br/>Task engine<br/>🚧 next"]
P2["<b>Phase 2</b><br/>Workflow DAG<br/>📋 designed"]
P3["<b>Phase 3</b><br/>Distribution<br/>📋 designed"]
P4["<b>Phase 4</b><br/>Production polish<br/>📋 designed"]
P0 --> P1 --> P2 --> P3 --> P4
classDef done fill:#1f6f43,stroke:#0d3d24,color:#fff
classDef next fill:#8a6d1f,stroke:#5c4813,color:#fff
classDef todo fill:#3a3a3a,stroke:#222,color:#ccc
class P0 done
class P1 next
class P2,P3,P4 todo
The phase where the queue becomes real. SELECT … FOR UPDATE SKIP LOCKED dequeue
with LISTEN/NOTIFY wake-up; retries with exponential backoff and jitter; cron
schedules, delayed tasks and priorities; heartbeats with lease extension; a
dead-letter queue for poison messages; the transactional outbox relayed to Kafka;
gRPC streaming dispatch to workers; and the Redis idempotency fast path in front of
the durable floor that already exists.
Topological scheduling with cycle detection; fan-out and exactly-once fan-in under
SERIALIZABLE isolation; CEL conditional branches; cascading cancellation;
saga-style compensation on failure; property-based tests over generated DAGs.
The phase this whole project exists for: etcd leader election with fencing tokens (a stale leader's writes get rejected by a monotonic epoch checked on every write), consistent-hash sharding with live rebalance touching ≤ 1/N of keys, active-active sharded schedulers, and chaos drills for leader loss, worker loss and network partitions.
Large payloads spilled to S3-compatible object storage; SLOs with burn-rate alerts and real dashboards; k6 load tests driven to the capacity targets; pprof-guided optimisation; distroless image and goreleaser.
Capacity targets (Phase 4 will measure, not assume): 5,000 tasks/s submit peak · 2,000 tasks/s sustained execution · p99 submit < 50 ms · p99 end-to-end < 5 s.
Phase 0 was supposed to be "just tooling". The lessons that actually stuck were mostly about failures that look like success — and that is the theme I would most want to talk about in an interview.
-
A gate you have never watched fail is not a gate. My architecture linter config had a schema error that made it exit non-zero before analysing a single import. It looked configured. It enforced nothing. Now the fix is verified by planting a violation and confirming the build breaks (how).
-
Silently wrong beats loudly broken — for the bug, never for you. An early
GEThandler checkedif err != nil && errors.Is(err, pgx.ErrNoRows), so any other database error fell through to the success path and returned200with a zero-valued task. A 500 is strictly better than a confident lie. -
Unit mismatches are invisible. My latency histogram matched instruments by name alone. Then a routine
go mod tidymovedotelhttppast a semantic-convention cutover, changing the instrument from milliseconds to seconds. Same name, same View, every measurement now 1000× off — and the dashboard would have looked better, with the SLO permanently "met". The selector now pins the unit. -
Test the seam, not just the sides. My service layer had 100% coverage handling a conflict error — supplied by a fake. Nothing proved the repository actually produced that error from a real driver. Deleting the SQLSTATE check would have kept every test green. Now a real Postgres proves it.
-
Two type systems meeting is where bugs live. My domain says
complete; my published API sayssucceeded. An unchecked string conversion compiles fine and emits contract violations at runtime. The mapper is explicit, with a louddefaultarm. -
Codegen gives you vocabulary, not implementation. Declaring
400and500in OpenAPI generates the types and nothing else — the default handler still writestext/plain, so the spec promised JSON the server never sent. Caught only by testing on the wire.
MIT © 2026 Roman Agaltsev
Built by Roman Agaltsev · GitHub
Questions about any design decision here are very welcome — every one of them has a reason, and I would enjoy defending or revising it.