module Ecluse.Service (
ServiceRuntime (..),
withServiceRuntime,
workerLiveness,
runWorker,
mountBindingFor,
) where
import Katip (LogEnv, Namespace (Namespace), SimpleLogPayload, katipAddNamespace, runKatipContextT)
import Network.HTTP.Client (Manager)
import Network.HTTP.Client.TLS (tlsManagerSettings)
import Ecluse.Boot (BootEnv (beLogEnv, beTelemetry), logBootWarning, logRuleBootOrder)
import Ecluse.Composition (BootWiring (bwBindings, bwPublishTargets))
import Ecluse.Composition.Executable (
ExecutablePlan (epBootPlan),
MirrorWiring (mwBootWiring, mwCveSync, mwDeferredMetrics, mwQueue, mwRole),
)
import Ecluse.Composition.MemoryPlan (
MemoryPlan (mpAdmissionCapacity, mpMirrorArtifactTenant, mpTransientBudget),
MirrorArtifactTenant (matMaxBytes),
TransientBudget (tbBootBytes),
mirrorArtifactBytesCap,
)
import Ecluse.Composition.MemoryPlan.Transient (brakeBounds, meterStepBytes)
import Ecluse.Composition.MirrorQueue (MirrorRuntimePlan (MirrorWith, NoMirroring))
import Ecluse.Composition.MirrorRole (enqueuesJobs, spawnsWorker)
import Ecluse.Composition.Plan (
BootPlan (bpCacheConfig, bpMemoryPlan, bpMirrorRuntime, bpPrivateConnections, bpPublicConnections, bpValidated),
)
import Ecluse.Composition.Sizing (newPooledManager)
import Ecluse.Composition.Sizing qualified as Composition
import Ecluse.Composition.Types (MirrorRole)
import Ecluse.Composition.Validate (ValidatedPlan (vpSettings))
import Ecluse.Composition.Worker (workerPoliciesFor)
import Ecluse.Config (AppConfig)
import Ecluse.Core.Credential.Refresh (CredentialError (Unconfigured))
import Ecluse.Core.Ecosystem (Ecosystem, prefixFor)
import Ecluse.Core.Queue (MirrorQueue)
import Ecluse.Core.Queue.Buffer (newEnqueueBuffer, reportWorthy)
import Ecluse.Core.Registry.Adapter (
RegistryAdapter,
adapterEcosystem,
adapterFor,
adapterServe,
serveCredential,
serveRouter,
)
import Ecluse.Core.Server.Admission (newServeAdmission)
import Ecluse.Core.Server.Admission.Brake (defaultBrakeMarks)
import Ecluse.Core.Server.Admission.Meter (MemoryMeter, MeterSettings (..), meterSnapshot, newMemoryMeter)
import Ecluse.Core.Server.Admission.Weighted (admissionWaitMicros)
import Ecluse.Core.Server.Cache (newMetadataCache)
import Ecluse.Core.Server.Context (PackumentDeps, PublishDeps)
import Ecluse.Core.Server.Readiness (Readiness)
import Ecluse.Core.Supervision (
FaultDisposition (Permanent, Transient),
SupervisionPolicy (SupervisionPolicy, spBackoff, spClassify, spLabel),
backgroundLoopBackoff,
superviseLoop,
transientPolicy,
)
import Ecluse.Core.Worker (Liveness, WorkerHeartbeat, WorkerPolicies, alwaysLive, heartbeatLivenessNow, runWorkerM, workerLoop)
import Ecluse.Cve.Sync (cveSyncReadiness, cveSyncScheduleFor, cveSyncTasks, registerAdvisoryAges)
import Ecluse.Rts (cgroupMemoryUse)
import Ecluse.Rts.Sampler (SamplerSettings (..), readCollector, runMemorySampler, samplerPeriodMicros, shareWindowPeriods)
import Ecluse.Runtime.Env (Env, envDdContext, envLogEnv, envMetrics, envTelemetry, newWorkerHeartbeat, withEnvWithAdmission, workerRuntimeOf)
import Ecluse.Runtime.Server (MountBinding (..))
import Ecluse.Runtime.Telemetry (Telemetry)
import Ecluse.Runtime.Telemetry.Correlation (ddPayloadNow)
import Ecluse.Runtime.Telemetry.Instruments (registerMemoryMeter)
import Ecluse.Runtime.Telemetry.Reporters (
DeferredMetrics,
deferredMirrorEnqueueFailure,
installMetrics,
)
import Ecluse.Runtime.Telemetry.Tracing (instrumentDataPlaneManagerSettings)
data ServiceRuntime = ServiceRuntime
{ ServiceRuntime -> MirrorRole
svcRole :: MirrorRole
, ServiceRuntime -> Bool
svcRunsWorker :: Bool
, ServiceRuntime -> Env
svcEnv :: Env
, ServiceRuntime -> AppConfig
svcAppConfig :: AppConfig
, ServiceRuntime -> [MountBinding]
svcBindings :: [MountBinding]
, ServiceRuntime -> WorkerPolicies
svcWorkerPolicies :: WorkerPolicies
, ServiceRuntime -> Maybe (IO ())
svcMirrorDrain :: Maybe (IO ())
, ServiceRuntime -> [IO ()]
svcSyncTasks :: [IO ()]
, ServiceRuntime -> IO ()
svcMemorySampler :: IO ()
, ServiceRuntime -> IO Readiness
svcCheckReady :: IO Readiness
, ServiceRuntime -> IO Liveness
svcCheckLive :: IO Liveness
}
withServiceRuntime :: BootEnv -> ExecutablePlan -> MirrorWiring -> (ServiceRuntime -> IO ()) -> IO ()
withServiceRuntime :: BootEnv
-> ExecutablePlan
-> MirrorWiring
-> (ServiceRuntime -> IO ())
-> IO ()
withServiceRuntime BootEnv
bootEnv ExecutablePlan
plan MirrorWiring
mirror ServiceRuntime -> IO ()
action = do
let logEnv :: LogEnv
logEnv = BootEnv -> LogEnv
beLogEnv BootEnv
bootEnv
telemetry :: Telemetry
telemetry = BootEnv -> Telemetry
beTelemetry BootEnv
bootEnv
bootPlan :: BootPlan
bootPlan = ExecutablePlan -> BootPlan
epBootPlan ExecutablePlan
plan
role :: MirrorRole
role = MirrorWiring -> MirrorRole
mwRole MirrorWiring
mirror
appConfig :: AppConfig
appConfig = ValidatedPlan -> AppConfig
vpSettings (BootPlan -> ValidatedPlan
bpValidated BootPlan
bootPlan)
mirrorRuntime :: MirrorRuntimePlan
mirrorRuntime = BootPlan -> MirrorRuntimePlan
bpMirrorRuntime BootPlan
bootPlan
memoryPlan :: MemoryPlan
memoryPlan = BootPlan -> MemoryPlan
bpMemoryPlan BootPlan
bootPlan
deferredMetrics :: DeferredMetrics
deferredMetrics = MirrorWiring -> DeferredMetrics
mwDeferredMetrics MirrorWiring
mirror
cveSyncPlan :: Map Ecosystem CveSyncHandle
cveSyncPlan = MirrorWiring -> Map Ecosystem CveSyncHandle
mwCveSync MirrorWiring
mirror
bindings :: [MountBinding]
bindings = BootWiring -> [MountBinding]
bwBindings (MirrorWiring -> BootWiring
mwBootWiring MirrorWiring
mirror)
serveAdmission <- Int -> IO ServeAdmission
newServeAdmission (MemoryPlan -> Int
mpAdmissionCapacity MemoryPlan
memoryPlan)
(memoryMeter, sampler) <- memoryAdmission memoryPlan
heartbeat <- newWorkerHeartbeat
let runsWorkerHere = MirrorRole -> MirrorRuntimePlan -> Bool
spawnsWorker MirrorRole
role MirrorRuntimePlan
mirrorRuntime
logRuleBootOrder logEnv bindings
(queue, mirrorDrain) <- mirrorHandOff role logEnv deferredMetrics mirrorRuntime (mwQueue mirror)
metadataCache <- newMetadataCache (bpCacheConfig bootPlan)
(manager, privateManager) <- dataPlaneManagers telemetry bootPlan
withEnvWithAdmission serveAdmission memoryMeter queue manager privateManager metadataCache logEnv telemetry heartbeat $ \Env
builtEnv -> do
DeferredMetrics -> Metrics -> IO ()
installMetrics DeferredMetrics
deferredMetrics (Env -> Metrics
envMetrics Env
builtEnv)
Metrics -> IO MeterSnapshot -> IO ()
registerMemoryMeter (Env -> Metrics
envMetrics Env
builtEnv) (MemoryMeter -> IO MeterSnapshot
meterSnapshot MemoryMeter
memoryMeter)
Metrics -> Map Ecosystem CveSyncHandle -> IO ()
registerAdvisoryAges (Env -> Metrics
envMetrics Env
builtEnv) Map Ecosystem CveSyncHandle
cveSyncPlan
let workerArtifactMaxBytes :: Int
workerArtifactMaxBytes = Int
-> (MirrorArtifactTenant -> Int)
-> Maybe MirrorArtifactTenant
-> Int
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Int
mirrorArtifactBytesCap MirrorArtifactTenant -> Int
matMaxBytes (MemoryPlan -> Maybe MirrorArtifactTenant
mpMirrorArtifactTenant MemoryPlan
memoryPlan)
ServiceRuntime -> IO ()
action
ServiceRuntime
{ svcRole :: MirrorRole
svcRole = MirrorRole
role
, svcRunsWorker :: Bool
svcRunsWorker = Bool
runsWorkerHere
, svcEnv :: Env
svcEnv = Env
builtEnv
, svcAppConfig :: AppConfig
svcAppConfig = AppConfig
appConfig
, svcBindings :: [MountBinding]
svcBindings = [MountBinding]
bindings
, svcWorkerPolicies :: WorkerPolicies
svcWorkerPolicies = Env -> [MountBinding] -> [PublishTarget] -> Int -> WorkerPolicies
workerPoliciesFor Env
builtEnv [MountBinding]
bindings (BootWiring -> [PublishTarget]
bwPublishTargets (MirrorWiring -> BootWiring
mwBootWiring MirrorWiring
mirror)) Int
workerArtifactMaxBytes
, svcMirrorDrain :: Maybe (IO ())
svcMirrorDrain = Text -> Env -> IO () -> IO ()
superviseBackground Text
"mirror-enqueue-drain" Env
builtEnv (IO () -> IO ()) -> Maybe (IO ()) -> Maybe (IO ())
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe (IO ())
mirrorDrain
, svcSyncTasks :: [IO ()]
svcSyncTasks =
LogEnv
-> Metrics
-> Telemetry
-> SyncSchedule
-> Map Ecosystem CveSyncHandle
-> [IO ()]
cveSyncTasks
(Env -> LogEnv
envLogEnv Env
builtEnv)
(Env -> Metrics
envMetrics Env
builtEnv)
(Env -> Telemetry
envTelemetry Env
builtEnv)
(AppConfig -> SyncSchedule
cveSyncScheduleFor AppConfig
appConfig)
Map Ecosystem CveSyncHandle
cveSyncPlan
, svcMemorySampler :: IO ()
svcMemorySampler = Text -> Env -> IO () -> IO ()
superviseBackground Text
"memory-sampler" Env
builtEnv IO ()
sampler
, svcCheckReady :: IO Readiness
svcCheckReady = Map Ecosystem CveSyncHandle -> IO Readiness
cveSyncReadiness Map Ecosystem CveSyncHandle
cveSyncPlan
, svcCheckLive :: IO Liveness
svcCheckLive = Bool -> WorkerHeartbeat -> IO Liveness
workerLiveness Bool
runsWorkerHere WorkerHeartbeat
heartbeat
}
memoryAdmission :: MemoryPlan -> IO (MemoryMeter, IO ())
memoryAdmission :: MemoryPlan -> IO (MemoryMeter, IO ())
memoryAdmission MemoryPlan
memoryPlan = do
meter <-
MeterSettings -> IO MemoryMeter
newMemoryMeter
MeterSettings
{ msBudgetBytes :: Int
msBudgetBytes = TransientBudget -> Int
tbBootBytes (MemoryPlan -> TransientBudget
mpTransientBudget MemoryPlan
memoryPlan)
, msStepBytes :: Int
msStepBytes = Int
meterStepBytes
, msEntryRoom :: Int
msEntryRoom = MemoryPlan -> Int
mpAdmissionCapacity MemoryPlan
memoryPlan
, msEntryWaitMicros :: Int
msEntryWaitMicros = Int
admissionWaitMicros
}
collector <- readCollector
kernelUse <- cgroupMemoryUse
pure
( meter
, runMemorySampler
SamplerSettings
{ ssMeter = meter
, ssMarks = defaultBrakeMarks
, ssBounds = brakeBounds (mpTransientBudget memoryPlan)
, ssCollector = collector
, ssKernelPermille = kernelUse
, ssPeriodMicros = samplerPeriodMicros
, ssWindowPeriods = shareWindowPeriods
}
)
dataPlaneManagers :: Telemetry -> BootPlan -> IO (Manager, Manager)
dataPlaneManagers :: Telemetry -> BootPlan -> IO (Manager, Manager)
dataPlaneManagers Telemetry
telemetry BootPlan
bootPlan = do
publicSettings <- Telemetry -> ManagerSettings -> IO ManagerSettings
instrumentDataPlaneManagerSettings Telemetry
telemetry ManagerSettings
tlsManagerSettings
privateSettings <- instrumentDataPlaneManagerSettings telemetry tlsManagerSettings
manager <- newPooledManager (bpPublicConnections bootPlan) publicSettings
privateManager <- newPooledManager (bpPrivateConnections bootPlan) privateSettings
pure (manager, privateManager)
workerLiveness :: Bool -> WorkerHeartbeat -> IO Liveness
workerLiveness :: Bool -> WorkerHeartbeat -> IO Liveness
workerLiveness Bool
runningWorker WorkerHeartbeat
heartbeat
| Bool
runningWorker = WorkerHeartbeat -> IO Liveness
heartbeatLivenessNow WorkerHeartbeat
heartbeat
| Bool
otherwise = Liveness -> IO Liveness
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Liveness
alwaysLive
mirrorHandOff :: MirrorRole -> LogEnv -> DeferredMetrics -> MirrorRuntimePlan -> MirrorQueue -> IO (MirrorQueue, Maybe (IO ()))
mirrorHandOff :: MirrorRole
-> LogEnv
-> DeferredMetrics
-> MirrorRuntimePlan
-> MirrorQueue
-> IO (MirrorQueue, Maybe (IO ()))
mirrorHandOff MirrorRole
role LogEnv
logEnv DeferredMetrics
deferredMetrics MirrorRuntimePlan
mirrorRuntime MirrorQueue
backendQueue = case MirrorRuntimePlan
mirrorRuntime of
MirrorRuntimePlan
NoMirroring -> (MirrorQueue, Maybe (IO ())) -> IO (MirrorQueue, Maybe (IO ()))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (MirrorQueue
backendQueue, Maybe (IO ())
forall a. Maybe a
Nothing)
MirrorWith MirrorQueuePlan
_
| Bool -> Bool
not (MirrorRole -> Bool
enqueuesJobs MirrorRole
role) -> (MirrorQueue, Maybe (IO ())) -> IO (MirrorQueue, Maybe (IO ()))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (MirrorQueue
backendQueue, Maybe (IO ())
forall a. Maybe a
Nothing)
| Bool
otherwise -> do
(queue, drainEnqueueBuffer) <-
(Text -> IO ()) -> IO () -> MirrorQueue -> IO (MirrorQueue, IO ())
bufferedMirrorHandOff (LogEnv -> Text -> IO ()
logBootWarning LogEnv
logEnv) (DeferredMetrics -> IO ()
deferredMirrorEnqueueFailure DeferredMetrics
deferredMetrics) MirrorQueue
backendQueue
pure (queue, Just drainEnqueueBuffer)
bufferedMirrorHandOff :: (Text -> IO ()) -> IO () -> MirrorQueue -> IO (MirrorQueue, IO ())
bufferedMirrorHandOff :: (Text -> IO ()) -> IO () -> MirrorQueue -> IO (MirrorQueue, IO ())
bufferedMirrorHandOff Text -> IO ()
warn IO ()
countEnqueueFailure =
Int
-> (Int -> IO ())
-> (Int -> Text -> IO ())
-> MirrorQueue
-> IO (MirrorQueue, IO ())
newEnqueueBuffer
Int
Composition.mirrorEnqueueBufferDepth
( \Int
drops -> do
Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Int -> Bool
enqueueReportWorthy Int
drops) (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$
Text -> IO ()
warn (Text
"mirror enqueue buffer full: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show Int
drops Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" job(s) dropped so far; each is re-enqueued on the next demand for its artifact")
IO ()
countEnqueueFailure
)
( \Int
failures Text
detail -> do
Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Int -> Bool
enqueueReportWorthy Int
failures) (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$
Text -> IO ()
warn (Text
"mirror enqueue delivery failed (" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show Int
failures Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" so far): " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
detail)
IO ()
countEnqueueFailure
)
enqueueReportWorthy :: Int -> Bool
enqueueReportWorthy :: Int -> Bool
enqueueReportWorthy Int
n = Int -> Int -> Bool
reportWorthy Int
n Int
Composition.mirrorEnqueueReportInterval
superviseBackground :: Text -> Env -> IO () -> IO ()
superviseBackground :: Text -> Env -> IO () -> IO ()
superviseBackground Text
label Env
builtEnv IO ()
task =
IO Void -> IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO Void -> IO ())
-> (KatipContextT IO Void -> IO Void)
-> KatipContextT IO Void
-> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. LogEnv
-> SimpleLogPayload
-> Namespace
-> KatipContextT IO Void
-> IO Void
forall c (m :: * -> *) a.
LogItem c =>
LogEnv -> c -> Namespace -> KatipContextT m a -> m a
runKatipContextT (Env -> LogEnv
envLogEnv Env
builtEnv) (SimpleLogPayload
forall a. Monoid a => a
mempty :: SimpleLogPayload) ([Text] -> Namespace
Namespace [Text
label]) (KatipContextT IO Void -> IO ()) -> KatipContextT IO Void -> IO ()
forall a b. (a -> b) -> a -> b
$
SupervisionPolicy -> KatipContextT IO () -> KatipContextT IO Void
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SupervisionPolicy -> m () -> m Void
superviseLoop (Text -> BackoffSchedule -> SupervisionPolicy
transientPolicy Text
label BackoffSchedule
backgroundLoopBackoff) (IO () -> KatipContextT IO ()
forall a. IO a -> KatipContextT IO a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO IO ()
task)
runWorker :: WorkerPolicies -> Env -> IO ()
runWorker :: WorkerPolicies -> Env -> IO ()
runWorker WorkerPolicies
policies Env
env = do
dd <- DdContext -> IO SimpleLogPayload
forall (m :: * -> *). MonadIO m => DdContext -> m SimpleLogPayload
ddPayloadNow (Env -> DdContext
envDdContext Env
env)
void (runWorkerM (envLogEnv env) dd (workerRuntimeOf policies env) (katipAddNamespace "worker" (workerLoop workerSupervision)))
workerSupervision :: SupervisionPolicy
workerSupervision :: SupervisionPolicy
workerSupervision =
SupervisionPolicy
{ spLabel :: Text
spLabel = Text
"worker"
, spClassify :: SomeException -> FaultDisposition
spClassify = SomeException -> FaultDisposition
classify
, spBackoff :: BackoffSchedule
spBackoff = BackoffSchedule
backgroundLoopBackoff
}
where
classify :: SomeException -> FaultDisposition
classify SomeException
fault
| Just (Unconfigured Text
_) <- SomeException -> Maybe CredentialError
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
fault = FaultDisposition
Permanent
| Bool
otherwise = FaultDisposition
Transient
mountBindingFor :: Ecosystem -> PackumentDeps -> Maybe PublishDeps -> Maybe MountBinding
mountBindingFor :: Ecosystem
-> PackumentDeps -> Maybe PublishDeps -> Maybe MountBinding
mountBindingFor Ecosystem
eco PackumentDeps
packumentDeps Maybe PublishDeps
publishDeps =
Ecosystem -> Maybe RegistryAdapter
adapterFor Ecosystem
eco Maybe RegistryAdapter
-> (RegistryAdapter -> MountBinding) -> Maybe MountBinding
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> \RegistryAdapter
adapter -> RegistryAdapter
-> PackumentDeps -> Maybe PublishDeps -> MountBinding
mountOf RegistryAdapter
adapter PackumentDeps
packumentDeps Maybe PublishDeps
publishDeps
mountOf :: RegistryAdapter -> PackumentDeps -> Maybe PublishDeps -> MountBinding
mountOf :: RegistryAdapter
-> PackumentDeps -> Maybe PublishDeps -> MountBinding
mountOf RegistryAdapter
adapter PackumentDeps
packumentDeps Maybe PublishDeps
publishDeps =
MountBinding
{ bindingPrefix :: NonEmpty Text
bindingPrefix = Ecosystem -> NonEmpty Text
prefixFor (RegistryAdapter -> Ecosystem
adapterEcosystem RegistryAdapter
adapter)
, bindingRouter :: MountRouter
bindingRouter = AdapterServe -> MountRouter
serveRouter (RegistryAdapter -> AdapterServe
adapterServe RegistryAdapter
adapter)
, bindingCredential :: CredentialMapping
bindingCredential = AdapterServe -> CredentialMapping
serveCredential (RegistryAdapter -> AdapterServe
adapterServe RegistryAdapter
adapter)
, bindingPackumentDeps :: PackumentDeps
bindingPackumentDeps = PackumentDeps
packumentDeps
, bindingPublishDeps :: Maybe PublishDeps
bindingPublishDeps = Maybe PublishDeps
publishDeps
}