module Ecluse.Runtime.Cve.Sync (
CveFetch (..),
DbEtag (..),
OsvDbFetchFault (..),
OsvDbCapExceeded (..),
S3CveSource,
newS3CveSource,
s3CveFetchFor,
cappedAt,
SyncEnv (..),
SyncOutcome (..),
syncStep,
SyncSchedule (..),
runCveSync,
bootBackoffDelays,
bootBurstPolicy,
) where
import Conduit (ConduitT, await, runResourceT, yield, (.|))
import Control.Retry (RetryPolicyM, RetryStatus (rsIterNumber), retryPolicy, retrying)
import Data.ByteString qualified as BS
import Data.Conduit.Combinators qualified as C
import Katip (KatipContext, Severity (DebugS, ErrorS, InfoS), logFM, ls)
import Network.HTTP.Types.Status (statusCode)
import System.Directory (removeFile, renameFile)
import UnliftIO (MonadUnliftIO)
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, DbEtag (..), openCveDb)
import Ecluse.Core.Cve.Slot (CveSlot, swapIn)
import Ecluse.Core.Ecosystem (Ecosystem)
import Ecluse.Core.Fault (TransportFault)
import Ecluse.Runtime.Aws.Fault (classifyAwsTransport)
import Ecluse.Runtime.Aws.S3 (buildS3Env)
data CveFetch = CveFetch
{ CveFetch -> IO (Either OsvDbFetchFault (Maybe DbEtag))
fetchHeadEtag :: IO (Either OsvDbFetchFault (Maybe DbEtag))
, CveFetch -> FilePath -> IO (Either OsvDbFetchFault DbEtag)
fetchDownload :: FilePath -> IO (Either OsvDbFetchFault DbEtag)
}
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 -> FilePath
syncDbPath :: FilePath
, SyncEnv -> CveSlot
syncSlot :: CveSlot
}
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 DbEtag))
fetchHeadEtag (SyncEnv -> CveFetch
syncFetch SyncEnv
env) IO (Either OsvDbFetchFault (Maybe DbEtag))
-> (Either OsvDbFetchFault (Maybe DbEtag) -> 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 DbEtag
Nothing -> SyncOutcome -> IO SyncOutcome
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure SyncOutcome
SyncAbsent
Right (Just DbEtag
remote)
| DbEtag -> Maybe DbEtag
forall a. a -> Maybe a
Just DbEtag
remote Maybe DbEtag -> Maybe DbEtag -> Bool
forall a. Eq a => a -> a -> Bool
== Maybe DbEtag
lastSeen -> 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 DbEtag)
fetchDownload (SyncEnv -> CveFetch
syncFetch SyncEnv
env) FilePath
temp IO (Either OsvDbFetchFault DbEtag)
-> IO () -> IO (Either OsvDbFetchFault DbEtag)
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 DbEtag
fetched -> do
opened <- Ecosystem -> FilePath -> IO (Either CveDbRejected CveDb)
openCveDb (SyncEnv -> Ecosystem
syncEcosystem 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 DbEtag
fetched CveDbRejected
rejection)
Right CveDb
db -> SyncEnv -> FilePath -> DbEtag -> CveDb -> IO SyncOutcome
publishVerified SyncEnv
env FilePath
temp DbEtag
fetched CveDb
db
publishVerified :: SyncEnv -> FilePath -> DbEtag -> CveDb -> IO SyncOutcome
publishVerified :: SyncEnv -> FilePath -> DbEtag -> CveDb -> IO SyncOutcome
publishVerified SyncEnv
env FilePath
temp DbEtag
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 -> CveDb -> IO ()
swapIn (SyncEnv -> CveSlot
syncSlot SyncEnv
env) DbEtag
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 DbEtag
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
}
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]
bootBurstPolicy :: (Monad m) => [Int] -> RetryPolicyM m
bootBurstPolicy :: forall (m :: * -> *). Monad m => [Int] -> RetryPolicyM m
bootBurstPolicy [Int]
delays = (RetryStatus -> Maybe Int) -> RetryPolicyM m
forall (m :: * -> *).
Monad m =>
(RetryStatus -> Maybe Int) -> RetryPolicyM m
retryPolicy (\RetryStatus
rs -> [Int]
delays [Int] -> Int -> Maybe Int
forall a. [a] -> Int -> Maybe a
!!? RetryStatus -> Int
rsIterNumber RetryStatus
rs)
runCveSync :: (MonadUnliftIO m, KatipContext m) => SyncEnv -> SyncSchedule -> IO () -> m ()
runCveSync :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncEnv -> SyncSchedule -> IO () -> m ()
runCveSync SyncEnv
env SyncSchedule
schedule IO ()
notifyFirstSync = do
seen <- m (Maybe DbEtag)
burst
poll seen
where
eco :: Text
eco = Ecosystem -> Text
forall b a. (Show a, IsString b) => a -> b
show (SyncEnv -> Ecosystem
syncEcosystem SyncEnv
env) :: Text
burst :: m (Maybe DbEtag)
burst = do
(settled, seen') <-
RetryPolicyM m
-> (RetryStatus -> (Bool, Maybe DbEtag) -> m Bool)
-> (RetryStatus -> m (Bool, Maybe DbEtag))
-> m (Bool, Maybe DbEtag)
forall (m :: * -> *) b.
MonadIO m =>
RetryPolicyM m
-> (RetryStatus -> b -> m Bool) -> (RetryStatus -> m b) -> m b
retrying
([Int] -> RetryPolicyM m
forall (m :: * -> *). Monad m => [Int] -> RetryPolicyM m
bootBurstPolicy (SyncSchedule -> [Int]
schedBootBackoff SyncSchedule
schedule))
(\RetryStatus
_ (Bool
done, Maybe DbEtag
_) -> Bool -> m Bool
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Bool -> Bool
not Bool
done))
(\RetryStatus
_ -> SyncEnv -> Text -> IO () -> Maybe DbEtag -> m (Bool, Maybe DbEtag)
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncEnv -> Text -> IO () -> Maybe DbEtag -> m (Bool, Maybe DbEtag)
loggedStep SyncEnv
env Text
eco IO ()
notifyFirstSync Maybe DbEtag
forall a. Maybe a
Nothing)
unless settled $
logFM ErrorS (ls ("cve-sync[" <> eco <> "]: boot fetch did not acquire an advisory database within the boot budget; this ecosystem stays not-ready and denies by default until one is acquired. Continuing to poll; investigate the bucket, object, or IAM if this persists."))
pure seen'
poll :: Maybe DbEtag -> m b
poll Maybe DbEtag
lastSeen = do
Int -> m ()
forall (m :: * -> *). MonadIO m => Int -> m ()
threadDelay (SyncSchedule -> Int
schedPollDelay SyncSchedule
schedule)
(_, seen') <- SyncEnv -> Text -> IO () -> Maybe DbEtag -> m (Bool, Maybe DbEtag)
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncEnv -> Text -> IO () -> Maybe DbEtag -> m (Bool, Maybe DbEtag)
loggedStep SyncEnv
env Text
eco IO ()
notifyFirstSync Maybe DbEtag
lastSeen
poll seen'
loggedStep :: (MonadUnliftIO m, KatipContext m) => SyncEnv -> Text -> IO () -> Maybe DbEtag -> m (Bool, Maybe DbEtag)
loggedStep :: forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SyncEnv -> Text -> IO () -> Maybe DbEtag -> m (Bool, Maybe DbEtag)
loggedStep 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 (Bool, Maybe DbEtag))
-> m (Bool, Maybe DbEtag)
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 -> 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
"]: 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))
(Bool, Maybe DbEtag) -> m (Bool, Maybe DbEtag)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Bool
False, Maybe DbEtag
lastSeen)
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
<> [(Text, Text)] -> Text
forall b a. (Show a, IsString b) => a -> b
show [(Text, Text)]
meta))
IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO IO ()
notifyFirstSync
(Bool, Maybe DbEtag) -> m (Bool, Maybe DbEtag)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Bool
True, DbEtag -> Maybe DbEtag
forall a. a -> Maybe a
Just DbEtag
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"))
(Bool, Maybe DbEtag) -> m (Bool, Maybe DbEtag)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (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"))
(Bool, Maybe DbEtag) -> m (Bool, Maybe DbEtag)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (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))
(Bool, Maybe DbEtag) -> m (Bool, Maybe DbEtag)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Bool
True, DbEtag -> Maybe DbEtag
forall a. a -> Maybe a
Just DbEtag
etag)
newtype S3CveSource = S3CveSource
{ S3CveSource -> Text -> Text -> Int -> CveFetch
s3CveFetchFor :: Text -> Text -> Int -> CveFetch
}
newS3CveSource :: Maybe (Bool, Text, Int) -> IO S3CveSource
newS3CveSource :: Maybe (Bool, Text, Int) -> IO S3CveSource
newS3CveSource Maybe (Bool, Text, Int)
mEndpoint = do
awsEnv <- Maybe (Bool, Text, Int) -> IO Env
buildS3Env Maybe (Bool, Text, Int)
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
{ fetchHeadEtag :: IO (Either OsvDbFetchFault (Maybe DbEtag))
fetchHeadEtag = Env -> Text -> Text -> IO (Either OsvDbFetchFault (Maybe DbEtag))
s3HeadEtag Env
awsEnv Text
bucket Text
key
, fetchDownload :: FilePath -> IO (Either OsvDbFetchFault DbEtag)
fetchDownload = Env
-> Text
-> Text
-> Int
-> FilePath
-> IO (Either OsvDbFetchFault DbEtag)
s3Download Env
awsEnv Text
bucket Text
key Int
maxBytes
}
s3HeadEtag :: AWS.Env -> Text -> Text -> IO (Either OsvDbFetchFault (Maybe DbEtag))
s3HeadEtag :: Env -> Text -> Text -> IO (Either OsvDbFetchFault (Maybe DbEtag))
s3HeadEtag 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 DbEtag))
-> IO (Either OsvDbFetchFault (Maybe DbEtag))
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> \case
Right HeadObjectResponse
resp -> Maybe DbEtag -> Either OsvDbFetchFault (Maybe DbEtag)
forall a b. b -> Either a b
Right (ETag -> DbEtag
dbEtag (ETag -> DbEtag) -> Maybe ETag -> Maybe DbEtag
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> 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 DbEtag -> Either OsvDbFetchFault (Maybe DbEtag)
forall a b. b -> Either a b
Right Maybe DbEtag
forall a. Maybe a
Nothing
| Bool
otherwise -> OsvDbFetchFault -> Either OsvDbFetchFault (Maybe DbEtag)
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 DbEtag)
s3Download :: Env
-> Text
-> Text
-> Int
-> FilePath
-> IO (Either OsvDbFetchFault DbEtag)
s3Download Env
awsEnv Text
bucket Text
key Int
maxBytes FilePath
dest = IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag)
classified (IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag))
-> (ResourceT IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag))
-> ResourceT IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ResourceT IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag)
forall (m :: * -> *) a. MonadUnliftIO m => ResourceT m a -> m a
runResourceT (ResourceT IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag))
-> ResourceT IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag)
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)
pure (maybe (Left OsvDbNoEtag) (Right . dbEtag) (resp ^. S3L.getObjectResponse_eTag))
where
classified :: IO (Either OsvDbFetchFault DbEtag) -> IO (Either OsvDbFetchFault DbEtag)
classified :: IO (Either OsvDbFetchFault DbEtag)
-> IO (Either OsvDbFetchFault DbEtag)
classified IO (Either OsvDbFetchFault DbEtag)
act =
IO (Either OsvDbFetchFault DbEtag)
act
IO (Either OsvDbFetchFault DbEtag)
-> (Error -> IO (Either OsvDbFetchFault DbEtag))
-> IO (Either OsvDbFetchFault DbEtag)
forall (m :: * -> *) e a.
(MonadUnliftIO m, Exception e) =>
m a -> (e -> m a) -> m a
`catch` (\(Error
err :: AWS.Error) -> Either OsvDbFetchFault DbEtag -> IO (Either OsvDbFetchFault DbEtag)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (OsvDbFetchFault -> Either OsvDbFetchFault DbEtag
forall a b. a -> Either a b
Left (TransportFault -> OsvDbFetchFault
OsvDbTransport (Error -> TransportFault
classifyAwsTransport Error
err))))
IO (Either OsvDbFetchFault DbEtag)
-> (OsvDbCapExceeded -> IO (Either OsvDbFetchFault DbEtag))
-> IO (Either OsvDbFetchFault DbEtag)
forall (m :: * -> *) e a.
(MonadUnliftIO m, Exception e) =>
m a -> (e -> m a) -> m a
`catch` (\(OsvDbCapExceeded Int
n) -> Either OsvDbFetchFault DbEtag -> IO (Either OsvDbFetchFault DbEtag)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (OsvDbFetchFault -> Either OsvDbFetchFault DbEtag
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 -> ConduitT ByteString ByteString m ()
forall (m :: * -> *).
MonadIO m =>
Int -> ConduitT ByteString ByteString m ()
go Int
0
where
go :: Int -> ConduitT ByteString ByteString m ()
go Int
seen =
ConduitT ByteString ByteString m (Maybe ByteString)
forall (m :: * -> *) i o. Monad m => ConduitT i o m (Maybe i)
await ConduitT ByteString ByteString m (Maybe ByteString)
-> (Maybe ByteString -> ConduitT ByteString ByteString m ())
-> ConduitT ByteString ByteString m ()
forall a b.
ConduitT ByteString ByteString m a
-> (a -> ConduitT ByteString ByteString m b)
-> ConduitT ByteString ByteString m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
Maybe ByteString
Nothing -> ConduitT ByteString ByteString m ()
forall (f :: * -> *). Applicative f => f ()
pass
Just ByteString
chunk -> do
let seen' :: Int
seen' = Int
seen Int -> Int -> Int
forall a. Num a => a -> a -> a
+ ByteString -> Int
BS.length ByteString
chunk
Bool
-> ConduitT ByteString ByteString m ()
-> ConduitT ByteString ByteString m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Int
seen' Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
maxBytes) (OsvDbCapExceeded -> ConduitT ByteString ByteString m ()
forall (m :: * -> *) e a. (MonadIO m, Exception e) => e -> m a
throwIO (Int -> OsvDbCapExceeded
OsvDbCapExceeded Int
maxBytes))
ByteString -> ConduitT ByteString ByteString m ()
forall (m :: * -> *) o i. Monad m => o -> ConduitT i o m ()
yield ByteString
chunk
Int -> ConduitT ByteString ByteString m ()
go Int
seen'