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

{- | The fetch transport, the one-cycle step and the scheduled task behind
"Ecluse.Runtime.Cve.Sync", which documents the sync and re-exports the curated surface.
Importing this module opts out of that stability promise, the convention @text@ and @bytestring@
use, so production code imports the public one.
-}
module Ecluse.Runtime.Cve.Sync.Internal (
    -- * The injected transport
    CveFetch (..),
    FetchedObject (..),
    DbEtag (..),
    OsvDbFetchFault (..),
    OsvDbCapExceeded (..),
    S3CveSource,
    newS3CveSource,
    s3CveFetchFor,
    cappedAt,

    -- * One sync cycle
    SyncEnv (..),
    SyncOutcome (..),
    syncStep,

    -- * The scheduled task
    SyncSchedule (..),
    SyncHooks (..),
    runCveSync,
    bootBackoffDelays,
    absentReportInterval,
) where

import Conduit (ConduitT, runResourceT, (.|))
import Data.Conduit.Combinators qualified as C
import Data.List (lookup)
import Data.Text qualified as T
import Data.Time (UTCTime)
import Data.Time.Format.ISO8601 (iso8601ParseM)
import Katip (KatipContext, Severity (DebugS, ErrorS, InfoS), logFM, ls)
import Network.HTTP.Types.Status (statusCode)
import System.Directory (removeFile, renameFile)
import UnliftIO (MonadUnliftIO, withRunInIO)
import UnliftIO.Concurrent (threadDelay)
import UnliftIO.Exception (catch, catchAny, mask, onException, throwIO)

import Amazonka qualified as AWS
import Amazonka.S3 qualified as S3
import Amazonka.S3.Lens qualified as S3L
import Lens.Micro ((^.))

import Ecluse.Core.Cve (CveDb (cveDbClose, cveDbMeta), CveDbRejected, openCveDb)
import Ecluse.Core.Cve.Slot (AdvisorySource (..), CveSlot, currentAdvisoryEtag, currentAdvisorySource, observeAdvisoryPublication, swapIn)
import Ecluse.Core.Cve.Types (DbEtag (..))
import Ecluse.Core.Ecosystem (Ecosystem)
import Ecluse.Core.Fault (TransportFault)
import Ecluse.Core.Osv.Provenance (AdvisoryProvenance (apEpssScoreDate, apOsvNewestModified, apOsvSource))
import Ecluse.Core.Osv.Schema (EpssRequirement, MetaKey (MetaBuiltAt, MetaRowCount), renderMetaKey)
import Ecluse.Core.Security.Authority (dialledAuthorityLabel)
import Ecluse.Core.Stream (boundBytes)
import Ecluse.Core.Telemetry.Metrics (
    AdvisorySyncResult (AdvisoryFetchFailed, AdvisoryNonePublished, AdvisoryRefused, AdvisorySwapped, AdvisoryUnchanged),
 )
import Ecluse.Core.Telemetry.Record (AdvisorySyncMetricsPort (asmpSyncAttempt, asmpSyncDuration), timedSeconds)
import Ecluse.Core.Telemetry.Span (AdvisorySyncTracingPort (astpSyncAttemptSpan))
import Ecluse.Core.Text (readDecimalText, renderIso8601Utc)
import Ecluse.Runtime.Aws.Env (AwsEndpoint)
import Ecluse.Runtime.Aws.Fault (classifyAwsTransport)
import Ecluse.Runtime.Aws.S3 (buildS3Env)

-- | The advisory transport supplied to 'syncStep' by 'newS3CveSource'.
data CveFetch = CveFetch
    { CveFetch -> IO (Either OsvDbFetchFault (Maybe FetchedObject))
fetchHead :: IO (Either OsvDbFetchFault (Maybe FetchedObject))
    {- ^ The object's ETag and publication time. @Right Nothing@ means it does not exist.
    Fetch failures use 'Left'.
    -}
    , CveFetch -> FilePath -> IO (Either OsvDbFetchFault FetchedObject)
fetchDownload :: FilePath -> IO (Either OsvDbFetchFault FetchedObject)
    {- ^ Download the artifact to the given path, byte-bounded. The ETag is the download's own, so a
    publish racing the poll is recorded truthfully. A 'Left' may leave a partial file at that path.
    -}
    }

{- | Metadata from one HEAD or GET response. A download carries its own metadata,
so a publication racing HEAD cannot mislabel the downloaded bytes.
-}
data FetchedObject = FetchedObject
    { FetchedObject -> DbEtag
foEtag :: DbEtag
    , FetchedObject -> Maybe UTCTime
foPushedAt :: Maybe UTCTime
    -- ^ The object's own timestamp, 'Nothing' when the store reported none.
    }
    deriving stock (FetchedObject -> FetchedObject -> Bool
(FetchedObject -> FetchedObject -> Bool)
-> (FetchedObject -> FetchedObject -> Bool) -> Eq FetchedObject
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: FetchedObject -> FetchedObject -> Bool
== :: FetchedObject -> FetchedObject -> Bool
$c/= :: FetchedObject -> FetchedObject -> Bool
/= :: FetchedObject -> FetchedObject -> Bool
Eq, Int -> FetchedObject -> ShowS
[FetchedObject] -> ShowS
FetchedObject -> FilePath
(Int -> FetchedObject -> ShowS)
-> (FetchedObject -> FilePath)
-> ([FetchedObject] -> ShowS)
-> Show FetchedObject
forall a.
(Int -> a -> ShowS) -> (a -> FilePath) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> FetchedObject -> ShowS
showsPrec :: Int -> FetchedObject -> ShowS
$cshow :: FetchedObject -> FilePath
show :: FetchedObject -> FilePath
$cshowList :: [FetchedObject] -> ShowS
showList :: [FetchedObject] -> ShowS
Show)

{- | Why an artifact fetch did not yield usable bytes. Every one is a value on the 'CveFetch'
channel, never an exception, and 'syncStep' folds it into its outcome.
-}
data OsvDbFetchFault
    = -- | The object exceeds the configured byte cap (carried, in bytes).
      OsvDbTooLarge Int
    | -- | The response carried no ETag, so there is nothing truthful to record.
      OsvDbNoEtag
    | -- | The transport could not deliver the object (carried, classified).
      OsvDbTransport TransportFault
    deriving stock (OsvDbFetchFault -> OsvDbFetchFault -> Bool
(OsvDbFetchFault -> OsvDbFetchFault -> Bool)
-> (OsvDbFetchFault -> OsvDbFetchFault -> Bool)
-> Eq OsvDbFetchFault
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: OsvDbFetchFault -> OsvDbFetchFault -> Bool
== :: OsvDbFetchFault -> OsvDbFetchFault -> Bool
$c/= :: OsvDbFetchFault -> OsvDbFetchFault -> Bool
/= :: OsvDbFetchFault -> OsvDbFetchFault -> Bool
Eq, Int -> OsvDbFetchFault -> ShowS
[OsvDbFetchFault] -> ShowS
OsvDbFetchFault -> FilePath
(Int -> OsvDbFetchFault -> ShowS)
-> (OsvDbFetchFault -> FilePath)
-> ([OsvDbFetchFault] -> ShowS)
-> Show OsvDbFetchFault
forall a.
(Int -> a -> ShowS) -> (a -> FilePath) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> OsvDbFetchFault -> ShowS
showsPrec :: Int -> OsvDbFetchFault -> ShowS
$cshow :: OsvDbFetchFault -> FilePath
show :: OsvDbFetchFault -> FilePath
$cshowList :: [OsvDbFetchFault] -> ShowS
showList :: [OsvDbFetchFault] -> ShowS
Show)

