A kit

Worker

A háttérmunka-chassis, asynq szerver, outbox relay és schedule runner egyetlen binárisban, közös graceful shutdownnal.

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:

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.

cmd/worker/main.go
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:

TaskMit 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:

  1. leállnak a termelő ciklusok: a schedule runner (elengedi a lease-t) és a relay poll-ciklusa;
  2. 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);
  3. 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);
  4. 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óDefaultJelentés
WORKER_CONCURRENCY10egyszerre ennyi task dolgozódik fel
WORKER_QUEUESdefault:1queue → súly párok (critical:6,default:3,low:1)
WORKER_STRICT_PRIORITYfalsesúlyozás helyett szigorú prioritás-sorrend
WORKER_SHUTDOWN_TIMEOUT30sa graceful drain korlátja; a be nem fejezett taskok visszakerülnek a sorba
WORKER_DEFAULT_MAX_RETRY25a 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:

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, attempt attribútumokkal;
  • worker: task done (Info): plusz duration;
  • worker: task failed (Error): plusz error; 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.

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.

Copyright © 2026