Kiroku can carry W3C trace context on every event so a consumer reading an
event later — a projection, a subscription, a process manager — can continue
the trace that produced it. The kiroku-otel package provides pure helpers
to write and read that context; the store's append/read hooks let you apply
them automatically.
kiroku-otel is a separate package so kiroku-store stays free of any
hs-opentelemetry dependency. Opt in by depending on kiroku-otel directly.
Trace context lives in the event's metadata JSONB column as W3C
trace-context strings:
{
"traceparent": "00-<32-hex traceId>-<16-hex spanId>-<2-hex flags>",
"tracestate": "<vendor entries, optional>"
}Any other keys in metadata are preserved. This is a plain JSON convention,
so consumers that do not use kiroku-otel (for example the
Shibuya Adapter, which copies traceparent into its
envelope) can still read it.
Kiroku.Otel.TraceContext exposes two pure functions:
injectTraceContext :: SpanContext -> EventData -> EventData
extractTraceContext :: RecordedEvent -> Maybe SpanContextinjectTraceContext encodes a SpanContext into traceparent / tracestate
and merges them into the event's metadata. Existing keys are preserved; any
existing traceparent / tracestate are overwritten (the W3C spec mandates
exactly one of each). If metadata is absent or not a JSON object, it starts
from an empty object.
extractTraceContext pulls a SpanContext back out of a RecordedEvent's
metadata. It returns Nothing — and never throws — when metadata is
absent, is not an object, lacks traceparent, or holds a traceparent that
fails W3C parsing.
import Kiroku.Otel.TraceContext (extractTraceContext)
handle :: RecordedEvent -> IO ()
handle event =
case extractTraceContext event of
Just ctx -> withLinkedSpan ctx (process event)
Nothing -> process eventYou can apply these helpers by hand on each event, but the store's
StoreSettings hooks apply them to every event automatically:
enrichEventfires on the append path, on the typedEventData, before the SQL encoder runs. Use it to inject the current span into outgoing events.decodeHookfires on the read and subscription paths, on theRecordedEventabout to be surfaced. Use it to attach derived metadata, decrypt, or redact.
import Control.Lens ((&), (.~))
import Kiroku.Store
import Kiroku.Otel.TraceContext (injectTraceContext)
settings :: ConnectionSettings
settings =
defaultConnectionSettings connStr
& #storeSettings .~
defaultStoreSettings
{ enrichEvent = Just $ \ed -> do
ctx <- captureCurrentSpanContext -- from your OTel setup
pure (injectTraceContext ctx ed)
}Wire the resulting StoreSettings into ConnectionSettings via its
storeSettings field; withStore copies it onto the store handle so the
interpreter reaches it for every event.
Both hooks default to Nothing, in which case the interpreter takes a pure
fast path that does no extra allocation and no traversal — there is no cost
when you do not use them.
The enrichEvent hook fires inside the Store interpreter
(runStorePool). Callers that bypass the interpreter and use
Kiroku.Store.Transaction.appendToStreamTx directly do not get
enrichment. Opt in manually with enrichEventsIO before building the
prepared event list:
import Kiroku.Store.Transaction (enrichEventsIO)
enriched <- enrichEventsIO storeSettings rawEvents
-- then prepare and append within your transactionThe high-level runTransactionAppending wrapper (see
Appending Events) goes through the normal path, so it
honors enrichEvent without extra work.
The helpers above propagate trace context on individual events. A separate opt-in handler turns a subscription worker's lifecycle into spans, so the finite state machine described in Subscriptions — catch up, go live, pause under backpressure, reconnect, retry, dead-letter — shows up on a trace timeline.
Kiroku.Otel.Subscription exposes a factory that builds a KirokuEvent
callback from an OpenTelemetry Tracer:
import Kiroku.Otel.Subscription (subscriptionTraceHandler)
import Control.Lens ((&), (.~))
import Kiroku.Store
-- tracer :: OpenTelemetry.Trace.Core.Tracer (from your TracerProvider)
handler <- subscriptionTraceHandler tracer
settings :: ConnectionSettings
settings =
defaultConnectionSettings connStr
& #eventHandler .~ Just handlerIt consumes the same KirokuEvent stream covered in
Observability; install it as the eventHandler and every
subscription emits spans thereafter. (Compose it with your own logging handler
by fanning the event out to both.)
Each span is short and ends promptly — the SDK only exports a span when it ends, so a single span held open for the worker's lifetime would be invisible while the worker runs and lost on a crash. The model is therefore:
| Span | Opened on | Ended on |
|---|---|---|
kiroku.subscription.catchup |
Started |
CaughtUp |
kiroku.subscription.deliver |
each delivered batch (every target, both phases) | immediately (per batch) |
kiroku.subscription.paused |
Paused |
Resumed |
kiroku.subscription.reconnecting |
first Reconnecting |
next CaughtUp (later attempts add a reconnect.attempt span event) |
kiroku.subscription.retrying |
first Retrying of a poison event |
DeadLettered (status Error) or the worker moving on (status Ok) |
kiroku.subscription.dead_letter |
an immediate dead-letter (no retry) | immediately |
kiroku.subscription.stopped |
Stopped (always) |
immediately |
kiroku.subscription.deliver is emitted once per non-empty batch handed to
the handler, on every target ($all, category, consumer group) and in
both phases — it carries kiroku.batch.rows,
messaging.batch.message_count, messaging.operation.type = "process", and a
kiroku.subscription.state of "catchup" or "live". This closes the previous blind spot where an $all
subscription's live phase emitted no per-batch span and so went dark on the
timeline once it caught up. (The worker still emits KirokuEventSubscriptionFetched
for the DB-driven live-fetch-rate signal, but the tracer intentionally does not
trace it — the deliver span subsumes it, so the category/consumer-group live path
does not produce two spans per batch.)
kiroku.subscription.stopped is always emitted when the worker stops,
carrying kiroku.subscription.stop_reason and kiroku.checkpoint.global_position,
so the terminal state is present in the trace even for a healthy worker that
catches up, runs live, and stops with no open episode span.
The honest limitation: an in-progress episode (a pause that has not resumed
yet) does not appear in the backend until it ends. For "what state is the worker
in right now," use the currentState handle accessor
(Observability) or the
store's subscription-state registry snapshot.
Every span carries a Kiroku-specific kiroku.* attribute set (the constants are
exported from Kiroku.Otel.Subscription):
kiroku.subscription.namekiroku.subscription.state—"catchup","live","paused","reconnecting","retrying"kiroku.consumer_group.member/kiroku.consumer_group.size— only for group memberskiroku.checkpoint.global_position,kiroku.subscription.attempt,kiroku.batch.rows,kiroku.event.global_position,kiroku.dead_letter.reason,kiroku.subscription.stop_reason
Where a current OpenTelemetry semantic-convention key exists, the tracer emits
that standard key as well. Subscription spans include messaging.system = "kiroku" and messaging.destination.name; delivery spans add
messaging.operation.type = "process" and messaging.batch.message_count.
Database-error spans add db.system.name = "postgresql" and
db.operation.name.
The eventHandler callback runs synchronously on the worker loop's thread
(one thread per consumer-group member). Opening and ending a span is cheap and
in-memory, but the export may block — so configure your TracerProvider with
a batch span processor (the SDK default), which exports on a background
thread. With a simple/synchronous processor a slow exporter would stall the
worker.
The same identity travels onto Shibuya's per-message spans. The
Shibuya Adapter stamps kiroku.subscription.name,
kiroku.consumer_group.member, kiroku.event.type, and
kiroku.event.global_position onto each Envelope. It also stamps
messaging.system = "kiroku" and messaging.destination.name with the
subscription name, which Shibuya merges into the span it opens for the message.
So a single trace is followable from a kiroku subscription into Shibuya's
processing, tagged with the same keys on both sides.
kiroku-otel depends on hs-opentelemetry-api (^>=1.0),
hs-opentelemetry-propagator-w3c (^>=1.0), and
hs-opentelemetry-semantic-conventions (^>=1.40). SpanContext,
wrapSpanContext, and the W3C propagator come from those packages; capturing
the current span context is the responsibility of your application's
OpenTelemetry tracer setup.
- Observability — operational events and metrics, which complement tracing.
- Shibuya Adapter — propagates
traceparentinto Shibuya envelopes. - Appending Events — where
enrichEventruns.