{- | 'cappedAt' sits in a conduit and has no value channel, so it reports an overstepped byte cap
by throwing this __confined__ exception. 's3Download' catches it and folds it into 'OsvDbTooLarge'.
-}
newtype OsvDbCapExceeded = OsvDbCapExceeded Int
    deriving stock (OsvDbCapExceeded -> OsvDbCapExceeded -> Bool
(OsvDbCapExceeded -> OsvDbCapExceeded -> Bool)
-> (OsvDbCapExceeded -> OsvDbCapExceeded -> Bool)
-> Eq OsvDbCapExceeded
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: OsvDbCapExceeded -> OsvDbCapExceeded -> Bool
== :: OsvDbCapExceeded -> OsvDbCapExceeded -> Bool
$c/= :: OsvDbCapExceeded -> OsvDbCapExceeded -> Bool
/= :: OsvDbCapExceeded -> OsvDbCapExceeded -> Bool
Eq, Int -> OsvDbCapExceeded -> ShowS
[OsvDbCapExceeded] -> ShowS
OsvDbCapExceeded -> FilePath
(Int -> OsvDbCapExceeded -> ShowS)
-> (OsvDbCapExceeded -> FilePath)
-> ([OsvDbCapExceeded] -> ShowS)
-> Show OsvDbCapExceeded
forall a.
(Int -> a -> ShowS) -> (a -> FilePath) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> OsvDbCapExceeded -> ShowS
showsPrec :: Int -> OsvDbCapExceeded -> ShowS
$cshow :: OsvDbCapExceeded -> FilePath
show :: OsvDbCapExceeded -> FilePath
$cshowList :: [OsvDbCapExceeded] -> ShowS
showList :: [OsvDbCapExceeded] -> ShowS
Show)

instance Exception OsvDbCapExceeded

-- | Everything one ecosystem's sync task operates on.
data SyncEnv = SyncEnv
    { SyncEnv -> CveFetch
syncFetch :: CveFetch
    -- ^ The transport for this ecosystem's object key.
    , SyncEnv -> Ecosystem
syncEcosystem :: Ecosystem
    -- ^ The ecosystem the artifact must verify as.
    , SyncEnv -> EpssRequirement
syncEpssRequirement :: EpssRequirement
    -- ^ Whether this ecosystem requires successful EPSS enrichment.
    , SyncEnv -> FilePath
syncDbPath :: FilePath
    -- ^ The canonical on-disk artifact path (the stable per-ecosystem name).
    , SyncEnv -> CveSlot
syncSlot :: CveSlot
    -- ^ The slot this task's swaps publish to.
    , SyncEnv -> Text
syncStoreRef :: Text
    -- ^ How the configured store reads back, for the reports that name where an artifact belongs.
    }

{- | What one 'syncStep' concluded. The caller ('runCveSync') logs it and decides
scheduling.
-}
data SyncOutcome
    = -- | Verification accepted a new artifact and it is now live (its ETag and provenance carried).
      SyncSwapped DbEtag [(Text, Text)]
    | -- | No database replacement, though an accepted republication can advance publication time.
      SyncUnchanged
    | -- | The object does not exist in the bucket (not yet published).
      SyncAbsent
    | {- | The artifact downloaded, and verification __refused__ it. The last-good
      generation keeps serving and the sync remembers the ETag.
      -}
      SyncRejected DbEtag CveDbRejected
    | {- | The fetch itself failed (carried). The step learned nothing about the
      remote artifact, so the last seen ETag stands and the schedule retries.
      -}
      SyncFetchFaulted OsvDbFetchFault
    deriving stock (Int -> SyncOutcome -> ShowS
[SyncOutcome] -> ShowS
SyncOutcome -> FilePath
(Int -> SyncOutcome -> ShowS)
-> (SyncOutcome -> FilePath)
-> ([SyncOutcome] -> ShowS)
-> Show SyncOutcome
forall a.
(Int -> a -> ShowS) -> (a -> FilePath) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> SyncOutcome -> ShowS
showsPrec :: Int -> SyncOutcome -> ShowS
$cshow :: SyncOutcome -> FilePath
show :: SyncOutcome -> FilePath
$cshowList :: [SyncOutcome] -> ShowS
showList :: [SyncOutcome] -> ShowS
Show)

{- | One detect-download-verify-swap cycle against the last seen ETag. Total over the fetch and
over verification: a failed fetch and a refused artifact are outcomes, not exceptions.
-}
syncStep :: SyncEnv -> Maybe DbEtag -> IO SyncOutcome
syncStep :: SyncEnv -> Maybe DbEtag -> IO SyncOutcome
syncStep SyncEnv
env Maybe DbEtag
lastSeen =
    CveFetch -> IO (Either OsvDbFetchFault (Maybe FetchedObject))
fetchHead (SyncEnv -> CveFetch
syncFetch SyncEnv
env) IO (Either OsvDbFetchFault (Maybe FetchedObject))
-> (Either OsvDbFetchFault (Maybe FetchedObject) -> IO SyncOutcome)
-> IO SyncOutcome
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
        Left OsvDbFetchFault
fault -> SyncOutcome -> IO SyncOutcome
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (OsvDbFetchFault -> SyncOutcome
SyncFetchFaulted OsvDbFetchFault
fault)
        Right Maybe FetchedObject
Nothing -> SyncOutcome -> IO SyncOutcome
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure SyncOutcome
SyncAbsent
        Right (Just FetchedObject
remote)
            | DbEtag -> Maybe DbEtag
forall a. a -> Maybe a
Just (FetchedObject -> DbEtag
foEtag FetchedObject
remote) Maybe DbEtag -> Maybe DbEtag -> Bool
forall a. Eq a => a -> a -> Bool
== Maybe DbEtag
lastSeen -> do
                CveSlot -> DbEtag -> Maybe UTCTime -> IO ()
observeAdvisoryPublication (SyncEnv -> CveSlot
syncSlot SyncEnv
env) (FetchedObject -> DbEtag
foEtag FetchedObject
remote) (FetchedObject -> Maybe UTCTime
foPushedAt FetchedObject
remote)
                SyncOutcome -> IO SyncOutcome
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure SyncOutcome
SyncUnchanged
            | Bool
otherwise -> SyncEnv -> IO SyncOutcome
syncNewArtifact SyncEnv
env

-- Nothing unverified is renamed onto the name the read path opens. The 'onException' guards
-- absorb nothing: they discard the temp file when a filesystem fault escapes, then re-propagate.
syncNewArtifact :: SyncEnv -> IO SyncOutcome
syncNewArtifact :: SyncEnv -> IO SyncOutcome
syncNewArtifact SyncEnv
env = do
    let temp :: FilePath
temp = SyncEnv -> FilePath
syncDbPath SyncEnv
env FilePath -> ShowS
forall a. Semigroup a => a -> a -> a
<> FilePath
".tmp"
    downloaded <- CveFetch -> FilePath -> IO (Either OsvDbFetchFault FetchedObject)
fetchDownload (SyncEnv -> CveFetch
syncFetch SyncEnv
env) FilePath
temp IO (Either OsvDbFetchFault FetchedObject)
-> IO () -> IO (Either OsvDbFetchFault FetchedObject)
forall (m :: * -> *) a b. MonadUnliftIO m => m a -> m b -> m a
`onException` FilePath -> IO ()
discardTemp FilePath
temp
    case downloaded of
        Left OsvDbFetchFault
fault -> do
            -- A byte-cap failure can leave a partial file.
            FilePath -> IO ()
discardTemp FilePath
temp
            SyncOutcome -> IO SyncOutcome
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (OsvDbFetchFault -> SyncOutcome
SyncFetchFaulted OsvDbFetchFault
fault)
        Right FetchedObject
fetched -> do
            opened <- Ecosystem
-> EpssRequirement -> FilePath -> IO (Either CveDbRejected CveDb)
openCveDb (SyncEnv -> Ecosystem
syncEcosystem SyncEnv
env) (SyncEnv -> EpssRequirement
syncEpssRequirement SyncEnv
env) FilePath
temp IO (Either CveDbRejected CveDb)
-> IO () -> IO (Either CveDbRejected CveDb)
forall (m :: * -> *) a b. MonadUnliftIO m => m a -> m b -> m a
`onException` FilePath -> IO ()
discardTemp FilePath
temp
            case opened of
                Left CveDbRejected
