SQL query interface for pgflow worker operations.
Provides functions for registering workers and managing their lifecycle.
Summary
Functions
Refreshes a worker heartbeat and reports whether it has been deprecated.
Marks a worker as stopped.
Registers a worker in the database.
Registers an edge/worker function for monitoring by pgflow.ensure_workers().
Functions
@spec heartbeat_worker(Ecto.Repo.t(), String.t()) :: {:ok, :alive | :deprecated} | {:error, term()}
Refreshes a worker heartbeat and reports whether it has been deprecated.
Missing registration is treated as deprecated, matching upstream edge-worker behavior when a worker row no longer exists.
@spec mark_worker_stopped(Ecto.Repo.t(), String.t()) :: {:ok, nil} | {:error, term()}
Marks a worker as stopped.
Sets the stopped_at timestamp for graceful shutdown signaling.
Parameters
repo- The Ecto repositoryworker_id- The worker identifier (UUID string)
Returns
{:ok, nil}- Success{:error, reason}- Error details if the operation fails
@spec register_worker(Ecto.Repo.t(), String.t(), String.t(), String.t()) :: {:ok, nil} | {:error, term()}
Registers a worker in the database.
Creates a new worker record or updates the heartbeat if the worker already exists.
Parameters
repo- The Ecto repositoryworker_id- The worker identifier (UUID string)queue_name- The queue name (flow_slug)function_name- The function name (e.g., "elixir:MyApp.Flows.MyFlow")
Returns
{:ok, nil}- Success{:error, reason}- Error details if the operation fails
@spec track_worker_function(Ecto.Repo.t(), String.t(), String.t()) :: {:ok, nil} | {:error, term()}
Registers an edge/worker function for monitoring by pgflow.ensure_workers().
Elixir workers call this on startup with "process" start mode after flow
compilation succeeds.