-
Notifications
You must be signed in to change notification settings - Fork 0
Add ccJobAttempts, retry strategy, and metrics for attempts for failed jobs #48
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
181d930
53b08a7
a211633
ae568bc
68dead9
330eb7e
0f07a8e
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -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. | ||
| 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-??-??) | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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) | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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. | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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 | ||
|
|
||
| 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 |
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Note: TODO release date