| Safe Haskell | None |
|---|---|
| Language | GHC2021 |
Ecluse.Core.Worker
Description
The mirror worker's public surface over the supervised loop that turns enqueued jobs into mirrored packages.
The loop long-polls the demand-driven mirror queue (Ecluse.Core.Queue) and resolves each job's
ecosystem bundle (WorkerPolicy). A job is fail-closed when its ecosystem carries no bundle, and
when its name is one the deployment owns: the queue outlives a namespace declaration, so that
privilege is read here before any public request. The per-job decision, the digest gate, the
receipt lease, verdict realisation and supervision live in the child modules re-exported below.
Synopsis
- data WorkerRuntime = WorkerRuntime {
- wrQueue :: MirrorQueue
- wrManager :: Manager
- wrHeartbeat :: WorkerHeartbeat
- wrMetrics :: WorkerMetricsPort
- wrTracing :: WorkerTracingPort
- wrInjectTraceContext :: forall (m :: Type -> Type) a. (KatipContext m, MonadIO m) => m a -> m a
- wrPolicies :: WorkerPolicies
- data WorkerPolicy = WorkerPolicy {
- wpFirstParty :: PackageName -> Bool
- wpResolveVersion :: PackageName -> Version -> IO VersionEvaluation
- wpRules :: [PreparedRule]
- wpMinIntegrity :: MinIntegrity
- wpArtifactHostHonoured :: Maybe HostPort -> Bool
- wpArtifact :: AdapterArtifact
- wpPublish :: MirrorPublish
- wpArtifactLimits :: Limits
- wpNow :: IO UTCTime
- type WorkerPolicies = Map Ecosystem WorkerPolicy
- data WorkerM a
- runWorkerM :: LogEnv -> SimpleLogPayload -> WorkerRuntime -> WorkerM a -> IO a
- workerLoop :: SupervisionPolicy -> WorkerM Void
- processBatch :: [QueueMessage] -> WorkerM ()
- processJob :: MirrorJob -> WorkerM JobOutcome
- data JobOutcome
- data RetryLeg
- data WorkerHeartbeat
- newWorkerHeartbeat :: IO WorkerHeartbeat
- recordPoll :: WorkerHeartbeat -> UTCTime -> IO ()
- lastPoll :: WorkerHeartbeat -> IO (Maybe UTCTime)
- workerJobStepAllowance :: NominalDiffTime
- workerHeartbeatStaleAfter :: NominalDiffTime
- heartbeatHealthy :: UTCTime -> UTCTime -> Maybe UTCTime -> Bool
- data Liveness = Liveness {}
- alwaysLive :: Liveness
- heartbeatLivenessNow :: WorkerHeartbeat -> IO Liveness
- data IntegrityResult
- verifyIntegrity :: NonEmpty Hash -> ByteString -> IntegrityResult
Worker runtime
data WorkerRuntime Source #
The effectful backends the mirror worker closes over, read through the WorkerM reader.
Constructors
| WorkerRuntime | |
Fields
| |
Instances
| MonadReader WorkerRuntime WorkerM Source # | |
Defined in Ecluse.Core.Worker.Types Methods ask :: WorkerM WorkerRuntime # local :: (WorkerRuntime -> WorkerRuntime) -> WorkerM a -> WorkerM a # reader :: (WorkerRuntime -> a) -> WorkerM a # | |
Per-ecosystem ingest re-evaluation
data WorkerPolicy Source #
The per-ecosystem bundle every job is dispatched through. Its resolver, rules, and gates are the serve path's own, so ingest and serve reach one decision over one policy.
Constructors
| WorkerPolicy | |
Fields
| |
type WorkerPolicies = Map Ecosystem WorkerPolicy Source #
The bundles keyed by a job's package ecosystem. A job whose ecosystem is absent is fail-closed: dropped, never mirrored unvetted.
The worker monad
A reader over the WorkerRuntime on katip's logging context. That base is a reader,
never a StateT, so the context behaves across the loop.
Instances
| MonadIO WorkerM Source # | |
Defined in Ecluse.Core.Worker.Types | |
| Applicative WorkerM Source # | |
| Functor WorkerM Source # | |
| Monad WorkerM Source # | |
| Katip WorkerM Source # | |
| KatipContext WorkerM Source # | |
Defined in Ecluse.Core.Worker.Types Methods getKatipContext :: WorkerM LogContexts # localKatipContext :: (LogContexts -> LogContexts) -> WorkerM a -> WorkerM a # getKatipNamespace :: WorkerM Namespace # localKatipNamespace :: (Namespace -> Namespace) -> WorkerM a -> WorkerM a # | |
| MonadUnliftIO WorkerM Source # | |
Defined in Ecluse.Core.Worker.Types | |
| MonadReader WorkerRuntime WorkerM Source # | |
Defined in Ecluse.Core.Worker.Types Methods ask :: WorkerM WorkerRuntime # local :: (WorkerRuntime -> WorkerRuntime) -> WorkerM a -> WorkerM a # reader :: (WorkerRuntime -> a) -> WorkerM a # | |
runWorkerM :: LogEnv -> SimpleLogPayload -> WorkerRuntime -> WorkerM a -> IO a Source #
Run a WorkerM at the caller's katip environment, so the application owns the log
stream and the trace-correlation identity every line carries.
Loop and job processing (exposed for direct testing)
workerLoop :: SupervisionPolicy -> WorkerM Void Source #
The continuous consume loop: long-poll, process, repeat, under the supervision policy. The
heartbeat advances only on progress, so a persistently faulting receive goes stale on /livez.
processBatch :: [QueueMessage] -> WorkerM () Source #
Process one batch sequentially under a lease on every receipt in it. The heartbeat advances
per job, so workerHeartbeatStaleAfter covers one job.
processJob :: MirrorJob -> WorkerM JobOutcome Source #
Decide one job, re-checking current policy before publishing, because the queue wait is unbounded and mirrored bytes bypass every later rule.
data JobOutcome Source #
The terminal outcome of processing one mirror job. It decides whether the worker acks the message or leaves it to redeliver.
Constructors
| Succeeded | The publish succeeded or the mirror already held the version. The worker acknowledges either result, including idempotent redelivery. |
| Dropped Text | A non-retryable rejection (a tampered artifact, an unformable request URL). Redelivery cannot help, so the job is acked to retire it after alarming. |
| SourceUnavailable Text | The source's own version object was unavailable, so no mirror write can reflect the
package. Retired like |
| DeadLettered Text | A terminal fault handed to |
| Retried RetryLeg Text | A transient fault: a fetch failure, or a registry rejection worth retrying. The message is left un-acked so it redelivers, carrying the leg it gave up on. |
Instances
| Show JobOutcome Source # | |
Defined in Ecluse.Core.Worker.Job Methods showsPrec :: Int -> JobOutcome -> ShowS # show :: JobOutcome -> String # showList :: [JobOutcome] -> ShowS # | |
| Eq JobOutcome Source # | |
Defined in Ecluse.Core.Worker.Job | |
Which leg a transient failure gave up on. The realisation half reads it to decide whether to reset the message's visibility, so the two legs cannot be conflated at the queue handle.
Constructors
| BeforePublish | The job gave up before it published: the inventory probe, the re-evaluation, or the artifact fetch. The message keeps its lease and redelivers when that window lapses. |
| AfterPublish | The publish itself failed transiently, after the bytes were fetched and verified. The message is released so its redelivery does not wait out the lease. |
Liveness
data WorkerHeartbeat Source #
Worker progress and its startup allowance, separate from HTTP readiness.
newWorkerHeartbeat :: IO WorkerHeartbeat Source #
Build a fresh WorkerHeartbeat with no poll yet recorded (lastPoll is
Nothing until the worker's first successful receive).
recordPoll :: WorkerHeartbeat -> UTCTime -> IO () Source #
Stamp the heartbeat with the given instant, recording a unit of worker progress.
The worker calls it through recordWorkerProgress.
lastPoll :: WorkerHeartbeat -> IO (Maybe UTCTime) Source #
The instant of the worker's last recorded progress, a successful poll or a completed
job, or Nothing before its first.
workerJobStepAllowance :: NominalDiffTime Source #
How long one long job step may run: uploading the largest artifact the memory plan admits (512 MiB) over a 2 MiB-per-second link. Renewing a receipt's visibility never extends it.
workerHeartbeatStaleAfter :: NominalDiffTime Source #
Startup and progress allowance: a fetch and a publish of that largest artifact plus a minute, so a slow job never restarts the process.
heartbeatHealthy :: UTCTime -> UTCTime -> Maybe UTCTime -> Bool Source #
Judge progress at now, using startup time only until the first successful progress.
What /livez answers from: the health verdict, plus the instant the checked loop last
recorded progress so an orchestrator can judge staleness rather than only pass or fail.
Constructors
| Liveness | |
Fields
| |
alwaysLive :: Liveness Source #
The verdict of a role with no background loop to stall: live, with nothing to report.
heartbeatLivenessNow :: WorkerHeartbeat -> IO Liveness Source #
Read the worker heartbeat and judge it against the current wall clock, keeping the
instant judged. Both the embedded and the dedicated worker answer /livez through this.
Integrity verification
data IntegrityResult Source #
Whether fetched bytes may enter the mirror, with a refusal detail for the operator.
Constructors
| IntegrityVerified | The bytes matched the selected digest or one of its SRI alternatives. |
| IntegrityMismatch Text | The selected digest mismatched or its algorithm cannot be computed. |
Instances
| Show IntegrityResult Source # | |
Defined in Ecluse.Core.Worker.Integrity Methods showsPrec :: Int -> IntegrityResult -> ShowS # show :: IntegrityResult -> String # showList :: [IntegrityResult] -> ShowS # | |
| Eq IntegrityResult Source # | |
Defined in Ecluse.Core.Worker.Integrity Methods (==) :: IntegrityResult -> IntegrityResult -> Bool # (/=) :: IntegrityResult -> IntegrityResult -> Bool # | |
verifyIntegrity :: NonEmpty Hash -> ByteString -> IntegrityResult Source #
Verify the strongest selected digest, allowing only its same-algorithm SRI alternatives. A weaker match or an uncomputable selected algorithm never permits publication.