module Ecluse.Pilot (
runPilot,
runExportLoop,
superviseExportCycles,
PilotCompileOptions (..),
runPilotCompile,
PilotUploadUnconfigured (..),
) where
import Conduit (MonadResource, runResourceT)
import Control.Monad.Catch (MonadMask)
import Katip (KatipContext, LogEnv, Severity (InfoS, WarningS), logFM, ls)
import UnliftIO (MonadUnliftIO)
import UnliftIO.Async (mapConcurrently_)
import UnliftIO.Concurrent (threadDelay)
import UnliftIO.Exception (throwIO)
import Ecluse.Boot (BootEnv (..), probeServerConfig)
import Ecluse.Composition.Plan (BootPlan (bpS3Endpoint))
import Ecluse.Config (
AdvisoriesSettings (advDataDir, advUrl),
AdvisoryStoreUrl,
AppConfig (cfgAdvisories),
Config (configApp, configMounts),
advisoryStoreUrlText,
)
import Ecluse.Core.Ecosystem (ecosystemName, parseEcosystem)
import Ecluse.Core.Osv.Compile (compileOsvToSqlite)
import Ecluse.Core.Osv.Ecosystem (osvEcosystemFor, osvEcosystemNamed)
import Ecluse.Core.Supervision (
BackoffSchedule (BackoffSchedule, bsBaseMicros, bsCapMicros),
superviseLoop,
transientPolicy,
)
import Ecluse.Pilot.Plan (
ExportLoopPlan (ExportIdle, ExportTo),
ExportTarget (etEcosystem, etEpss),
PilotCompileOptions (..),
PilotUploadUnconfigured (..),
UploadPlan (UploadSkipped, UploadTo),
compileEpssRequirement,
configuredSources,
exportCadenceMicros,
idleCadenceMicros,
quietTimeFor,
unmountedCompileWarning,
uploadPlan,
uploadTarget,
)
import Ecluse.Runtime.Aws.Env (AwsEndpoint)
import Ecluse.Runtime.Log (moduleContext)
import Ecluse.Runtime.Pilot.Export (exportToS3)
import Ecluse.Runtime.Server (probeOnlyApplication, raceServerAgainstLoop, runWarp)
import Ecluse.Runtime.Telemetry (Telemetry, telemetryTracerProvider)
import Ecluse.Runtime.Telemetry.Instruments (Metrics, advisoryCompileMetricsPortOf, newMetrics)
runPilot :: BootEnv -> ExportLoopPlan -> IO ()
runPilot :: BootEnv -> ExportLoopPlan -> IO ()
runPilot BootEnv
bootEnv ExportLoopPlan
exportPlan = do
let cfg :: ServerConfig
cfg = AppConfig -> ServerConfig
probeServerConfig (Config -> AppConfig
configApp (BootEnv -> Config
beConfig BootEnv
bootEnv))
LogEnv -> Text -> KatipContextT IO () -> IO ()
forall (m :: * -> *) a. LogEnv -> Text -> KatipContextT m a -> m a
moduleContext (BootEnv -> LogEnv
beLogEnv BootEnv
bootEnv) Text
"Ecluse.Pilot" (KatipContextT IO () -> IO ()) -> KatipContextT IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
Severity -> LogStr -> KatipContextT IO ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
InfoS LogStr
"Pilot mode starting up"
KatipContextT IO () -> KatipContextT IO () -> KatipContextT IO ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
raceServerAgainstLoop
(IO () -> KatipContextT IO ()
forall a. IO a -> KatipContextT IO a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> KatipContextT IO ()) -> IO () -> KatipContextT IO ()
forall a b. (a -> b) -> a -> b
$ LogEnv
-> Text
-> ServerConfig
-> (ServerConfig -> IO Application)
-> IO ()
runWarp (BootEnv -> LogEnv
beLogEnv BootEnv
bootEnv) Text
"Pilot health probes" ServerConfig
cfg ServerConfig -> IO Application
probeOnlyApplication)
(Telemetry
-> Maybe AwsEndpoint
-> Config
-> ExportLoopPlan
-> KatipContextT IO ()
forall (m :: * -> *).
(MonadMask m, MonadUnliftIO m, KatipContext m) =>
Telemetry -> Maybe AwsEndpoint -> Config -> ExportLoopPlan -> m ()
runExportLoop (BootEnv -> Telemetry
beTelemetry BootEnv
bootEnv) (BootPlan -> Maybe AwsEndpoint
bpS3Endpoint (BootEnv -> BootPlan
beBootPlan BootEnv
bootEnv)) (BootEnv -> Config
beConfig BootEnv
bootEnv) ExportLoopPlan
exportPlan)
runExportLoop :: (MonadMask m, MonadUnliftIO m, KatipContext m) => Telemetry -> Maybe AwsEndpoint -> Config -> ExportLoopPlan -> m ()
runExportLoop :: forall (m :: * -> *).
(MonadMask m, MonadUnliftIO m, KatipContext m) =>
Telemetry -> Maybe AwsEndpoint -> Config -> ExportLoopPlan -> m ()
runExportLoop Telemetry
telemetry Maybe AwsEndpoint
s3Endpoint Config
config = \case
ExportLoopPlan
ExportIdle -> do
Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
InfoS LogStr
"No advisory store configured for OSV database export; export loop disabled."
m () -> m ()
forall (f :: * -> *) a b. Applicative f => f a -> f b
forever (Int -> m ()
forall (m :: * -> *). MonadIO m => Int -> m ()
threadDelay Int
idleCadenceMicros)
ExportTo AdvisoryStoreUrl
store NonEmpty ExportTarget
targets -> do
Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
InfoS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"Export loop starting up. Target store: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> AdvisoryStoreUrl -> Text
advisoryStoreUrlText AdvisoryStoreUrl
store))
metrics <- IO Metrics -> m Metrics
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Telemetry -> IO Metrics
newMetrics Telemetry
telemetry)
superviseExportCycles schedule targets $ \ExportTarget
target -> do
ResourceT m () -> m ()
forall (m :: * -> *) a. MonadUnliftIO m => ResourceT m a -> m a
runResourceT (Metrics
-> ExportTarget
-> Telemetry
-> Maybe AwsEndpoint
-> AdvisoriesSettings
-> AdvisoryStoreUrl
-> ResourceT m ()
forall (m :: * -> *).
(MonadResource m, MonadMask m, MonadUnliftIO m, KatipContext m) =>
Metrics
-> ExportTarget
-> Telemetry
-> Maybe AwsEndpoint
-> AdvisoriesSettings
-> AdvisoryStoreUrl
-> m ()
exportEcosystem Metrics
metrics ExportTarget
target Telemetry
telemetry Maybe AwsEndpoint
s3Endpoint AdvisoriesSettings
advisories AdvisoryStoreUrl
store)
Int -> m ()
forall (m :: * -> *). MonadIO m => Int -> m ()
threadDelay Int
cadence
where
advisories :: AdvisoriesSettings
advisories = AppConfig -> AdvisoriesSettings
cfgAdvisories (Config -> AppConfig
configApp Config
config)
cadence :: Int
cadence = AdvisoriesSettings -> Int
exportCadenceMicros AdvisoriesSettings
advisories
schedule :: BackoffSchedule
schedule = BackoffSchedule{bsBaseMicros :: Int
bsBaseMicros = Int
cadence, bsCapMicros :: Int
bsCapMicros = Int
cadence}
superviseExportCycles :: (MonadUnliftIO m, KatipContext m) => BackoffSchedule -> NonEmpty ExportTarget -> (ExportTarget -> m ()) -> m ()
superviseExportCycles :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
BackoffSchedule
-> NonEmpty ExportTarget -> (ExportTarget -> m ()) -> m ()
superviseExportCycles BackoffSchedule
schedule NonEmpty ExportTarget
targets ExportTarget -> m ()
runCycle =
(ExportTarget -> m Void) -> NonEmpty ExportTarget -> m ()
forall (m :: * -> *) (f :: * -> *) a b.
(MonadUnliftIO m, Foldable f) =>
(a -> m b) -> f a -> m ()
mapConcurrently_ (\ExportTarget
target -> SupervisionPolicy -> m () -> m Void
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SupervisionPolicy -> m () -> m Void
superviseLoop (ExportTarget -> SupervisionPolicy
policyFor ExportTarget
target) (ExportTarget -> m ()
runCycle ExportTarget
target)) NonEmpty ExportTarget
targets
where
policyFor :: ExportTarget -> SupervisionPolicy
policyFor ExportTarget
target = Text -> BackoffSchedule -> SupervisionPolicy
transientPolicy (Text
"pilot-export-" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Ecosystem -> Text
ecosystemName (ExportTarget -> Ecosystem
etEcosystem ExportTarget
target)) BackoffSchedule
schedule
exportEcosystem :: (MonadResource m, MonadMask m, MonadUnliftIO m, KatipContext m) => Metrics -> ExportTarget -> Telemetry -> Maybe AwsEndpoint -> AdvisoriesSettings -> AdvisoryStoreUrl -> m ()
exportEcosystem :: forall (m :: * -> *).
(MonadResource m, MonadMask m, MonadUnliftIO m, KatipContext m) =>
Metrics
-> ExportTarget
-> Telemetry
-> Maybe AwsEndpoint
-> AdvisoriesSettings
-> AdvisoryStoreUrl
-> m ()
exportEcosystem Metrics
metrics ExportTarget
target Telemetry
telemetry Maybe AwsEndpoint
s3Endpoint AdvisoriesSettings
advisories AdvisoryStoreUrl
store = do
Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
InfoS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"Starting " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Ecosystem -> Text
ecosystemName Ecosystem
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" OSV database compilation"))
dbPath <-
AdvisoryCompileMetricsPort
-> Maybe TracerProvider
-> String
-> OsvEcosystem
-> EpssRequirement
-> CompileSources
-> QuietTime
-> m String
forall (m :: * -> *).
(MonadResource m, MonadMask m, MonadUnliftIO m, KatipContext m) =>
AdvisoryCompileMetricsPort
-> Maybe TracerProvider
-> String
-> OsvEcosystem
-> EpssRequirement
-> CompileSources
-> QuietTime
-> m String
compileOsvToSqlite
(Metrics -> Maybe Ecosystem -> AdvisoryCompileMetricsPort
advisoryCompileMetricsPortOf Metrics
metrics (Ecosystem -> Maybe Ecosystem
forall a. a -> Maybe a
Just Ecosystem
eco))
(Telemetry -> Maybe TracerProvider
telemetryTracerProvider Telemetry
telemetry)
(AdvisoriesSettings -> String
advDataDir AdvisoriesSettings
advisories)
OsvEcosystem
osvEco
(ExportTarget -> EpssRequirement
etEpss ExportTarget
target)
(AdvisoriesSettings -> OsvEcosystem -> CompileSources
configuredSources AdvisoriesSettings
advisories OsvEcosystem
osvEco)
(AdvisoriesSettings -> Maybe Ecosystem -> QuietTime
quietTimeFor AdvisoriesSettings
advisories (Ecosystem -> Maybe Ecosystem
forall a. a -> Maybe a
Just Ecosystem
eco))
uploadToStore telemetry s3Endpoint store dbPath
where
eco :: Ecosystem
eco = ExportTarget -> Ecosystem
etEcosystem ExportTarget
target
osvEco :: OsvEcosystem
osvEco = Ecosystem -> OsvEcosystem
osvEcosystemFor Ecosystem
eco
uploadToStore :: (MonadResource m, MonadUnliftIO m, MonadMask m, KatipContext m) => Telemetry -> Maybe AwsEndpoint -> AdvisoryStoreUrl -> FilePath -> m ()
uploadToStore :: forall (m :: * -> *).
(MonadResource m, MonadUnliftIO m, MonadMask m, KatipContext m) =>
Telemetry
-> Maybe AwsEndpoint -> AdvisoryStoreUrl -> String -> m ()
uploadToStore Telemetry
telemetry Maybe AwsEndpoint
s3Endpoint AdvisoryStoreUrl
store String
dbPath =
Maybe TracerProvider
-> Maybe AwsEndpoint -> Text -> Text -> String -> m ()
forall (m :: * -> *).
(MonadResource m, MonadUnliftIO m, MonadThrow m, KatipContext m) =>
Maybe TracerProvider
-> Maybe AwsEndpoint -> Text -> Text -> String -> m ()
exportToS3 (Telemetry -> Maybe TracerProvider
telemetryTracerProvider Telemetry
telemetry) Maybe AwsEndpoint
s3Endpoint Text
bucket Text
key String
dbPath
where
(Text
bucket, Text
key) = AdvisoryStoreUrl -> String -> (Text, Text)
uploadTarget AdvisoryStoreUrl
store String
dbPath
runUploadPlan :: (MonadResource m, MonadUnliftIO m, MonadMask m, KatipContext m) => Telemetry -> Maybe AwsEndpoint -> UploadPlan -> FilePath -> m ()
runUploadPlan :: forall (m :: * -> *).
(MonadResource m, MonadUnliftIO m, MonadMask m, KatipContext m) =>
Telemetry -> Maybe AwsEndpoint -> UploadPlan -> String -> m ()
runUploadPlan Telemetry
telemetry Maybe AwsEndpoint
s3Endpoint UploadPlan
plan String
dbPath = case UploadPlan
plan of
UploadPlan
UploadSkipped -> m ()
forall (f :: * -> *). Applicative f => f ()
pass
UploadTo AdvisoryStoreUrl
store -> Telemetry
-> Maybe AwsEndpoint -> AdvisoryStoreUrl -> String -> m ()
forall (m :: * -> *).
(MonadResource m, MonadUnliftIO m, MonadMask m, KatipContext m) =>
Telemetry
-> Maybe AwsEndpoint -> AdvisoryStoreUrl -> String -> m ()
uploadToStore Telemetry
telemetry Maybe AwsEndpoint
s3Endpoint AdvisoryStoreUrl
store String
dbPath
runPilotCompile :: LogEnv -> Telemetry -> Maybe AwsEndpoint -> Config -> PilotCompileOptions -> IO FilePath
runPilotCompile :: LogEnv
-> Telemetry
-> Maybe AwsEndpoint
-> Config
-> PilotCompileOptions
-> IO String
runPilotCompile LogEnv
logEnv Telemetry
telemetry Maybe AwsEndpoint
s3Endpoint Config
config PilotCompileOptions
opts = do
metrics <- Telemetry -> IO Metrics
newMetrics Telemetry
telemetry
let compileMetrics = Metrics -> Maybe Ecosystem -> AdvisoryCompileMetricsPort
advisoryCompileMetricsPortOf Metrics
metrics (Text -> Maybe Ecosystem
parseEcosystem (PilotCompileOptions -> Text
pcoEcosystem PilotCompileOptions
opts))
moduleContext logEnv "Ecluse.Pilot" $
runResourceT $ do
plan <- either throwIO pure planned
traverse_ (logFM WarningS . ls) (unmountedCompileWarning (configMounts config) opts)
dbFile <-
compileOsvToSqlite
compileMetrics
(telemetryTracerProvider telemetry)
(pcoOutDir opts)
osvEco
(compileEpssRequirement (configMounts config) opts)
(configuredSources advisories osvEco)
(quietTimeFor advisories (parseEcosystem (pcoEcosystem opts)))
runUploadPlan telemetry s3Endpoint plan dbFile
pure dbFile
where
advisories :: AdvisoriesSettings
advisories = AppConfig -> AdvisoriesSettings
cfgAdvisories (Config -> AppConfig
configApp Config
config)
osvEco :: OsvEcosystem
osvEco = Text -> OsvEcosystem
osvEcosystemNamed (PilotCompileOptions -> Text
pcoEcosystem PilotCompileOptions
opts)
planned :: Either PilotUploadUnconfigured UploadPlan
planned = PilotCompileOptions
-> Maybe AdvisoryStoreUrl
-> Either PilotUploadUnconfigured UploadPlan
uploadPlan PilotCompileOptions
opts (AdvisoriesSettings -> Maybe AdvisoryStoreUrl
advUrl AdvisoriesSettings
advisories)