From 181d930ba4f7d09b7d3e14a5003423805115f1f7 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 22 Sep 2026 12:49:02 +0000 Subject: [PATCH 1/6] Revive PR #5: wrap consumer jobs in a Job idx info record Rebase Jonathan's 2021 proof of concept (github.com/scrive/consumers/pull/5) onto current master. ConsumerConfig's job type is now always Job idx info, exposing jobIndex, jobRunAt, jobFinishedAt and jobAttempts (the queue bookkeeping columns already maintained by reserveJobs) alongside the caller-supplied payload as jobInfo, instead of leaving attempts opaque to everything outside the reservation query. Update ccJobFetcher/ccProcessJob/ccOnException/ccJobLogData signatures accordingly, and add exponentialBackoff, a ready-made ccOnException handler built on jobAttempts. Update the example and test consumers to select the new columns and construct Job values; fix a crash in the ported diffTimeToInterval (the original used `^` with a negative Integer exponent). This only touches the consumers package; consumers-metrics-prometheus is deliberately left as-is for now, since it doesn't need any changes to keep compiling (it only threads job values through opaquely). Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_016NeDfW7w3TwyS1UM82uvkL --- consumers/CHANGELOG.md | 11 ++- consumers/consumers.cabal | 2 +- consumers/example/Example.hs | 23 ++++-- .../PostgreSQL/Consumers/Components.hs | 22 +++--- .../Database/PostgreSQL/Consumers/Config.hs | 76 ++++++++++++++++--- .../Database/PostgreSQL/Consumers/Consumer.hs | 4 +- consumers/test/Test.hs | 22 ++++-- 7 files changed, 122 insertions(+), 38 deletions(-) diff --git a/consumers/CHANGELOG.md b/consumers/CHANGELOG.md index 0de28c3..1aa4ded 100644 --- a/consumers/CHANGELOG.md +++ b/consumers/CHANGELOG.md @@ -1,4 +1,13 @@ -# consumers-2.3.5.0 (2026-??-??) +# consumers-2.4.0.0 (2026-??-??) +* **Breaking:** `ConsumerConfig`'s job parameter is now always wrapped in the + new `Job` type, which exposes `jobIndex`, `jobRunAt`, `jobFinishedAt` and + `jobAttempts` (the queue bookkeeping columns) alongside the caller-supplied + payload as `jobInfo`. `ccJobFetcher`, `ccJobIndex`, `ccProcessJob`, + `ccOnException` and `ccJobLogData` all change shape accordingly; existing + consumers need to select `run_at`, `finished_at` and `attempts` in + `ccJobSelectors` and construct a `Job` in `ccJobFetcher`. +* Add `exponentialBackoff`, a ready-made `ccOnException` handler built on + `jobAttempts`. * 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 b9a3e96..e2db5df 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 diff --git a/consumers/example/Example.hs b/consumers/example/Example.hs index 47f563a..bab8707 100644 --- a/consumers/example/Example.hs +++ b/consumers/example/Example.hs @@ -13,6 +13,7 @@ import Control.Monad import Control.Monad.IO.Class import Data.Int import Data.Text qualified as T +import Data.Time (UTCTime) import Database.PostgreSQL.Consumers import Database.PostgreSQL.PQTypes import Database.PostgreSQL.PQTypes.Checks @@ -101,15 +102,23 @@ main = do ConsumerConfig { ccJobsTable = "consumers_example_jobs" , ccConsumersTable = "consumers_example_consumers" - , ccJobSelectors = ["id", "message"] - , ccJobFetcher = id - , ccJobIndex = \(i :: Int64, _msg :: T.Text) -> i + , ccJobSelectors = ["id", "run_at", "finished_at", "attempts", "message"] + , ccJobFetcher = + \(i :: Int64, runAt :: Maybe UTCTime, finishedAt :: Maybe UTCTime, attempts :: Int32, msg :: T.Text) -> + Job + { jobIndex = i + , jobRunAt = runAt + , jobFinishedAt = finishedAt + , jobAttempts = fromIntegral attempts + , jobInfo = msg + } + , ccJobIndex = jobIndex , ccNotificationChannel = Just "consumers_example_chan" , ccNotificationTimeout = 10 * 1000000 -- 10 sec , ccMaxRunningJobs = 1 , ccProcessJob = processJob , ccOnException = handleException - , ccJobLogData = \(i, _) -> ["job_id" .= i] + , ccJobLogData = \job -> ["job_id" .= jobIndex job] } -- Add a job to the consumer's queue. @@ -124,15 +133,15 @@ main = do commit -- Invoked when a job is ready to be processed. - processJob :: (Int64, T.Text) -> AppM Result - processJob (_idx, msg) = do + processJob :: Job Int64 T.Text -> AppM Result + processJob Job {jobInfo = 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 :: SomeException -> Job Int64 T.Text -> AppM Action handleException _ _ = pure . RerunAfter $ imicroseconds 500000 -- | Table where jobs are stored. See diff --git a/consumers/src/Database/PostgreSQL/Consumers/Components.hs b/consumers/src/Database/PostgreSQL/Consumers/Components.hs index 5302fad..155dee8 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/Components.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/Components.hs @@ -46,7 +46,7 @@ runConsumer , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx job + => ConsumerConfig m idx info -- ^ The consumer. -> ConnectionSourceM m -> m (m ()) @@ -62,7 +62,7 @@ runConsumerWithIdleSignal , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx job + => ConsumerConfig m idx info -- ^ The consumer. -> ConnectionSourceM m -> TMVar Bool @@ -81,7 +81,7 @@ runConsumerWithMaybeIdleSignal , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx job + => ConsumerConfig m idx info -> ConnectionSourceM m -> Maybe (TMVar Bool) -> m (m ()) @@ -182,7 +182,7 @@ runConsumerWithMaybeIdleSignal cc0 cs mIdleSignal -- database for incoming jobs. spawnListener :: (MonadBaseControl IO m, MonadMask m) - => ConsumerConfig m idx job + => ConsumerConfig m idx info -> ConnectionSourceM m -> MVar () -> m ThreadId @@ -211,7 +211,7 @@ spawnListener cc cs semaphore = -- | Spawn a thread that monitors working consumers for activity and -- periodically updates its own. spawnMonitor - :: forall m idx job + :: forall m idx info . ( MonadBaseControl IO m , MonadLog m , MonadMask m @@ -220,7 +220,7 @@ spawnMonitor , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx job + => ConsumerConfig m idx info -> ConnectionSourceM m -> ConsumerID -> m ThreadId @@ -294,7 +294,7 @@ spawnMonitor ConsumerConfig {..} cs cid = forkP "monitor" . forever $ do -- | Spawn a thread that reserves and processes jobs. spawnDispatcher - :: forall m idx job + :: forall m idx info . ( MonadBaseControl IO m , MonadLog m , MonadMask m @@ -302,7 +302,7 @@ spawnDispatcher , Show idx , ToSQL idx ) - => ConsumerConfig m idx job + => ConsumerConfig m idx info -> ConnectionSourceM m -> ConsumerID -> MVar () @@ -357,7 +357,7 @@ spawnDispatcher ConsumerConfig {..} cs cid semaphore runningJobsInfo runningJobs pure (batchSize > 0) - reserveJobs :: Int -> m ([job], Int) + reserveJobs :: Int -> m ([Job idx info], Int) reserveJobs limit = runDBT cs ts $ do now <- currentTime n <- @@ -389,7 +389,7 @@ spawnDispatcher ConsumerConfig {..} cs cid semaphore runningJobsInfo runningJobs ] -- Spawn each job in a separate thread. - startJob :: job -> m (job, m (T.Result Result)) + startJob :: Job idx info -> m (Job idx info, m (T.Result Result)) startJob job = do (_, joinFork) <- mask $ \restore -> T.fork $ do tid <- myThreadId @@ -403,7 +403,7 @@ spawnDispatcher ConsumerConfig {..} cs cid semaphore runningJobsInfo runningJobs modifyTVar' runningJobsInfo $ M.delete tid -- Wait for all the jobs and collect their results. - joinJob :: (job, m (T.Result Result)) -> m (idx, Result) + joinJob :: (Job idx info, m (T.Result Result)) -> m (idx, Result) joinJob (job, joinFork) = joinFork >>= \case Right result -> pure (ccJobIndex job, result) diff --git a/consumers/src/Database/PostgreSQL/Consumers/Config.hs b/consumers/src/Database/PostgreSQL/Consumers/Config.hs index b759594..9d9d1d4 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/Config.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/Config.hs @@ -1,8 +1,10 @@ module Database.PostgreSQL.Consumers.Config ( Action (..) , Result (..) + , Job (..) , ConsumerConfig (..) , hoistConsumer + , exponentialBackoff ) where import Control.Exception (SomeException) @@ -26,8 +28,21 @@ data Action data Result = Ok Action | Failed Action deriving (Eq, Ord, Show) +-- | A job fetched from the queue, together with its queue bookkeeping +-- fields (as opposed to just the caller-supplied 'jobInfo' payload). +data Job idx info = Job + { jobIndex :: !idx + , jobRunAt :: !(Maybe UTCTime) + , jobFinishedAt :: !(Maybe UTCTime) + , jobAttempts :: !Int + -- ^ Number of consecutive failed processing attempts, including the + -- current one. Reset to 1 once a job has succeeded, so this is a streak + -- of failures rather than a lifetime total. + , jobInfo :: !info + } + -- | Config of a consumer. -data ConsumerConfig m idx job = forall row. FromRow row => ConsumerConfig +data ConsumerConfig m idx info = forall row. FromRow row => ConsumerConfig { ccJobsTable :: !(RawSQL ()) -- ^ Name of the database table where jobs are stored. The table needs to have -- the following columns in order to be suitable for acting as a job queue: @@ -78,10 +93,12 @@ data ConsumerConfig m idx job = forall row. FromRow row => ConsumerConfig -- and these jobs stay locked forever, yet are never processed. , ccJobSelectors :: ![SQL] -- ^ Fields needed to be selected from the jobs table in order to assemble a - -- job. - , ccJobFetcher :: !(row -> job) + -- job. Needs to match what 'ccJobFetcher' expects, and typically includes + -- @id@, @run_at@, @finished_at@ and @attempts@ alongside whatever columns + -- make up the job's 'jobInfo' payload. + , ccJobFetcher :: !(row -> Job idx info) -- ^ Function that transforms the list of fields into a job. - , ccJobIndex :: !(job -> idx) + , ccJobIndex :: !(Job idx info -> idx) -- ^ Selector for taking out job ID from the job object. , ccNotificationChannel :: !(Maybe Channel) -- ^ Notification channel used for listening for incoming jobs. Whenever the @@ -107,26 +124,67 @@ data ConsumerConfig m idx job = forall row. FromRow row => ConsumerConfig -- it needs to be a positive number. , ccMaxRunningJobs :: !Int -- ^ Maximum amount of jobs that can be processed in parallel. - , ccProcessJob :: !(job -> m Result) + , ccProcessJob :: !(Job idx info -> m Result) -- ^ Function that processes a job. It's recommended to process each job in a -- separate DB transaction, otherwise you'll have to remember to commit your -- changes to the database manually. - , ccOnException :: !(SomeException -> job -> m Action) + , ccOnException :: !(SomeException -> Job idx info -> m Action) -- ^ Action taken if a job processing function throws an exception. For -- robustness it's best to ensure that it doesn't throw. If it does, the -- exception will be logged and the job in question postponed by a day. - , ccJobLogData :: !(job -> [A.Pair]) + , ccJobLogData :: !(Job idx info -> [A.Pair]) -- ^ Data to attach to each log message while processing a job. } -- | Change the monad the consumer lives in. hoistConsumer :: (forall r. m r -> n r) - -> ConsumerConfig m idx job - -> ConsumerConfig n idx job + -> ConsumerConfig m idx info + -> ConsumerConfig n idx info hoistConsumer f ConsumerConfig {..} = ConsumerConfig { ccProcessJob = f . ccProcessJob , ccOnException = \ex -> f . ccOnException ex , .. } + +-- | A ready-made 'ccOnException' handler: reruns a job with exponentially +-- increasing delay, then gives up (removes the job) once 'jobAttempts' +-- exceeds @maxAttempts@. +-- +-- /Note:/ this mirrors the growth curve from the original proof of concept +-- (delay ^ attempts), so it only grows the delay if @startInterval@ is +-- greater than one second; below that it shrinks towards zero instead. +-- Pick @startInterval@ accordingly, or write a custom handler if you need a +-- more conventional @startInterval * base ^ attempts@ curve. +exponentialBackoff + :: Monad m + => Int + -- ^ Maximum number of attempts before the job is removed. + -> Interval + -- ^ Delay before the first retry. + -> SomeException + -> Job idx info + -> m Action +exponentialBackoff maxAttempts startInterval _ Job {..} = + pure $ case (jobAttempts > maxAttempts, jobAttempts > 2) of + (True, _) -> Remove + (False, False) -> RerunAfter startInterval + (False, True) -> + RerunAfter . diffTimeToInterval $ + intervalToDiffTime startInterval ^ jobAttempts + +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/src/Database/PostgreSQL/Consumers/Consumer.hs b/consumers/src/Database/PostgreSQL/Consumers/Consumer.hs index a836db9..5e568e9 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/Consumer.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/Consumer.hs @@ -33,7 +33,7 @@ instance Show ConsumerID where -- acquired ID. registerConsumer :: (MonadBase IO m, MonadMask m, MonadTime m) - => ConsumerConfig n idx job + => ConsumerConfig n idx info -> ConnectionSourceM m -> m ConsumerID registerConsumer ConsumerConfig {..} cs = runDBT cs defaultTransactionSettings $ do @@ -49,7 +49,7 @@ registerConsumer ConsumerConfig {..} cs = runDBT cs defaultTransactionSettings $ -- | Unregister consumer with a given ID. unregisterConsumer :: (MonadBase IO m, MonadMask m) - => ConsumerConfig n idx job + => ConsumerConfig n idx info -> ConnectionSourceM m -> ConsumerID -> m () diff --git a/consumers/test/Test.hs b/consumers/test/Test.hs index fd2b61d..9677d51 100644 --- a/consumers/test/Test.hs +++ b/consumers/test/Test.hs @@ -144,16 +144,24 @@ test = do ConsumerConfig { ccJobsTable = "consumers_test_jobs" , ccConsumersTable = "consumers_test_consumers" - , ccJobSelectors = ["id", "countdown"] - , ccJobFetcher = id - , ccJobIndex = \(i :: Int64, _ :: Int32) -> i + , ccJobSelectors = ["id", "run_at", "finished_at", "attempts", "countdown"] + , ccJobFetcher = + \(i :: Int64, runAt :: Maybe UTCTime, finishedAt :: Maybe UTCTime, attempts :: Int32, countdown :: Int32) -> + Job + { jobIndex = i + , jobRunAt = runAt + , jobFinishedAt = finishedAt + , jobAttempts = fromIntegral attempts + , jobInfo = countdown + } + , ccJobIndex = jobIndex , 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 = \job -> ["job_id" .= jobIndex job] } putJob :: Int32 -> TestEnv () @@ -167,15 +175,15 @@ test = do <> ")" notify "consumers_test_chan" "" - processJob :: (Int64, Int32) -> TestEnv Result - processJob (_idx, countdown) = do + processJob :: Job Int64 Int32 -> TestEnv Result + processJob Job {jobInfo = countdown} = do when (countdown > 0) $ do putJob (countdown - 1) putJob (countdown - 1) commit pure (Ok Remove) - handleException :: SomeException -> (Int64, Int32) -> TestEnv Action + handleException :: SomeException -> Job Int64 Int32 -> TestEnv Action handleException _ _ = pure . RerunAfter $ imicroseconds 500000 jobsTable :: Table From 53b08a73f937c648cd55b2f5c19d435cf0f2cfa0 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 22 Sep 2026 15:07:45 +0000 Subject: [PATCH 2/6] Add ccJobAttempts and RetryStrategy module, supersede Job wrapper attempt Revert the previous Job idx info restructuring (63e226f) in favour of a much smaller approach worked out with a colleague: add a single new selector ccJobAttempts :: job -> Int to ConsumerConfig, mirroring ccJobIndex. Existing job types don't change shape at all, only need to also select the attempts column and expose it via this new field. Add Database.PostgreSQL.Consumers.RetryStrategy as a separate module with ready-made ccOnException handlers built on ccJobAttempts: - constantBackoff: fixed delay every retry - linearBackoff: delay grows linearly with attempt count - exponentialBackoff: delay doubles each attempt, capped at a max delay - exponentialBackoffWithJitter: like exponentialBackoff, but with random jitter to avoid multiple consumer instances retrying in lockstep after a shared dependency recovers Update the example and test consumers to select attempts, wire up ccJobAttempts, and use exponentialBackoff/constantBackoff respectively. consumers-metrics-prometheus is still untouched; ccJobAttempts is what a future failed-attempts histogram there would read from. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_016NeDfW7w3TwyS1UM82uvkL --- consumers/CHANGELOG.md | 16 +-- consumers/consumers.cabal | 2 + consumers/example/Example.hs | 33 +++-- .../PostgreSQL/Consumers/Components.hs | 22 +-- .../Database/PostgreSQL/Consumers/Config.hs | 84 +++--------- .../Database/PostgreSQL/Consumers/Consumer.hs | 4 +- .../PostgreSQL/Consumers/RetryStrategy.hs | 128 ++++++++++++++++++ consumers/test/Test.hs | 27 ++-- 8 files changed, 194 insertions(+), 122 deletions(-) create mode 100644 consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs diff --git a/consumers/CHANGELOG.md b/consumers/CHANGELOG.md index 1aa4ded..e8624b3 100644 --- a/consumers/CHANGELOG.md +++ b/consumers/CHANGELOG.md @@ -1,13 +1,11 @@ # consumers-2.4.0.0 (2026-??-??) -* **Breaking:** `ConsumerConfig`'s job parameter is now always wrapped in the - new `Job` type, which exposes `jobIndex`, `jobRunAt`, `jobFinishedAt` and - `jobAttempts` (the queue bookkeeping columns) alongside the caller-supplied - payload as `jobInfo`. `ccJobFetcher`, `ccJobIndex`, `ccProcessJob`, - `ccOnException` and `ccJobLogData` all change shape accordingly; existing - consumers need to select `run_at`, `finished_at` and `attempts` in - `ccJobSelectors` and construct a `Job` in `ccJobFetcher`. -* Add `exponentialBackoff`, a ready-made `ccOnException` handler built on - `jobAttempts`. +* **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 e2db5df..82027af 100644 --- a/consumers/consumers.cabal +++ b/consumers/consumers.cabal @@ -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 bab8707..edff381 100644 --- a/consumers/example/Example.hs +++ b/consumers/example/Example.hs @@ -13,8 +13,8 @@ import Control.Monad import Control.Monad.IO.Class import Data.Int import Data.Text qualified as T -import Data.Time (UTCTime) import Database.PostgreSQL.Consumers +import Database.PostgreSQL.Consumers.RetryStrategy (exponentialBackoff) import Database.PostgreSQL.PQTypes import Database.PostgreSQL.PQTypes.Checks import Database.PostgreSQL.PQTypes.Model @@ -102,23 +102,16 @@ main = do ConsumerConfig { ccJobsTable = "consumers_example_jobs" , ccConsumersTable = "consumers_example_consumers" - , ccJobSelectors = ["id", "run_at", "finished_at", "attempts", "message"] - , ccJobFetcher = - \(i :: Int64, runAt :: Maybe UTCTime, finishedAt :: Maybe UTCTime, attempts :: Int32, msg :: T.Text) -> - Job - { jobIndex = i - , jobRunAt = runAt - , jobFinishedAt = finishedAt - , jobAttempts = fromIntegral attempts - , jobInfo = msg - } - , ccJobIndex = jobIndex + , ccJobSelectors = ["id", "attempts", "message"] + , ccJobFetcher = id + , 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 = \job -> ["job_id" .= jobIndex job] + , ccJobLogData = \(i, _, _) -> ["job_id" .= i] } -- Add a job to the consumer's queue. @@ -133,16 +126,22 @@ main = do commit -- Invoked when a job is ready to be processed. - processJob :: Job Int64 T.Text -> AppM Result - processJob Job {jobInfo = 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 -> Job 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/Components.hs b/consumers/src/Database/PostgreSQL/Consumers/Components.hs index 155dee8..5302fad 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/Components.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/Components.hs @@ -46,7 +46,7 @@ runConsumer , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx info + => ConsumerConfig m idx job -- ^ The consumer. -> ConnectionSourceM m -> m (m ()) @@ -62,7 +62,7 @@ runConsumerWithIdleSignal , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx info + => ConsumerConfig m idx job -- ^ The consumer. -> ConnectionSourceM m -> TMVar Bool @@ -81,7 +81,7 @@ runConsumerWithMaybeIdleSignal , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx info + => ConsumerConfig m idx job -> ConnectionSourceM m -> Maybe (TMVar Bool) -> m (m ()) @@ -182,7 +182,7 @@ runConsumerWithMaybeIdleSignal cc0 cs mIdleSignal -- database for incoming jobs. spawnListener :: (MonadBaseControl IO m, MonadMask m) - => ConsumerConfig m idx info + => ConsumerConfig m idx job -> ConnectionSourceM m -> MVar () -> m ThreadId @@ -211,7 +211,7 @@ spawnListener cc cs semaphore = -- | Spawn a thread that monitors working consumers for activity and -- periodically updates its own. spawnMonitor - :: forall m idx info + :: forall m idx job . ( MonadBaseControl IO m , MonadLog m , MonadMask m @@ -220,7 +220,7 @@ spawnMonitor , FromSQL idx , ToSQL idx ) - => ConsumerConfig m idx info + => ConsumerConfig m idx job -> ConnectionSourceM m -> ConsumerID -> m ThreadId @@ -294,7 +294,7 @@ spawnMonitor ConsumerConfig {..} cs cid = forkP "monitor" . forever $ do -- | Spawn a thread that reserves and processes jobs. spawnDispatcher - :: forall m idx info + :: forall m idx job . ( MonadBaseControl IO m , MonadLog m , MonadMask m @@ -302,7 +302,7 @@ spawnDispatcher , Show idx , ToSQL idx ) - => ConsumerConfig m idx info + => ConsumerConfig m idx job -> ConnectionSourceM m -> ConsumerID -> MVar () @@ -357,7 +357,7 @@ spawnDispatcher ConsumerConfig {..} cs cid semaphore runningJobsInfo runningJobs pure (batchSize > 0) - reserveJobs :: Int -> m ([Job idx info], Int) + reserveJobs :: Int -> m ([job], Int) reserveJobs limit = runDBT cs ts $ do now <- currentTime n <- @@ -389,7 +389,7 @@ spawnDispatcher ConsumerConfig {..} cs cid semaphore runningJobsInfo runningJobs ] -- Spawn each job in a separate thread. - startJob :: Job idx info -> m (Job idx info, m (T.Result Result)) + startJob :: job -> m (job, m (T.Result Result)) startJob job = do (_, joinFork) <- mask $ \restore -> T.fork $ do tid <- myThreadId @@ -403,7 +403,7 @@ spawnDispatcher ConsumerConfig {..} cs cid semaphore runningJobsInfo runningJobs modifyTVar' runningJobsInfo $ M.delete tid -- Wait for all the jobs and collect their results. - joinJob :: (Job idx info, m (T.Result Result)) -> m (idx, Result) + joinJob :: (job, m (T.Result Result)) -> m (idx, Result) joinJob (job, joinFork) = joinFork >>= \case Right result -> pure (ccJobIndex job, result) diff --git a/consumers/src/Database/PostgreSQL/Consumers/Config.hs b/consumers/src/Database/PostgreSQL/Consumers/Config.hs index 9d9d1d4..b37ae05 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/Config.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/Config.hs @@ -1,10 +1,8 @@ module Database.PostgreSQL.Consumers.Config ( Action (..) , Result (..) - , Job (..) , ConsumerConfig (..) , hoistConsumer - , exponentialBackoff ) where import Control.Exception (SomeException) @@ -28,21 +26,8 @@ data Action data Result = Ok Action | Failed Action deriving (Eq, Ord, Show) --- | A job fetched from the queue, together with its queue bookkeeping --- fields (as opposed to just the caller-supplied 'jobInfo' payload). -data Job idx info = Job - { jobIndex :: !idx - , jobRunAt :: !(Maybe UTCTime) - , jobFinishedAt :: !(Maybe UTCTime) - , jobAttempts :: !Int - -- ^ Number of consecutive failed processing attempts, including the - -- current one. Reset to 1 once a job has succeeded, so this is a streak - -- of failures rather than a lifetime total. - , jobInfo :: !info - } - -- | Config of a consumer. -data ConsumerConfig m idx info = forall row. FromRow row => ConsumerConfig +data ConsumerConfig m idx job = forall row. FromRow row => ConsumerConfig { ccJobsTable :: !(RawSQL ()) -- ^ Name of the database table where jobs are stored. The table needs to have -- the following columns in order to be suitable for acting as a job queue: @@ -93,13 +78,19 @@ data ConsumerConfig m idx info = forall row. FromRow row => ConsumerConfig -- and these jobs stay locked forever, yet are never processed. , ccJobSelectors :: ![SQL] -- ^ Fields needed to be selected from the jobs table in order to assemble a - -- job. Needs to match what 'ccJobFetcher' expects, and typically includes - -- @id@, @run_at@, @finished_at@ and @attempts@ alongside whatever columns - -- make up the job's 'jobInfo' payload. - , ccJobFetcher :: !(row -> Job idx info) + -- job. + , ccJobFetcher :: !(row -> job) -- ^ Function that transforms the list of fields into a job. - , ccJobIndex :: !(Job idx info -> idx) + , ccJobIndex :: !(job -> idx) -- ^ Selector for taking out job ID from the job object. + , ccJobAttempts :: !(job -> Int) + -- ^ 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 @@ -124,67 +115,26 @@ data ConsumerConfig m idx info = forall row. FromRow row => ConsumerConfig -- it needs to be a positive number. , ccMaxRunningJobs :: !Int -- ^ Maximum amount of jobs that can be processed in parallel. - , ccProcessJob :: !(Job idx info -> m Result) + , ccProcessJob :: !(job -> m Result) -- ^ Function that processes a job. It's recommended to process each job in a -- separate DB transaction, otherwise you'll have to remember to commit your -- changes to the database manually. - , ccOnException :: !(SomeException -> Job idx info -> m Action) + , ccOnException :: !(SomeException -> job -> m Action) -- ^ Action taken if a job processing function throws an exception. For -- robustness it's best to ensure that it doesn't throw. If it does, the -- exception will be logged and the job in question postponed by a day. - , ccJobLogData :: !(Job idx info -> [A.Pair]) + , ccJobLogData :: !(job -> [A.Pair]) -- ^ Data to attach to each log message while processing a job. } -- | Change the monad the consumer lives in. hoistConsumer :: (forall r. m r -> n r) - -> ConsumerConfig m idx info - -> ConsumerConfig n idx info + -> ConsumerConfig m idx job + -> ConsumerConfig n idx job hoistConsumer f ConsumerConfig {..} = ConsumerConfig { ccProcessJob = f . ccProcessJob , ccOnException = \ex -> f . ccOnException ex , .. } - --- | A ready-made 'ccOnException' handler: reruns a job with exponentially --- increasing delay, then gives up (removes the job) once 'jobAttempts' --- exceeds @maxAttempts@. --- --- /Note:/ this mirrors the growth curve from the original proof of concept --- (delay ^ attempts), so it only grows the delay if @startInterval@ is --- greater than one second; below that it shrinks towards zero instead. --- Pick @startInterval@ accordingly, or write a custom handler if you need a --- more conventional @startInterval * base ^ attempts@ curve. -exponentialBackoff - :: Monad m - => Int - -- ^ Maximum number of attempts before the job is removed. - -> Interval - -- ^ Delay before the first retry. - -> SomeException - -> Job idx info - -> m Action -exponentialBackoff maxAttempts startInterval _ Job {..} = - pure $ case (jobAttempts > maxAttempts, jobAttempts > 2) of - (True, _) -> Remove - (False, False) -> RerunAfter startInterval - (False, True) -> - RerunAfter . diffTimeToInterval $ - intervalToDiffTime startInterval ^ jobAttempts - -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/src/Database/PostgreSQL/Consumers/Consumer.hs b/consumers/src/Database/PostgreSQL/Consumers/Consumer.hs index 5e568e9..a836db9 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/Consumer.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/Consumer.hs @@ -33,7 +33,7 @@ instance Show ConsumerID where -- acquired ID. registerConsumer :: (MonadBase IO m, MonadMask m, MonadTime m) - => ConsumerConfig n idx info + => ConsumerConfig n idx job -> ConnectionSourceM m -> m ConsumerID registerConsumer ConsumerConfig {..} cs = runDBT cs defaultTransactionSettings $ do @@ -49,7 +49,7 @@ registerConsumer ConsumerConfig {..} cs = runDBT cs defaultTransactionSettings $ -- | Unregister consumer with a given ID. unregisterConsumer :: (MonadBase IO m, MonadMask m) - => ConsumerConfig n idx info + => ConsumerConfig n idx job -> ConnectionSourceM m -> ConsumerID -> m () diff --git a/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs b/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs new file mode 100644 index 0000000..dbc727b --- /dev/null +++ b/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs @@ -0,0 +1,128 @@ +-- | +-- 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.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. + -> Int + -- ^ The job's current attempt count (see 'ccJobAttempts'). + -> Action +constantBackoff maxAttempts delay attempts + | 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@. + -> Int + -- ^ The job's current attempt count (see 'ccJobAttempts'). + -> Action +linearBackoff maxAttempts delayUnit attempts + | attempts >= maxAttempts = Remove + | otherwise = RerunAfter $ scaleInterval 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. + -> Int + -- ^ The job's current attempt count (see 'ccJobAttempts'). + -> Action +exponentialBackoff maxAttempts baseDelay maxDelay attempts + | attempts >= maxAttempts = Remove + | otherwise = RerunAfter $ nextDelay baseDelay maxDelay 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. + -> Int + -- ^ The job's current attempt count (see 'ccJobAttempts'). + -> m Action +exponentialBackoffWithJitter maxAttempts baseDelay maxDelay attempts + | attempts >= maxAttempts = pure Remove + | otherwise = do + jitter <- liftIO $ randomRIO (0.5, 1.5 :: Double) + pure . RerunAfter . scaleIntervalD jitter $ nextDelay baseDelay maxDelay 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 9677d51..812d8c9 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,24 +145,17 @@ test = do ConsumerConfig { ccJobsTable = "consumers_test_jobs" , ccConsumersTable = "consumers_test_consumers" - , ccJobSelectors = ["id", "run_at", "finished_at", "attempts", "countdown"] - , ccJobFetcher = - \(i :: Int64, runAt :: Maybe UTCTime, finishedAt :: Maybe UTCTime, attempts :: Int32, countdown :: Int32) -> - Job - { jobIndex = i - , jobRunAt = runAt - , jobFinishedAt = finishedAt - , jobAttempts = fromIntegral attempts - , jobInfo = countdown - } - , ccJobIndex = jobIndex + , ccJobSelectors = ["id", "attempts", "countdown"] + , ccJobFetcher = id + , 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 = \job -> ["job_id" .= jobIndex job] + , ccJobLogData = \(i, _, _) -> ["job_id" .= i] } putJob :: Int32 -> TestEnv () @@ -175,16 +169,17 @@ test = do <> ")" notify "consumers_test_chan" "" - processJob :: Job Int64 Int32 -> TestEnv Result - processJob Job {jobInfo = 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 -> Job 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 = From a21163312ef1a2732ee22bff67dbe1c2eb7f5340 Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 22 Sep 2026 20:58:15 +0000 Subject: [PATCH 3/6] Add consumers_job_failed_attempts metric Add a new histogram, consumers_job_failed_attempts, labelled by job_name, observing ccJobAttempts for jobs whose result wasn't Ok (i.e. Failed, an exception, or an abort). This is the metric Jan asked about: it lets you tell one-off failures (attempt 1) apart from jobs stuck in a persistent retry loop, to denoise alerting/dashboards for e.g. the document sealing queue. Reuses the existing reportJob release action from generalBracket alongside consumers_job_execution_seconds, now also passed the job value so it can read ccJobAttempts. Bucket boundaries are configurable via the new jobFailedAttemptsBuckets field on ConsumerMetricsConfig, same pattern as jobExecutionBuckets. Requires consumers >= 2.4.0.0 for ccJobAttempts. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_016NeDfW7w3TwyS1UM82uvkL --- consumers-metrics-prometheus/CHANGELOG.md | 8 +++++ .../consumers-metrics-prometheus.cabal | 4 +-- .../PostgreSQL/Consumers/Instrumented.hs | 32 +++++++++++++++++-- 3 files changed, 40 insertions(+), 4 deletions(-) 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 e023d87..88f7f12 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..c3d15cc 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,16 @@ 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] From ae568bca89decf43bf541155695d4d09204489f1 Mon Sep 17 00:00:00 2001 From: Jan Sipr Date: Wed, 23 Sep 2026 13:26:33 +0200 Subject: [PATCH 4/6] Haskell CI and Postgres-15 --- .github/workflows/haskell-ci.yml | 55 +++++++++++++++++++++----------- cabal.haskell-ci | 1 + postgres-15.patch | 11 +++++++ 3 files changed, 48 insertions(+), 19 deletions(-) create mode 100644 postgres-15.patch diff --git a/.github/workflows/haskell-ci.yml b/.github/workflows/haskell-ci.yml index 862fa02..7595667 100644 --- a/.github/workflows/haskell-ci.yml +++ b/.github/workflows/haskell-ci.yml @@ -8,9 +8,9 @@ # # For more information, see https://github.com/haskell-CI/haskell-ci # -# version: 0.19.20260104 +# version: 0.19.20260923 # -# REGENDATA ("0.19.20260104",["github","cabal.project","--config=cabal.haskell-ci"]) +# REGENDATA ("0.19.20260923",["github","cabal.project","--config=cabal.haskell-ci"]) # name: Haskell-CI on: @@ -23,6 +23,8 @@ on: merge_group: branches: - master + workflow_dispatch: + {} jobs: linux: name: Haskell-CI - Linux - ${{ matrix.compiler }} @@ -33,7 +35,7 @@ jobs: image: buildpack-deps:jammy services: postgres: - image: postgres:14 + image: postgres:15 env: POSTGRES_PASSWORD: postgres options: --health-cmd pg_isready --health-interval 10s --health-timeout 5s --health-retries 5 @@ -44,7 +46,7 @@ jobs: - compiler: ghc-9.14.1 compilerKind: ghc compilerVersion: 9.14.1 - setup-method: ghcup + setup-method: ghcup-prerelease allow-failure: false - compiler: ghc-9.12.2 compilerKind: ghc @@ -105,6 +107,21 @@ jobs: HCKIND: ${{ matrix.compilerKind }} HCNAME: ${{ matrix.compiler }} HCVER: ${{ matrix.compilerVersion }} + - name: Install GHC (GHCup prerelease) + if: matrix.setup-method == 'ghcup-prerelease' + run: | + "$HOME/.ghcup/bin/ghcup" config add-release-channel prereleases + "$HOME/.ghcup/bin/ghcup" install ghc "$HCVER" || (cat "$HOME"/.ghcup/logs/*.* && false) + HC=$("$HOME/.ghcup/bin/ghcup" whereis ghc "$HCVER") + HCPKG=$(echo "$HC" | sed 's#ghc$#ghc-pkg#') + HADDOCK=$(echo "$HC" | sed 's#ghc$#haddock#') + echo "HC=$HC" >> "$GITHUB_ENV" + echo "HCPKG=$HCPKG" >> "$GITHUB_ENV" + echo "HADDOCK=$HADDOCK" >> "$GITHUB_ENV" + env: + HCKIND: ${{ matrix.compilerKind }} + HCNAME: ${{ matrix.compiler }} + HCVER: ${{ matrix.compilerVersion }} - name: Set PATH and environment variables run: | echo "$HOME/.cabal/bin" >> $GITHUB_PATH @@ -166,14 +183,14 @@ jobs: chmod a+x $HOME/.cabal/bin/cabal-plan cabal-plan --version - name: checkout - uses: actions/checkout@v5 + uses: actions/checkout@v7 with: path: source - name: initial cabal.project for sdist run: | touch cabal.project - echo "packages: $GITHUB_WORKSPACE/source/consumers" >> cabal.project echo "packages: $GITHUB_WORKSPACE/source/consumers-metrics-prometheus" >> cabal.project + echo "packages: $GITHUB_WORKSPACE/source/consumers" >> cabal.project cat cabal.project - name: sdist run: | @@ -185,27 +202,27 @@ jobs: find sdist -maxdepth 1 -type f -name '*.tar.gz' -exec tar -C $GITHUB_WORKSPACE/unpacked -xzvf {} \; - name: generate cabal.project run: | - PKGDIR_consumers="$(find "$GITHUB_WORKSPACE/unpacked" -maxdepth 1 -type d -regex '.*/consumers-[0-9.]*')" - echo "PKGDIR_consumers=${PKGDIR_consumers}" >> "$GITHUB_ENV" PKGDIR_consumers_metrics_prometheus="$(find "$GITHUB_WORKSPACE/unpacked" -maxdepth 1 -type d -regex '.*/consumers-metrics-prometheus-[0-9.]*')" echo "PKGDIR_consumers_metrics_prometheus=${PKGDIR_consumers_metrics_prometheus}" >> "$GITHUB_ENV" + PKGDIR_consumers="$(find "$GITHUB_WORKSPACE/unpacked" -maxdepth 1 -type d -regex '.*/consumers-[0-9.]*')" + echo "PKGDIR_consumers=${PKGDIR_consumers}" >> "$GITHUB_ENV" rm -f cabal.project cabal.project.local touch cabal.project touch cabal.project.local - echo "packages: ${PKGDIR_consumers}" >> cabal.project echo "packages: ${PKGDIR_consumers_metrics_prometheus}" >> cabal.project - echo "package consumers" >> cabal.project - echo " ghc-options: -Werror=missing-methods -Werror=missing-fields" >> cabal.project + echo "packages: ${PKGDIR_consumers}" >> cabal.project echo "package consumers-metrics-prometheus" >> cabal.project echo " ghc-options: -Werror=missing-methods -Werror=missing-fields" >> cabal.project - if [ $((HCNUMVER >= 90400)) -ne 0 ] ; then echo "package consumers" >> cabal.project ; fi - if [ $((HCNUMVER >= 90400)) -ne 0 ] ; then echo " ghc-options: -Werror=unused-packages" >> cabal.project ; fi + echo "package consumers" >> cabal.project + echo " ghc-options: -Werror=missing-methods -Werror=missing-fields" >> cabal.project if [ $((HCNUMVER >= 90400)) -ne 0 ] ; then echo "package consumers-metrics-prometheus" >> cabal.project ; fi if [ $((HCNUMVER >= 90400)) -ne 0 ] ; then echo " ghc-options: -Werror=unused-packages" >> cabal.project ; fi - echo "package consumers" >> cabal.project - echo " ghc-options: -Werror=incomplete-patterns -Werror=incomplete-uni-patterns" >> cabal.project + if [ $((HCNUMVER >= 90400)) -ne 0 ] ; then echo "package consumers" >> cabal.project ; fi + if [ $((HCNUMVER >= 90400)) -ne 0 ] ; then echo " ghc-options: -Werror=unused-packages" >> cabal.project ; fi echo "package consumers-metrics-prometheus" >> cabal.project echo " ghc-options: -Werror=incomplete-patterns -Werror=incomplete-uni-patterns" >> cabal.project + echo "package consumers" >> cabal.project + echo " ghc-options: -Werror=incomplete-patterns -Werror=incomplete-uni-patterns" >> cabal.project cat >> cabal.project < Date: Wed, 23 Sep 2026 10:42:28 +0000 Subject: [PATCH 5/6] hlint redundant \$ --- .../src/Database/PostgreSQL/Consumers/Instrumented.hs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs b/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs index c3d15cc..5596fa1 100644 --- a/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs +++ b/consumers-metrics-prometheus/src/Database/PostgreSQL/Consumers/Instrumented.hs @@ -287,7 +287,6 @@ instrumentConsumerConfig ConsumerMetrics {..} ConsumerConfig {..} = ExitCaseSuccess (Ok _) -> pure () _ -> liftBase $ - Prom.withLabel jobsFailedAttempts jobName $ - (`Prom.observe` fromIntegral (ccJobAttempts job)) + Prom.withLabel jobsFailedAttempts jobName (`Prom.observe` fromIntegral (ccJobAttempts job)) handleEx e = logAttention "Exception while instrumenting job" $ object ["exception" .= show e] From 0f07a8e1531fe228c149957030903724a4118cdc Mon Sep 17 00:00:00 2001 From: Jonathan Jouty Date: Wed, 30 Sep 2026 14:34:10 +0100 Subject: [PATCH 6/6] Use Int32 for ccJobAttempts --- .../Database/PostgreSQL/Consumers/Config.hs | 3 ++- .../PostgreSQL/Consumers/RetryStrategy.hs | 23 ++++++++++--------- 2 files changed, 14 insertions(+), 12 deletions(-) diff --git a/consumers/src/Database/PostgreSQL/Consumers/Config.hs b/consumers/src/Database/PostgreSQL/Consumers/Config.hs index b37ae05..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,7 +84,7 @@ 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 -> Int) + , 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 diff --git a/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs b/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs index dbc727b..029e44c 100644 --- a/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs +++ b/consumers/src/Database/PostgreSQL/Consumers/RetryStrategy.hs @@ -22,6 +22,7 @@ module Database.PostgreSQL.Consumers.RetryStrategy ) where import Control.Monad.IO.Class +import Data.Int (Int32) import Data.Time import Database.PostgreSQL.Consumers.Config import Database.PostgreSQL.PQTypes.Interval @@ -34,11 +35,11 @@ constantBackoff -- ^ Maximum number of attempts before the job is removed. -> Interval -- ^ Delay before every retry. - -> Int + -> Int32 -- ^ The job's current attempt count (see 'ccJobAttempts'). -> Action constantBackoff maxAttempts delay attempts - | attempts >= maxAttempts = Remove + | fromIntegral attempts >= maxAttempts = Remove | otherwise = RerunAfter delay -- | Retry with a delay that grows linearly with the attempt count @@ -48,12 +49,12 @@ linearBackoff -- ^ Maximum number of attempts before the job is removed. -> Interval -- ^ Delay unit; the delay before retry number @n@ is @n * delayUnit@. - -> Int + -> Int32 -- ^ The job's current attempt count (see 'ccJobAttempts'). -> Action linearBackoff maxAttempts delayUnit attempts - | attempts >= maxAttempts = Remove - | otherwise = RerunAfter $ scaleInterval attempts delayUnit + | 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 @@ -65,12 +66,12 @@ exponentialBackoff -- ^ Base delay, used for the first retry. -> Interval -- ^ Delay cap; the computed delay never exceeds this. - -> Int + -> Int32 -- ^ The job's current attempt count (see 'ccJobAttempts'). -> Action exponentialBackoff maxAttempts baseDelay maxDelay attempts - | attempts >= maxAttempts = Remove - | otherwise = RerunAfter $ nextDelay 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. @@ -88,14 +89,14 @@ exponentialBackoffWithJitter -- ^ Base delay, used for the first retry (before jitter). -> Interval -- ^ Delay cap, applied before jitter is added. - -> Int + -> Int32 -- ^ The job's current attempt count (see 'ccJobAttempts'). -> m Action exponentialBackoffWithJitter maxAttempts baseDelay maxDelay attempts - | attempts >= maxAttempts = pure Remove + | fromIntegral attempts >= maxAttempts = pure Remove | otherwise = do jitter <- liftIO $ randomRIO (0.5, 1.5 :: Double) - pure . RerunAfter . scaleIntervalD jitter $ nextDelay baseDelay maxDelay attempts + pure . RerunAfter . scaleIntervalD jitter $ nextDelay baseDelay maxDelay (fromIntegral attempts) ----------------------------------------