rejection -> do
                    FilePath -> IO ()
discardTemp FilePath
temp
                    SyncOutcome -> IO SyncOutcome
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (DbEtag -> CveDbRejected -> SyncOutcome
SyncRejected (FetchedObject -> DbEtag
foEtag FetchedObject
fetched) CveDbRejected
rejection)
                Right CveDb
db -> SyncEnv -> FilePath -> FetchedObject -> CveDb -> IO SyncOutcome
publishVerified SyncEnv
env FilePath
temp FetchedObject
fetched CveDb
db

publishVerified :: SyncEnv -> FilePath -> FetchedObject -> CveDb -> IO SyncOutcome
publishVerified :: SyncEnv -> FilePath -> FetchedObject -> CveDb -> IO SyncOutcome
publishVerified SyncEnv
env FilePath
temp FetchedObject
fetched CveDb
db = ((forall a. IO a -> IO a) -> IO SyncOutcome) -> IO SyncOutcome
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> m a) -> m b) -> m b
mask (((forall a. IO a -> IO a) -> IO SyncOutcome) -> IO SyncOutcome)
-> ((forall a. IO a -> IO a) -> IO SyncOutcome) -> IO SyncOutcome
forall a b. (a -> b) -> a -> b
$ \forall a. IO a -> IO a
restore -> do
    -- The verified connection follows the inode through the rename. This side still owns it,
    -- so a failure closes the connection and discards the download.
    IO () -> IO ()
forall a. IO a -> IO a
restore (FilePath -> FilePath -> IO ()
renameFile FilePath
temp (SyncEnv -> FilePath
syncDbPath SyncEnv
env))
        IO () -> IO () -> IO ()
forall (m :: * -> *) a b. MonadUnliftIO m => m a -> m b -> m a
`onException` (CveDb -> IO ()
cveDbClose CveDb
db IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> FilePath -> IO ()
discardTemp FilePath
temp)
    -- 'swapIn' owns the connection from entry, so nothing wraps it: a failure while the displaced
    -- generation drains must never close the newly live database. The mask pins the handoff.
    CveSlot -> DbEtag -> Maybe UTCTime -> CveDb -> IO ()
swapIn (SyncEnv -> CveSlot
syncSlot SyncEnv
env) (FetchedObject -> DbEtag
foEtag FetchedObject
fetched) (FetchedObject -> Maybe UTCTime
foPushedAt FetchedObject
fetched) CveDb
db
    SyncOutcome -> IO SyncOutcome
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (DbEtag -> [(Text, Text)] -> SyncOutcome
SyncSwapped (FetchedObject -> DbEtag
foEtag FetchedObject
fetched) (CveDb -> [(Text, Text)]
cveDbMeta CveDb
db))

-- Best-effort: the temp may already be renamed away or never created.
discardTemp :: FilePath -> IO ()
discardTemp :: FilePath -> IO ()
discardTemp FilePath
temp = FilePath -> IO ()
removeFile FilePath
temp IO () -> (SomeException -> IO ()) -> IO ()
forall (m :: * -> *) a.
MonadUnliftIO m =>
m a -> (SomeException -> m a) -> m a
`catchAny` IO () -> SomeException -> IO ()
forall a b. a -> b -> a
const IO ()
forall (f :: * -> *). Applicative f => f ()
pass

{- | The task's timing: the boot burst's backoff delays and the steady poll interval, both in
microseconds. The composition root ships 'bootBackoffDelays' and the configured poll interval.
-}
data SyncSchedule = SyncSchedule
    { SyncSchedule -> [Int]
schedBootBackoff :: [Int]
    -- ^ Delays before each boot-burst retry. The list's length is the budget.
    , SyncSchedule -> Int
schedPollDelay :: Int
    -- ^ The steady ETag-poll interval.
    , SyncSchedule -> Int
schedAbsentReport :: Int
    -- ^ How long between repeats of the unloaded-database and fetch-failure reports.
    }

{- | The shipped boot-burst backoff: an immediate first attempt, then a retry after each delay,
then the burst concedes to the steady poll. The poll interval, not this, is the operator's knob.
-}
bootBackoffDelays :: [Int]
bootBackoffDelays :: [Int]
bootBackoffDelays = [Int
1_000_000, Int
2_000_000, Int
4_000_000, Int
8_000_000, Int
16_000_000]

{- | The shipped gap, in microseconds, between repeats of the unloaded-database and fetch-failure
reports. The rules' outage reminder paces on the same gap.
-}
absentReportInterval :: Int
absentReportInterval :: Int
absentReportInterval = Int
900_000_000

{- | What the shell hangs off one sync task. Both run inside the task, so neither may block it,
and both must tolerate being called again.
-}
data SyncHooks = SyncHooks
    { SyncHooks -> IO ()
hookFirstSync :: IO ()
    -- ^ Runs after every swap, so it must be idempotent.
    , SyncHooks -> IO ()
hookPushAge :: IO ()
    {- ^ Runs after every step, settled or not, so the push age is read on a poll that
    changed nothing.
    -}
    }

-- | Retry at boot, then poll forever. A refused artifact ends the boot burst.
runCveSync ::
    (MonadUnliftIO m, KatipContext m) =>
    AdvisorySyncMetricsPort ->
    AdvisorySyncTracingPort ->
    SyncEnv ->
    SyncSchedule ->
    SyncHooks ->
    m ()
runCveSync :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
AdvisorySyncMetricsPort
-> AdvisorySyncTracingPort
-> SyncEnv
-> SyncSchedule
-> SyncHooks
-> m ()
runCveSync AdvisorySyncMetricsPort
metrics AdvisorySyncTracingPort
tracing SyncEnv
env SyncSchedule
schedule SyncHooks
hooks =
    SyncLoop -> Pacing -> Int -> [Int] -> m (Pacing, Maybe DbEtag)
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Pacing -> Int -> [Int] -> m (Pacing, Maybe DbEtag)
burstCycle SyncLoop
loop Pacing
initialPacing Int
0 (SyncSchedule -> [Int]
schedBootBackoff SyncSchedule
schedule) m (Pacing, Maybe DbEtag)
-> ((Pacing, Maybe DbEtag) -> m ()) -> m ()
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= (Pacing -> Maybe DbEtag -> m ()) -> (Pacing, Maybe DbEtag) -> m ()
forall a b c. (a -> b -> c) -> (a, b) -> c
uncurry (SyncLoop -> Pacing -> Maybe DbEtag -> m ()
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Pacing -> Maybe DbEtag -> m ()
pollCycle SyncLoop
loop)
  where
    loop :: SyncLoop
loop =
        SyncLoop
            { slMetrics :: AdvisorySyncMetricsPort
slMetrics = AdvisorySyncMetricsPort
metrics
            , slTracing :: AdvisorySyncTracingPort
slTracing = AdvisorySyncTracingPort
tracing
            , slEnv :: SyncEnv
slEnv = SyncEnv
env
            , slSchedule :: SyncSchedule
slSchedule = SyncSchedule
schedule
            , slHooks :: SyncHooks
slHooks = SyncHooks
hooks
            , slEcosystem :: Text
slEcosystem = Ecosystem -> Text
forall b a. (Show a, IsString b) => a -> b
show (SyncEnv -> Ecosystem
syncEcosystem SyncEnv
env)
            }

-- Everything the loop's arms read. The ecosystem label is rendered once, at the top of the task.
data SyncLoop = SyncLoop
    { SyncLoop -> AdvisorySyncMetricsPort
slMetrics :: AdvisorySyncMetricsPort
    , SyncLoop -> AdvisorySyncTracingPort
slTracing :: AdvisorySyncTracingPort
    , SyncLoop -> SyncEnv
slEnv :: SyncEnv
    , SyncLoop -> SyncSchedule
slSchedule :: SyncSchedule
    , SyncLoop -> SyncHooks
slHooks :: SyncHooks
    , SyncLoop -> Text
slEcosystem :: Text
    }

