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
8 changes: 8 additions & 0 deletions consumers-metrics-prometheus/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,2 +1,10 @@
# consumers-metrics-prometheus-1.1.0.0 (2026-??-??)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note: TODO release date

* Add `consumers_job_failed_attempts`, a histogram of the number of
processing attempts made so far by jobs whose execution didn't succeed
(`Failed`, an exception, or an abort), by `job_name`. Distinguishes
one-off failures from jobs that keep failing. Configurable via the new
`jobFailedAttemptsBuckets` field on `ConsumerMetricsConfig`.
* Requires `consumers` >= 2.4.0.0, for `ccJobAttempts`.

# consumers-metrics-prometheus-1.0.0.0 (2025-03-03)
* Initial release.
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
cabal-version: 3.0
name: consumers-metrics-prometheus
version: 1.0.0.0
version: 1.1.0.0
synopsis: Prometheus metrics for the consumers library

description: Provides seamless instrumentation of your existing
Expand Down Expand Up @@ -43,7 +43,7 @@ library
exposed-modules: Database.PostgreSQL.Consumers.Instrumented

build-depends: base >= 4.16 && < 5
, consumers >= 2.3
, consumers >= 2.4
, exceptions >= 0.10
, hpqtypes >= 1.13
, lifted-base >= 0.2
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ data ConsumerMetricsConfig = ConsumerMetricsConfig
-- ^ Collection interval in seconds
, jobExecutionBuckets :: [Prom.Bucket]
-- ^ Buckets to use for the 'jobExecution' histogram
, jobFailedAttemptsBuckets :: [Prom.Bucket]
-- ^ Buckets to use for the 'jobsFailedAttempts' histogram
, collectDegradeThresholdSeconds :: Double
-- ^ @logAttention@ and graceful degrade if collection takes longer than @x@ seconds
, collectDegradeSeconds :: Int
Expand All @@ -43,6 +45,7 @@ data ConsumerMetricsConfig = ConsumerMetricsConfig
-- ConsumerMetricsConfig
-- { collectSeconds = 15
-- , jobExecutionBuckets = [0.01, 0.05, 0.1, 0.5, 1, 2, 4, 8, 16, 32, 64, 128, 256, 512]
-- , jobFailedAttemptsBuckets = [1, 2, 3, 5, 10, 20, 50]
-- , collectDegradeThresholdSeconds = 0.1
-- , collectDegradeSeconds = 60
-- }
Expand All @@ -52,6 +55,7 @@ defaultConsumerMetricsConfig =
ConsumerMetricsConfig
{ collectSeconds = 15
, jobExecutionBuckets = [0.01, 0.05, 0.1, 0.5, 1, 2, 4, 8, 16, 32, 64, 128, 256, 512]
, jobFailedAttemptsBuckets = [1, 2, 3, 5, 10, 20, 50]
, collectDegradeThresholdSeconds = 0.1
, collectDegradeSeconds = 60
}
Expand All @@ -62,6 +66,9 @@ defaultConsumerMetricsConfig =
-- # HELP consumers_job_execution_seconds Execution time of jobs in seconds, by job_name, includes the job_result
-- # TYPE consumers_job_execution_seconds histogram
--
-- # HELP consumers_job_failed_attempts Number of processing attempts made so far by jobs whose execution didn't succeed, by job_name
-- # TYPE consumers_job_failed_attempts histogram
--
-- # HELP consumers_jobs_reserved_total The total number of job reserved, by job_name
-- # TYPE consumers_jobs_reserved_total counter
--
Expand All @@ -79,6 +86,7 @@ data ConsumerMetrics = ConsumerMetrics
, jobsOverdue :: Prom.Vector Prom.Label1 Prom.Gauge
, jobsReserved :: Prom.Vector Prom.Label1 Prom.Counter
, jobsExecution :: Prom.Vector Prom.Label2 Prom.Histogram
, jobsFailedAttempts :: Prom.Vector Prom.Label1 Prom.Histogram
}

registerConsumerMetrics :: MonadBaseControl IO m => ConsumerMetricsConfig -> m ConsumerMetrics
Expand Down Expand Up @@ -116,6 +124,15 @@ registerConsumerMetrics ConsumerMetricsConfig {..} = liftBase $ do
, metricHelp = "Execution time of jobs in seconds, by job_name, includes the job_result"
}
jobExecutionBuckets
jobsFailedAttempts <-
Prom.register
. Prom.vector "job_name"
$ Prom.histogram
Prom.Info
{ metricName = "consumers_job_failed_attempts"
, metricHelp = "Number of processing attempts made so far by jobs whose execution didn't succeed, by job_name"
}
jobFailedAttemptsBuckets
pure $ ConsumerMetrics {..}

