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

{- | The advisory-sync plan: one 'CveSyncHandle' per mount ecosystem ('planCveSync'), the
projections the composition root reads off it, and one supervised sync task per handle.
-}
module Ecluse.Cve.Sync (
    CveSyncHandle (..),
    AdvisoryNeed (..),
    planCveSync,
    sweepStaleTemps,
    sweepStep,
    cveRuleDepsFor,
    advisoryFreshnessFor,
    reportPushAge,
    katipOutageReporter,
    outageReportPeriod,
    cveSyncReadiness,
    cveSyncScheduleFor,
    cveSyncTasks,
    registerAdvisoryAges,
) where

import Data.Map.Strict qualified as Map
import Data.Text qualified as T
import Data.Time (NominalDiffTime, UTCTime, getCurrentTime)
import Katip (LogEnv, Severity (ErrorS, InfoS, WarningS), SimpleLogPayload, runKatipContextT, sl)
import System.Directory (createDirectoryIfMissing, listDirectory, removeFile)
import System.FilePath (isExtensionOf, (</>))
import System.IO.Error (IOError, catchIOError)

import Ecluse.Config (
    AdvisoriesSettings (advDataDir, advPollInterval, advUrl),
    AdvisoryStoreUrl,
    AppConfig (cfgAdvisories, cfgLimits),
    LimitsSettings (limMaxAdvisoryDatabaseBytes),
    advisoryObjectKey,
    advisoryStoreBucket,
    advisoryStoreUrlText,
 )
import Ecluse.Core.Breaker (BreakerReporter)
import Ecluse.Core.Clock (secondsToMicros)
import Ecluse.Core.Cve.Slot (AdvisorySource (asPushedAt), CveSlot, currentAdvisoryEtag, currentAdvisorySource, generationInstalledAt, newCveSlot, withSlotGeneration)
import Ecluse.Core.Ecosystem (Ecosystem, ecosystemName)
import Ecluse.Core.Osv.Schema (EpssRequirement, osvDbFileName)
import Ecluse.Core.Rules (AdvisoryDatabase (AdvisoryDatabase, NoAdvisoryDatabase), RuleDeps (..), SourceReporter, noSourceReporter)
import Ecluse.Core.Rules.Freshness (
    AdvisoryAge (advisoryAge, advisoryMaxAge, advisoryPushedAt),
    AdvisoryFreshness (AdvisoryFresh),
    AdvisoryPublication (NoGeneration, PublishedAt, UndatedGeneration),
    MaxAdvisoryAge,
    ageAlarmStep,
    assessAdvisoryAge,
 )
import Ecluse.Core.Rules.Outage (OutageReport (..), OutageState (Healthy), sourceReporter, tvarOutageStore)
import Ecluse.Core.Server.Readiness (
    DatabaseRequirement,
    MountReadiness,
    Readiness,
    mountReadiness,
    mountStateFor,
 )
import Ecluse.Core.Supervision (
    backgroundLoopBackoff,
    superviseLoop,
    transientPolicy,
 )
import Ecluse.Core.Text (renderIso8601Utc)
import Ecluse.Runtime.Aws.Env (AwsEndpoint)
import Ecluse.Runtime.Cve.Sync (S3CveSource, SyncEnv (..), SyncHooks (SyncHooks, hookFirstSync, hookPushAge), SyncSchedule (SyncSchedule, schedAbsentReport, schedBootBackoff, schedPollDelay), absentReportInterval, bootBackoffDelays, newS3CveSource, runCveSync, s3CveFetchFor)
import Ecluse.Runtime.Log (logLine, moduleField)
import Ecluse.Runtime.Telemetry (Telemetry)
import Ecluse.Runtime.Telemetry.Instruments (Metrics, advisorySyncMetricsPortOf, registerAdvisoryDatabaseAge, registerAdvisorySourceAge)
import Ecluse.Runtime.Telemetry.Tracing (advisorySyncTracingPortOf)