{- The loop's pacing of its two repeating reports, both on 'schedAbsentReport': the time since the
unloaded-database report, and the time since the fetch-failure report while fetches keep failing. -}
data Pacing = Pacing
    { Pacing -> Int
pacUnloaded :: Int
    , Pacing -> Maybe Int
pacFetchFailure :: Maybe Int
    }

initialPacing :: Pacing
initialPacing :: Pacing
initialPacing = Pacing{pacUnloaded :: Int
pacUnloaded = Int
0, pacFetchFailure :: Maybe Int
pacFetchFailure = Maybe Int
forall a. Maybe a
Nothing}

loopStep :: (MonadUnliftIO m, KatipContext m) => SyncLoop -> Maybe DbEtag -> m Stepped
loopStep :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Maybe DbEtag -> m Stepped
loopStep SyncLoop
loop Maybe DbEtag
lastSeen = do
    stepped <-
        AdvisorySyncMetricsPort
-> AdvisorySyncTracingPort
-> SyncEnv
-> Text
-> IO ()
-> Maybe DbEtag
-> m Stepped
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
AdvisorySyncMetricsPort
-> AdvisorySyncTracingPort
-> SyncEnv
-> Text
-> IO ()
-> Maybe DbEtag
-> m Stepped
observedStep
            (SyncLoop -> AdvisorySyncMetricsPort
slMetrics SyncLoop
loop)
            (SyncLoop -> AdvisorySyncTracingPort
slTracing SyncLoop
loop)
            (SyncLoop -> SyncEnv
slEnv SyncLoop
loop)
            (SyncLoop -> Text
slEcosystem SyncLoop
loop)
            (SyncHooks -> IO ()
hookFirstSync (SyncLoop -> SyncHooks
slHooks SyncLoop
loop))
            Maybe DbEtag
lastSeen
    liftIO (hookPushAge (slHooks loop))
    pure stepped

-- Each attempt reads 'Nothing' as last seen, because a not-settled outcome never advances it.
-- The burst concedes to the steady poll once its delays are spent.
burstCycle :: (MonadUnliftIO m, KatipContext m) => SyncLoop -> Pacing -> Int -> [Int] -> m (Pacing, Maybe DbEtag)
burstCycle :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Pacing -> Int -> [Int] -> m (Pacing, Maybe DbEtag)
burstCycle SyncLoop
loop Pacing
pacing Int
delta [Int]
delays = do
    stepped <- SyncLoop -> Maybe DbEtag -> m Stepped
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Maybe DbEtag -> m Stepped
loopStep SyncLoop
loop Maybe DbEtag
forall a. Maybe a
Nothing
    pacing' <- reportFetch loop delta stepped pacing
    case delays of
        [Int]
_ | Stepped -> Bool
stSettled Stepped
stepped -> (Pacing, Maybe DbEtag) -> m (Pacing, Maybe DbEtag)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Pacing
pacing', Stepped -> Maybe DbEtag
stSeen Stepped
stepped)
        [] -> (Pacing
pacing', Stepped -> Maybe DbEtag
stSeen Stepped
stepped) (Pacing, Maybe DbEtag) -> m () -> m (Pacing, Maybe DbEtag)
forall a b. a -> m b -> m a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ SyncLoop -> AdvisorySyncResult -> m ()
forall (m :: * -> *).
KatipContext m =>
SyncLoop -> AdvisorySyncResult -> m ()
reportLoopUnloaded SyncLoop
loop (Stepped -> AdvisorySyncResult
stResult Stepped
stepped)
        Int
delay : [Int]
rest -> Int -> m ()
forall (m :: * -> *). MonadIO m => Int -> m ()
threadDelay Int
delay m () -> m (Pacing, Maybe DbEtag) -> m (Pacing, Maybe DbEtag)
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> SyncLoop -> Pacing -> Int -> [Int] -> m (Pacing, Maybe DbEtag)
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Pacing -> Int -> [Int] -> m (Pacing, Maybe DbEtag)
burstCycle SyncLoop
loop Pacing
pacing' Int
delay [Int]
rest

pollCycle :: (MonadUnliftIO m, KatipContext m) => SyncLoop -> Pacing -> Maybe DbEtag -> m ()
pollCycle :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Pacing -> Maybe DbEtag -> m ()
pollCycle SyncLoop
loop Pacing
pacing Maybe DbEtag
lastSeen = do
    Int -> m ()
