PgFlow.Queries.Workers (PgFlow v0.4.0)

Copy Markdown View Source

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 an edge/worker function for monitoring by pgflow.ensure_workers().

Functions

heartbeat_worker(repo, worker_id)

@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.

mark_worker_stopped(repo, worker_id)

@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 repository
  • worker_id - The worker identifier (UUID string)

Returns

  • {:ok, nil} - Success
  • {:error, reason} - Error details if the operation fails

register_worker(repo, worker_id, queue_name, function_name)

@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 repository
  • worker_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

track_worker_function(repo, function_name, start_mode)

@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.