Események
Az eventek dispatcheléséhez és a rájuk feliratkozáshoz add hozzá az events-et a queue mellé:
import (
"github.com/gp-system/events"
"github.com/gp-system/queue"
)
Az events egy önálló event- és listener-modul, a queue csomagra épülve: minden listener eleve queue-zott, alapból aszinkron módon fut. Az alkalmazáskód eldob egy eventet, a listenerek típus szerint feliratkoznak, és minden listener önálló asynq taskként fut, saját retry-kerettel: egy event N listener-taskra terül szét, és egy hibázó listener újrapróbálkozása soha nem futtatja újra a többit. Bármilyen Go programban működik, kit nélkül is: egy sima net/http service, egy CLI-eszköz vagy egy Lambda ugyanígy tud eventet eldobni és listenereket regisztrálni, csak ez a modul és egy Valkey kell hozzá.
A kézbesítés végig legalább-egyszeres (outbox → relay → asynq → listener), ezért két szabály tartja össze az egészet: állapotváltó folyamatban outboxon át dobd el az eventet, a listenereidet pedig írd idempotensre.
Telepítés
go get github.com/gp-system/events@v0.1.0 # a kit is ezt a taget használja
# vagy: go get github.com/gp-system/events@latest
Az events közvetlenül behúzza a gp-system/errs-t és a gp-system/queue-t; a tranzakciós events/outbox alcsomaghoz ezen felül kell a gp-system/dbx, a jackc/pgx/v5, és ha a bun-alapú store-t használod, az uptrace/bun is; az events/scheduler alcsomagnak a gyökér-modulon túl semmi extra nem kell. A hibiken/asynq, a redis/go-redis/v9 és a google/uuid a queue-n keresztül tranzitívan érkezik. Ehhez semmihez nem kell a kit: importáld azt az alcsomagot, ami kell, és semmi mást.
Event definiálása
Az event egy sima struct stabil EventName()-mel. A shop orderPlaced eventjét az add event shop orderPlaced generálta, a payload-mezőket te adod hozzá:
package events
type OrderPlaced struct {
OrderID string `json:"orderId"`
Email string `json:"email"`
}
func (OrderPlaced) EventName() string { return "shop.orderPlaced" }
A generátor névkonvenciója <modul>.<eventNév> camelCase-ben, ebből lesz a event:shop.orderPlaced task-típus is. A név perzisztálódik (az outbox-táblában és a Valkeyben), ezért tartsd stabilan a deployok közt; ha átnevezed, a már bent lévő taskok az új deploy alatt nem találnak handlert.
EventName()-nek a típus zero value-ján is működnie kell: érték típust regisztrálj, ne pointert. A Listen a zero value-ból olvassa ki a nevet, és startupkor azonnal panikol, ha az üres vagy a hívás elszáll, így a hibás regisztráció nem jut el prodig.Dispatch
A Dispatcher a modul Dependencies-ébe injektálódik (az első add event drótozza be mindkét belépési pontba). A shop rendelési szabálya (a modul core/ csomagjában) a tranzakción belül dobja el az eventet:
func (s *Service) PlaceOrder(ctx context.Context, in PlaceOrderInput) (Order, error) {
var order Order
err := s.tx.WithinTransaction(ctx, func(ctx context.Context) error {
var err error
if order, err = s.orders.Insert(ctx, in); err != nil {
return err
}
if err = s.stock.Decrement(ctx, in.Items); err != nil {
return err
}
return s.dispatcher.Dispatch(ctx, events.OrderPlaced{
OrderID: order.ID,
Email: order.Email,
})
})
return order, err
}
Az events.Dispatcher interfészt három implementáció elégíti ki:
| Konstruktor | Viselkedés | Mikor |
|---|---|---|
outbox.NewDispatcher(store) | az eventet az outbox_events táblába írja a hívó tranzakciójában; egy relay továbbítja Valkeybe | állapotváltó folyamatok, a generált alapértelmezés, a shop is ezt használja |
events.NewDispatcher(client) | közvetlenül Valkeybe enqueue-l (queue.Client-en át) | fire-and-forget: ha a folyamat a DB-commit és az enqueue közt hal le, az event elvész, ott jó, ahol ez elfogadható |
events.NewSyncDispatcher(reg) | inline futtatja a regisztrált listenereket, a hívó goroutine-jában és tranzakciójában; a listener-hibákat errors.Join-nal adja vissza | unit tesztek és lokális tooling |
Hogy miért az outbox az alapértelmezés, és mi a baj a közvetlen enqueue-val állapotváltásnál, azt az outbox oldal fejti ki.
A boríték és a kézbesítési metaadat
Dispatchkor az events.NewEnvelope(ctx, ev) épít egy queue.Envelope-ot: friss UUID (ID), az event neve, a JSON-payload, a context trace-kontextusa és az időbélyeg. Bármelyik dispatchert használod, a listener ugyanazt a borítékot kapja.
A worker a listener hívása előtt a borítékból kézbesítési metaadatot tesz a contextre: a listener az events.MetaFromContext(ctx)-szel olvassa:
| Mező | Jelentés |
|---|---|
ID | a boríték id-ja, stabil idempotencia-kulcs a retryk és az event összes listenere közt |
Name | az event neve |
OccurredAt | mikor dobták el az eventet |
Attempt | hányadik kísérlet ez (0 az első kézbesítéskor) |
Listenerek
A listenereket a generikus events.Listen[T] regisztrálja: az event típusát a handler szignatúrája adja, kézi típus-assertion nincs. Az add listener shop orderPlaced sendOrderConfirmation --queue mail két dolgot generál. A regisztrációs sort a modul listeners/register.go-jába:
package listeners
import (
"github.com/acme/shop/internal/modules/shop"
"github.com/gp-system/events"
"github.com/gp-system/queue"
// gpsystem:worker-imports
)
func Register(reg *events.Registry, deps shop.Dependencies) {
events.Listen(reg, "sendOrderConfirmation", sendOrderConfirmation(deps), queue.OnQueue("mail"))
// gpsystem:listeners
}
És a listener-stubot külön fájlba: a törzs a tiéd, itt a shop kitöltött változata:
package listeners
import (
"context"
"github.com/gp-system/mail"
"github.com/gp-system/mail/mjml"
"github.com/acme/shop/internal/modules/shop"
shopevents "github.com/acme/shop/internal/modules/shop/events"
)
func sendOrderConfirmation(deps shop.Dependencies) func(context.Context, shopevents.OrderPlaced) error {
return func(ctx context.Context, ev shopevents.OrderPlaced) error {
msg := mail.NewMessage().
WithTo(ev.Email).
WithSubject("Rendelésed visszaigazolása").
WithBody(mjml.Template(templates, "templates/order_confirmation.mjml.tmpl",
map[string]any{"OrderID": ev.OrderID}))
return deps.Mailer.Send(ctx, msg)
}
}
(A mail/mjml API-hoz lásd az E-mail oldalt.)
Amit a regisztrációról tudni kell:
- A név (a
Listen2. argumentuma) eventenként egyedi; ebből lesz alistener:shop.orderPlaced:sendOrderConfirmationtask-típus, ami a worker logjaiban és az Asynqmonban is látszik. Ugyanarra a (event, név) párra kétszer regisztrálni startup-panik. - A további argumentumok listenerenkénti enqueue-opciók:
queue.OnQueuea queue-hoz (a shop levelei amailqueue-ba mennek, hogy a worker súlyozhassa őket),queue.MaxRetry,queue.Timeouta retry-kerethez. - Másik modul eventjére a
add listener <modul> <producer>.<event> <név>formával iratkozol fel: az event típusát a producereventscsomagjából importálod.
A Registry ezen felül EventNames() []string-et, FanoutTask(env queue.Envelope) ([]FanoutTask, error)-t és Handler(taskType string) (func(context.Context, queue.Envelope) error, bool)-t is ad: ezekkel bontja ki a worker chassis az event: taskot listener-taskokra, és irányítja a listener: taskot a megfelelő függvényhez. Ezekhez csak akkor nyúlj közvetlenül, ha a kit worker csomagja nélkül drótozol be egy fogyasztót; a kittel a worker.Run mindezt elvégzi helyetted.
Fan-out: hogyan lesz egy eventből N task
A worker az event:shop.orderPlaced taskot kibontja, és listenerenként egy listener:... taskot enqueue-l. Minden listener-task determinisztikus TaskID-t kap (<boríték-id>:<listener-név>), így ha a fan-out félbeszakad és újrafut, a már beadott taskok ErrDuplicate-ként no-opok: a fan-out maga is idempotens.
Innentől minden listener a saját életét éli: külön queue-ban lehet, külön retry-keretet fogyaszt, és a hibája csak őt futtatja újra. Ez az alapból-mindig-aszinkron viselkedés (nem egy szinkron event-loop) a csomag meghatározó tulajdonsága, és az ára az idempotencia-követelmény.
Idempotencia
Legalább-egyszeres kézbesítés mellett egy listener többször is lefuthat ugyanarra az eventre (retry, worker-crash a feldolgozás közben, relay-újraküldés). A Meta.ID erre a stabil kulcs: ugyanannál a dispatchnél minden kísérletben ugyanaz.
meta, _ := events.MetaFromContext(ctx)
// pl. egyedi kulcs a "megtörtént már?" táblában, vagy a mail-szolgáltató
// idempotency kulcsa:
return deps.Mailer.SendOnce(ctx, meta.ID, msg)
Praktikus minták:
- Természetes kulcs: ha a mellékhatás DB-írás,
INSERT ... ON CONFLICT DO NOTHINGaz event id-ra (vagy az üzleti kulcsra). - Külső API: add át a
meta.ID-t idempotency-key headerként, ha a szolgáltató támogatja. - Olvasás + feltétel: a "már elküldtük?" ellenőrzés és a küldés közt maradó versenyt csak kulcs-alapú megoldás zárja be, a puszta
SELECT-tel ne elégedj meg.
Hol futnak a listenerek
A listenerek a worker binárisban futnak: a modulonkénti listeners alcsomagok Register függvényeit a generált cmd/worker/main.go komponálja egyetlen events.Registry-be. Ütemezett event-tüzeléshez (cron szerint dispatchelt eventek) lásd az ütemezőt; a generátorparancsok (add event, add listener) részleteihez a CLI worker-generátorokat.
A kittel
A new project a kezdetektől bedrótozza az events-et: az add event, az add listener és az add job generálja a fenti Registry/Schedule-kompozíciót, és a worker chassis az, ami ténylegesen futtatja (a saját asynq szervere mellett az outbox relay-t és az ütemező-futtatót is). Ehhez semmi nem szükséges a kitből: kit nélkül saját asynq szervert építesz, és a reg.Handler(taskType)-pel regisztrálod a rád tartozó task-típusokat.
Használt patternek
Legalább-egyszeres kézbesítés + idempotens listener-szerződés (Meta.ID mint dedup-kulcs, determinisztikus TaskID a fan-outban): lásd a Design patterneket.