forall (m :: * -> *). MonadIO m => Int -> m ()
threadDelay Int
pollDelay
    stepped <- SyncLoop -> Maybe DbEtag -> m Stepped
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncLoop -> Maybe DbEtag -> m Stepped
loopStep SyncLoop
loop Maybe DbEtag
lastSeen
    pacing' <- reportFetch loop pollDelay stepped pacing
    unloaded <- repeatUnloaded loop (stResult stepped) (pacUnloaded pacing' + pollDelay)
    pollCycle loop pacing'{pacUnloaded = unloaded} (stSeen stepped)
  where
    pollDelay :: Int
pollDelay = SyncSchedule -> Int
schedPollDelay (SyncLoop -> SyncSchedule
slSchedule SyncLoop
loop)

reportFetch :: (KatipContext m) => SyncLoop -> Int -> Stepped -> Pacing -> m Pacing
reportFetch :: forall (m :: * -> *).
KatipContext m =>
SyncLoop -> Int -> Stepped -> Pacing -> m Pacing
reportFetch SyncLoop
loop Int
delta Stepped
stepped Pacing
pacing = do
    let (Maybe Int
fetching, Maybe FetchHealth
report) =
            Int
-> Int
-> Maybe OsvDbFetchFault
-> Maybe Int
-> (Maybe Int, Maybe FetchHealth)
paceFetchFailure (SyncSchedule -> Int
schedAbsentReport (SyncLoop -> SyncSchedule
slSchedule SyncLoop
loop)) Int
delta (Stepped -> Maybe OsvDbFetchFault
stFault Stepped
stepped) (Pacing -> Maybe Int
pacFetchFailure Pacing
pacing)
    (FetchHealth -> m ()) -> Maybe FetchHealth -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (Text -> FetchHealth -> m ()
forall (m :: * -> *). KatipContext m => Text -> FetchHealth -> m ()
reportFetchHealth (SyncLoop -> Text
slEcosystem SyncLoop
loop)) Maybe FetchHealth
report
    Pacing -> m Pacing
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Pacing
pacing{pacFetchFailure = fetching}

-- The report repeats only while the slot has never been filled, so the first swap ends it
-- and a later outage starts the interval again.
repeatUnloaded :: (KatipContext m) => SyncLoop -> AdvisorySyncResult -> Int -> m Int
repeatUnloaded :: forall (m :: * -> *).
KatipContext m =>
SyncLoop -> AdvisorySyncResult -> Int -> m Int
repeatUnloaded SyncLoop
loop AdvisorySyncResult
result Int
elapsed =
    IO (Maybe DbEtag) -> m (Maybe DbEtag)
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (CveSlot -> IO (Maybe DbEtag)
currentAdvisoryEtag (SyncEnv -> CveSlot
syncSlot (SyncLoop -> SyncEnv
slEnv SyncLoop
loop))) m (Maybe DbEtag) -> (Maybe DbEtag -> m Int) -> m Int
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
        Just DbEtag
_ -> Int -> m Int
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Int
0
        Maybe DbEtag
Nothing
            | Int
elapsed Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
< SyncSchedule -> Int
schedAbsentReport (SyncLoop -> SyncSchedule
slSchedule SyncLoop
loop) -> Int -> m Int
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Int
elapsed
            | Bool
otherwise -> Int
0 Int -> m () -> m Int
forall a b. a -> m b -> m a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ SyncLoop -> AdvisorySyncResult -> m ()
forall (m :: * -> *).
KatipContext m =>
SyncLoop -> AdvisorySyncResult -> m ()
reportLoopUnloaded SyncLoop
loop AdvisorySyncResult
result

reportLoopUnloaded :: (KatipContext m) => SyncLoop -> AdvisorySyncResult -> m ()
reportLoopUnloaded :: forall (m :: * -> *).
KatipContext m =>
SyncLoop -> AdvisorySyncResult -> m ()
reportLoopUnloaded SyncLoop
loop = Text -> Text -> AdvisorySyncResult -> m ()
forall (m :: * -> *).
KatipContext m =>
Text -> Text -> AdvisorySyncResult -> m ()
reportUnloaded (SyncLoop -> Text
slEcosystem SyncLoop
loop) (SyncEnv -> Text
syncStoreRef (SyncLoop -> SyncEnv
slEnv SyncLoop
loop))

-- What one paced fetch outcome reports, if anything.
data FetchHealth
    = FetchFailing OsvDbFetchFault
    | FetchStillFailing OsvDbFetchFault
    | FetchRecovered

{- Pace the fetch-failure report: the first failure reports at once, a later one only after
@interval@, and the first fetch that succeeds after a failure reports the recovery. -}
paceFetchFailure :: Int -> Int -> Maybe OsvDbFetchFault -> Maybe Int -> (Maybe Int, Maybe FetchHealth)
paceFetchFailure :: Int
-> Int
-> Maybe OsvDbFetchFault
-> Maybe Int
-> (Maybe Int, Maybe FetchHealth)
paceFetchFailure Int
interval Int
delta Maybe OsvDbFetchFault
fault Maybe Int
failing = case (Maybe OsvDbFetchFault
fault, Maybe Int
failing) of
    (Maybe OsvDbFetchFault
Nothing, Maybe Int
Nothing) -> (Maybe Int
forall a. Maybe a
Nothing, Maybe FetchHealth
forall a. Maybe a
Nothing)
    (Maybe OsvDbFetchFault
Nothing, Just Int
_) -> (Maybe Int
forall a. Maybe a
Nothing, FetchHealth -> Maybe FetchHealth
forall a. a -> Maybe a
Just FetchHealth
FetchRecovered)
    (Just OsvDbFetchFault
f, Maybe Int
Nothing) -> (Int -> Maybe Int
forall a. a -> Maybe a
Just Int
0, FetchHealth -> Maybe FetchHealth
forall a. a -> Maybe a
Just (OsvDbFetchFault -> FetchHealth
FetchFailing OsvDbFetchFault
f))
    (Just OsvDbFetchFault
f, Just Int
elapsed)
        | Int
elapsed Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
delta Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= Int
interval -> (Int -> Maybe Int
forall a. a -> Maybe a
Just Int
0, FetchHealth -> Maybe FetchHealth
forall a. a -> Maybe a
Just (OsvDbFetchFault -> FetchHealth
FetchStillFailing OsvDbFetchFault
f))
        | Bool
otherwise -> (Int -> Maybe Int
forall a. a -> Maybe a
Just (Int
elapsed Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
delta), Maybe FetchHealth
forall a. Maybe a
Nothing)

{- The store that keeps failing ages the serving artifact towards its maximum, so the failure logs at
the level an operator pages on. The fault names the transport cause, never a credential. -}
reportFetchHealth :: (KatipContext m) => Text -> FetchHealth -> m ()
reportFetchHealth :: forall (m :: * -> *). KatipContext m => Text -> FetchHealth -> m ()
reportFetchHealth Text
eco = \case
    FetchFailing OsvDbFetchFault
fault -> Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
ErrorS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: sync fetch failed: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> OsvDbFetchFault -> Text
forall b a. (Show a, IsString b) => a -> b
show OsvDbFetchFault
fault))
    FetchStillFailing OsvDbFetchFault
fault -> Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
ErrorS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: sync fetch still failing: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> OsvDbFetchFault -> Text
forall b a. (Show a, IsString b) => a -> b
show OsvDbFetchFault
fault))
    FetchHealth
FetchRecovered -> 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
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: sync fetch recovered"))

{- The line an operator alerts on while nothing is loaded. Its cause separates an artifact never
published from one verification refused and from an access that keeps failing. -}
reportUnloaded :: (KatipContext m) => Text -> Text -> AdvisorySyncResult -> m ()
reportUnloaded :: forall (m :: * -> *).
KatipContext m =>
Text -> Text -> AdvisorySyncResult -> m ()
reportUnloaded Text
eco Text
store AdvisorySyncResult
result =
    Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
ErrorS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> AdvisorySyncResult -> Text
unloadedCause Text
store AdvisorySyncResult
result))

-- Why nothing is loaded, worded to name the role or the access that has to change.
unloadedCause :: Text -> AdvisorySyncResult -> Text
unloadedCause :: Text -> AdvisorySyncResult -> Text
unloadedCause Text
store = \case
    AdvisorySyncResult
AdvisoryNonePublished ->
        Text
"no advisory artifact has ever been published to " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
store Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
", and ecluse pilot is what compiles and publishes one. This ecosystem stays not-ready and its advisory denies refuse until an artifact lands."
    AdvisorySyncResult
AdvisoryRefused -> Text
refusedArtifact
    -- A poll that finds nothing changed while nothing is loaded is the refused artifact standing.
    AdvisorySyncResult
AdvisoryUnchanged -> Text
refusedArtifact
    AdvisorySyncResult
AdvisoryFetchFailed ->
        Text
"no advisory database could be fetched from " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
store Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
": this ecosystem stays not-ready and denies by default until one loads. Continuing to poll; investigate the bucket, object, or IAM if this persists."
    AdvisorySyncResult
AdvisorySwapped ->
        Text
"no advisory database is loaded from " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
store Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
", so this ecosystem stays not-ready and denies by default until one is."
  where
    refusedArtifact :: Text
refusedArtifact =
        Text
"the advisory artifact at " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
store Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" was refused by verification, so nothing is loaded and this ecosystem stays not-ready. The refusal line names what failed, and ecluse pilot must publish an artifact that verifies."

-- One observed step as the loop reads it. A fetch fault rides along for the loop's paced report.
data Stepped = Stepped
    { Stepped -> AdvisorySyncResult
stResult :: AdvisorySyncResult
    , Stepped -> Bool
stSettled :: Bool
    -- ^ Whether the boot burst may stop.
    , Stepped -> Maybe DbEtag
stSeen :: Maybe DbEtag
    -- ^ The ETag now last seen.
    , Stepped -> Maybe OsvDbFetchFault
stFault :: Maybe OsvDbFetchFault
    }

steppedOf :: AdvisorySyncResult -> Bool -> Maybe DbEtag -> Stepped
steppedOf :: AdvisorySyncResult -> Bool -> Maybe DbEtag -> Stepped
steppedOf AdvisorySyncResult
result Bool
settled Maybe DbEtag
seen =
    Stepped{stResult :: AdvisorySyncResult
stResult = AdvisorySyncResult
result, stSettled :: Bool
stSettled = Bool
settled, stSeen :: Maybe DbEtag
stSeen = Maybe DbEtag
seen, stFault :: Maybe OsvDbFetchFault
stFault = Maybe OsvDbFetchFault
forall a. Maybe a
Nothing}

-- One observed step: the attempt, timed and labelled, inside this ecosystem's attempt span.
observedStep ::
    (MonadUnliftIO m, KatipContext m) =>
    AdvisorySyncMetricsPort ->
    AdvisorySyncTracingPort ->
    SyncEnv ->
    Text ->
    IO () ->
    Maybe DbEtag ->
    m Stepped
observedStep :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
AdvisorySyncMetricsPort
-> AdvisorySyncTracingPort
-> SyncEnv
-> Text
-> IO ()
-> Maybe DbEtag
-> m Stepped
observedStep AdvisorySyncMetricsPort
metrics AdvisorySyncTracingPort
tracing SyncEnv
env Text
eco IO ()
notifyFirstSync Maybe DbEtag
lastSeen =
    ((forall a. m a -> IO a) -> IO Stepped) -> m Stepped