{- | The rules' boot-bound capabilities for one mount ecosystem. A mount's rules read only their own
ecosystem's advisory database. One the plan carries no handle for has none configured, and reports nowhere.
-}
cveRuleDepsFor :: Map.Map Ecosystem CveSyncHandle -> BreakerReporter -> (Ecosystem -> OutageReport -> IO ()) -> Ecosystem -> RuleDeps
cveRuleDepsFor :: Map Ecosystem CveSyncHandle
-> BreakerReporter
-> (Ecosystem -> OutageReport -> IO ())
-> Ecosystem
-> RuleDeps
cveRuleDepsFor Map Ecosystem CveSyncHandle
plan BreakerReporter
reporter Ecosystem -> OutageReport -> IO ()
reportOutage Ecosystem
eco =
    RuleDeps
        { rdAdvisoryDatabase :: AdvisoryDatabase
rdAdvisoryDatabase = AdvisoryDatabase
-> (CveSyncHandle -> AdvisoryDatabase)
-> Maybe CveSyncHandle
-> AdvisoryDatabase
forall b a. b -> (a -> b) -> Maybe a -> b
maybe AdvisoryDatabase
NoAdvisoryDatabase CveSyncHandle -> AdvisoryDatabase
slotDatabase Maybe CveSyncHandle
handle
        , rdCurrentAdvisoryEtag :: IO (Maybe DbEtag)
rdCurrentAdvisoryEtag = IO (Maybe DbEtag)
-> (CveSyncHandle -> IO (Maybe DbEtag))
-> Maybe CveSyncHandle
-> IO (Maybe DbEtag)
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (Maybe DbEtag -> IO (Maybe DbEtag)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe DbEtag
forall a. Maybe a
Nothing) (CveSlot -> IO (Maybe DbEtag)
currentAdvisoryEtag (CveSlot -> IO (Maybe DbEtag))
-> (CveSyncHandle -> CveSlot) -> CveSyncHandle -> IO (Maybe DbEtag)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. SyncEnv -> CveSlot
syncSlot (SyncEnv -> CveSlot)
-> (CveSyncHandle -> SyncEnv) -> CveSyncHandle -> CveSlot
forall b c a. (b -> c) -> (a -> b) -> a -> c
. CveSyncHandle -> SyncEnv
csEnv) Maybe CveSyncHandle
handle
        , rdBreakerReporter :: BreakerReporter
rdBreakerReporter = BreakerReporter
reporter
        , rdSourceReporter :: SourceReporter
rdSourceReporter = SourceReporter
-> (CveSyncHandle -> SourceReporter)
-> Maybe CveSyncHandle
-> SourceReporter
forall b a. b -> (a -> b) -> Maybe a -> b
maybe SourceReporter
noSourceReporter ((OutageReport -> IO ()) -> CveSyncHandle -> SourceReporter
sourceReporterOf (Ecosystem -> OutageReport -> IO ()
reportOutage Ecosystem
eco)) Maybe CveSyncHandle
handle
        , rdAdvisoryFreshness :: IO AdvisoryFreshness
rdAdvisoryFreshness = Map Ecosystem CveSyncHandle -> Ecosystem -> IO AdvisoryFreshness
advisoryFreshnessFor Map Ecosystem CveSyncHandle
plan Ecosystem
eco
        }
  where
    handle :: Maybe CveSyncHandle
handle = Ecosystem -> Map Ecosystem CveSyncHandle -> Maybe CveSyncHandle
forall k a. Ord k => k -> Map k a -> Maybe a
Map.lookup Ecosystem
eco Map Ecosystem CveSyncHandle
plan
    slotDatabase :: CveSyncHandle -> AdvisoryDatabase
slotDatabase CveSyncHandle
h = (forall a. (Maybe (DbEtag, CveLookup) -> IO a) -> IO a)
-> AdvisoryDatabase
AdvisoryDatabase (CveSlot -> (Maybe (DbEtag, CveLookup) -> IO a) -> IO a
forall a. CveSlot -> (Maybe (DbEtag, CveLookup) -> IO a) -> IO a
withSlotGeneration (SyncEnv -> CveSlot
syncSlot (CveSyncHandle -> SyncEnv
csEnv CveSyncHandle
h)))

-- One handle's reporter, over the outage state every mount of the ecosystem shares.
sourceReporterOf :: (OutageReport -> IO ()) -> CveSyncHandle -> SourceReporter
sourceReporterOf :: (OutageReport -> IO ()) -> CveSyncHandle -> SourceReporter
sourceReporterOf OutageReport -> IO ()
emit CveSyncHandle
handle = NominalDiffTime
-> IO UTCTime
-> OutageStore
-> (OutageReport -> IO ())
-> SourceReporter
sourceReporter NominalDiffTime
outageReportPeriod (CveSyncHandle -> IO UTCTime
csClock CveSyncHandle
handle) (TVar OutageState -> OutageStore
tvarOutageStore (CveSyncHandle -> TVar OutageState
csOutage CveSyncHandle
handle)) OutageReport -> IO ()
emit

