Event Bus
Want to react when data changes — push to a search index, notify another service,
warm a cache? Wire an EventBus with Client.UseEventBus, and every committed
Create, Update, and Delete publishes a typed event you can route anywhere:
client.UseEventBus(quark.NewLoggerEventBus(slog.Default()))
quark.For[Order](ctx, client).Create(&order)
// => level=INFO msg=... event=created table=orders (a "created" event is published)
The interfaces
type Event interface {
Kind() string // "created" | "updated" | "deleted"
Table() string
Payload() any // the model value the operation acted on
}
type EventBus interface {
Publish(ctx context.Context, event Event) error
}
Payload() is the model pointer, so a subscriber type-switches to recover the
struct:
func (b *myBus) Publish(ctx context.Context, e quark.Event) error {
if o, ok := e.Payload().(*Order); ok && e.Kind() == "created" {
return b.broker.Publish(ctx, "orders.created", o)
}
return nil
}
Two buses ship in-tree
quark.NewLoggerEventBus(logger)— logs each event to anslog.Loggerat Info (nil →slog.Default()). Never errors.quark.NewOTelEventBus(logger)— writes a correlation-tagged log record (event=quark.event.emit). It deliberately doesn't pull the OTel SDK into the core package; for real spans, implementEventBusagainst your tracer. (The OTel link only happens if your process installs an slog→OTel bridge such asotelslog.)
Both are starting points — production wires EventBus to a broker.
Delivery semantics
Events publish synchronously after the write commits, and there's no transactional outbox. If the process dies between commit and publish, the data is durable but the event is lost. For guaranteed delivery, write your own outbox and publish from a poller, or make subscribers idempotent and accept at-least-once. Quark publishes synchronously on purpose: a built-in outbox would need its own table, poller and delivery guarantees, and those belong to your infrastructure, not to the ORM.
Inside client.Tx (with ForTx[T]), the publish is registered via
Tx.OnCommit: it
fires after a durable commit and is discarded on rollback — a rolled-back
write emits nothing. Outside a transaction (For[T]), it publishes inline right
after the statement.
What emits: the per-row methods — Create (created),
Update/UpdateFields (updated), Delete (deleted). Bulk and WHERE-based
methods (CreateBatch, UpdateBatch, DeleteBatch, DeleteBy, UpdateMap) do
not emit — there's no single *T to attach. Emission isn't gated on rows
affected: treat an event as "the operation committed", not "a row definitely
changed" (whether a no-op update counts varies by engine, so gating would be
engine-dependent).
When Publish fails, the write is already saved — an emit failure never rolls
anything back:
| Path | On Publish error |
|---|---|
Non-transactional (For[T]) | The method returns the error wrapped in quark.ErrEventEmitFailed. The row stays written — retry the emit, don't re-create. |
Transactional (ForTx[T]) | Logged (event=quark.event.emit_failure), not propagated — the commit already succeeded. |
if err := quark.For[Order](ctx, client).Create(&order); err != nil {
if errors.Is(err, quark.ErrEventEmitFailed) {
// The order IS saved; only the event failed. Re-publish, don't re-create.
}
}
Connecting a broker
type NATSBus struct{ nc *nats.Conn }
func (b *NATSBus) Publish(ctx context.Context, e quark.Event) error {
data, err := json.Marshal(e.Payload())
if err != nil {
return err
}
return b.nc.Publish("quark."+e.Table()+"."+e.Kind(), data)
}
client.UseEventBus(&NATSBus{nc: nc})
Kafka or Redis Streams follow the same shape: serialize e.Payload(), derive a
topic from e.Table() + e.Kind(), hand off to the broker.
Inbound: PostgreSQL LISTEN/NOTIFY
The event bus above is the outbound side. Separately, on PostgreSQL you can
consume NOTIFY messages. ListenerFactory.CreateListener returns a listener
that pins one dedicated connection (a LISTEN registration lives on the physical
connection, and the pool rotates connections freely). Every non-PostgreSQL
dialect returns ErrDialectNotSupported.
The listener needs the pgx driver in database/sql mode — open the client
with quark.New("pgx", dsn) after importing
_ "github.com/jackc/pgx/v5/stdlib". It reaches under the pool for the raw
*pgx.Conn, so a client opened with lib/pq (the driver the installation
guide suggests as a general default) fails with ErrDialectNotSupported at
Listen time even though the dialect is PostgreSQL.
listener, err := quark.NewListenerFactory(client).CreateListener()
if err != nil {
return err // ErrDialectNotSupported off PostgreSQL
}
defer listener.Close()
if err := listener.Listen(ctx, "orders"); err != nil {
return err
}
// Emit from anywhere: quark.Notify(ctx, client, "orders", `{"id":42}`)
for {
payload, err := listener.Receive(ctx) // blocks until a NOTIFY arrives
if err != nil {
return err // ctx cancelled or connection dropped — reconnect
}
log.Printf("channel=%s payload=%s", payload.Channel, payload.Payload)
}
Sentinels (match with errors.Is): ErrDialectNotSupported (CreateListener off
PG), ErrNoSubscription (Receive before any Listen), ErrListenerClosed
(any call after Close).
Listen/Receive/Close are serialized over the one pinned connection, and
Receive blocks holding it — register every channel first, then loop Receive in
one goroutine; to stop, cancel the Receive context, then Close. PostgreSQL
LISTEN/NOTIFY has no durable buffer: notifications emitted while the
connection is down are lost, so reconnect on error and reconcile state yourself.
The listener holds one connection from the pool for its lifetime — size the pool
accordingly or give it its own Client. A dedicated connection is required:
LISTEN registers on a specific session, so a listener sharing the general pool
would silently stop receiving notifications the moment its connection was
recycled.