forall b. ((forall a. m a -> IO a) -> IO b) -> m b
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> IO a) -> IO b) -> m b
withRunInIO (((forall a. m a -> IO a) -> IO Stepped) -> m Stepped)
-> ((forall a. m a -> IO a) -> IO Stepped) -> m Stepped
forall a b. (a -> b) -> a -> b
$ \forall a. m a -> IO a
runInIO ->
        AdvisorySyncTracingPort
-> forall a. Ecosystem -> (a -> AdvisorySyncResult) -> IO a -> IO a
astpSyncAttemptSpan
            AdvisorySyncTracingPort
tracing
            Ecosystem
ecosystem
            Stepped -> AdvisorySyncResult
stResult
            (AdvisorySyncMetricsPort -> Ecosystem -> IO Stepped -> IO Stepped
meteredStep AdvisorySyncMetricsPort
metrics Ecosystem
ecosystem (m Stepped -> IO Stepped
forall a. m a -> IO a
runInIO (SyncEnv -> Text -> IO () -> Maybe DbEtag -> m Stepped
forall (m :: * -> *).
KatipContext m =>
SyncEnv -> Text -> IO () -> Maybe DbEtag -> m Stepped
attemptStep SyncEnv
env Text
eco IO ()
notifyFirstSync Maybe DbEtag
lastSeen)))
  where
    ecosystem :: Ecosystem
ecosystem = SyncEnv -> Ecosystem
syncEcosystem SyncEnv
env

-- Residue escaping the attempt bypasses these records. An attempt that never concluded has no
-- result to label, and the supervision above reports it.
meteredStep :: AdvisorySyncMetricsPort -> Ecosystem -> IO Stepped -> IO Stepped
meteredStep :: AdvisorySyncMetricsPort -> Ecosystem -> IO Stepped -> IO Stepped
meteredStep AdvisorySyncMetricsPort
metrics Ecosystem
ecosystem IO Stepped
act = do
    (attempted, seconds) <- IO Stepped -> IO (Stepped, Double)
forall (m :: * -> *) a. MonadIO m => m a -> m (a, Double)
timedSeconds IO Stepped
act
    asmpSyncAttempt metrics ecosystem (stResult attempted)
    asmpSyncDuration metrics ecosystem (stResult attempted) seconds
    pure attempted

attemptStep :: (KatipContext m) => SyncEnv -> Text -> IO () -> Maybe DbEtag -> m Stepped
attemptStep :: forall (m :: * -> *).
KatipContext m =>
SyncEnv -> Text -> IO () -> Maybe DbEtag -> m Stepped
attemptStep SyncEnv
env Text
eco IO ()
notifyFirstSync Maybe DbEtag
lastSeen =
    IO SyncOutcome -> m SyncOutcome
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (SyncEnv -> Maybe DbEtag -> IO SyncOutcome
syncStep SyncEnv
env Maybe DbEtag
lastSeen) m SyncOutcome -> (SyncOutcome -> m Stepped) -> m Stepped
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
        SyncFetchFaulted OsvDbFetchFault
fault ->
            -- The step learned nothing about the remote artifact, so the last seen ETag and
            -- the last good database both stand and the next poll retries.
            Stepped -> m Stepped
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (AdvisorySyncResult -> Bool -> Maybe DbEtag -> Stepped
steppedOf AdvisorySyncResult
AdvisoryFetchFailed Bool
False Maybe DbEtag
lastSeen){stFault = Just fault}
        SyncSwapped DbEtag
etag [(Text, Text)]
meta -> 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
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: advisory database swapped in: etag=" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> DbEtag -> Text
forall b a. (Show a, IsString b) => a -> b
show DbEtag
etag Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" meta=" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> (Maybe UTCTime, Maybe Word64) -> Text
forall b a. (Show a, IsString b) => a -> b
show ([(Text, Text)] -> (Maybe UTCTime, Maybe Word64)
metadataSummary [(Text, Text)]
meta)))
            source <- IO (Maybe AdvisorySource) -> m (Maybe AdvisorySource)
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (CveSlot -> IO (Maybe AdvisorySource)
currentAdvisorySource (SyncEnv -> CveSlot
syncSlot SyncEnv
env))
            logFM InfoS (ls ("cve-sync[" <> eco <> "]: serving artifact source: " <> maybe unrecordedValue renderAdvisorySource source))
            whenNothing_ (asPushedAt =<< source) (undatedArtifact eco etag)
            liftIO notifyFirstSync
            pure (steppedOf AdvisorySwapped True (Just etag))
        SyncOutcome
SyncUnchanged -> do
            Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
DebugS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: advisory database unchanged"))
            Stepped -> m Stepped
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (AdvisorySyncResult -> Bool -> Maybe DbEtag -> Stepped
steppedOf AdvisorySyncResult
AdvisoryUnchanged Bool
True Maybe DbEtag
lastSeen)
        SyncOutcome
SyncAbsent -> do
            Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
DebugS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: no advisory database published yet"))
            Stepped -> m Stepped
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (AdvisorySyncResult -> Bool -> Maybe DbEtag -> Stepped
steppedOf AdvisorySyncResult
AdvisoryNonePublished Bool
False Maybe DbEtag
lastSeen)
        SyncRejected DbEtag
etag CveDbRejected
rejection -> do
            Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
ErrorS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"cve-sync[" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: downloaded artifact refused (keeping last good): " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> CveDbRejected -> Text
forall b a. (Show a, IsString b) => a -> b
show CveDbRejected
rejection))
            -- Remember the ETag so the same refused artifact is not re-downloaded.
            -- A fixed re-publish carries a new one. Identical bytes cannot end differently.
            Stepped -> m Stepped
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (AdvisorySyncResult -> Bool -> Maybe DbEtag -> Stepped
steppedOf AdvisorySyncResult
AdvisoryRefused Bool
True (DbEtag -> Maybe DbEtag
forall a. a -> Maybe a
Just DbEtag
etag))

{- An artifact the object store gave no publication time for: its age cannot be established, so
CVE-based denial refuses on it. One line per swap, because only a swap can install one. -}
undatedArtifact :: (KatipContext m) => Text -> DbEtag -> m ()
undatedArtifact :: forall (m :: * -> *). KatipContext m => Text -> DbEtag -> m ()
undatedArtifact Text
eco DbEtag
etag =
    Severity -> LogStr -> m ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM
        Severity
ErrorS
        ( Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls
            ( Text
"cve-sync["
                Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
eco
                Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"]: the object store reported no publication time for the artifact it served (etag="
                Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> DbEtag -> Text
forall b a. (Show a, IsString b) => a -> b
show DbEtag
etag
                Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"), so its age cannot be established and CVE-based denial refuses until a push carries one"
            )
        )

{- Where the serving artifact came from, for the swap line. The source renders as its authority
alone, on the same rule as 'metadataSummary' below: artifact text never reaches a log verbatim. -}
renderAdvisorySource :: AdvisorySource -> Text
renderAdvisorySource :: AdvisorySource -> Text
renderAdvisorySource AdvisorySource
source =
    Text
"pushed_at="
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Maybe UTCTime -> Text
renderStamp (AdvisorySource -> Maybe UTCTime
asPushedAt AdvisorySource
source)
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" osv_source="
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text -> (Text -> Text) -> Maybe Text -> Text
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Text
unrecordedValue Text -> Text
dialledAuthorityLabel (AdvisoryProvenance -> Maybe Text
apOsvSource AdvisoryProvenance
prov)
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" osv_newest_modified="
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Maybe UTCTime -> Text
renderStamp (AdvisoryProvenance -> Maybe UTCTime
apOsvNewestModified AdvisoryProvenance
prov)
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" epss_score_date="
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Maybe UTCTime -> Text
renderStamp (AdvisoryProvenance -> Maybe UTCTime
apEpssScoreDate AdvisoryProvenance
prov)
  where
    prov :: AdvisoryProvenance
