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

{- | Bounded registry exchanges and their transport-fault classification.
Read exchanges retain explicit access refusals before reading an error body. Every exchange runs
under the "Ecluse.Core.Registry.Progress" watchdog: one that moves fewer than the floor's body
bytes, up or down, in a window of waiting fails with the transport timeout, and its connection
closes rather than returning to the pool. The serve path also wraps its exchanges in
'withinServeCap', which ends each one before the request timeout.
-}
module Ecluse.Core.Registry.Exchange (
    -- * The bounded exchange
    boundedExchange,
    singleAttemptSettings,
    boundedFetch,
    boundedJsonFetch,
    withSuccessBody,
    boundedRelay,

    -- * The serve-path cap
    withinServeCap,

    -- * Source digests
    digestingRead,

    -- * Metered reads
    chargedRead,

    -- * Request formation
    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)

-- | Destructive clients must return uncertain transport failures for reassessment before retry.
singleAttemptSettings :: ManagerSettings -> ManagerSettings
singleAttemptSettings :: ManagerSettings -> ManagerSettings
singleAttemptSettings ManagerSettings
settings = ManagerSettings
settings{managerRetryableException = const False}

-- | Project status, decompressed byte count, and body. Transport failures retain their typed cause.
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)}

{- | Fail an action that outlives the floor's serve-path cap with the transport timeout. The serve
path wraps its exchanges in it, and the mirror worker and the Dredger do not.
-}
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"

-- Whole seconds as an operator wrote them, and a fraction only where one exists.
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)

-- | Preserve explicit auth refusals without reading their untrusted error bodies.
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

-- | The exchange keeping the answered status alongside the body, for the first-party relay.
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}

-- | Report request-formation and exchange failures through the same error channel.
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))

-- | Extract selected values from a 2xx body within the response lifetime. Other statuses are never parsed.
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)

-- | Consume a 2xx body within the response lifetime. A status outside 2xx never reaches the consumer.
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))

{- | Run a consumer over a source that hashes each chunk it passes on. A successful result carries
the digest of every chunk the consumer read.
-}
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)

-- | Pay for each chunk's length before the consumer sees it.
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