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

{- | Bounded-memory artifact streaming: the constant-memory serve path.

The proxy serves an artifact by __streaming it through__ from upstream, never buffering it
whole, so a multi-hundred-megabyte tarball never becomes a local memory spike. The mirror
worker's whole-artifact fetch ('Ecluse.Core.Worker.Fetch.fetchArtifactBytes') is the
separate, buffered mirroring concern.

== Resource lifetime

A WAI streaming body __runs after the handler returns__, so an upstream connection released
lexically is already gone by the time the body streams. Raw WAI avoids that: the relay opens
the connection with @responseOpen@ before it commits the response, and closes it through
@finally@ around the whole relay. The open-to-@finally@ handoff is masked, so the request
timeout's kill or Warp's teardown on client disconnect cannot strand the connection between
@responseOpen@ returning and @finally@ arming @responseClose@.

== Backpressure

'pumpBody' writes each chunk through the sink's bounded output buffer before pulling the
next, and the write blocks once that buffer spills. The proxy therefore reads upstream only
as fast as the client drains, in __constant memory whatever the artifact's size__. Only the
first chunk is flushed, for a prompt first byte: at relay byte rates a per-chunk flush
degenerates into one socket send per upstream read (see @docs\/architecture\/web-layer.md@ →
"Streaming and resource lifetime").

== Progress

The pump reads under the "Ecluse.Core.Registry.Progress" watchdog, which counts only the time
spent waiting on upstream. An upstream that falls below the floor aborts the relay mid-stream,
and a client that drains slowly does not.
-}
module Ecluse.Core.Server.Stream (
    -- * A typed relay responder
    RelayResponder (..),

    -- * Relaying an upstream response through
    UpstreamBody (..),
    withUpstreamWhen,

    -- * The pump
    pumpBody,
) where

import Data.ByteString qualified as BS
import Data.ByteString.Builder (Builder, byteString)
import Network.HTTP.Client (BodyReader, Manager, Request, brRead, responseClose, responseHeaders, responseOpen, responseStatus)
import Network.HTTP.Client qualified as HTTP
import Network.HTTP.Types (ResponseHeaders, Status)
import Network.Wai (StreamingBody)
import UnliftIO.Exception (finally, mask, tryAny)

import Ecluse.Core.Registry.Progress (meteredReader, watchedRaising)
import Ecluse.Core.Security (ProgressFloor)
import Ecluse.Core.Server.Conditional (isNotModified)

{- | The two ways an upstream relay can answer, over the caller's route-scoped response value. WAI
construction stays out of this module, so no pipeline holds an unrestricted WAI responder.
-}
data RelayResponder response = RelayResponder
    { forall response.
RelayResponder response
-> Status -> ResponseHeaders -> StreamingBody -> IO response
relayStreamResponse :: Status -> ResponseHeaders -> StreamingBody -> IO response
    -- ^ Commit a status, headers, and bounded-memory streaming body.
    , forall response.
RelayResponder response -> Status -> ResponseHeaders -> IO response
relayEmptyResponse :: Status -> ResponseHeaders -> IO response
    -- ^ Commit the same response without a body (a @304@ or @HEAD@).
    }

-- | Whether a relay pumps the upstream body through, or answers bodiless.
data UpstreamBody
    = {- | Stream the body through. A @304@ still answers bodiless, because it carries no body
      (RFC 9110 15.4.5) and upstream's reader is never read.
      -}
      StreamBody
    | {- | Answer bodiless and never read upstream's body reader, so a @HEAD@ cannot make the
      proxy stream a whole artifact to nowhere.
      -}
      NoBody
    deriving stock (UpstreamBody -> UpstreamBody -> Bool
(UpstreamBody -> UpstreamBody -> Bool)
-> (UpstreamBody -> UpstreamBody -> Bool) -> Eq UpstreamBody
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: UpstreamBody -> UpstreamBody -> Bool
== :: UpstreamBody -> UpstreamBody -> Bool
$c/= :: UpstreamBody -> UpstreamBody -> Bool
/= :: UpstreamBody -> UpstreamBody -> Bool
Eq, Int -> UpstreamBody -> ShowS
[UpstreamBody] -> ShowS
UpstreamBody -> String
(Int -> UpstreamBody -> ShowS)
-> (UpstreamBody -> String)
-> ([UpstreamBody] -> ShowS)
-> Show UpstreamBody
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> UpstreamBody -> ShowS
showsPrec :: Int -> UpstreamBody -> ShowS
$cshow :: UpstreamBody -> String
show :: UpstreamBody -> String
$cshowList :: [UpstreamBody] -> ShowS
showList :: [UpstreamBody] -> ShowS
Show)

{- | Relay an upstream response when its status passes @accept@. 'Nothing' is a recoverable
miss that commits nothing, and a failure after the commit propagates rather than re-answering.
-}
withUpstreamWhen ::
    Manager ->
    ProgressFloor ->
    Request ->
    UpstreamBody ->
    -- | Whether upstream's status is a hit. A rejected status is a clean miss.
    (Status -> Bool) ->
    {- | Run once, pre-commit, on the accepted status and headers: the client-facing status and
    headers, plus the caller's own verdict on the relay.
    -}
    (Status -> ResponseHeaders -> IO (Status, ResponseHeaders, verdict)) ->
    RelayResponder response ->
    IO (Maybe (verdict, response))