prov = AdvisorySource -> AdvisoryProvenance
asProvenance AdvisorySource
source

renderStamp :: Maybe UTCTime -> Text
renderStamp :: Maybe UTCTime -> Text
renderStamp = Text -> (UTCTime -> Text) -> Maybe UTCTime -> Text
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Text
unrecordedValue UTCTime -> Text
renderIso8601Utc

-- What a value the artifact never recorded reads as, so absence is not read as a zero.
unrecordedValue :: Text
unrecordedValue :: Text
unrecordedValue = Text
"<unrecorded>"

-- Legacy artifacts contain arbitrary text. Only parsed, bounded values reach the log.
metadataSummary :: [(Text, Text)] -> (Maybe UTCTime, Maybe Word64)
metadataSummary :: [(Text, Text)] -> (Maybe UTCTime, Maybe Word64)
metadataSummary [(Text, Text)]
meta =
    ( MetaKey -> Int -> Maybe Text
boundedValue MetaKey
MetaBuiltAt Int
64 Maybe Text -> (Text -> Maybe UTCTime) -> Maybe UTCTime
forall a b. Maybe a -> (a -> Maybe b) -> Maybe b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= FilePath -> Maybe UTCTime
forall (m :: * -> *) t. (MonadFail m, ISO8601 t) => FilePath -> m t
iso8601ParseM (FilePath -> Maybe UTCTime)
-> (Text -> FilePath) -> Text -> Maybe UTCTime
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> FilePath
forall a. ToString a => a -> FilePath
toString
    , MetaKey -> Int -> Maybe Text
boundedValue MetaKey
MetaRowCount Int
20 Maybe Text -> (Text -> Maybe Word64) -> Maybe Word64
forall a b. Maybe a -> (a -> Maybe b) -> Maybe b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= Text -> Maybe Word64
parseMetadataCount
    )
  where
    boundedValue :: MetaKey -> Int -> Maybe Text
boundedValue MetaKey
key Int
limit = do
        value <- Text -> [(Text, Text)] -> Maybe Text
forall a b. Eq a => a -> [(a, b)] -> Maybe b
lookup (MetaKey -> Text
renderMetaKey MetaKey
key) [(Text, Text)]
meta
        guard (T.compareLength value limit /= GT)
        pure value

parseMetadataCount :: Text -> Maybe Word64
parseMetadataCount :: Text -> Maybe Word64
parseMetadataCount Text
value = do
    count <- Text -> Maybe Integer
forall a. Integral a => Text -> Maybe a
readDecimalText Text
value :: Maybe Integer
    guard (count <= toInteger (maxBound :: Word64))
    pure (fromInteger count)

{- | An S3-backed advisory-fetch source. 'newS3CveSource' captures one @amazonka@ 'AWS.Env', so
every mount's 'CveFetch' shares one credential discovery. The composition shell never sees it.
-}
newtype S3CveSource = S3CveSource
    { S3CveSource -> Text -> Text -> Int -> CveFetch
s3CveFetchFor :: Text -> Text -> Int -> CveFetch
    -- ^ A 'CveFetch' against one bucket, object key, and byte cap, over the captured env.
    }

-- | Build an 'S3CveSource' over one S3 @amazonka@ env, honouring the resolved endpoint override.
newS3CveSource :: Maybe AwsEndpoint -> IO S3CveSource
newS3CveSource :: Maybe AwsEndpoint -> IO S3CveSource
newS3CveSource Maybe AwsEndpoint
mEndpoint = do
    awsEnv <- Maybe AwsEndpoint -> IO Env
buildS3Env Maybe AwsEndpoint
mEndpoint
    pure (S3CveSource (s3CveFetch awsEnv))

s3CveFetch :: AWS.Env -> Text -> Text -> Int -> CveFetch
s3CveFetch :: Env -> Text -> Text -> Int -> CveFetch
s3CveFetch Env
awsEnv Text
bucket Text
key Int
maxBytes =
    CveFetch
        { fetchHead :: IO (Either OsvDbFetchFault (Maybe FetchedObject))
fetchHead = Env
-> Text
-> Text
-> IO (Either OsvDbFetchFault (Maybe FetchedObject))
s3Head Env
awsEnv Text
bucket Text
key
        , fetchDownload :: FilePath -> IO (Either OsvDbFetchFault FetchedObject)
fetchDownload = Env
-> Text
-> Text
-> Int
-> FilePath
-> IO (Either OsvDbFetchFault FetchedObject)
s3Download Env
awsEnv Text
bucket Text
key Int
maxBytes
        }

s3Head :: AWS.Env -> Text -> Text -> IO (Either OsvDbFetchFault (Maybe FetchedObject))
s3Head :: Env
-> Text
-> Text
-> IO (Either OsvDbFetchFault (Maybe FetchedObject))
s3Head Env
awsEnv Text
bucket Text
key =
    ResourceT IO (Either Error HeadObjectResponse)
-> IO (Either Error HeadObjectResponse)
forall (m :: * -> *) a. MonadUnliftIO m => ResourceT m a -> m a
runResourceT (Env
-> HeadObject
-> ResourceT IO (Either Error (AWSResponse HeadObject))
forall (m :: * -> *) a.
(MonadResource m, AWSRequest a) =>
Env -> a -> m (Either Error (AWSResponse a))
AWS.sendEither Env
awsEnv (BucketName -> ObjectKey -> HeadObject
S3.newHeadObject (Text -> BucketName
S3.BucketName Text
bucket) (Text -> ObjectKey
S3.ObjectKey Text
key))) IO (Either Error HeadObjectResponse)
-> (Either Error HeadObjectResponse
    -> Either OsvDbFetchFault (Maybe FetchedObject))
-> IO (Either OsvDbFetchFault (Maybe FetchedObject))
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> \case
        Right HeadObjectResponse
resp ->
            let observed :: ETag -> FetchedObject
observed ETag
etag = DbEtag -> Maybe UTCTime -> FetchedObject
FetchedObject (ETag -> DbEtag
dbEtag ETag
etag) (HeadObjectResponse
resp HeadObjectResponse
-> Getting (Maybe UTCTime) HeadObjectResponse (Maybe UTCTime)
-> Maybe UTCTime
forall s a. s -> Getting a s a -> a
^. Getting (Maybe UTCTime) HeadObjectResponse (Maybe UTCTime)
Lens' HeadObjectResponse (Maybe UTCTime)
S3L.headObjectResponse_lastModified)
             in Either OsvDbFetchFault (Maybe FetchedObject)
-> (ETag -> Either OsvDbFetchFault (Maybe FetchedObject))
-> Maybe ETag
-> Either OsvDbFetchFault (Maybe FetchedObject)
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (OsvDbFetchFault -> Either OsvDbFetchFault (Maybe FetchedObject)
forall a b. a -> Either a b
Left OsvDbFetchFault
OsvDbNoEtag) (Maybe FetchedObject -> Either OsvDbFetchFault (Maybe FetchedObject)
forall a b. b -> Either a b
Right (Maybe FetchedObject
 -> Either OsvDbFetchFault (Maybe FetchedObject))
-> (ETag -> Maybe FetchedObject)
-> ETag
-> Either OsvDbFetchFault (Maybe FetchedObject)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. FetchedObject -> Maybe FetchedObject
forall a. a -> Maybe a
Just (FetchedObject -> Maybe FetchedObject)
-> (ETag -> FetchedObject) -> ETag -> Maybe FetchedObject
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ETag -> FetchedObject
observed) (HeadObjectResponse
resp HeadObjectResponse
-> Getting (Maybe ETag) HeadObjectResponse (Maybe ETag)
-> Maybe ETag
forall s a. s -> Getting a s a -> a
^. Getting (Maybe ETag) HeadObjectResponse (Maybe ETag)
Lens' HeadObjectResponse (Maybe ETag)
S3L.headObjectResponse_eTag)
        Left Error