{- | How often a continuing outage reminds the operator: the unloaded-database report's own gap, so
an outage costs the log one line per interval on either path.
-}
outageReportPeriod :: NominalDiffTime
outageReportPeriod :: NominalDiffTime
outageReportPeriod = Int -> NominalDiffTime
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
absentReportInterval NominalDiffTime -> NominalDiffTime -> NominalDiffTime
forall a. Fractional a => a -> a -> a
/ NominalDiffTime
1_000_000

{- | How old one mount's serving artifact's push is. An ecosystem the plan carries no handle for
has no advisory stack at all, so nothing ages and the absent-database path decides instead.
-}
advisoryFreshnessFor :: Map.Map Ecosystem CveSyncHandle -> Ecosystem -> IO AdvisoryFreshness
advisoryFreshnessFor :: Map Ecosystem CveSyncHandle -> Ecosystem -> IO AdvisoryFreshness
advisoryFreshnessFor Map Ecosystem CveSyncHandle
plan Ecosystem
eco = IO AdvisoryFreshness
-> (CveSyncHandle -> IO AdvisoryFreshness)
-> Maybe CveSyncHandle
-> IO AdvisoryFreshness
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (AdvisoryFreshness -> IO AdvisoryFreshness
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure AdvisoryFreshness
AdvisoryFresh) CveSyncHandle -> IO AdvisoryFreshness
advisoryFreshnessOf (Ecosystem -> Map Ecosystem CveSyncHandle -> Maybe CveSyncHandle
forall k a. Ord k => k -> Map k a -> Maybe a
Map.lookup Ecosystem
eco Map Ecosystem CveSyncHandle
plan)

{- | One handle's reading: the slot's publication time against this mount's maximum, on the
handle's own clock. A failed poll never swaps, so a warm process keeps the last time it read.
-}
advisoryFreshnessOf :: CveSyncHandle -> IO AdvisoryFreshness
advisoryFreshnessOf :: CveSyncHandle -> IO AdvisoryFreshness
advisoryFreshnessOf CveSyncHandle
handle = do
    now <- CveSyncHandle -> IO UTCTime
csClock CveSyncHandle
handle
    assessAdvisoryAge (csMaxAge handle) now . publicationOf <$> currentAdvisorySource (syncSlot (csEnv handle))

{- Nothing serving, a dated push, or a generation the store gave no publication time for. The
third is unverified evidence rather than an absent database, so it is kept distinct here. -}
publicationOf :: Maybe AdvisorySource -> AdvisoryPublication
publicationOf :: Maybe AdvisorySource -> AdvisoryPublication
publicationOf = AdvisoryPublication
-> (AdvisorySource -> AdvisoryPublication)
-> Maybe AdvisorySource
-> AdvisoryPublication
forall b a. b -> (a -> b) -> Maybe a -> b
maybe AdvisoryPublication
NoGeneration (AdvisoryPublication
-> (UTCTime -> AdvisoryPublication)
-> Maybe UTCTime
-> AdvisoryPublication
forall b a. b -> (a -> b) -> Maybe a -> b
maybe AdvisoryPublication
UndatedGeneration UTCTime -> AdvisoryPublication
PublishedAt (Maybe UTCTime -> AdvisoryPublication)
-> (AdvisorySource -> Maybe UTCTime)
-> AdvisorySource
-> AdvisoryPublication
forall b c a. (b -> c) -> (a -> b) -> a -> c
. AdvisorySource -> Maybe UTCTime
asPushedAt)

{- | Report one ecosystem's push age when it passes half its maximum, once per crossing. The latch
re-arms when a fresh push brings the age back under, so a long outage does not repeat every poll.
-}
reportPushAge :: LogEnv -> Ecosystem -> CveSyncHandle -> IO ()
reportPushAge :: LogEnv -> Ecosystem -> CveSyncHandle -> IO ()
reportPushAge LogEnv
logEnv Ecosystem
eco CveSyncHandle
handle = do
    freshness <- CveSyncHandle -> IO AdvisoryFreshness