-- | Run a 'ConsumerConfig', but with instrumentation added.
Expand Down Expand Up @@ -250,9 +267,9 @@ instrumentConsumerConfig ConsumerMetrics {..} ConsumerConfig {..} =
-- result of the job).
ccProcessJob' job = do
handleAny handleEx . liftBase $ Prom.withLabel jobsReserved jobName Prom.incCounter
fst <$> generalBracket monotonicTime reportJob (const $ ccProcessJob job)
fst <$> generalBracket monotonicTime (reportJob job) (const $ ccProcessJob job)

reportJob t1 jobExit = handleAny handleEx $ do
reportJob job t1 jobExit = handleAny handleEx $ do
t2 <- monotonicTime
let duration = t2 - t1
resultLabel = case jobExit of
Expand All @@ -261,5 +278,15 @@ instrumentConsumerConfig ConsumerMetrics {..} ConsumerConfig {..} =
ExitCaseException _ -> "exception"
ExitCaseAbort -> "abort"
liftBase $ Prom.withLabel jobsExecution (jobName, resultLabel) (`Prom.observe` duration)
-- Only observe the attempt count for jobs that didn't succeed: for
-- 'Ok' results it carries no information (denoising one-off failures
-- from persistent retries is the whole point of this metric), and it
-- would otherwise dominate the histogram with a flood of "attempt 1"
-- observations.
case jobExit of
ExitCaseSuccess (Ok _) -> pure ()
_ ->
liftBase $
Prom.withLabel jobsFailedAttempts jobName (`Prom.observe` fromIntegral (ccJobAttempts job))

handleEx e = logAttention "Exception while instrumenting job" $ object ["exception" .= show e]
9 changes: 8 additions & 1 deletion consumers/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,4 +1,11 @@
# consumers-2.3.5.0 (2026-??-??)
# consumers-2.4.0.0 (2026-??-??)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note: TODO release date

* **Breaking:** add `ccJobAttempts :: job -> Int` to `ConsumerConfig`, a
selector for the job's current (consecutive-failure) attempt count. Needs
`ccJobSelectors`/`ccJobFetcher` to expose the `attempts` column.
* Add `Database.PostgreSQL.Consumers.RetryStrategy`, ready-made
`ccOnException` retry strategies (`constantBackoff`, `linearBackoff`,
`exponentialBackoff`, `exponentialBackoffWithJitter`) built on
`ccJobAttempts`.
* Add `hoistConsumer` to `Database.PostgreSQL.Consumers.Config`.

# consumers-2.3.4.0 (2025-11-27)
Expand Down
4 changes: 3 additions & 1 deletion consumers/consumers.cabal
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
cabal-version: 3.0
name: consumers
version: 2.3.5.0
version: 2.4.0.0
synopsis: Concurrent PostgreSQL data consumers

description: Library for setting up concurrent consumers of data
Expand Down Expand Up @@ -51,6 +51,7 @@ library
Database.PostgreSQL.Consumers.Config,
Database.PostgreSQL.Consumers.Consumer,
Database.PostgreSQL.Consumers.Components,
Database.PostgreSQL.Consumers.RetryStrategy,
Database.PostgreSQL.Consumers.Utils

build-depends: base >= 4.16 && < 5
Expand All @@ -64,6 +65,7 @@ library
, monad-control >= 1.0
, monad-time >= 0.4
, mtl >= 2.2
, random >= 1.2
, safe-exceptions >= 0.1.7
, stm >= 2.4
, text >= 1.2
Expand Down
22 changes: 15 additions & 7 deletions consumers/example/Example.hs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import Control.Monad.IO.Class
import Data.Int
import Data.Text qualified as T
import Database.PostgreSQL.Consumers
import Database.PostgreSQL.Consumers.RetryStrategy (exponentialBackoff)
import Database.PostgreSQL.PQTypes
import Database.PostgreSQL.PQTypes.Checks
import Database.PostgreSQL.PQTypes.Model
Expand Down Expand Up @@ -101,15 +102,16 @@ main = do
ConsumerConfig
{ ccJobsTable = "consumers_example_jobs"
, ccConsumersTable = "consumers_example_consumers"
, ccJobSelectors = ["id", "message"]
, ccJobSelectors = ["id", "attempts", "message"]
, ccJobFetcher = id
, ccJobIndex = \(i :: Int64, _msg :: T.Text) -> i
, ccJobIndex = \(i :: Int64, _attempts :: Int32, _msg :: T.Text) -> i
, ccJobAttempts = \(_i, attempts :: Int32, _msg) -> fromIntegral attempts
, ccNotificationChannel = Just "consumers_example_chan"
, ccNotificationTimeout = 10 * 1000000 -- 10 sec
, ccMaxRunningJobs = 1
, ccProcessJob = processJob
, ccOnException = handleException
, ccJobLogData = \(i, _) -> ["job_id" .= i]
, ccJobLogData = \(i, _, _) -> ["job_id" .= i]
}

