Telemetry events emitted by PgFlow.
PgFlow uses :telemetry to emit events at key points in the workflow lifecycle.
These events can be used for monitoring, logging, and metrics collection.
Events
Worker Lifecycle
[:pgflow, :worker, :start]— Worker process started[:pgflow, :worker, :stop]— Worker process stopped
Poll Cycles
[:pgflow, :worker, :poll, :start]— Poll cycle started[:pgflow, :worker, :poll, :stop]— Poll cycle completed
Task Execution
[:pgflow, :worker, :task, :start]— Task execution started[:pgflow, :worker, :task, :stop]— Task execution completed successfully[:pgflow, :worker, :task, :exception]— Task execution failed
Step Lifecycle
[:pgflow, :step, :skipped]— Step skipped without a worker (e.g. unmetif/if_not)
Run Lifecycle
[:pgflow, :run, :started]— Flow run created (emitted byPgFlow.Client)[:pgflow, :run, :completed]— Flow run completed (emitted by worker after task cascades)[:pgflow, :run, :failed]— Flow run failed (emitted by worker after task cascades)
Attaching Handlers
:telemetry.attach_many(
"my-handler",
[
[:pgflow, :worker, :task, :stop],
[:pgflow, :run, :completed],
[:pgflow, :run, :failed]
],
&MyModule.handle_event/4,
nil
)Default Logger
PgFlow includes a default logger handler that can be attached by setting
attach_default_logger: true in the configuration.
Note: The default logger is disabled by default since PgFlow.Logger provides
structured logging directly in the worker. Enable this if you need telemetry-based
logging for specific use cases like metrics collection or external log aggregation.
Summary
Functions
Attaches the default telemetry handlers for logging.
Detaches the default telemetry handlers.
Emits [:pgflow, :step, :skipped] for every skipped step on a run.
Emits [:pgflow, :step, :skipped] for skips not already announced.
Emits [:pgflow, :step, :skipped] when a step is skipped without a worker.
Functions
@spec attach_default_logger() :: :ok | {:error, :already_exists}
Attaches the default telemetry handlers for logging.
This is called automatically on application start if attach_default_logger: true
is set in the configuration.
@spec detach_default_logger() :: :ok | {:error, :not_found}
Detaches the default telemetry handlers.
@spec emit_skipped_steps(Ecto.Repo.t(), String.t(), String.t()) :: :ok
Emits [:pgflow, :step, :skipped] for every skipped step on a run.
One-shot form, for callers that observe a run exactly once (such as
PgFlow.Client.start_flow/2, which sees the skips the run was born with).
Callers that sweep the same run repeatedly must use
emit_skipped_steps/4 instead, or they will re-announce every skip on
every sweep.
@spec emit_skipped_steps(Ecto.Repo.t(), String.t(), String.t(), MapSet.t(String.t())) :: MapSet.t(String.t())
Emits [:pgflow, :step, :skipped] for skips not already announced.
Looks up skipped steps via Flows.list_skipped_steps/2 (dependency
ordered, so a parent's skip is always emitted before its cascaded
children). already_emitted is a MapSet of step slugs this caller has
already announced for run_id; the union of it and the slugs emitted by
this call is returned, ready to be passed back in on the next sweep.
Query errors are swallowed so callers can treat this as fire-and-forget; the set is returned unchanged in that case and the skips are picked up on a later sweep.
Delivery contract
Skips are decided in PostgreSQL, and PgFlow discovers them by polling
step_states after each complete_task/fail_task. Passing the returned
set back in makes a single emitter announce each skip exactly once — the
guarantee non-idempotent handlers need.
It is a per-emitter guarantee, not a global one. Two workers processing
the same run each sweep it independently, so a skip can be announced once
per worker that touches the run (and once by Client.start_flow/2 for
skips decided at run start). Handlers that must be globally exactly-once
should dedupe on {run_id, step_slug}; PgFlow.LiveClient does this
structurally by treating step:skipped as an idempotent state transition
rather than a counter. Closing the gap properly means having the SQL
return newly transitioned rows, which is a core pgflow change.
@spec emit_step_skipped(%{ flow_slug: String.t(), run_id: String.t(), step_slug: String.t(), skip_reason: String.t() | nil }) :: :ok
Emits [:pgflow, :step, :skipped] when a step is skipped without a worker.
Metadata
:flow_slug- Flow identifier:run_id- Run UUID:step_slug- Skipped step identifier:skip_reason- Why the step was skipped, ornil