advisoryFreshnessOf CveSyncHandle
handle
    crossing <- atomically $ do
        latched <- readTVar (csAgeAlarmed handle)
        let (latched', crossed) = ageAlarmStep latched freshness
        writeTVar (csAgeAlarmed handle) latched'
        pure crossed
    whenJust crossing (logPushAge logEnv eco)

-- What the crossing line carries: enough to tell an update outage from a maximum set too short.
logPushAge :: LogEnv -> Ecosystem -> AdvisoryAge -> IO ()
logPushAge :: LogEnv -> Ecosystem -> AdvisoryAge -> IO ()
logPushAge LogEnv
logEnv Ecosystem
eco AdvisoryAge
observed =
    LogEnv -> SimpleLogPayload -> Severity -> Text -> IO ()
logLine
        LogEnv
logEnv
        ( Text -> SimpleLogPayload
moduleField Text
"Ecluse.Cve.Sync"
            SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"ecosystem" (Ecosystem -> Text
ecosystemName Ecosystem
eco)
            SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"pushed_at" (UTCTime -> Text
renderIso8601Utc (AdvisoryAge -> UTCTime
advisoryPushedAt AdvisoryAge
observed))
            SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Integer -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"age_seconds" (NominalDiffTime -> Integer
forall b. Integral b => NominalDiffTime -> b
forall a b. (RealFrac a, Integral b) => a -> b
round (AdvisoryAge -> NominalDiffTime
advisoryAge AdvisoryAge
observed) :: Integer)
            SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Integer -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"max_age_seconds" (NominalDiffTime -> Integer
forall b. Integral b => NominalDiffTime -> b
forall a b. (RealFrac a, Integral b) => a -> b
round (AdvisoryAge -> NominalDiffTime
advisoryMaxAge AdvisoryAge
observed) :: Integer)
        )
        Severity
ErrorS
        Text
"the advisory push age has passed half its maximum; past the maximum, CVE-based denial refuses"

-- The publication time the serving artifact carries, or nothing before the first sync.
advisoryPushTime :: CveSyncHandle -> IO (Maybe UTCTime)
advisoryPushTime :: CveSyncHandle -> IO (Maybe UTCTime)
advisoryPushTime CveSyncHandle
handle = (AdvisorySource -> Maybe UTCTime
asPushedAt (AdvisorySource -> Maybe UTCTime)
-> Maybe AdvisorySource -> Maybe UTCTime
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<<) (Maybe AdvisorySource -> Maybe UTCTime)
-> IO (Maybe AdvisorySource) -> IO (Maybe UTCTime)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> CveSlot -> IO (Maybe AdvisorySource)
currentAdvisorySource (SyncEnv -> CveSlot
syncSlot (CveSyncHandle -> SyncEnv
csEnv CveSyncHandle
handle))

{- | Log one ecosystem's advisory-source outage reports: the start and each reminder at ERROR, the
level an operator pages on, and the recovery at INFO. A fault's detail rides along and reaches no client.
-}
katipOutageReporter :: LogEnv -> Ecosystem -> OutageReport -> IO ()
katipOutageReporter :: LogEnv -> Ecosystem -> OutageReport -> IO ()
katipOutageReporter LogEnv
logEnv Ecosystem
eco = \case
    OutageBegan Text
rule Text
cause ->
        LogEnv -> SimpleLogPayload -> Severity -> Text -> IO ()
logLine LogEnv
logEnv (SimpleLogPayload
payload SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"rule" Text
rule SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"cause" Text
cause) Severity
ErrorS Text
"advisory source outage began: a rule cannot consult it"
    OutageContinues UTCTime
since Map Text Text
rules ->
        LogEnv -> SimpleLogPayload -> Severity -> Text -> IO ()
logLine LogEnv
logEnv (SimpleLogPayload
payload SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"since" (UTCTime -> Text
renderIso8601Utc UTCTime
since) SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"rules" (Map Text Text -> Text
renderCauses Map Text Text
rules)) Severity
ErrorS Text
"advisory source outage continues"
    OutageRecovered UTCTime
since ->
        LogEnv -> SimpleLogPayload -> Severity -> Text -> IO ()
logLine LogEnv
logEnv (SimpleLogPayload
payload SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"since" (UTCTime -> Text
renderIso8601Utc UTCTime
since)) Severity
InfoS Text
"advisory source outage recovered: every rule consults it again"
  where
    payload :: SimpleLogPayload
payload = Text -> SimpleLogPayload
moduleField Text
"Ecluse.Core.Rules" SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"ecosystem" (Ecosystem -> Text
ecosystemName Ecosystem
eco)
    renderCauses :: Map Text Text -> Text
renderCauses = Text -> [Text] -> Text
T.intercalate Text
"; " ([Text] -> Text)
-> (Map Text Text -> [Text]) -> Map Text Text -> Text
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ((Text, Text) -> Text) -> [(Text, Text)] -> [Text]
forall a b. (a -> b) -> [a] -> [b]
map (\(Text
rule, Text
cause) -> Text
rule Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
": " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
cause) ([(Text, Text)] -> [Text])
-> (Map Text Text -> [(Text, Text)]) -> Map Text Text -> [Text]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Map Text Text -> [(Text, Text)]
forall k a. Map k a -> [(k, a)]
Map.toList

{- | The readiness verdict over the sync plan. Only a mount whose rules deny on the database waits
for its first sync, and one ecosystem's missing artifact leaves the others routable.
-}
cveSyncReadiness :: Map.Map Ecosystem CveSyncHandle -> IO Readiness
cveSyncReadiness :: Map Ecosystem CveSyncHandle -> IO Readiness
cveSyncReadiness Map Ecosystem CveSyncHandle
plan = Map Ecosystem MountReadiness -> Readiness
mountReadiness (Map Ecosystem MountReadiness -> Readiness)
-> IO (Map Ecosystem MountReadiness) -> IO Readiness
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> (CveSyncHandle -> IO MountReadiness)
-> Map Ecosystem CveSyncHandle -> IO (Map Ecosystem MountReadiness)
forall (t :: * -> *) (f :: * -> *) a b.
(Traversable t, Applicative f) =>
(a -> f b) -> t a -> f (t b)
forall (f :: * -> *) a b.
Applicative f =>
(a -> f b) -> Map Ecosystem a -> f (Map Ecosystem b)
traverse CveSyncHandle -> IO MountReadiness
mountStateOf Map Ecosystem CveSyncHandle
plan

-- One mount's advisory state, from what its rules need and the one-way flag its sync task flips.
mountStateOf :: CveSyncHandle -> IO MountReadiness
mountStateOf :: CveSyncHandle -> IO MountReadiness
mountStateOf CveSyncHandle
handle = DatabaseRequirement -> Bool -> MountReadiness
mountStateFor (CveSyncHandle -> DatabaseRequirement
csDatabase CveSyncHandle
handle) (Bool -> MountReadiness) -> IO Bool -> IO MountReadiness
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> TVar Bool -> IO Bool
forall (m :: * -> *) a. MonadIO m => TVar a -> m a
readTVarIO (CveSyncHandle -> TVar Bool
csReady CveSyncHandle
handle)

{- | The sync tasks' timing: the shipped boot burst over the configured poll interval. The microsecond
conversion cannot wrap: the config decoder bounds the interval to @[1, maxBound div 1_000_000]@ seconds.
-}
cveSyncScheduleFor :: AppConfig -> SyncSchedule
cveSyncScheduleFor :: AppConfig -> SyncSchedule
cveSyncScheduleFor AppConfig
env =
    SyncSchedule
        { schedBootBackoff :: [Int]
schedBootBackoff = [Int]
bootBackoffDelays
        , schedPollDelay :: Int
schedPollDelay = NominalDiffTime -> Int
secondsToMicros (AdvisoriesSettings -> NominalDiffTime
advPollInterval (AppConfig -> AdvisoriesSettings
cfgAdvisories AppConfig
env))
        , schedAbsentReport :: Int
schedAbsentReport = Int
absentReportInterval
        }

{- | One supervised sync task per configured ecosystem, each flipping its own one-way readiness
flag once its first sync lands. Every role that evaluates rules runs these.
-}
cveSyncTasks :: LogEnv -> Metrics -> Telemetry -> SyncSchedule -> Map.Map Ecosystem CveSyncHandle -> [IO ()]
cveSyncTasks :: LogEnv
-> Metrics
-> Telemetry
-> SyncSchedule
-> Map Ecosystem CveSyncHandle
-> [IO ()]
cveSyncTasks LogEnv
logEnv Metrics
metrics Telemetry
telemetry SyncSchedule
schedule Map Ecosystem CveSyncHandle
plan =
    [ 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 LogEnv
logEnv (SimpleLogPayload
forall a. Monoid a => a
mempty :: SimpleLogPayload) Namespace
"cve-sync" (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
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Ecosystem -> Text
forall b a. (Show a, IsString b) => a -> b
show (SyncEnv -> Ecosystem
syncEcosystem (CveSyncHandle -> SyncEnv
csEnv CveSyncHandle
handle)) Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]") BackoffSchedule
backgroundLoopBackoff)
            (AdvisorySyncMetricsPort
-> AdvisorySyncTracingPort
-> SyncEnv
-> SyncSchedule
-> SyncHooks
-> KatipContextT IO ()
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
AdvisorySyncMetricsPort
-> AdvisorySyncTracingPort
-> SyncEnv
-> SyncSchedule
-> SyncHooks
-> m ()
runCveSync AdvisorySyncMetricsPort
syncMetrics AdvisorySyncTracingPort
syncTracing (CveSyncHandle -> SyncEnv
csEnv CveSyncHandle
handle) SyncSchedule
schedule (Ecosystem -> CveSyncHandle -> SyncHooks
hooksFor Ecosystem
eco CveSyncHandle
handle))
    | (Ecosystem
eco, CveSyncHandle
handle) <- Map Ecosystem CveSyncHandle -> [(Ecosystem, CveSyncHandle)]
forall k a. Map k a -> [(k, a)]
Map.toList Map Ecosystem CveSyncHandle
plan
    ]
  where
    syncMetrics :: AdvisorySyncMetricsPort
