Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 33 additions & 12 deletions lib/prom_ex/plugins/oban.ex
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,14 @@ if Code.ensure_loaded?(Oban) do
- `duration_unit`: This is an OPTIONAL option and is a `Telemetry.Metrics.time_unit()`. It can be one of:
`:second | :millisecond | :microsecond | :nanosecond`. It is `:millisecond` by default.

- `job_attempt_buckets`: OPTIONAL. Buckets for job attempt distributions. Defaults to `[1, 5, 10]`.

- `job_duration_buckets`: OPTIONAL. Buckets for job duration and queue time distributions
(in the configured `duration_unit`). Defaults to `[10, 100, 500, 1_000, 5_000, 20_000]`.

- `producer_duration_buckets`: OPTIONAL. Buckets for producer duration distributions
(in the configured `duration_unit`). Defaults to `[10, 100, 500, 1_000, 5_000, 10_000]`.

This plugin exposes the following metric groups:
- `:oban_init_event_metrics`
- `:oban_job_event_metrics`
Expand Down Expand Up @@ -74,6 +82,9 @@ if Code.ensure_loaded?(Oban) do
otp_app = Keyword.fetch!(opts, :otp_app)
metric_prefix = Keyword.get(opts, :metric_prefix, PromEx.metric_prefix(otp_app, :oban))
duration_unit = Keyword.get(opts, :duration_unit, :millisecond)
job_attempt_buckets = Keyword.get(opts, :job_attempt_buckets, [1, 5, 10])
job_duration_buckets = Keyword.get(opts, :job_duration_buckets, [10, 100, 500, 1_000, 5_000, 20_000])
producer_duration_buckets = Keyword.get(opts, :producer_duration_buckets, [10, 100, 500, 1_000, 5_000, 10_000])

oban_supervisors = get_oban_supervisors(opts)
keep_function_filter = keep_oban_instance_metrics(oban_supervisors)
Expand All @@ -83,8 +94,14 @@ if Code.ensure_loaded?(Oban) do

[
oban_supervisor_init_event_metrics(metric_prefix, keep_function_filter, duration_unit),
oban_job_event_metrics(metric_prefix, keep_function_filter, duration_unit),
oban_producer_event_metrics(metric_prefix, keep_function_filter, duration_unit),
oban_job_event_metrics(
metric_prefix,
keep_function_filter,
duration_unit,
job_attempt_buckets,
job_duration_buckets
),
oban_producer_event_metrics(metric_prefix, keep_function_filter, duration_unit, producer_duration_buckets),
oban_circuit_breaker_event_metrics(metric_prefix, keep_function_filter)
]
end
Expand Down Expand Up @@ -186,9 +203,13 @@ if Code.ensure_loaded?(Oban) do
}
end

defp oban_job_event_metrics(metric_prefix, keep_function_filter, duration_unit) do
job_attempt_buckets = [1, 5, 10]
job_duration_buckets = [10, 100, 500, 1_000, 5_000, 20_000]
defp oban_job_event_metrics(
metric_prefix,
keep_function_filter,
duration_unit,
job_attempt_buckets,
job_duration_buckets
) do
duration_unit_plural = Utils.make_plural_atom(duration_unit)

Event.build(
Expand Down Expand Up @@ -279,7 +300,7 @@ if Code.ensure_loaded?(Oban) do
)
end

defp oban_producer_event_metrics(metric_prefix, keep_function_filter, duration_unit) do
defp oban_producer_event_metrics(metric_prefix, keep_function_filter, duration_unit, producer_duration_buckets) do
duration_unit_plural = Utils.make_plural_atom(duration_unit)

Event.build(
Expand All @@ -291,7 +312,7 @@ if Code.ensure_loaded?(Oban) do
measurement: :duration,
description: "How long it took to dispatch the job.",
reporter_options: [
buckets: [10, 100, 500, 1_000, 5_000, 10_000]
buckets: producer_duration_buckets
],
unit: {:native, duration_unit},
tag_values: &producer_tag_values/1,
Expand All @@ -318,7 +339,7 @@ if Code.ensure_loaded?(Oban) do
measurement: :duration,
description: "How long it took for the producer to raise an exception.",
reporter_options: [
buckets: [10, 100, 500, 1_000, 5_000, 10_000]
buckets: producer_duration_buckets
],
unit: {:native, duration_unit},
tag_values: &producer_tag_values/1,
Expand Down Expand Up @@ -438,7 +459,7 @@ if Code.ensure_loaded?(Oban) do

config
|> Oban.Repo.all(query)
|> include_zeros_for_missing_queue_states()
|> include_zeros_for_missing_queue_states(config)
|> Enum.each(fn {{queue, state}, count} ->
measurements = %{count: count}
metadata = %{name: normalize_module_name(oban_supervisor), queue: queue, state: state}
Expand All @@ -447,10 +468,10 @@ if Code.ensure_loaded?(Oban) do
end)
end

defp include_zeros_for_missing_queue_states(query_result) do
defp include_zeros_for_missing_queue_states(query_result, config) do
{_, opts} =
Oban.config().plugins
|> Enum.find({nil, [queues: Oban.config().queues]}, fn {plugin, _} ->
config.plugins
|> Enum.find({nil, [queues: config.queues]}, fn {plugin, _} ->
plugin == Oban.Pro.Plugins.DynamicQueues
end)

Expand Down
Loading