module Ecluse.Core.Registry.Exchange (
boundedExchange,
singleAttemptSettings,
boundedFetch,
boundedJsonFetch,
withSuccessBody,
boundedRelay,
withinServeCap,
digestingRead,
chargedRead,
formThen,
) where
import Crypto.Hash (hashInit, hashUpdate)
import Data.ByteString qualified as BS
import Data.ByteString.Lazy qualified as LBS
import Data.JsonStream.Parser qualified as J
import Network.HTTP.Client (
BodyReader,
Manager,
ManagerSettings (managerRetryableException),
Request (requestBody),
Response (responseStatus),
brRead,
responseBody,
withResponse,
)
import Network.HTTP.Types.Status (statusCode)
import UnliftIO (try)
import UnliftIO.Timeout (timeout)
import Ecluse.Core.Fault (TransportCause (TransportTimeout), transportFault)
import Ecluse.Core.Fault.Http (classifyTransport)
import Ecluse.Core.Registry (
BodyOutcome (SuccessBody, UnreadStatus),
FetchFault (FetchBoundExceeded, FetchTransport),
PublishRelayResponse (..),
RegistryResponse (RegistryResponse),
UrlFormationError,
isAuthorisationFailure,
isSuccessStatus,
)
import Ecluse.Core.Registry.JsonStream (StreamResult, readJsonStream)
import Ecluse.Core.Registry.Progress (meteredReader, meteredUpload, watched)
import Ecluse.Core.Security (
BodyLimit,
LimitError,
ProgressFloor,
boundedRead,
floorMinBytes,
floorServeCapMicros,
floorWindowMicros,
)
import Ecluse.Core.Snapshot (ContentDigest, digestFromContext)
singleAttemptSettings :: ManagerSettings -> ManagerSettings
singleAttemptSettings :: ManagerSettings -> ManagerSettings
singleAttemptSettings ManagerSettings
settings = ManagerSettings
settings{managerRetryableException = const False}
boundedExchange :: (Int -> Int -> ByteString -> a) -> Manager -> ProgressFloor -> BodyLimit -> Request -> IO (Either FetchFault a)
boundedExchange :: forall a.
(Int -> Int -> ByteString -> a)
-> Manager
-> ProgressFloor
-> BodyLimit
-> Request
-> IO (Either FetchFault a)
boundedExchange Int -> Int -> ByteString -> a
project Manager
manager ProgressFloor
progress BodyLimit
limits Request
request =
Manager
-> ProgressFloor
-> Request
-> (Response BodyReader -> IO (Either LimitError a))
-> IO (Either FetchFault a)
forall a.
Manager
-> ProgressFloor
-> Request
-> (Response BodyReader -> IO (Either LimitError a))
-> IO (Either FetchFault a)
runExchange Manager
manager ProgressFloor
progress Request
request ((Int -> Int -> ByteString -> a)
-> BodyLimit -> Response BodyReader -> IO (Either LimitError a)
forall a.
(Int -> Int -> ByteString -> a)
-> BodyLimit -> Response BodyReader -> IO (Either LimitError a)
readBounded Int -> Int -> ByteString -> a
project BodyLimit
limits)
runExchange :: Manager -> ProgressFloor -> Request -> (Response BodyReader -> IO (Either LimitError a)) -> IO (Either FetchFault a)
runExchange :: forall a.
Manager
-> ProgressFloor
-> Request
-> (Response BodyReader -> IO (Either LimitError a))
-> IO (Either FetchFault a)
runExchange Manager
manager ProgressFloor
progress Request
request Response BodyReader -> IO (Either LimitError a)
readResponse =
ProgressFloor
-> (Watch -> IO (Either HttpException (Either LimitError a)))
-> IO (Maybe (Either HttpException (Either LimitError a)))
forall a. ProgressFloor -> (Watch -> IO a) -> IO (Maybe a)
watched ProgressFloor
progress (\Watch
watch -> IO (Either LimitError a)
-> IO (Either HttpException (Either LimitError a))
forall (m :: * -> *) e a.
(MonadUnliftIO m, Exception e) =>
m a -> m (Either e a)
try (Request
-> Manager
-> (Response BodyReader -> IO (Either LimitError a))
-> IO (Either LimitError a)
forall a.
Request -> Manager -> (Response BodyReader -> IO a) -> IO a
withResponse (Watch -> Request
metered Watch
watch) Manager
manager (Response BodyReader -> IO (Either LimitError a)
readResponse (Response BodyReader -> IO (Either LimitError a))
-> (Response BodyReader -> Response BodyReader)
-> Response BodyReader
-> IO (Either LimitError a)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (BodyReader -> BodyReader)
-> Response BodyReader -> Response BodyReader
forall a b. (a -> b) -> Response a -> Response b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap (Watch -> BodyReader -> BodyReader
meteredReader Watch
watch))))
IO (Maybe (Either HttpException (Either LimitError a)))
-> (Maybe (Either HttpException (Either LimitError a))
-> Either FetchFault a)
-> IO (Either FetchFault a)
forall (f :: * -> *) a b. Functor f => f a -> (a -> b) -> f b
<&> \case
Maybe (Either HttpException (Either LimitError a))
Nothing -> FetchFault -> Either FetchFault a
forall a b. a -> Either a b
Left (ProgressFloor -> FetchFault
belowFloor ProgressFloor
progress)
Just (Left HttpException
httpErr) -> FetchFault -> Either FetchFault a
forall a b. a -> Either a b
Left (TransportFault -> FetchFault
FetchTransport (HttpException -> TransportFault
classifyTransport HttpException
httpErr))
Just (Right (Left LimitError
limitErr)) -> FetchFault -> Either FetchFault a
forall a b. a -> Either a b
Left (LimitError -> FetchFault
FetchBoundExceeded LimitError
limitErr)
Just (Right (Right a
projected)) -> a -> Either FetchFault a
forall a b. b -> Either a b
Right a
projected
where
metered :: Watch -> Request
metered Watch
watch = Request
request{requestBody = meteredUpload watch (requestBody request)}
withinServeCap :: ProgressFloor -> (FetchFault -> e) -> IO (Either e a) -> IO (Either e a)
withinServeCap :: forall e a.
ProgressFloor
-> (FetchFault -> e) -> IO (Either e a) -> IO (Either e a)
withinServeCap ProgressFloor
progress FetchFault -> e
inject IO (Either e a)
action =
Either e a -> Maybe (Either e a) -> Either e a
forall a. a -> Maybe a -> a
fromMaybe (e -> Either e a
forall a b. a -> Either a b
Left (FetchFault -> e
inject FetchFault
capExceeded)) (Maybe (Either e a) -> Either e a)
-> IO (Maybe (Either e a)) -> IO (Either e a)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Int -> IO (Either e a) -> IO (Maybe (Either e a))
forall (m :: * -> *) a.
MonadUnliftIO m =>
Int -> m a -> m (Maybe a)
timeout (ProgressFloor -> Int
floorServeCapMicros ProgressFloor
progress) IO (Either e a)
action
where
capExceeded :: FetchFault
capExceeded = TransportFault -> FetchFault
FetchTransport (TransportCause -> Text -> TransportFault
transportFault TransportCause
TransportTimeout (Text
"the upstream exchange outlived its " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
seconds (ProgressFloor -> Int
floorServeCapMicros ProgressFloor
progress) Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"-second serve-path cap"))
belowFloor :: ProgressFloor -> FetchFault
belowFloor :: ProgressFloor -> FetchFault
belowFloor ProgressFloor
progress =
TransportFault -> FetchFault
FetchTransport (TransportFault -> FetchFault)
-> (Text -> TransportFault) -> Text -> FetchFault
forall b c a. (b -> c) -> (a -> b) -> a -> c
. TransportCause -> Text -> TransportFault
transportFault TransportCause
TransportTimeout (Text -> FetchFault) -> Text -> FetchFault
forall a b. (a -> b) -> a -> b
$
Text
"the exchange moved fewer than "
Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show (ProgressFloor -> Int
floorMinBytes ProgressFloor
progress)
Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" body bytes, counting both directions together, in a "
Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
seconds (ProgressFloor -> Int
floorWindowMicros ProgressFloor
progress)
Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"-second progress window"
seconds :: Int -> Text
seconds :: Int -> Text
seconds Int
micros = case Int
micros Int -> Int -> (Int, Int)
forall a. Integral a => a -> a -> (a, a)
`divMod` Int
1_000_000 of
(Int
whole, Int
0) -> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show Int
whole
(Int, Int)
_ -> Double -> Text
forall b a. (Show a, IsString b) => a -> b
show (Int -> Double
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
micros Double -> Double -> Double
forall a. Fractional a => a -> a -> a
/ Double
1_000_000 :: Double)
boundedFetch :: Manager -> ProgressFloor -> BodyLimit -> Request -> IO (Either FetchFault RegistryResponse)
boundedFetch :: Manager
-> ProgressFloor
-> BodyLimit
-> Request
-> IO (Either FetchFault RegistryResponse)
boundedFetch Manager
manager ProgressFloor
progress BodyLimit
limits Request
request = Manager
-> ProgressFloor
-> Request
-> (Response BodyReader -> IO (Either LimitError RegistryResponse))
-> IO (Either FetchFault RegistryResponse)
forall a.
Manager
-> ProgressFloor
-> Request
-> (Response BodyReader -> IO (Either LimitError a))
-> IO (Either FetchFault a)
runExchange Manager
manager ProgressFloor
progress Request
request ((Response BodyReader -> IO (Either LimitError RegistryResponse))
-> IO (Either FetchFault RegistryResponse))
-> (Response BodyReader -> IO (Either LimitError RegistryResponse))
-> IO (Either FetchFault RegistryResponse)
forall a b. (a -> b) -> a -> b
$ \Response BodyReader
response ->
let code :: Int
code = Status -> Int
statusCode (Response BodyReader -> Status
forall body. Response body -> Status
responseStatus Response BodyReader
response)
in if Int -> Bool
isAuthorisationFailure Int
code
then Either LimitError RegistryResponse
-> IO (Either LimitError RegistryResponse)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (RegistryResponse -> Either LimitError RegistryResponse
forall a b. b -> Either a b
Right (Int -> Int -> ByteString -> RegistryResponse
RegistryResponse Int
code Int
0 ByteString
""))
else (Int -> Int -> ByteString -> RegistryResponse)
-> BodyLimit
-> Response BodyReader
-> IO (Either LimitError RegistryResponse)
forall a.
(Int -> Int -> ByteString -> a)
-> BodyLimit -> Response BodyReader -> IO (Either LimitError a)
readBounded Int -> Int -> ByteString -> RegistryResponse
RegistryResponse BodyLimit
limits Response BodyReader
response
boundedRelay :: Manager -> ProgressFloor -> BodyLimit -> Request -> IO (Either FetchFault PublishRelayResponse)
boundedRelay :: Manager
-> ProgressFloor
-> BodyLimit
-> Request
-> IO (Either FetchFault PublishRelayResponse)
boundedRelay =
(Int -> Int -> ByteString -> PublishRelayResponse)
-> Manager
-> ProgressFloor
-> BodyLimit
-> Request
-> IO (Either FetchFault PublishRelayResponse)
forall a.
(Int -> Int -> ByteString -> a)
-> Manager
-> ProgressFloor
-> BodyLimit
-> Request
-> IO (Either FetchFault a)
boundedExchange ((Int -> Int -> ByteString -> PublishRelayResponse)
-> Manager
-> ProgressFloor
-> BodyLimit
-> Request
-> IO (Either FetchFault PublishRelayResponse))
-> (Int -> Int -> ByteString -> PublishRelayResponse)
-> Manager
-> ProgressFloor
-> BodyLimit
-> Request
-> IO (Either FetchFault PublishRelayResponse)
forall a b. (a -> b) -> a -> b
$ \Int
status Int
_ ByteString
body ->
PublishRelayResponse{relayStatus :: Int
relayStatus = Int
status, relayBody :: LByteString
relayBody = ByteString -> LByteString
LBS.fromStrict ByteString
body}
formThen ::
(UrlFormationError -> fault) ->
(Request -> IO (Either fault a)) ->
Either UrlFormationError Request ->
IO (Either fault a)
formThen :: forall fault a.
(UrlFormationError -> fault)
-> (Request -> IO (Either fault a))
-> Either UrlFormationError Request
-> IO (Either fault a)
formThen UrlFormationError -> fault
unformable = (UrlFormationError -> IO (Either fault a))
-> (Request -> IO (Either fault a))
-> Either UrlFormationError Request
-> IO (Either fault a)
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either (Either fault a -> IO (Either fault a)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Either fault a -> IO (Either fault a))
-> (UrlFormationError -> Either fault a)
-> UrlFormationError
-> IO (Either fault a)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. fault -> Either fault a
forall a b. a -> Either a b
Left (fault -> Either fault a)
-> (UrlFormationError -> fault)
-> UrlFormationError
-> Either fault a
forall b c a. (b -> c) -> (a -> b) -> a -> c
. UrlFormationError -> fault
unformable)
readBounded :: (Int -> Int -> ByteString -> a) -> BodyLimit -> Response BodyReader -> IO (Either LimitError a)
readBounded :: forall a.
(Int -> Int -> ByteString -> a)
-> BodyLimit -> Response BodyReader -> IO (Either LimitError a)
readBounded Int -> Int -> ByteString -> a
project BodyLimit
limits Response BodyReader
response =
((Int, ByteString) -> a)
-> Either LimitError (Int, ByteString) -> Either LimitError a
forall a b. (a -> b) -> Either LimitError a -> Either LimitError b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap ((Int -> ByteString -> a) -> (Int, ByteString) -> a
forall a b c. (a -> b -> c) -> (a, b) -> c
uncurry (Int -> Int -> ByteString -> a
project (Status -> Int
statusCode (Response BodyReader -> Status
forall body. Response body -> Status
responseStatus Response BodyReader
response))))
(Either LimitError (Int, ByteString) -> Either LimitError a)
-> IO (Either LimitError (Int, ByteString))
-> IO (Either LimitError a)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> BodyLimit -> BodyReader -> IO (Either LimitError (Int, ByteString))
forall (m :: * -> *).
Monad m =>
BodyLimit
-> m ByteString -> m (Either LimitError (Int, ByteString))
boundedRead BodyLimit
limits (BodyReader -> BodyReader
brRead (Response BodyReader -> BodyReader
forall body. Response body -> body
responseBody Response BodyReader
response))
boundedJsonFetch :: Manager -> ProgressFloor -> BodyLimit -> J.Parser a -> (s -> a -> Either LimitError s) -> s -> Request -> IO (Either FetchFault (BodyOutcome (StreamResult s)))
boundedJsonFetch :: forall a s.
Manager
-> ProgressFloor
-> BodyLimit
-> Parser a
-> (s -> a -> Either LimitError s)
-> s
-> Request
-> IO (Either FetchFault (BodyOutcome (StreamResult s)))
boundedJsonFetch Manager
manager ProgressFloor
progress BodyLimit
limits Parser a
parser s -> a -> Either LimitError s
step s
initial = Manager
-> ProgressFloor
-> (BodyReader -> IO (Either LimitError (StreamResult s)))
-> Request
-> IO (Either FetchFault (BodyOutcome (StreamResult s)))
forall a.
Manager
-> ProgressFloor
-> (BodyReader -> IO (Either LimitError a))
-> Request
-> IO (Either FetchFault (BodyOutcome a))
withSuccessBody Manager
manager ProgressFloor
progress (BodyLimit
-> Parser a
-> (s -> a -> Either LimitError s)
-> s
-> BodyReader
-> IO (Either LimitError (StreamResult s))
forall (m :: * -> *) a s.
Monad m =>
BodyLimit
-> Parser a
-> (s -> a -> Either LimitError s)
-> s
-> m ByteString
-> m (Either LimitError (StreamResult s))
readJsonStream BodyLimit
limits Parser a
parser s -> a -> Either LimitError s
step s
initial)
withSuccessBody :: Manager -> ProgressFloor -> (IO ByteString -> IO (Either LimitError a)) -> Request -> IO (Either FetchFault (BodyOutcome a))
withSuccessBody :: forall a.
Manager
-> ProgressFloor
-> (BodyReader -> IO (Either LimitError a))
-> Request
-> IO (Either FetchFault (BodyOutcome a))
withSuccessBody Manager
manager ProgressFloor
progress BodyReader -> IO (Either LimitError a)
consume Request
request = Manager
-> ProgressFloor
-> Request
-> (Response BodyReader -> IO (Either LimitError (BodyOutcome a)))
-> IO (Either FetchFault (BodyOutcome a))
forall a.
Manager
-> ProgressFloor
-> Request
-> (Response BodyReader -> IO (Either LimitError a))
-> IO (Either FetchFault a)
runExchange Manager
manager ProgressFloor
progress Request
request ((Response BodyReader -> IO (Either LimitError (BodyOutcome a)))
-> IO (Either FetchFault (BodyOutcome a)))
-> (Response BodyReader -> IO (Either LimitError (BodyOutcome a)))
-> IO (Either FetchFault (BodyOutcome a))
forall a b. (a -> b) -> a -> b
$ \Response BodyReader
response -> do
let code :: Int
code = Status -> Int
statusCode (Response BodyReader -> Status
forall body. Response body -> Status
responseStatus Response BodyReader
response)
if Int -> Bool
isSuccessStatus Int
code
then (a -> BodyOutcome a)
-> Either LimitError a -> Either LimitError (BodyOutcome a)
forall a b. (a -> b) -> Either LimitError a -> Either LimitError b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap (Int -> a -> BodyOutcome a
forall a. Int -> a -> BodyOutcome a
SuccessBody Int
code) (Either LimitError a -> Either LimitError (BodyOutcome a))
-> IO (Either LimitError a)
-> IO (Either LimitError (BodyOutcome a))
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> BodyReader -> IO (Either LimitError a)
consume (BodyReader -> BodyReader
brRead (Response BodyReader -> BodyReader
forall body. Response body -> body
responseBody Response BodyReader
response))
else Either LimitError (BodyOutcome a)
-> IO (Either LimitError (BodyOutcome a))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (BodyOutcome a -> Either LimitError (BodyOutcome a)
forall a b. b -> Either a b
Right (Int -> BodyOutcome a
forall a. Int -> BodyOutcome a
UnreadStatus Int
code))
digestingRead :: (IO ByteString -> IO (Either e a)) -> IO ByteString -> IO (Either e (a, ContentDigest))
digestingRead :: forall e a.
(BodyReader -> IO (Either e a))
-> BodyReader -> IO (Either e (a, ContentDigest))
digestingRead BodyReader -> IO (Either e a)
consume BodyReader
readChunk = do
context <- Context SHA256 -> IO (IORef (Context SHA256))
forall (m :: * -> *) a. MonadIO m => a -> m (IORef a)
newIORef Context SHA256
forall a. HashAlgorithm a => Context a
hashInit
let next = do
chunk <- BodyReader
readChunk
modifyIORef' context (`hashUpdate` chunk)
pure chunk
consume next >>= traverse (\a
result -> (a
result,) (ContentDigest -> (a, ContentDigest))
-> (Context SHA256 -> ContentDigest)
-> Context SHA256
-> (a, ContentDigest)
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Context SHA256 -> ContentDigest
digestFromContext (Context SHA256 -> (a, ContentDigest))
-> IO (Context SHA256) -> IO (a, ContentDigest)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> IORef (Context SHA256) -> IO (Context SHA256)
forall (m :: * -> *) a. MonadIO m => IORef a -> m a
readIORef IORef (Context SHA256)
context)
chargedRead :: (Int -> IO ()) -> IO ByteString -> IO ByteString
chargedRead :: (Int -> IO ()) -> BodyReader -> BodyReader
chargedRead Int -> IO ()
charge BodyReader
readChunk = do
chunk <- BodyReader
readChunk
charge (BS.length chunk)
pure chunk