withUpstreamWhen :: forall verdict response.
Manager
-> ProgressFloor
-> Request
-> UpstreamBody
-> (Status -> Bool)
-> (Status
    -> ResponseHeaders -> IO (Status, ResponseHeaders, verdict))
-> RelayResponder response
-> IO (Maybe (verdict, response))
withUpstreamWhen Manager
manager ProgressFloor
progress Request
request UpstreamBody
body Status -> Bool
accept Status -> ResponseHeaders -> IO (Status, ResponseHeaders, verdict)
relay RelayResponder response
respond =
    -- Masked from 'responseOpen' to 'finally' arming 'responseClose', so an async exception
    -- between the two cannot strand the connection. 'restore' keeps the relay interruptible.
    ((forall a. IO a -> IO a) -> IO (Maybe (verdict, response)))
-> IO (Maybe (verdict, response))
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> m a) -> m b) -> m b
mask (((forall a. IO a -> IO a) -> IO (Maybe (verdict, response)))
 -> IO (Maybe (verdict, response)))
-> ((forall a. IO a -> IO a) -> IO (Maybe (verdict, response)))
-> IO (Maybe (verdict, response))
forall a b. (a -> b) -> a -> b
$ \forall a. IO a -> IO a
restore ->
        IO (Response BodyReader)
-> IO (Either SomeException (Response BodyReader))
forall (m :: * -> *) a.
MonadUnliftIO m =>
m a -> m (Either SomeException a)
tryAny (IO (Response BodyReader) -> IO (Response BodyReader)
forall a. IO a -> IO a
restore (Request -> Manager -> IO (Response BodyReader)
responseOpen Request
request Manager
manager)) IO (Either SomeException (Response BodyReader))
-> (Either SomeException (Response BodyReader)
    -> IO (Maybe (verdict, response)))
-> IO (Maybe (verdict, response))
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 SomeException
_ -> Maybe (verdict, response) -> IO (Maybe (verdict, response))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe (verdict, response)
forall a. Maybe a
Nothing
            Right Response BodyReader
upstream -> IO (Maybe (verdict, response)) -> IO (Maybe (verdict, response))
forall a. IO a -> IO a
restore (Response BodyReader -> IO (Maybe (verdict, response))
answer Response BodyReader
upstream) IO (Maybe (verdict, response))
-> IO () -> IO (Maybe (verdict, response))
forall (m :: * -> *) a b. MonadUnliftIO m => m a -> m b -> m a
`finally` Response BodyReader -> IO ()
forall a. Response a -> IO ()
responseClose Response BodyReader
upstream
  where
    answer :: Response BodyReader -> IO (Maybe (verdict, response))
answer Response BodyReader
upstream
        | Bool -> Bool
not (Status -> Bool
accept Status
upstreamStatus) = Maybe (verdict, response) -> IO (Maybe (verdict, response))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe (verdict, response)
forall a. Maybe a
Nothing
        | Bool
otherwise = do
            (status, headers, verdict) <- Status -> ResponseHeaders -> IO (Status, ResponseHeaders, verdict)
relay Status
upstreamStatus (Response BodyReader -> ResponseHeaders
forall body. Response body -> ResponseHeaders
responseHeaders Response BodyReader
upstream)
            received <-
                if bodiless
                    then relayEmptyResponse respond status headers
                    else relayStreamResponse respond status headers pump
            pure (Just (verdict, received))
      where
        upstreamStatus :: Status
upstreamStatus = Response BodyReader -> Status
forall body. Response body -> Status
responseStatus Response BodyReader
upstream
        bodiless :: Bool
bodiless = UpstreamBody
body UpstreamBody -> UpstreamBody -> Bool
forall a. Eq a => a -> a -> Bool
== UpstreamBody
NoBody Bool -> Bool -> Bool
|| Status -> Bool
isNotModified Status
upstreamStatus
        pump :: StreamingBody
pump Builder -> IO ()
write IO ()
flush =
            ProgressFloor -> (Watch -> IO ()) -> IO ()
forall a. ProgressFloor -> (Watch -> IO a) -> IO a
watchedRaising ProgressFloor
progress (\Watch
watch -> BodyReader -> StreamingBody
pumpBody (Watch -> BodyReader -> BodyReader
meteredReader Watch
watch (BodyReader -> BodyReader
brRead (Response BodyReader -> BodyReader
forall body. Response body -> body
HTTP.responseBody Response BodyReader
upstream))) Builder -> IO ()
write IO ()
flush)

{- | Pump a chunked body from a reader to a WAI stream sink in constant memory. An empty chunk
is @http-client@'s 'BodyReader' end-of-body terminator, and the pump never writes it.
-}
pumpBody :: BodyReader -> (Builder -> IO ()) -> IO () -> IO ()
pumpBody :: BodyReader -> StreamingBody
pumpBody BodyReader
readChunk Builder -> IO ()
write IO ()
flush = do
    opening <- BodyReader
readChunk
    unless (BS.null opening) $ do
        write (byteString opening)
        flush
        rest
  where
    rest :: IO ()
    rest :: IO ()
rest = do
        chunk <- BodyReader
readChunk
        unless (BS.null chunk) $ do
            write (byteString chunk)
            rest