| Safe Haskell | None |
|---|---|
| Language | GHC2024 |
Arbiter.Core.Job.Types
Description
Job records, the enqueue setters, and the lifecycle hook types.
Synopsis
- type Job payload key q insertedAt adm = JobRecord payload key q insertedAt adm
- class HasKind payload where
- constructorKind :: (GKindOf (Rep a), Generic a) => a -> Text
- constructorKinds :: GKindsOf (Rep a) => [Text]
- data PayloadKeys = PayloadKeys {}
- data PayloadColumns = PayloadColumns {}
- type JobRead payload = Job payload Int64 Text UTCTime PayloadKeys
- type JobWrite payload = Job payload () () () ()
- primaryKey :: JobRecord payload key q insertedAt adm -> key
- payload :: JobRecord payload key q insertedAt adm -> payload
- queueName :: JobRecord payload key q insertedAt adm -> q
- groupKey :: JobRecord payload key q insertedAt adm -> Maybe Text
- insertedAt :: JobRecord payload key q insertedAt adm -> insertedAt
- updatedAt :: JobRecord payload Int64 q insertedAt adm -> Maybe UTCTime
- attempts :: JobRecord payload Int64 q insertedAt adm -> Int32
- lastError :: JobRecord payload Int64 q insertedAt adm -> Maybe Text
- priority :: JobRecord payload key q insertedAt adm -> Int32
- lastAttemptedAt :: JobRecord payload Int64 q insertedAt adm -> Maybe UTCTime
- notVisibleUntil :: JobRecord payload key q insertedAt adm -> Maybe UTCTime
- dedupKey :: JobRecord payload key q insertedAt adm -> Maybe DedupKey
- maxAttempts :: JobRecord payload key q insertedAt adm -> Maybe Int32
- parentId :: JobRecord payload Int64 q insertedAt adm -> Maybe Int64
- parentState :: JobRecord payload Int64 q insertedAt adm -> Maybe Value
- traceContext :: JobRecord payload key q insertedAt adm -> Maybe TraceContext
- suspended :: JobRecord payload Int64 q insertedAt adm -> Bool
- claimedBy :: JobRecord payload Int64 q insertedAt adm -> Maybe UUID
- claimSeq :: JobRecord payload Int64 q insertedAt adm -> Int64
- archiveFor :: JobRecord payload key q insertedAt adm -> Maybe Int32
- payloadKeys :: JobRecord payload key q insertedAt adm -> adm
- defaultJob :: payload -> JobWrite payload
- defaultGroupedJob :: Text -> payload -> JobWrite payload
- setPayload :: payload' -> JobWrite payload -> JobWrite payload'
- setGroupKey :: Maybe Text -> JobWrite payload -> JobWrite payload
- setPriority :: Int32 -> JobWrite payload -> JobWrite payload
- setNotVisibleUntil :: Maybe UTCTime -> JobWrite payload -> JobWrite payload
- setDedupKey :: Maybe DedupKey -> JobWrite payload -> JobWrite payload
- setMaxAttempts :: Maybe Int32 -> JobWrite payload -> JobWrite payload
- setTraceContext :: Maybe TraceContext -> JobWrite payload -> JobWrite payload
- setArchiveFor :: Maybe Int32 -> JobWrite payload -> JobWrite payload
- mapPayload :: (payload -> payload') -> Job payload key q insertedAt adm -> Job payload' key q insertedAt adm
- defaultMaxAttempts :: Int32
- dayRetention :: Int32
- isRollup :: Job p Int64 q t adm -> Bool
- data JobStatus
- jobStatusToText :: JobStatus -> Text
- jobStatusFromText :: Text -> Either Text JobStatus
- type JobPayload payload = (FromJSON payload, ToJSON payload, HasKind payload, HasRateLimit payload, HasConcurrency payload)
- type RegistryAdmissionPolicies (registry :: JobPayloadRegistry) = (RegistryConcurrencyPolicies registry, RegistryRateLimitPolicies registry)
- data DedupKey
- dedupParts :: Maybe DedupKey -> (Maybe Text, Maybe Text)
- data TraceContext = TraceContext {
- traceparent :: Text
- tracestate :: Maybe Text
- toTraceContext :: Maybe Text -> Maybe Text -> Maybe TraceContext
- data ObservabilityHooks (m :: Type -> Type) payload = ObservabilityHooks {
- onJobClaimed :: JobPayload payload => JobRead payload -> ClaimTime -> m ()
- onJobSuccess :: JobPayload payload => JobRead payload -> StartTime -> EndTime -> m ()
- onJobFailure :: JobPayload payload => JobRead payload -> ErrorMsg -> StartTime -> EndTime -> m ()
- onJobRetry :: JobPayload payload => JobRead payload -> BackoffDelay -> m ()
- onJobFailedAndMovedToDLQ :: JobPayload payload => ErrorMsg -> JobRead payload -> m ()
- onJobCancelled :: JobPayload payload => JobRead payload -> ErrorMsg -> m ()
- onJobUnavailable :: JobPayload payload => JobRead payload -> ErrorMsg -> m ()
- onJobHeartbeat :: JobPayload payload => JobRead payload -> CurrentTime -> StartTime -> m ()
- defaultObservabilityHooks :: forall (m :: Type -> Type) payload. Applicative m => ObservabilityHooks m payload
- andThen :: MonadUnliftIO m => m () -> m () -> m ()
- type JobId = Int64
- type ClaimSeq = Int64
- type ClaimTime = UTCTime
- type CurrentTime = UTCTime
- type StartTime = UTCTime
- type EndTime = UTCTime
- type ErrorMsg = Text
- type BackoffDelay = NominalDiffTime
Core Job Type
Source #type Job payload key q insertedAt adm = JobRecord payload key q insertedAt adm
A job parametrized over payload, primary key, queue name, insertion timestamp, and the columns derived from the payload. The constructor is internal.
Source #class HasKind payload where
A payload's per-job variant label. Defaults to unlabelled.
Minimal complete definition
Nothing
Source #constructorKind :: (GKindOf (Rep a), Generic a) => a -> Text
The constructor name of a value, for a payload that wraps the sum it wants
labelled. Needs Generic on the wrapped type.
instance HasKind Envelope where kindOf = Just . constructorKind . envelopePayload kindsFor = constructorKinds @EmailPayload
Source #constructorKinds :: GKindsOf (Rep a) => [Text]
Every constructor name of a type, in declaration order.
Source #data PayloadKeys
The labels and keys a stored job carries from its payload, one field each.
Constructors
| PayloadKeys | |
Fields
| |
Instances
Source #data PayloadColumns
The writable columns resolved from a payload at enqueue. The kind, key
and prefix columns round-trip via PayloadKeys. cost is write-only.
Constructors
| PayloadColumns | |
Fields | |
Instances
| Eq PayloadColumns Source # | |||||
Defined in Arbiter.Core.Job.Types Methods #(==) :: PayloadColumns -> PayloadColumns -> Bool #(/=) :: PayloadColumns -> PayloadColumns -> Bool | |||||
| Generic PayloadColumns Source # | |||||
Defined in Arbiter.Core.Job.Types Associated Types
| |||||
| Show PayloadColumns Source # | |||||
Defined in Arbiter.Core.Job.Types Methods #showsPrec :: Int -> PayloadColumns -> ShowS #show :: PayloadColumns -> String #showList :: [PayloadColumns] -> ShowS | |||||
| type Rep PayloadColumns Source # | |||||
Defined in Arbiter.Core.Job.Types type Rep PayloadColumns = D1 ('MetaData "PayloadColumns" "Arbiter.Core.Job.Types" "arbiter-core-0.1.0.0-inplace" 'False) (C1 ('MetaCons "PayloadColumns" 'PrefixI 'True) ((S1 ('MetaSel ('Just "pcKind") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 (Maybe Text)) :*: (S1 ('MetaSel ('Just "pcRateLimitKey") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 (Maybe Text)) :*: S1 ('MetaSel ('Just "pcRateLimitPrefix") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 (Maybe Text)))) :*: (S1 ('MetaSel ('Just "pcRateLimitCost") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 Double) :*: (S1 ('MetaSel ('Just "pcConcurrencyKey") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 (Maybe Text)) :*: S1 ('MetaSel ('Just "pcConcurrencyPrefix") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 (Maybe Text)))))) | |||||
Source #type JobRead payload = Job payload Int64 Text UTCTime PayloadKeys
A job read from the database.
Source #type JobWrite payload = Job payload () () () ()
A job ready to enqueue. Arbiter owns claim, retry, parent, rollup, and suspension state. Use the exported setters to configure enqueue fields.
Source #primaryKey :: JobRecord payload key q insertedAt adm -> key
Database-assigned identifier for a stored job.
Source #payload :: JobRecord payload key q insertedAt adm -> payload
User-defined payload stored as JSONB.
Source #groupKey :: JobRecord payload key q insertedAt adm -> Maybe Text
Serial-processing group, or Nothing for an ungrouped job.
Source #insertedAt :: JobRecord payload key q insertedAt adm -> insertedAt
Time at which the job was inserted.
Source #updatedAt :: JobRecord payload Int64 q insertedAt adm -> Maybe UTCTime
Time at which the job was last updated.
Source #attempts :: JobRecord payload Int64 q insertedAt adm -> Int32
Number of attempts made so far.
Source #lastError :: JobRecord payload Int64 q insertedAt adm -> Maybe Text
Error message from the last failed attempt.
Source #priority :: JobRecord payload key q insertedAt adm -> Int32
Claim priority. Lower numbers have higher priority.
Source #lastAttemptedAt :: JobRecord payload Int64 q insertedAt adm -> Maybe UTCTime
Time at which a worker last claimed the job.
Source #notVisibleUntil :: JobRecord payload key q insertedAt adm -> Maybe UTCTime
Earliest time at which the job can be claimed.
Source #dedupKey :: JobRecord payload key q insertedAt adm -> Maybe DedupKey
Deduplication strategy and key.
Source #maxAttempts :: JobRecord payload key q insertedAt adm -> Maybe Int32
Attempt limit before the job moves to the DLQ.
Source #parentId :: JobRecord payload Int64 q insertedAt adm -> Maybe Int64
Identifier of this job's parent in a job tree.
Source #parentState :: JobRecord payload Int64 q insertedAt adm -> Maybe Value
Snapshot of accumulated child results for a rollup finalizer.
Source #traceContext :: JobRecord payload key q insertedAt adm -> Maybe TraceContext
W3C trace context captured at enqueue.
Source #suspended :: JobRecord payload Int64 q insertedAt adm -> Bool
Whether the job is currently ineligible for claiming.
Source #claimedBy :: JobRecord payload Int64 q insertedAt adm -> Maybe UUID
Worker pool that most recently claimed the job.
Source #claimSeq :: JobRecord payload Int64 q insertedAt adm -> Int64
Monotonically increasing claim identifier.
Source #archiveFor :: JobRecord payload key q insertedAt adm -> Maybe Int32
Completed-job archive retention in seconds.
Source #payloadKeys :: JobRecord payload key q insertedAt adm -> adm
The labels and keys stamped from the payload at enqueue.
Source #defaultJob :: payload -> JobWrite payload
Ungrouped JobWrite with default values. For serial processing within a
group, use defaultGroupedJob.
Source #defaultGroupedJob :: Text -> payload -> JobWrite payload
defaultJob with a group key. Jobs sharing a group key are processed serially.
Source #setPayload :: payload' -> JobWrite payload -> JobWrite payload'
Replace the payload.
Source #setPriority :: Int32 -> JobWrite payload -> JobWrite payload
Set the claim priority. Lower numbers claim first.
Source #setNotVisibleUntil :: Maybe UTCTime -> JobWrite payload -> JobWrite payload
Delay the job's first visibility.
Source #setMaxAttempts :: Maybe Int32 -> JobWrite payload -> JobWrite payload
Override the queue's attempt limit.
Source #setTraceContext :: Maybe TraceContext -> JobWrite payload -> JobWrite payload
Attach a trace context.
Source #setArchiveFor :: Maybe Int32 -> JobWrite payload -> JobWrite payload
Set the archive retention, in seconds.
Source #mapPayload :: (payload -> payload') -> Job payload key q insertedAt adm -> Job payload' key q insertedAt adm
Transform a job's payload without changing its stored metadata.
Source #defaultMaxAttempts :: Int32
Default attempt limit stamped onto jobs whose maxAttempts is unset.
24h in seconds, a convenience value for archiveFor.
Source #isRollup :: Job p Int64 q t adm -> Bool
A rollup finalizer is any job whose parentState snapshot is present
(an empty object on insert, the merged child results before a DLQ move).
Derived status
Effective job status. Arbiter derives status from the stored fields. The status SQL in Arbiter.Core.Sql.Jobs is its source of truth.
Instances
| FromJSON JobStatus Source # | |||||
Defined in Arbiter.Core.Job.Status | |||||
| ToJSON JobStatus Source # | |||||
| Eq JobStatus Source # | |||||
| Bounded JobStatus Source # | |||||
| Enum JobStatus Source # | |||||
Defined in Arbiter.Core.Job.Status | |||||
| Generic JobStatus Source # | |||||
Defined in Arbiter.Core.Job.Status Associated Types
| |||||
| Show JobStatus Source # | |||||
| type Rep JobStatus Source # | |||||
Defined in Arbiter.Core.Job.Status type Rep JobStatus = D1 ('MetaData "JobStatus" "Arbiter.Core.Job.Status" "arbiter-core-0.1.0.0-inplace" 'False) ((C1 ('MetaCons "Ready" 'PrefixI 'False) (U1 :: Type -> Type) :+: (C1 ('MetaCons "InFlight" 'PrefixI 'False) (U1 :: Type -> Type) :+: C1 ('MetaCons "Backoff" 'PrefixI 'False) (U1 :: Type -> Type))) :+: ((C1 ('MetaCons "Scheduled" 'PrefixI 'False) (U1 :: Type -> Type) :+: C1 ('MetaCons "Suspended" 'PrefixI 'False) (U1 :: Type -> Type)) :+: (C1 ('MetaCons "Throttled" 'PrefixI 'False) (U1 :: Type -> Type) :+: C1 ('MetaCons "Cancelled" 'PrefixI 'False) (U1 :: Type -> Type)))) | |||||
Source #jobStatusToText :: JobStatus -> Text
The wire name for a status.
Source #jobStatusFromText :: Text -> Either Text JobStatus
Strict inverse of jobStatusToText. Unknown values are rejected.
Type Constraints
Source #type JobPayload payload = (FromJSON payload, ToJSON payload, HasKind payload, HasRateLimit payload, HasConcurrency payload)
The full payload contract. JSON round-trip for JSONB storage plus the rate-limit and concurrency declarations. Both default to unlimited.
Source #type RegistryAdmissionPolicies (registry :: JobPayloadRegistry) = (RegistryConcurrencyPolicies registry, RegistryRateLimitPolicies registry)
The registry declares both admission policy kinds.
Deduplication
Deduplication strategy, checked on INSERT via ON CONFLICT on the dedup key.
Constructors
| IgnoreDuplicate Text | Skip if a job with this key exists ( |
| ReplaceDuplicate Text | Replace the existing job with this key ( |
Instances
| FromJSON DedupKey Source # | |||||
Defined in Arbiter.Core.Job.Dedup | |||||
| ToJSON DedupKey Source # | |||||
| Eq DedupKey Source # | |||||
| Generic DedupKey Source # | |||||
Defined in Arbiter.Core.Job.Dedup Associated Types
| |||||
| Show DedupKey Source # | |||||
| type Rep DedupKey Source # | |||||
Defined in Arbiter.Core.Job.Dedup type Rep DedupKey = D1 ('MetaData "DedupKey" "Arbiter.Core.Job.Dedup" "arbiter-core-0.1.0.0-inplace" 'False) (C1 ('MetaCons "IgnoreDuplicate" 'PrefixI 'False) (S1 ('MetaSel ('Nothing :: Maybe Symbol) 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 Text)) :+: C1 ('MetaCons "ReplaceDuplicate" 'PrefixI 'False) (S1 ('MetaSel ('Nothing :: Maybe Symbol) 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 Text))) | |||||
Source #dedupParts :: Maybe DedupKey -> (Maybe Text, Maybe Text)
The dedup_key and dedup_strategy column values for a DedupKey.
Trace context
Source #data TraceContext
A job's W3C trace context.
Constructors
| TraceContext | |
Fields
| |
Instances
| Eq TraceContext Source # | |||||
Defined in Arbiter.Core.Job.TraceContext | |||||
| Generic TraceContext Source # | |||||
Defined in Arbiter.Core.Job.TraceContext Associated Types
| |||||
| Show TraceContext Source # | |||||
Defined in Arbiter.Core.Job.TraceContext Methods #showsPrec :: Int -> TraceContext -> ShowS #show :: TraceContext -> String #showList :: [TraceContext] -> ShowS | |||||
| type Rep TraceContext Source # | |||||
Defined in Arbiter.Core.Job.TraceContext type Rep TraceContext = D1 ('MetaData "TraceContext" "Arbiter.Core.Job.TraceContext" "arbiter-core-0.1.0.0-inplace" 'False) (C1 ('MetaCons "TraceContext" 'PrefixI 'True) (S1 ('MetaSel ('Just "traceparent") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 Text) :*: S1 ('MetaSel ('Just "tracestate") 'NoSourceUnpackedness 'NoSourceStrictness 'DecidedLazy) (Rec0 (Maybe Text)))) | |||||
Source #toTraceContext :: Maybe Text -> Maybe Text -> Maybe TraceContext
A trace context from its two stored halves. An orphan tracestate is dropped.
Observability
Source #data ObservabilityHooks (m :: Type -> Type) payload
Callbacks fired at each point of a job's lifecycle, for metrics, logging or tracing. An exception thrown inside one is caught and dropped.
Constructors
| ObservabilityHooks | |
Fields
| |
Instances
| MonadUnliftIO m => Monoid (ObservabilityHooks m payload) Source # | |
Defined in Arbiter.Core.Job.Types Methods #mempty :: ObservabilityHooks m payload #mappend :: ObservabilityHooks m payload -> ObservabilityHooks m payload -> ObservabilityHooks m payload #mconcat :: [ObservabilityHooks m payload] -> ObservabilityHooks m payload | |
| MonadUnliftIO m => Semigroup (ObservabilityHooks m payload) Source # | Run both hooks at each lifecycle point, left before right. The right one runs however the left ended. When both throw, the right's failure propagates. |
Defined in Arbiter.Core.Job.Types Methods #(<>) :: ObservabilityHooks m payload -> ObservabilityHooks m payload -> ObservabilityHooks m payload #sconcat :: NonEmpty (ObservabilityHooks m payload) -> ObservabilityHooks m payload #stimes :: Integral b => b -> ObservabilityHooks m payload -> ObservabilityHooks m payload | |
Source #defaultObservabilityHooks :: forall (m :: Type -> Type) payload. Applicative m => ObservabilityHooks m payload
No-op hooks. Override fields to add observability:
myHooks = defaultObservabilityHooks
{ onJobSuccess = \job startTime endTime -> do
let duration = diffUTCTime endTime startTime
logInfo $ "Job " <> show (primaryKey job) <> " took " <> show duration
}
Source #andThen :: MonadUnliftIO m => m () -> m () -> m ()
base's finally. The second action stays interruptible.
Source #type CurrentTime = UTCTime
The transaction's clock reading.
Source #type BackoffDelay = NominalDiffTime
How long to wait before the next attempt.