-- Add a job to the consumer's queue.
Expand All @@ -124,16 +126,22 @@ main = do
commit

-- Invoked when a job is ready to be processed.
processJob :: (Int64, T.Text) -> AppM Result
processJob (_idx, msg) = do
processJob :: (Int64, Int32, T.Text) -> AppM Result
processJob (_idx, _attempts, msg) = do
logInfo_ msg
pure (Ok Remove)

-- Invoked when 'processJob' throws an exception. Can handle
-- failure in different ways, such as: remove the job from the
-- queue, mark it as processed, or schedule it for rerun.
handleException :: SomeException -> (Int64, T.Text) -> AppM Action
handleException _ _ = pure . RerunAfter $ imicroseconds 500000
--
-- This uses 'exponentialBackoff' from
-- 'Database.PostgreSQL.Consumers.RetryStrategy': give up (remove the
-- job) after 5 attempts, otherwise retry with a delay that doubles each
-- time, capped at 10 seconds.
handleException :: SomeException -> (Int64, Int32, T.Text) -> AppM Action
handleException _ (_idx, attempts, _msg) =
pure $ exponentialBackoff 5 (imicroseconds 500000) (iseconds 10) (fromIntegral attempts)

-- | Table where jobs are stored. See
-- 'Database.PostgreSQL.Consumers.Config.ConsumerConfig'.
Expand Down
9 changes: 9 additions & 0 deletions consumers/src/Database/PostgreSQL/Consumers/Config.hs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ module Database.PostgreSQL.Consumers.Config

