Worker
A worker-chassis a háttérfeldolgozásodat (eventek, listenerek, ütemezett jobok) saját binárisként futtatja, cmd/worker néven, a HTTP API mellett. Ugyanúgy importálod, mint a server-t:
import "github.com/gp-system/gpsystem/worker"
A worker.Run a HTTP-engine-ek Run-jának háttérmunka-párja: ugyanaz az app életciklus (jelkezelés, graceful shutdown, closerek, telemetria-flush), csak HTTP-szerver helyett három dolgot futtat egy binárisban:
- az asynq szervert: az event-, listener-, ütemezett- és job-taskok feldolgozása;
- az outbox relay-t: a commitált eventek továbbítása a queue-ba;
- a schedule runnert: a cron-bejegyzések tüzelése, leader-electionnel.
Ez a task-feldolgozást és az ütemezést egyetlen folyamatban futtatja: külön worker- és scheduler-deployment helyett egy cmd/worker binárist skálázol. A bináris alapból a new project-tel készül; ha egy projekt new project --no-worker-rel jött létre, az add worker adja hozzá. Akárhogy is, az API-tól függetlenül deployolod és skálázod.
A generált cmd/worker
A shop worker-binárisa: minden sorát a generátor írta, a modulokat a new module / add listener / add job drótozza be a // gpsystem:* anchoroknál.
package main
import (
"context"
"log"
"github.com/gp-system/dbx/pg"
"github.com/gp-system/envconf"
"github.com/gp-system/events"
"github.com/gp-system/events/outbox"
"github.com/gp-system/events/scheduler"
"github.com/gp-system/gpsystem/worker"
"github.com/acme/shop/internal/modules/shop"
shopjobs "github.com/acme/shop/internal/modules/shop/jobs"
shoplisteners "github.com/acme/shop/internal/modules/shop/listeners"
"github.com/acme/shop/internal/platform/config"
)
func main() {
cfg := envconf.MustLoad[config.Config]()
ctx := context.Background()
pool := pg.MustNewPool(ctx, cfg.DB)
relay := outbox.MustNewRelay(ctx, cfg.Outbox, pool, cfg.Worker.Valkey)
err := worker.Run(ctx, cfg.Worker, func(reg *events.Registry, sched *scheduler.Schedule) error {
shopDeps := shop.Dependencies{
DB: pg.NewDB(pool),
Transactor: pg.NewTransactor(pool),
Dispatcher: outbox.NewDispatcher(outbox.NewStore(pg.NewDB(pool))),
}
shoplisteners.Register(reg, shopDeps)
shopjobs.Register(sched, shopDeps)
return nil
},
worker.WithOutboxRelay(relay),
worker.WithCloser("pgxpool", func(context.Context) error {
pool.Close()
return nil
}),
)
if err != nil {
log.Fatal(err)
}
}
A minta ugyanaz, mint a HTTP-oldalon: a függőségek explicit Dependencies structban utaznak, a register callback pedig megkapja az event-Registry-t és a Schedule-t, amiket a modulok generált listeners / jobs alcsomagjainak Register függvényei töltenek fel. A worker Dispatcher-e is az outboxos: egy listener által dispatchelt esemény ugyanazt a garanciát kapja, mint egy HTTP-kérésből dobott.
Run
func Run(ctx context.Context, cfg Config, register RegisterFunc, opts ...Option) error
type RegisterFunc func(reg *events.Registry, sched *scheduler.Schedule) error
A Run sorban: telemetria-bootstrap (telemetry.Setup) → a register felépíti a registryt és a schedule-t → queue.Client nyitása (PING) → asynq szerver indítása, mellette a relay (ha WithOutboxRelay-jel átadtad) és a schedule runner (ha a schedule nem üres; üres schedule mellett leader-election sincs) → blokkol SIGINT/SIGTERM-ig. Tiszta leállásnál nil-t ad vissza.
Mi dolgozza fel a taskokat
A worker négy beépített handlert regisztrál a task-típus-prefixekre:
| Task | Mit csinál a handler |
|---|---|
event:<név> | fan-out: listenerenként egy listener: taskot enqueue-l, determinisztikus TaskID-vel |
listener:<event>:<név> | dekódolja a borítékot, contextre teszi a Meta-t, meghívja a listeneredet |
schedule:<event> | ütemezett tüzelés: friss borítékot épít (id = a trigger task-id-ja), és fan-outol |
job:<név> | meghívja az ütemezett job handlerét |
Minden handler három middleware-en át fut: Sentry-hub taskonként (breadcrumb-izoláció), panic-recovery (a panikoló handler errs-hibává válik és retry-zik, a worker nem hal bele) és tracing: a consumer span a boríték trace-kontextusát folytatja, így a teljes út egy trace. Olyan taskot, aminek nincs regisztrált handlere (pl. egy deploy kivette a listenert), a worker warning mellett átugrik, nem retry-zik örökké.
Shutdown-sorrend
SIGINT/SIGTERM-re a Run determinisztikus sorrendben áll le:
- leállnak a termelő ciklusok: a schedule runner (elengedi a lease-t) és a relay poll-ciklusa;
- kiürülnek a folyamatban lévő taskok: legfeljebb
WORKER_SHUTDOWN_TIMEOUT-ig; ami nem fér bele, azt az asynq visszateszi a sorba, és később újra kézbesíti (ezért is: idempotens listenerek); - a closerek LIFO-ban futnak: előbb a kit sajátjai (scheduler, relay, queue-kliens), aztán a te
WithCloser-eid (a shopban a pgx-pool); - telemetria-flush: az utolsó spanok és logok még elmennek.
Ez ugyanaz a garancia, mint a HTTP-oldalon: egy deploy nem vág el félbe levélküldést. A részletekhez lásd az alkalmazás-életciklust.
Opciók
| Opció | Hatás |
|---|---|
WithOutboxRelay(*outbox.Relay) | az outbox relay futtatása ebben a workerben (N replikán biztonságos) |
WithCloser(name, fn) | takarítás graceful shutdownkor, a taskok kiürítése után (LIFO) |
WithoutTelemetry() | a telemetry.Setup kihagyása; teszthez, vagy ha a folyamat maga konfigurálja |
WithAsynqConfig(func(*asynq.Config)) | a nyers asynq-konfig módosítása a szerver felépítése előtt; escape hatch a kit által nem exponált beállításokhoz |
WithMux(func(*asynq.ServeMux)) | extra nyers asynq handlerek regisztrálása, olyan task-típusokhoz, amiket a kit nem modellez |
A két escape hatch a kit "centralizál, nem absztrahál" elvének worker-oldali fele: ha az asynq egy képességére van szükséged, nem kell kilépned a chassis-ból.
worker.Run(ctx, cfg.Worker, register,
worker.WithAsynqConfig(func(c *asynq.Config) {
c.HealthCheckFunc = reportHealth
}),
worker.WithMux(func(mux *asynq.ServeMux) {
mux.HandleFunc("import:products", handleProductImport) // saját task-típus
}),
)
Queue-prioritások
A worker a queue-kat súlyozottan szolgálja ki: a WORKER_QUEUES egy név:súly lista, a Horizon balance/queue-súlyozásának megfelelője. A shop a leveleit külön mail queue-ba teszi (a listener queue.OnQueue("mail") regisztrációja), így egy levél-burst nem szoríthatja ki a többi munkát:
WORKER_QUEUES=default:3,mail:1
Így a worker a kapacitása ~75%-át a default, ~25%-át a mail queue-ra fordítja, amikor mindkettőben van munka. A WORKER_STRICT_PRIORITY=true súlyozás helyett szigorú sorrendet ad: előbb a legnagyobb súlyú queue ürül ki teljesen, csak aztán jön a következő, ami hasznos, ha van egy critical sorod, ami mindent megelőz.
Konfiguráció
A worker.Config-ot prefix nélkül ágyazod a projekt-konfigba (a generált config.go-ban készen van):
type Config struct {
Worker worker.Config
DB dbx.Config `envPrefix:"DB_"`
Outbox outbox.Config `envPrefix:"OUTBOX_"`
}
| Változó | Default | Jelentés |
|---|---|---|
WORKER_CONCURRENCY | 10 | egyszerre ennyi task dolgozódik fel |
WORKER_QUEUES | default:1 | queue → súly párok (critical:6,default:3,low:1) |
WORKER_STRICT_PRIORITY | false | súlyozás helyett szigorú prioritás-sorrend |
WORKER_SHUTDOWN_TIMEOUT | 30s | a graceful drain korlátja; a be nem fejezett taskok visszakerülnek a sorba |
WORKER_DEFAULT_MAX_RETRY | 25 | a saját queue.MaxRetry nélküli listener-/job-taskok retry-kerete |
A worker.Config három további configot ágyaz be, így a worker-bináris egyetlen env-készletből konfigurálódik:
queue.ConfigVALKEY_prefixszel (VALKEY_ADDR, ...),scheduler.ConfigSCHEDULER_prefixszel,telemetry.Configprefix nélkül (a standardOTEL_*,LOG_*,SENTRY_*nevek), lásd az observability áttekintőt.
A teljes env-referencia: konfiguráció.
Logging
A worker taskonként logol az alapértelmezett slog loggeren át, a consumer spanhez trace-korreláltan:
worker: task started(Info):task_type,task_id,attemptattribútumokkal;worker: task done(Info): pluszduration;worker: task failed(Error): pluszerror; a szándékos skip (SkipRetry, pl. nincs regisztrált handler) csak egyszer, warningként jelenik meg.
Az asynq saját belső logjai is ugyanebbe a pipeline-ba folynak. Hogy a sorok hova kerülnek (dev-konzol, JSON-stdout, OTLP), azt a logging-konfig dönti el; LOG_LEVEL=WARN elnémítja a taskonkénti Info-sorokat.
Skálázás
A worker állapotmentes: annyi replikát futtatsz, amennyi munkád van, és mindhárom komponens replikabiztos.
- az asynq szerver természeténél fogva: egy taskot egyszerre egy worker kap;
- az outbox relay a
FOR UPDATE SKIP LOCKED+ determinisztikus TaskID párossal; - a schedule runner a Valkey-lease leader-electionnel: N replikából egy tüzel.
Semmit nem kell kikapcsolni vagy külön deployolni a skálázáshoz. Az ára egyetlen szerződés: a kézbesítés legalább-egyszeres, tehát a listenerek idempotensek.
Helyi futtatás
go run ./cmd/migrate up # az outbox_events migráció is lefut
mise run worker # go run ./cmd/worker
A Valkey a projekt compose.yml-jéből jön (mise run dev). A generátorparancsokhoz (add event, add listener, add job, add worker) lásd a CLI worker-generátorokat.