module Ecluse.Runtime.Cve.Sync.Internal (
CveFetch (..),
FetchedObject (..),
DbEtag (..),
OsvDbFetchFault (..),
OsvDbCapExceeded (..),
S3CveSource,
newS3CveSource,
s3CveFetchFor,
cappedAt,
SyncEnv (..),
SyncOutcome (..),
syncStep,
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)
data CveFetch = CveFetch
{ CveFetch -> IO (Either OsvDbFetchFault (Maybe FetchedObject))
fetchHead :: IO (Either OsvDbFetchFault (Maybe FetchedObject))
, CveFetch -> FilePath -> IO (Either OsvDbFetchFault FetchedObject)
fetchDownload :: FilePath -> IO (Either OsvDbFetchFault FetchedObject)
}
data FetchedObject = FetchedObject
{ FetchedObject -> DbEtag
foEtag :: DbEtag
, FetchedObject -> Maybe UTCTime
foPushedAt :: Maybe UTCTime
}
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)
data OsvDbFetchFault
=
OsvDbTooLarge Int
|
OsvDbNoEtag
|
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)
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
data SyncEnv = SyncEnv
{ SyncEnv -> CveFetch
syncFetch :: CveFetch
, SyncEnv -> Ecosystem
syncEcosystem :: Ecosystem
, SyncEnv -> EpssRequirement
syncEpssRequirement :: EpssRequirement
, SyncEnv -> FilePath
syncDbPath :: FilePath
, SyncEnv -> CveSlot
syncSlot :: CveSlot
, SyncEnv -> Text
syncStoreRef :: Text
}
data SyncOutcome
=
SyncSwapped DbEtag [(Text, Text)]
|
SyncUnchanged
|
SyncAbsent
|
SyncRejected DbEtag CveDbRejected
|
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)
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
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
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
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)
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))
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
data SyncSchedule = SyncSchedule
{ SyncSchedule -> [Int]
schedBootBackoff :: [Int]
, SyncSchedule -> Int
schedPollDelay :: Int
, SyncSchedule -> Int
schedAbsentReport :: Int
}
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]
absentReportInterval :: Int
absentReportInterval :: Int
absentReportInterval = Int
900_000_000
data SyncHooks = SyncHooks
{ SyncHooks -> IO ()
hookFirstSync :: IO ()
, SyncHooks -> IO ()
hookPushAge :: IO ()
}
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)
}
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
}
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
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}
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))
data FetchHealth
= FetchFailing OsvDbFetchFault
| FetchStillFailing OsvDbFetchFault
| FetchRecovered
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)
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"))
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))
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
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."
data Stepped = Stepped
{ Stepped -> AdvisorySyncResult
stResult :: AdvisorySyncResult
, Stepped -> Bool
stSettled :: Bool
, Stepped -> Maybe DbEtag
stSeen :: Maybe DbEtag
, 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}
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
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 ->
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))
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))
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"
)
)
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
unrecordedValue :: Text
unrecordedValue :: Text
unrecordedValue = Text
"<unrecorded>"
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)
newtype S3CveSource = S3CveSource
{ S3CveSource -> Text -> Text -> Int -> CveFetch
s3CveFetchFor :: Text -> Text -> Int -> CveFetch
}
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))
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))
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
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)))