syncMetrics = Metrics -> AdvisorySyncMetricsPort
advisorySyncMetricsPortOf Metrics
metrics
    syncTracing :: AdvisorySyncTracingPort
syncTracing = Telemetry -> AdvisorySyncTracingPort
advisorySyncTracingPortOf Telemetry
telemetry
    hooksFor :: Ecosystem -> CveSyncHandle -> SyncHooks
hooksFor Ecosystem
eco CveSyncHandle
handle =
        SyncHooks
            { hookFirstSync :: IO ()
hookFirstSync = STM () -> IO ()
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (TVar Bool -> Bool -> STM ()
forall a. TVar a -> a -> STM ()
writeTVar (CveSyncHandle -> TVar Bool
csReady CveSyncHandle
handle) Bool
True)
            , hookPushAge :: IO ()
hookPushAge = LogEnv -> Ecosystem -> CveSyncHandle -> IO ()
reportPushAge LogEnv
logEnv Ecosystem
eco CveSyncHandle
handle
            }

-- | Register once per role. Callbacks read the slots, so observations survive sync-task restarts.
registerAdvisoryAges :: Metrics -> Map.Map Ecosystem CveSyncHandle -> IO ()
registerAdvisoryAges :: Metrics -> Map Ecosystem CveSyncHandle -> IO ()
registerAdvisoryAges Metrics
metrics Map Ecosystem CveSyncHandle
plan =
    [(Ecosystem, CveSyncHandle)]