err
            | Error -> Bool
isNotFound Error
err -> Maybe FetchedObject -> Either OsvDbFetchFault (Maybe FetchedObject)
forall a b. b -> Either a b
Right Maybe FetchedObject
forall a. Maybe a
Nothing
            | Bool
otherwise -> OsvDbFetchFault -> Either OsvDbFetchFault (Maybe FetchedObject)
forall a b. a -> Either a b
Left (TransportFault -> OsvDbFetchFault
OsvDbTransport (Error -> TransportFault
classifyAwsTransport Error
err))

s3Download :: AWS.Env -> Text -> Text -> Int -> FilePath -> IO (Either OsvDbFetchFault FetchedObject)
s3Download :: Env
-> Text
-> Text
-> Int
-> FilePath
-> IO (Either OsvDbFetchFault FetchedObject)
s3Download Env
awsEnv Text
bucket Text
key Int
maxBytes FilePath
dest = IO (Either OsvDbFetchFault FetchedObject)
-> IO (Either OsvDbFetchFault FetchedObject)
foldFetchEscapes (IO (Either OsvDbFetchFault FetchedObject)
 -> IO (Either OsvDbFetchFault FetchedObject))
-> (ResourceT IO (Either OsvDbFetchFault FetchedObject)
    -> IO (Either OsvDbFetchFault FetchedObject))
-> ResourceT IO (Either OsvDbFetchFault FetchedObject)
-> IO (Either OsvDbFetchFault FetchedObject)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ResourceT IO (Either OsvDbFetchFault FetchedObject)
-> IO (Either OsvDbFetchFault FetchedObject)
forall (m :: * -> *) a. MonadUnliftIO m => ResourceT m a -> m a
runResourceT (ResourceT IO (Either OsvDbFetchFault FetchedObject)
 -> IO (Either OsvDbFetchFault FetchedObject))
-> ResourceT IO (Either OsvDbFetchFault FetchedObject)
-> IO (Either OsvDbFetchFault FetchedObject)
forall a b. (a -> b) -> a -> b
$ do
    resp <- Env -> GetObject -> ResourceT IO (AWSResponse GetObject)
forall (m :: * -> *) a.
(MonadResource m, AWSRequest a) =>
Env -> a -> m (AWSResponse a)
AWS.send Env
awsEnv (BucketName -> ObjectKey -> GetObject
S3.newGetObject (Text -> BucketName
S3.BucketName Text
bucket) (Text -> ObjectKey
S3.ObjectKey Text
key))
    -- The declared length fails fast. The streaming cap is the enforcement: a
    -- declared length is not a guarantee.
    for_ (resp ^. S3L.getObjectResponse_contentLength) $ \Integer
len ->
        Bool -> ResourceT IO () -> ResourceT IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Integer
len Integer -> Integer -> Bool
forall a. Ord a => a -> a -> Bool
> Int -> Integer
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
maxBytes) (OsvDbCapExceeded -> ResourceT IO ()
forall (m :: * -> *) e a. (MonadIO m, Exception e) => e -> m a
throwIO (Int -> OsvDbCapExceeded
OsvDbCapExceeded Int
maxBytes))
    AWS.sinkBody (resp ^. S3L.getObjectResponse_body) (cappedAt maxBytes .| C.sinkFile dest)
    let fetched ETag
etag = FetchedObject{foEtag :: DbEtag
foEtag = ETag -> DbEtag
dbEtag ETag
etag, foPushedAt :: Maybe UTCTime
foPushedAt = GetObjectResponse
resp GetObjectResponse
-> Getting (Maybe UTCTime) GetObjectResponse (Maybe UTCTime)
-> Maybe UTCTime
forall s a. s -> Getting a s a -> a
^. Getting (Maybe UTCTime) GetObjectResponse (Maybe UTCTime)
Lens' GetObjectResponse (Maybe UTCTime)
S3L.getObjectResponse_lastModified}
    pure (maybe (Left OsvDbNoEtag) (Right . fetched) (resp ^. S3L.getObjectResponse_eTag))

-- The adapter boundary: fold the two typed escapes into the value channel. Nothing else is
-- caught, so a filesystem fault writing the destination propagates as residue.
foldFetchEscapes :: IO (Either OsvDbFetchFault FetchedObject) -> IO (Either OsvDbFetchFault FetchedObject)
foldFetchEscapes :: IO (Either OsvDbFetchFault FetchedObject)
-> IO (Either OsvDbFetchFault FetchedObject)
foldFetchEscapes IO (Either OsvDbFetchFault FetchedObject)
act =
    IO (Either OsvDbFetchFault FetchedObject)
act
        IO (Either OsvDbFetchFault FetchedObject)
-> (Error -> IO (Either OsvDbFetchFault FetchedObject))
-> IO (Either OsvDbFetchFault FetchedObject)
forall (m :: * -> *) e a.
(MonadUnliftIO m, Exception e) =>
m a -> (e -> m a) -> m a
`catch` (\(Error
err :: AWS.Error) -> Either OsvDbFetchFault FetchedObject
-> IO (Either OsvDbFetchFault FetchedObject)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (OsvDbFetchFault -> Either OsvDbFetchFault FetchedObject
forall a b. a -> Either a b
Left (TransportFault -> OsvDbFetchFault
OsvDbTransport (Error -> TransportFault
classifyAwsTransport Error
err))))
        IO (Either OsvDbFetchFault FetchedObject)
-> (OsvDbCapExceeded -> IO (Either OsvDbFetchFault FetchedObject))
-> IO (Either OsvDbFetchFault FetchedObject)
forall (m :: * -> *) e a.
(MonadUnliftIO m, Exception e) =>
m a -> (e -> m a) -> m a
`catch` (\(OsvDbCapExceeded Int
n) -> Either OsvDbFetchFault FetchedObject
-> IO (Either OsvDbFetchFault FetchedObject)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (OsvDbFetchFault -> Either OsvDbFetchFault FetchedObject
forall a b. a -> Either a b
Left (Int -> OsvDbFetchFault
OsvDbTooLarge Int
n)))

dbEtag :: S3.ETag -> DbEtag
dbEtag :: ETag -> DbEtag
dbEtag (S3.ETag ByteString
bytes) = Text -> DbEtag
DbEtag (ByteString -> Text
forall a b. ConvertUtf8 a b => b -> a
decodeUtf8 ByteString
bytes)

isNotFound :: AWS.Error -> Bool
isNotFound :: Error -> Bool
isNotFound = \case
    AWS.ServiceError ServiceError
se -> Status -> Int
statusCode (ServiceError
se ServiceError -> Getting Status ServiceError Status -> Status
forall s a. s -> Getting a s a -> a
^. Getting Status ServiceError Status
Lens' ServiceError Status
AWS.serviceError_status) Int -> Int -> Bool
forall a. Eq a => a -> a -> Bool
== Int
404
    Error
_ -> Bool
False

{- | A breach throws 'OsvDbCapExceeded' before yielding the excess chunk.
The S3 adapter folds it into 'OsvDbTooLarge'.
-}
cappedAt :: (MonadIO m) => Int -> ConduitT ByteString ByteString m ()
cappedAt :: forall (m :: * -> *).
MonadIO m =>
Int -> ConduitT ByteString ByteString m ()
cappedAt Int
maxBytes = Int -> (Int -> m ()) -> ConduitT ByteString ByteString m ()
forall (m :: * -> *).
Monad m =>
Int -> (Int -> m ()) -> ConduitT ByteString ByteString m ()
boundBytes Int
maxBytes (m () -> Int -> m ()
forall a b. a -> b -> a
const (OsvDbCapExceeded -> m ()
forall (m :: * -> *) e a. (MonadIO m, Exception e) => e -> m a
throwIO (Int -> OsvDbCapExceeded
OsvDbCapExceeded Int
maxBytes)))