module Ecluse.Core.Server.Stream (
RelayResponder (..),
UpstreamBody (..),
withUpstreamWhen,
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)
data RelayResponder response = RelayResponder
{ forall response.
RelayResponder response
-> Status -> ResponseHeaders -> StreamingBody -> IO response
relayStreamResponse :: Status -> ResponseHeaders -> StreamingBody -> IO response
, forall response.
RelayResponder response -> Status -> ResponseHeaders -> IO response
relayEmptyResponse :: Status -> ResponseHeaders -> IO response
}
data UpstreamBody
=
StreamBody
|
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)
withUpstreamWhen ::
Manager ->
ProgressFloor ->
Request ->
UpstreamBody ->
(Status -> Bool) ->
(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 =
((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)
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