-- SPDX-FileCopyrightText: 2026 Alexandra de Wit
--
-- SPDX-License-Identifier: MIT

{- | The Pilot role's front door. Every decision it acts on is planned before it starts, in
"Ecluse.Pilot.Plan" and the boot arm that spends the plan's refusal, so what is left here is
the effect each plan names.
-}
module Ecluse.Pilot (
    runPilot,
    runExportLoop,
    superviseExportCycles,

    -- * One-shot compilation
    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)

{- | The entry point for the Pilot worker mode. The export loop never returns, so the
server's graceful return on shutdown must cancel it, resuming from the remote artifact next boot.
-}
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)

{- | Run the loop the boot planned, never returning. Every fault inside a cycle is transient,
because a cycle has no wiring fault to fail up on.
-}
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))
        -- One instrument set for the whole loop. Rebuilding it per cycle would register
        -- the catalogue again and split each signal across two streams.
        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}

{- | Run one supervised cycle loop per ecosystem, side by side. Each keeps its own backoff, so a
feed that fails costs its own cadence and holds back no other ecosystem's artifact.
-}
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

-- One full cycle for one ecosystem: compile its OSV artifact and upload it.
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

-- Upload one artifact to the address 'uploadTarget' derives, so it lands where the sync reads.
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

-- Act on a planned upload. A skipped one is the caller's success path, not an error.
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

{- | Run a single OSV compilation, optionally upload the artifact, and return its path. An OSV
failure or a required EPSS failure propagates, so the command exits non-zero and stays scriptable.
-}
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
    -- The metric label domain is the closed 'Ecosystem' enum. A one-shot compile of a name
    -- outside it still writes its artifact, and records no series.
    let compileMetrics = Metrics -> Maybe Ecosystem -> AdvisoryCompileMetricsPort
advisoryCompileMetricsPortOf Metrics
metrics (Text -> Maybe Ecosystem
parseEcosystem (PilotCompileOptions -> Text
pcoEcosystem PilotCompileOptions
opts))
    moduleContext logEnv "Ecluse.Pilot" $
        runResourceT $ do
            -- Raised before the compile, so an upload the wiring cannot satisfy fails
            -- without first compiling the artifact it could never publish.
            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)