import Control.Exception (SomeException)
import Data.Aeson.Types qualified as A
import Data.Int (Int32)
import Data.Time
import Database.PostgreSQL.PQTypes.FromRow
import Database.PostgreSQL.PQTypes.Interval
Expand Down Expand Up @@ -83,6 +84,14 @@ data ConsumerConfig m idx job = forall row. FromRow row => ConsumerConfig
-- ^ Function that transforms the list of fields into a job.
, ccJobIndex :: !(job -> idx)
-- ^ Selector for taking out job ID from the job object.
, ccJobAttempts :: !(job -> Int32)
-- ^ Selector for taking out the number of processing attempts made so far
-- (i.e. the job's @attempts@ column) from the job object. Needs
-- 'ccJobSelectors'/'ccJobFetcher' to expose it. This is the number of
-- consecutive failed attempts, including the current one: it's reset to 1
-- once a job has succeeded, so it's a streak of failures rather than a
-- lifetime total. See "Database.PostgreSQL.Consumers.RetryStrategy" for
-- ready-made 'ccOnException' handlers built on it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: seems like dear Claude made this a bit longer than it needs to be

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Try to use my rules wrt. code comments etc. from https://github.com/arybczak/claude/blob/master/CLAUDE.md and tell it to rewrite the text in this PR using current rules, that should help.

, ccNotificationChannel :: !(Maybe Channel)
-- ^ Notification channel used for listening for incoming jobs. Whenever the
-- consumer receives a notification, it checks the database for any pending
Expand Down
129 changes: 129 additions & 0 deletions consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
-- |
-- Ready-made retry strategies for 'ccOnException', built around
-- 'ccJobAttempts' (the number of consecutive failed processing attempts for
-- a job, including the current one).
--
-- Each strategy takes the maximum number of attempts to allow, some
-- delay-shaping parameters, and the job's current attempt count, and
-- returns the 'Action' to apply: 'Remove' once @maxAttempts@ has been
-- exceeded, otherwise 'RerunAfter' some delay. Wire one into
-- 'ccOnException' using 'ccJobAttempts' to get the attempt count out of the
-- job:
--
-- @
-- ccOnException = \\_ex job -> pure $
-- exponentialBackoff 10 (iseconds 1) (iminutes 10) (ccJobAttempts config job)
-- @
module Database.PostgreSQL.Consumers.RetryStrategy
( constantBackoff
, linearBackoff
, exponentialBackoff
, exponentialBackoffWithJitter
) where

import Control.Monad.IO.Class
import Data.Int (Int32)
import Data.Time
import Database.PostgreSQL.Consumers.Config
import Database.PostgreSQL.PQTypes.Interval
import System.Random (randomRIO)

-- | Retry after the same fixed delay every time, until @maxAttempts@ is
-- exceeded.
constantBackoff
:: Int
-- ^ Maximum number of attempts before the job is removed.
-> Interval
-- ^ Delay before every retry.
-> Int32
-- ^ The job's current attempt count (see 'ccJobAttempts').
-> Action
constantBackoff maxAttempts delay attempts
| fromIntegral attempts >= maxAttempts = Remove
| otherwise = RerunAfter delay

-- | Retry with a delay that grows linearly with the attempt count
-- (@delayUnit * attempts@), until @maxAttempts@ is exceeded.
linearBackoff
:: Int
-- ^ Maximum number of attempts before the job is removed.
-> Interval
-- ^ Delay unit; the delay before retry number @n@ is @n * delayUnit@.
-> Int32
-- ^ The job's current attempt count (see 'ccJobAttempts').
-> Action
linearBackoff maxAttempts delayUnit attempts
| fromIntegral attempts >= maxAttempts = Remove
| otherwise = RerunAfter $ scaleInterval (fromIntegral attempts) delayUnit

-- | Retry with a delay that doubles on every attempt
-- (@baseDelay * 2 ^ (attempts - 1)@), capped at @maxDelay@ so it doesn't
-- grow without bound, until @maxAttempts@ is exceeded.
exponentialBackoff
:: Int
-- ^ Maximum number of attempts before the job is removed.
-> Interval
-- ^ Base delay, used for the first retry.
-> Interval
-- ^ Delay cap; the computed delay never exceeds this.
-> Int32
-- ^ The job's current attempt count (see 'ccJobAttempts').
-> Action
exponentialBackoff maxAttempts baseDelay maxDelay attempts
| fromIntegral attempts >= maxAttempts = Remove
| otherwise = RerunAfter $ nextDelay baseDelay maxDelay (fromIntegral attempts)

-- | Like 'exponentialBackoff', but adds up to +/-50% random jitter to the
-- computed delay.
--
-- Useful when several consumer instances can end up retrying the same kind
-- of job at once, e.g. after a shared dependency (a downstream API, a
-- database) comes back up from an outage: without jitter, every instance
-- backs off on the same schedule and they all retry in lockstep, hitting
-- the recovering dependency again simultaneously.
exponentialBackoffWithJitter
:: MonadIO m
=> Int
-- ^ Maximum number of attempts before the job is removed.
-> Interval
-- ^ Base delay, used for the first retry (before jitter).
-> Interval
-- ^ Delay cap, applied before jitter is added.
-> Int32
-- ^ The job's current attempt count (see 'ccJobAttempts').
-> m Action
exponentialBackoffWithJitter maxAttempts baseDelay maxDelay attempts
| fromIntegral attempts >= maxAttempts = pure Remove
| otherwise = do
jitter <- liftIO $ randomRIO (0.5, 1.5 :: Double)
pure . RerunAfter . scaleIntervalD jitter $ nextDelay baseDelay maxDelay (fromIntegral attempts)

----------------------------------------

nextDelay :: Interval -> Interval -> Int -> Interval
nextDelay baseDelay maxDelay attempts = min' maxDelay $ scaleInterval (2 ^ (attempts - 1)) baseDelay
where
min' a b = if intervalToDiffTime a < intervalToDiffTime b then a else b

-- | Scale an 'Interval' by an integer factor.
scaleInterval :: Int -> Interval -> Interval
scaleInterval factor = diffTimeToInterval . (fromIntegral factor *) . intervalToDiffTime

-- | Scale an 'Interval' by a fractional factor.
scaleIntervalD :: Double -> Interval -> Interval
scaleIntervalD factor = diffTimeToInterval . (realToFrac factor *) . intervalToDiffTime

intervalToDiffTime :: Interval -> DiffTime
intervalToDiffTime Interval {..} = secondsToDiffTime seconds
where
seconds =
(toInteger intYears * 365 * 24 * 60 * 60)
+ (toInteger intMonths * 30 * 24 * 60 * 60)
+ (toInteger intDays * 24 * 60 * 60)
+ (toInteger intHours * 60 * 60)
+ (toInteger intMinutes * 60)
+ toInteger intSeconds
+ (toInteger intMicroseconds `div` 1000000)

diffTimeToInterval :: DiffTime -> Interval
diffTimeToInterval = imicroseconds . fromInteger . (`div` 1000000) . diffTimeToPicoseconds
Loading
Loading