diff --git a/consumers-metrics-prometheus/CHANGELOG.md b/consumers-metrics-prometheus/CHANGELOG.md index 2e4eec1..f77b044 100644 --- a/consumers-metrics-prometheus/CHANGELOG.md +++ b/consumers-metrics-prometheus/CHANGELOG.md @@ -1,2 +1,10 @@ +# consumers-metrics-prometheus-1.1.0.0 (2026-??-??) +* 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. diff --git a/consumers-metrics-prometheus/consumers-metrics-prometheus.cabal b/consumers-metrics-prometheus/consumers-metrics-prometheus.cabal index 60100a7..a1d5c3b 100644 --- a/consumers-metrics-prometheus/consumers-metrics-prometheus.cabal +++ b/consumers-metrics-prometheus/consumers-metrics-prometheus.cabal @@ -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 @@ -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 diff --git a/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs b/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs index 3a6a357..5596fa1 100644 --- a/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs +++ b/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs @@ -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 @@ -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 -- } @@ -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 } @@ -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 -- @@ -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 @@ -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. @@ -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 @@ -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] diff --git a/consumers/CHANGELOG.md b/consumers/CHANGELOG.md index 0de28c3..e8624b3 100644 --- a/consumers/CHANGELOG.md +++ b/consumers/CHANGELOG.md @@ -1,4 +1,11 @@ -# consumers-2.3.5.0 (2026-??-??) +# consumers-2.4.0.0 (2026-??-??) +* **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) diff --git a/consumers/consumers.cabal b/consumers/consumers.cabal index c67003c..cd1d4a1 100644 --- a/consumers/consumers.cabal +++ b/consumers/consumers.cabal @@ -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 @@ -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 @@ -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 diff --git a/consumers/example/Example.hs b/consumers/example/Example.hs index c25af82..5f8082b 100644 --- a/consumers/example/Example.hs +++ b/consumers/example/Example.hs @@ -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 @@ -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. @@ -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'. diff --git a/consumers/src/Database/PostgreSQL/Consumers/Config.hs b/consumers/src/Database/PostgreSQL/Consumers/Config.hs index b759594..108dc98 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/Config.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/Config.hs @@ -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 @@ -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. , ccNotificationChannel :: !(Maybe Channel) -- ^ Notification channel used for listening for incoming jobs. Whenever the -- consumer receives a notification, it checks the database for any pending diff --git a/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs b/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs new file mode 100644 index 0000000..029e44c --- /dev/null +++ b/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs @@ -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 diff --git a/consumers/test/Test.hs b/consumers/test/Test.hs index 1c619ac..a45ae3d 100644 --- a/consumers/test/Test.hs +++ b/consumers/test/Test.hs @@ -13,6 +13,7 @@ import Data.Int import Data.Text qualified as T import Data.Time import Database.PostgreSQL.Consumers +import Database.PostgreSQL.Consumers.RetryStrategy (constantBackoff) import Database.PostgreSQL.PQTypes import Database.PostgreSQL.PQTypes.Checks import Database.PostgreSQL.PQTypes.Model @@ -144,16 +145,17 @@ test = do ConsumerConfig { ccJobsTable = "consumers_test_jobs" , ccConsumersTable = "consumers_test_consumers" - , ccJobSelectors = ["id", "countdown"] + , ccJobSelectors = ["id", "attempts", "countdown"] , ccJobFetcher = id - , ccJobIndex = \(i :: Int64, _ :: Int32) -> i + , ccJobIndex = \(i :: Int64, _attempts :: Int32, _countdown :: Int32) -> i + , ccJobAttempts = \(_i, attempts :: Int32, _countdown) -> fromIntegral attempts , ccNotificationChannel = Just "consumers_test_chan" , -- select some small timeout ccNotificationTimeout = 100 * 1000 -- 100 msec , ccMaxRunningJobs = 20 , ccProcessJob = processJob , ccOnException = handleException - , ccJobLogData = \(i, _) -> ["job_id" .= i] + , ccJobLogData = \(i, _, _) -> ["job_id" .= i] } putJob :: Int32 -> TestEnv () @@ -167,16 +169,17 @@ test = do <> ")" notify "consumers_test_chan" "" - processJob :: (Int64, Int32) -> TestEnv Result - processJob (_idx, countdown) = do + processJob :: (Int64, Int32, Int32) -> TestEnv Result + processJob (_idx, _attempts, countdown) = do when (countdown > 0) $ do putJob (countdown - 1) putJob (countdown - 1) commit pure (Ok Remove) - handleException :: SomeException -> (Int64, Int32) -> TestEnv Action - handleException _ _ = pure . RerunAfter $ imicroseconds 500000 + handleException :: SomeException -> (Int64, Int32, Int32) -> TestEnv Action + handleException _ (_idx, attempts, _countdown) = + pure $ constantBackoff 5 (imicroseconds 500000) (fromIntegral attempts) jobsTable :: Table jobsTable = diff --git a/postgres-15.patch b/postgres-15.patch new file mode 100644 index 0000000..e8d8cee --- /dev/null +++ b/postgres-15.patch @@ -0,0 +1,11 @@ +--- a/.github/workflows/haskell-ci.yml ++++ b/.github/workflows/haskell-ci.yml +@@ -33,7 +33,7 @@ + container: + image: buildpack-deps:jammy + services: + postgres: +- image: postgres:14 ++ image: postgres:15 + env: + POSTGRES_PASSWORD: postgres