ecluse:ecluse-core
Safe HaskellNone
LanguageGHC2021

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

Worker runtime

data WorkerRuntime Source #

The effectful backends the mirror worker closes over, read through the WorkerM reader.

Constructors

WorkerRuntime 

Fields

Instances

Instances details
MonadReader WorkerRuntime WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

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

  • wpFirstParty :: PackageName -> Bool

    Whether a name belongs to a namespace this deployment owns.

  • wpResolveVersion :: PackageName -> Version -> IO VersionEvaluation

    Resolve one version's metadata through the guarded public origin. Total by type: every failure, transport included, classifies as a VersionMetadataUnavailable value.

  • wpRules :: [PreparedRule]

    The prepared rule set re-evaluated against the resolved version.

  • wpMinIntegrity :: MinIntegrity

    The mount's public-integrity floor, re-applied through the shared admission gate.

  • wpArtifactHostHonoured :: Maybe HostPort -> Bool

    The mount's tarball-host gate, re-checked on the job's fetch URL. An unextractable authority (Nothing) is refused, because the queue payload is a trust boundary.

  • wpArtifact :: AdapterArtifact

    The mount ecosystem's artifact capability. A job's GET rides its by-URL member.

  • wpPublish :: MirrorPublish

    The mirror write bound to the mount's declared target, so a job's presence probe and publish reach only its own ecosystem's mirror.

  • wpArtifactLimits :: Limits

    The bounded-fetch budget for the artifact download, set from the memory plan's mirror-artifact tenant so a publish envelope cannot breach the heap ceiling.

  • wpNow :: IO UTCTime

    Wall-clock now for the rules, injected so the quarantine rule is deterministic.

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

data WorkerM a Source #

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

Instances details
MonadIO WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

Methods

liftIO :: IO a -> WorkerM a #

Applicative WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

Methods

pure :: a -> WorkerM a #

(<*>) :: WorkerM (a -> b) -> WorkerM a -> WorkerM b #

liftA2 :: (a -> b -> c) -> WorkerM a -> WorkerM b -> WorkerM c #

(*>) :: WorkerM a -> WorkerM b -> WorkerM b #

(<*) :: WorkerM a -> WorkerM b -> WorkerM a #

Functor WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

Methods

fmap :: (a -> b) -> WorkerM a -> WorkerM b #

(<$) :: a -> WorkerM b -> WorkerM a #

Monad WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

Methods

(>>=) :: WorkerM a -> (a -> WorkerM b) -> WorkerM b #

(>>) :: WorkerM a -> WorkerM b -> WorkerM b #

return :: a -> WorkerM a #

Katip WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

KatipContext WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

MonadUnliftIO WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

Methods

withRunInIO :: ((forall a. WorkerM a -> IO a) -> IO b) -> WorkerM b #

MonadReader WorkerRuntime WorkerM Source # 
Instance details

Defined in Ecluse.Core.Worker.Types

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 Dropped, but reported apart from a policy deny.

DeadLettered Text

A terminal fault handed to deadLetter rather than acked, because a plain delete would silently discard it on a durable queue.

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

Instances details
Show JobOutcome Source # 
Instance details

Defined in Ecluse.Core.Worker.Job

Eq JobOutcome Source # 
Instance details

Defined in Ecluse.Core.Worker.Job

data RetryLeg Source #

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.

Instances

Instances details
Show RetryLeg Source # 
Instance details

Defined in Ecluse.Core.Worker.Job

Eq RetryLeg Source # 
Instance details

Defined in Ecluse.Core.Worker.Job

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.

data Liveness Source #

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

Instances

Instances details
Show Liveness Source # 
Instance details

Defined in Ecluse.Core.Worker.Liveness

Eq Liveness Source # 
Instance details

Defined in Ecluse.Core.Worker.Liveness

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.

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.