-> ((Ecosystem, CveSyncHandle) -> IO ()) -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
t a -> (a -> f b) -> f ()
for_ (Map Ecosystem CveSyncHandle -> [(Ecosystem, CveSyncHandle)]
forall k a. Map k a -> [(k, a)]
Map.toList Map Ecosystem CveSyncHandle
plan) (((Ecosystem, CveSyncHandle) -> IO ()) -> IO ())
-> ((Ecosystem, CveSyncHandle) -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \(Ecosystem
eco, CveSyncHandle
handle) -> do
        Metrics -> Ecosystem -> IO (Maybe Double) -> IO ()
registerAdvisoryDatabaseAge Metrics
metrics Ecosystem
eco (CveSlot -> IO (Maybe Double)
generationInstalledAt (SyncEnv -> CveSlot
syncSlot (CveSyncHandle -> SyncEnv
csEnv CveSyncHandle
handle)))
        Metrics -> Ecosystem -> IO (Maybe UTCTime) -> IO ()
registerAdvisorySourceAge Metrics
metrics Ecosystem
eco (CveSyncHandle -> IO (Maybe UTCTime)
advisoryPushTime CveSyncHandle
handle)

-- | One configured ecosystem's advisory-sync wiring.
data CveSyncHandle = CveSyncHandle
    { CveSyncHandle -> TVar Bool
csReady :: TVar Bool
    -- ^ The one-way first-sync readiness flag.
    , CveSyncHandle -> SyncEnv
csEnv :: SyncEnv
    -- ^ The sync task's environment. Its 'syncSlot' is the slot this ecosystem's rules borrow through.
    , CveSyncHandle -> MaxAdvisoryAge
csMaxAge :: MaxAdvisoryAge
    -- ^ This mount's effective maximum push age, derived once at boot from its own rules.
    , CveSyncHandle -> IO UTCTime
csClock :: IO UTCTime
    -- ^ The wall clock the push age is read on, injected so a suite can fix it.
    , CveSyncHandle -> TVar Bool
csAgeAlarmed :: TVar Bool
    -- ^ Whether the half-maximum crossing has already been reported for the current push.
    , CveSyncHandle -> TVar OutageState
csOutage :: TVar OutageState
    -- ^ The rules-side outage state every mount of this ecosystem reports through.
    , CveSyncHandle -> DatabaseRequirement
csDatabase :: DatabaseRequirement
    -- ^ Whether this mount's own rules deny on the database, which is what its readiness turns on.
    }

-- | What one vetted mount's rules ask of the advisory stack, read off its own policy at boot.
data AdvisoryNeed = AdvisoryNeed
    { AdvisoryNeed -> Ecosystem
anEcosystem :: Ecosystem
    , AdvisoryNeed -> MaxAdvisoryAge
anMaxAge :: MaxAdvisoryAge
    , AdvisoryNeed -> EpssRequirement
anEpss :: EpssRequirement
    , AdvisoryNeed -> DatabaseRequirement
anDatabase :: DatabaseRequirement
    }

{- | Build the advisory-sync plan, one 'CveSyncHandle' per vetted mount ecosystem, or nothing with
no store. A mount the build does not ship awaits an artifact that never comes, so it stays unready.
-}
planCveSync :: LogEnv -> Maybe AwsEndpoint -> AppConfig -> [AdvisoryNeed] -> IO (Map.Map Ecosystem CveSyncHandle)
planCveSync :: LogEnv
-> Maybe AwsEndpoint
-> AppConfig
-> [AdvisoryNeed]
-> IO (Map Ecosystem CveSyncHandle)
planCveSync LogEnv
logEnv Maybe AwsEndpoint
s3Endpoint AppConfig
appCfg [AdvisoryNeed]
needs = case AdvisoriesSettings -> Maybe AdvisoryStoreUrl
advUrl (AppConfig -> AdvisoriesSettings
cfgAdvisories AppConfig
appCfg) of
    Maybe AdvisoryStoreUrl
Nothing -> Map Ecosystem CveSyncHandle -> IO (Map Ecosystem CveSyncHandle)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Map Ecosystem CveSyncHandle
forall k a. Map k a
Map.empty
    Just AdvisoryStoreUrl
store -> do
        let dataDir :: String
dataDir = AdvisoriesSettings -> String
advDataDir (AppConfig -> AdvisoriesSettings
cfgAdvisories AppConfig
appCfg)
        Bool -> String -> IO ()
createDirectoryIfMissing Bool
True String
dataDir
        LogEnv -> String -> IO ()
sweepStaleTemps LogEnv
logEnv String
dataDir
        cveSource <- Maybe AwsEndpoint -> IO S3CveSource
newS3CveSource Maybe AwsEndpoint
s3Endpoint
        Map.fromList <$> traverse (cveSyncHandleFor appCfg cveSource store) needs

cveSyncHandleFor :: AppConfig -> S3CveSource -> AdvisoryStoreUrl -> AdvisoryNeed -> IO (Ecosystem, CveSyncHandle)
cveSyncHandleFor :: AppConfig
-> S3CveSource
-> AdvisoryStoreUrl
-> AdvisoryNeed
-> IO (Ecosystem, CveSyncHandle)
cveSyncHandleFor AppConfig
appCfg S3CveSource
cveSource AdvisoryStoreUrl
store AdvisoryNeed
need = do
    slot <- IO CveSlot
newCveSlot
    ready <- newTVarIO False
    alarmed <- newTVarIO False
    outage <- newTVarIO Healthy
    pure
        ( anEcosystem need
        , CveSyncHandle
            { csReady = ready
            , csEnv = syncEnvFor appCfg cveSource store need slot
            , csMaxAge = anMaxAge need
            , csClock = getCurrentTime
            , csAgeAlarmed = alarmed
            , csOutage = outage
            , csDatabase = anDatabase need
            }
        )

-- 'cveSource' captures the S3 environment once, so every ecosystem's transport shares one
-- credential discovery. The store addresses the remote object, the local copy its bare file name.
syncEnvFor :: AppConfig -> S3CveSource -> AdvisoryStoreUrl -> AdvisoryNeed -> CveSlot -> SyncEnv
syncEnvFor :: AppConfig
-> S3CveSource
-> AdvisoryStoreUrl
-> AdvisoryNeed
-> CveSlot
-> SyncEnv
syncEnvFor AppConfig
appCfg S3CveSource
cveSource AdvisoryStoreUrl
store AdvisoryNeed
need CveSlot
slot =
    SyncEnv
        { syncFetch :: CveFetch
syncFetch =
            S3CveSource -> Text -> Text -> Int -> CveFetch
s3CveFetchFor
                S3CveSource
cveSource
                (AdvisoryStoreUrl -> Text
advisoryStoreBucket AdvisoryStoreUrl
store)
                (AdvisoryStoreUrl -> String -> Text
advisoryObjectKey AdvisoryStoreUrl
store String
fileName)
                (LimitsSettings -> Int
limMaxAdvisoryDatabaseBytes (AppConfig -> LimitsSettings
cfgLimits AppConfig
appCfg))
        , syncEcosystem :: Ecosystem
syncEcosystem = Ecosystem
eco
        , syncEpssRequirement :: EpssRequirement
syncEpssRequirement = AdvisoryNeed -> EpssRequirement
anEpss AdvisoryNeed
need
        , syncDbPath :: String
syncDbPath = AdvisoriesSettings -> String
advDataDir (AppConfig -> AdvisoriesSettings
cfgAdvisories AppConfig
appCfg) String -> String -> String
</> String
fileName
        , syncSlot :: CveSlot
syncSlot = CveSlot
slot
        , syncStoreRef :: Text
syncStoreRef = AdvisoryStoreUrl -> Text
advisoryStoreUrlText AdvisoryStoreUrl
store
        }
  where
    eco :: Ecosystem
eco = AdvisoryNeed -> Ecosystem
anEcosystem AdvisoryNeed
need
    fileName :: String
fileName = Text -> String
osvDbFileName (Ecosystem -> Text
ecosystemName Ecosystem
eco)

{- | Sweep the in-progress downloads an interrupted run left behind, which an @emptyDir@ keeps
across a container restart. The sweep is best effort, per 'sweepStep'.
-}
sweepStaleTemps :: LogEnv -> FilePath -> IO ()
sweepStaleTemps :: LogEnv -> String -> IO ()
sweepStaleTemps LogEnv
logEnv String
dataDir =
    LogEnv -> String -> IO () -> IO ()
sweepStep LogEnv
logEnv String
dataDir (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
        entries <- String -> IO [String]
listDirectory String
dataDir
        traverse_ (removeStaleTemp logEnv dataDir) (filter (isExtensionOf "tmp") entries)

-- Remove one stray @.tmp@ entry, tolerating a per-entry filesystem fault so a single
-- unremovable file does not abort the rest of the sweep.
removeStaleTemp :: LogEnv -> FilePath -> FilePath -> IO ()
removeStaleTemp :: LogEnv -> String -> String -> IO ()
removeStaleTemp LogEnv
logEnv String
dataDir String
entry =
    let path :: String
path = String
dataDir String -> String -> String
</> String
entry in LogEnv -> String -> IO () -> IO ()
sweepStep LogEnv
logEnv String
path (String -> IO ()
removeFile String
path)

{- | Run one best-effort step of the stale-temp sweep. It logs and swallows an 'IOError', so a
read-only or mispermissioned data dir does not stop the boot, and any other exception propagates.
-}
sweepStep :: LogEnv -> FilePath -> IO () -> IO ()
sweepStep :: LogEnv -> String -> IO () -> IO ()
sweepStep LogEnv
logEnv String
path IO ()
step = IO ()
step IO () -> (IOError -> IO ()) -> IO ()
forall a. IO a -> (IOError -> IO a) -> IO a
`catchIOError` LogEnv -> String -> IOError -> IO ()
logSweepFailure LogEnv
logEnv String
path

-- The logged OS error detail is the operator's own filesystem, not untrusted input.
logSweepFailure :: LogEnv -> FilePath -> IOError -> IO ()
logSweepFailure :: LogEnv -> String -> IOError -> IO ()
logSweepFailure LogEnv
logEnv String
path IOError
err =
    LogEnv -> SimpleLogPayload -> Severity -> Text -> IO ()
logLine LogEnv
logEnv SimpleLogPayload
payload Severity
WarningS (Text
"could not sweep stale advisory temp files: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> IOError -> Text
forall b a. (Show a, IsString b) => a -> b
show IOError
err)
  where
    payload :: SimpleLogPayload
payload = Text -> SimpleLogPayload
moduleField Text
"Ecluse.Cve.Sync" SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"path" (String -> Text
forall a. ToText a => a -> Text
toText String
path)