{-# LANGUAGE BangPatterns #-} {-# LANGUAGE CPP #-} {-# LANGUAGE MultiWayIf #-} {-# LANGUAGE NamedFieldPuns #-} {-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE ScopedTypeVariables #-} {-# LANGUAGE TupleSections #-} {-# OPTIONS_GHC -fno-warn-deprecations #-} module Network.Wai.Handler.Warp.Run where import Control.Arrow (first) import Control.Concurrent.STM ( TVar, atomically, check, modifyTVar', newTVarIO, readTVar, writeTVar, ) import qualified Control.Exception as E import qualified Data.ByteString as S import Data.Functor (($>)) import Data.IORef (IORef, newIORef, readIORef, writeIORef) import Data.Streaming.Network (bindPortTCP) import Foreign.C.Error ( Errno (..), eBADF, eCONNABORTED, eHOSTDOWN, eHOSTUNREACH, eMFILE, eNETDOWN, eNETUNREACH, eNONET, eNOPROTOOPT, ePROTO, ) import GHC.Conc.Sync (labelThread, myThreadId) import GHC.IO.Exception (IOErrorType (..), IOException (..)) import Network.Socket ( SockAddr, Socket, SocketOption (..), close, #if !WINDOWS fdSocket, #if MIN_VERSION_network(3,2,2) waitAndCancelReadSocketSTM, waitReadSocketSTM, #endif #endif getSocketName, setSocketOption, withSocketsDo, ) #if MIN_VERSION_network(3,1,1) import Network.Socket (gracefulClose) #endif import Network.Socket.BufferPool import qualified Network.Socket.ByteString as Sock import Network.Wai import System.Environment (lookupEnv) import System.IO.Error (ioeGetErrorType) import qualified System.TimeManager as T import System.Timeout (timeout) import Network.Wai.Handler.Warp.Buffer (createWriteBuffer) import Network.Wai.Handler.Warp.Counter import qualified Network.Wai.Handler.Warp.Date as D import qualified Network.Wai.Handler.Warp.FdCache as F import qualified Network.Wai.Handler.Warp.FileInfoCache as I import Network.Wai.Handler.Warp.HTTP1 (http1) import Network.Wai.Handler.Warp.HTTP2 (http2) import Network.Wai.Handler.Warp.HTTP2.Types (isHTTP2) import Network.Wai.Handler.Warp.Imports hiding (readInt) import Network.Wai.Handler.Warp.SendFile (sendFile) import Network.Wai.Handler.Warp.Settings import Network.Wai.Handler.Warp.ShuttingDown (readShuttingDown, writeShuttingDown) import Network.Wai.Handler.Warp.Types -- | Creating 'Connection' for plain HTTP based on a given socket. -- -- (N.B. make sure the 'Settings' have an initialized 'ServerState' to guarantee -- a graceful shutdown) socketConnection :: Settings -> Socket -> IO Connection socketConnection set s = do (ss, _) <- makeServerState set bufferPool <- newBufferPool 2048 16384 writeBuffer <- createWriteBuffer 16384 writeBufferRef <- newIORef writeBuffer isH2 <- newIORef False -- HTTP/1.x mysa <- getSocketName s appsInProgress <- newTVarIO 0 return Connection { connSendMany = Sock.sendMany s , connSendAll = sendall , connSendFile = sendfile writeBufferRef #if MIN_VERSION_network(3,1,1) , connClose = do h2 <- readIORef isH2 let tm = if h2 then settingsGracefulCloseTimeout2 set else settingsGracefulCloseTimeout1 set if tm <= 0 then close s else gracefulClose s tm `E.catch` throughAsync (return ()) #else , connClose = close s #endif , connRecv = receive' bufferPool ss appsInProgress , connRecvBuf = \_ _ -> return True -- obsoleted , connWriteBuffer = writeBufferRef , connHTTP2 = isH2 , connMySockAddr = mysa , connAppsInProgress = appsInProgress } where receive' bufferPool ss appsInProgress = E.handle handler $ makeGracefulRecv s bufferPool ss appsInProgress where handler :: E.IOException -> IO ByteString handler e | ioeGetErrorType e == InvalidArgument = return "" | otherwise = E.throwIO e sendfile writeBufferRef fid offset len hook headers = do writeBuffer <- readIORef writeBufferRef sendFile s writeBuffer sendall fid offset len hook headers sendall bs = E.handleJust ( \e -> if ioeGetErrorType e == ResourceVanished then Just ConnectionClosedByPeer else Nothing ) E.throwIO $ Sock.sendAll s bs -- | Create a 'Recv' using 'Network.Socket.BufferPool.Recv.receive', but make -- it non-blocking with 'waitReadSocketSTM' /AND/ cut off receiving any bytes -- when the server is shutting down and there are no more 'Application's -- actively using this 'Socket'. makeGracefulRecv :: Socket -> BufferPool -> ServerState -> TVar Int -> Recv makeGracefulRecv sock pool ss appsInProgress = do tryFastPath <- not <$> readShuttingDown (serverShuttingDown ss) if tryFastPath then do mbs <- receiveNoWait sock pool case mbs of Just bs -> return bs Nothing -> slowPath else slowPath where slowPath = makeGracefulRecvSlow sock pool ss appsInProgress makeGracefulRecvSlow :: Socket -> BufferPool -> ServerState -> TVar Int -> Recv makeGracefulRecvSlow sock pool ss appsInProgress = do sockWait <- #if !WINDOWS && MIN_VERSION_network(3,2,2) waitReadSocketSTM sock #else -- FIXME: 'waitReadSocketSTM' doesn't work on WINDOWS, and actually -- blocks indefinitely, so we fall back to going straight to 'recv'. pure (pure ()) #endif isShuttingDown <- atomically $ -- when shutting down (checkShutdown $> True) <|> -- else wait for socket readiness and do non-blocking read (sockWait $> False) if isShuttingDown then pure "" else recv where recv = receive sock pool checkShutdown = do check =<< currentShuttingDownStateSTM ss check . (<= 0) =<< readTVar appsInProgress -- | Run an 'Application' on the given port. -- This calls 'runSettings' with 'defaultSettings'. run :: Port -> Application -> IO () run p = runSettings defaultSettings{settingsPort = p} -- | Run an 'Application' on the port present in the @PORT@ -- environment variable. Uses the 'Port' given when the variable is unset. -- This calls 'runSettings' with 'defaultSettings'. -- -- @since 3.0.9 runEnv :: Port -> Application -> IO () runEnv p app = do mp <- lookupEnv "PORT" maybe (run p app) runReadPort mp where runReadPort :: String -> IO () runReadPort sp = case reads sp of ((p', _) : _) -> run p' app _ -> fail $ "Invalid value in $PORT: " ++ sp -- | Run an 'Application' with the given 'Settings'. -- This opens a listen socket on the port defined in 'Settings' and -- calls 'runSettingsSocket'. runSettings :: Settings -> Application -> IO () runSettings set app = withSocketsDo $ E.bracket (bindPortTCP (settingsPort set) (settingsHost set)) close ( \socket -> do setSocketCloseOnExec socket runSettingsSocket set socket app ) -- | What the accept loop waits on at the top of each turn, and what it -- leaves behind when it stops. -- -- A server used to be stopped by closing the listening socket under the -- thread waiting in @accept@ on it. That takes an IO manager that can wake -- a thread out of a wait by closing the file descriptor under it, which is -- what 'GHC.Conc.closeFdWith' is for, and the IO manager that provides it -- is not the only one there will be. So the loop waits for the socket and -- for the shutdown at once and ends on whichever comes first; the socket is -- closed after it, by which time nothing is waiting on it. data Listener = Listener { waitAcceptable :: IO Bool -- ^ Waits until there is something to accept. 'False' instead when the -- server has been asked to stop. , closeListener :: IO () -- ^ Run once the accept loop has ended and before the connections are -- waited for: the port is a successor's to take, and a successor should -- not have to wait out someone else's drain for it. } -- | For a caller that handed warp a connection maker and no socket. There -- is nothing here to wait on and nothing to close, and such a server ends -- the way it always did: by whatever its maker accepts on being closed. noListener :: Listener noListener = Listener{waitAcceptable = return True, closeListener = return ()} makeListener :: Settings -> Socket -> IO Listener #if !WINDOWS && MIN_VERSION_network(3,2,2) makeListener set socket = do stopping <- newTVarIO False settingsInstallShutdownHandler set $ atomically $ writeTVar stopping True return Listener { waitAcceptable = handleClosedListener $ do -- Cancelled however this ends, and not only when it ends -- because the socket became readable: an IO manager built on -- an interface like io_uring holds a reference on the socket -- for as long as the poll it was asked for is outstanding, -- so a wait left behind by a server that has stopped is a -- listening socket that does not close. (acceptable, cancelWait) <- waitAndCancelReadSocketSTM socket flip E.finally cancelWait $ atomically $ -- when shutting down ((check =<< readTVar stopping) $> False) <|> -- else wait for a connection to accept (acceptable $> True) , closeListener = close socket } -- Closing the listening socket is how a server was stopped before there was -- anything to tell, and callers that do it are still out there. Waiting on -- a descriptor that has been closed is an EBADF, under whichever name the IO -- manager gives it, and it means the same thing the shutdown handler means: -- stop accepting. Ending the loop quietly is what happened before, and -- 'closeListener' closing an already-closed socket is no error. -- -- 'acceptNewConnection' says the same of the EBADF from accept() itself. handleClosedListener :: IO Bool -> IO Bool handleClosedListener = E.handle $ \e -> if ioeGetErrorType e == InvalidArgument then return False else E.throwIO e #else makeListener set socket = do -- Two ways to get here. As in 'makeGracefulRecvSlow', -- 'waitReadSocketSTM' doesn't work on WINDOWS and blocks indefinitely; -- and before network 3.2.2 there is no 'waitAndCancelReadSocketSTM' to -- call at all. Either way a shutdown ends the accept loop the old way, -- by closing the listening socket under it. settingsInstallShutdownHandler set $ close socket return noListener #endif -- | This installs a shutdown handler for the given socket and runs the -- default connection setup action, which handles plain (non-cipher) HTTP. -- Running the handler stops the server accepting, closes the listening -- socket, and gracefully shuts the live connections down. -- -- The supplied socket can be a Unix named socket, which -- can be used when reverse HTTP proxying into your application. -- -- Note that the 'settingsPort' will still be passed to 'Application's via the -- 'serverPort' record. runSettingsSocket :: Settings -> Socket -> Application -> IO () runSettingsSocket oldSettings@Settings{settingsAccept = accept'} socket app = do listener <- makeListener oldSettings socket (_, newSettings) <- makeServerState oldSettings runSettingsConnectionMakerSecureWith newSettings listener (first ((,TCP) <$>) <$> getConnMaker newSettings) app where getConnMaker set = do (conn, sa) <- getConn set return (return conn, sa) getConn set = do (s, sa) <- accept' socket setSocketCloseOnExec s -- NoDelay causes an error for AF_UNIX. setSocketOption s NoDelay 1 `E.catch` throughAsync (return ()) conn <- socketConnection set s return (conn, sa) -- | The connection setup action would be expensive. A good example -- is initialization of TLS. -- So, this converts the connection setup action to the connection maker -- which will be executed after forking a new worker thread. -- Then this calls 'runSettingsConnectionMaker' with the connection maker. -- This allows the expensive computations to be performed -- in a separate worker thread instead of the main server loop. -- -- @since 1.3.5 runSettingsConnection :: Settings -> IO (Connection, SockAddr) -> Application -> IO () runSettingsConnection set getConn app = runSettingsConnectionMaker set getConnMaker app where getConnMaker = do (conn, sa) <- getConn return (return conn, sa) -- | This modifies the connection maker so that it returns 'TCP' for 'Transport' -- (i.e. plain HTTP) then calls 'runSettingsConnectionMakerSecure'. runSettingsConnectionMaker :: Settings -> IO (IO Connection, SockAddr) -> Application -> IO () runSettingsConnectionMaker x y = runSettingsConnectionMakerSecure x (toTCP <$> y) where toTCP = first ((,TCP) <$>) ---------------------------------------------------------------- -- | The core run function which takes 'Settings', -- a connection maker and 'Application'. -- The connection maker can return a connection of either plain HTTP -- or HTTP over TLS. -- -- @since 2.1.4 runSettingsConnectionMakerSecure :: Settings -> IO (IO (Connection, Transport), SockAddr) -> Application -> IO () runSettingsConnectionMakerSecure set = runSettingsConnectionMakerSecureWith set noListener {-# DEPRECATED runSettingsConnectionMakerSecure "use runSettingsConnectionMakerSecureWith instead" #-} -- | 'runSettingsConnectionMakerSecure' for a caller that has the listening -- socket, and so can say how the accept loop waits on it and what becomes -- of it when the loop stops. runSettingsConnectionMakerSecureWith :: Settings -> Listener -> IO (IO (Connection, Transport), SockAddr) -> Application -> IO () runSettingsConnectionMakerSecureWith oldSettings listener getConnMaker app = do settingsBeforeMainLoop oldSettings (ServerState{serverConnectionCounter}, newSettings) <- makeServerState oldSettings withII newSettings $ \ii -> initFdExhaustionRef >>= acceptConnection newSettings listener getConnMaker app serverConnectionCounter ii -- | Running an action with internal info. -- -- @since 3.3.11 withII :: Settings -> (InternalInfo -> IO a) -> IO a withII set action = withTimeoutManager $ \tm -> D.withDateCache $ \dc -> F.withFdCache fdCacheDurationInMicroseconds $ \fdc -> I.withFileInfoCache fdFileInfoDurationInMicroseconds $ \fic -> do let ii = InternalInfo tm dc fdc fic action ii where !fdCacheDurationInMicroseconds = settingsFdCacheDuration set * 1000000 !fdFileInfoDurationInMicroseconds = settingsFileInfoCacheDuration set * 1000000 !timeoutInMicroseconds = settingsTimeout set * 1000000 withTimeoutManager f = case settingsManager set of Just tm -> f tm Nothing -> E.bracket (T.initialize timeoutInMicroseconds) T.stopManager f -- Note that there is a thorough discussion of the exception safety of the -- following code at: https://github.com/yesodweb/wai/issues/146 -- -- We need to make sure of two things: -- -- 1. Asynchronous exceptions are not blocked entirely in the main loop. -- Doing so would make it impossible to kill the Warp thread. -- -- 2. Once a connection maker is received via acceptNewConnection, the -- connection is guaranteed to be closed, even in the presence of -- async exceptions. -- -- Our approach is explained in the comments below. acceptConnection :: Settings -> Listener -> IO (IO (Connection, Transport), SockAddr) -> Application -> Counter -> InternalInfo -> IORef FdExhaustion -- ^ This ref will be used to "debounce" the call to 'settingsOnException' -- when we hit an 'IOError' with 'eMFILE' in the case that Warp is not -- the reason the file descriptors are exhausted. -> IO () acceptConnection set listener getConnMaker app counter ii fdRef = do -- First mask all exceptions in acceptLoop. This is necessary to -- ensure that no async exception is throw between the call to -- acceptNewConnection and the registering of connClose. -- -- acceptLoop ends when the listener says the server has been asked to -- stop, or, for a caller that gave us no socket to wait on, when the -- socket it does accept on is closed. void $ E.mask_ acceptLoop -- Nothing is waiting on the listening socket any more, so it can be -- closed and the port left to a successor rather than held for as long -- as the connections below take to finish. closeListener listener -- In some cases, we want to stop Warp here without graceful shutdown. -- So, async exceptions are allowed here. -- That's why `finally` is not used. gracefulShutdown set counter where acceptLoop = do -- Allow async exceptions before receiving the next connection maker. E.allowInterrupt -- Wait for the socket to have something to accept and for the -- server to be asked to stop, whichever comes first. acceptable <- waitAcceptable listener when acceptable $ do -- acceptNewConnection will try to receive the next incoming -- request. It returns a /connection maker/, not a connection, -- since in some circumstances creating a working connection -- from a raw socket may be an expensive operation, and this -- expensive work should not be performed in the main event -- loop. An example of something expensive would be TLS -- negotiation. mx <- acceptNewConnection case mx of Nothing -> return () Just (mkConn, addr) -> do fork set mkConn addr app counter ii acceptLoop acceptNewConnection = do ex <- E.try getConnMaker case ex of Right x -> do -- Important to mark the exhaustion issue to be resolved -- when we get connections again. resetFdExhaustion fdRef return $ Just x Left e -> do let getErrno (Errno cInt) = cInt isErrno err = ioe_errno e == Just (getErrno err) -- Errors about one queued connection rather than about the -- listening socket. eCONNABORTED is the familiar one, a -- peer that went away before it could be accepted. Linux -- also reports the new socket's already-pending network -- errors through accept(), and accept(2) asks for those to -- be treated the same way, retried like EAGAIN. Either -- way the connection has left the queue, so the retry -- blocks for a new one rather than spinning on the same -- failure. -- -- eOPNOTSUPP is on that list in accept(2) and is left off -- this one on purpose: it also means the listening socket -- is not SOCK_STREAM, which is a permanent condition that -- retrying would spin on forever. Throwing tells whoever -- passed that socket, which is the only thing that helps. isQueuedConnectionError = any isErrno [ eCONNABORTED , eNETDOWN , ePROTO , eNOPROTOOPT , eHOSTDOWN , eNONET , eHOSTUNREACH , eNETUNREACH ] isFdExhaustion = isErrno eMFILE isIntentionallyClosedSocket = isErrno eBADF if | isQueuedConnectionError -> do -- Important to mark the exhaustion issue to be resolved resetFdExhaustion fdRef acceptNewConnection -- Keep in mind to reset the ref when anything other -- than this branch runs | isFdExhaustion -> do handleFdExhaustion e acceptNewConnection -- A graceful shutdown ends this loop by closing the -- listening socket, and that always arrives here as -- EBADF: 'close' replaces the descriptor with -1 before -- closing it, so every later accept() is handed -1 and -- fails that way. Matching EBADF alone therefore cannot -- miss a deliberate shutdown. | isIntentionallyClosedSocket -> do resetFdExhaustion fdRef settingsOnException set Nothing $ E.toException e return Nothing -- Everything left is something the listening socket will -- keep giving: descriptors exhausted system-wide, no -- memory for a socket, a socket that cannot accept. -- Ending the loop for those would return the same () a -- graceful shutdown returns, so a server that died and a -- server that was asked to stop would be reported -- identically and the caller could not tell which had -- happened. Throw instead, so it can. | otherwise -> do -- Maybe not important to mark the exhaustion issue -- as resolved here, but just for completeness' sake. resetFdExhaustion fdRef settingsOnException set Nothing $ E.toException e #if WINDOWS -- None of the guards above can match on Windows, where -- network reports a socket error with no errno on it, -- so every accept() failure arrives here including the -- EBADF of a deliberate shutdown. Throwing would turn -- an ordinary shutdown into an exception, so Windows -- keeps ending the loop the way it always has. return Nothing #else E.throwIO e #endif handleFdExhaustion e = do fdExhaustion <- readIORef fdRef -- If file descriptors are exhausted while Warp has -- no current connections, 'settingsOnException' would -- get called an enormous amount of times per second. when (fdExhaustion /= FdExhausted) $ settingsOnException set Nothing $ E.toException e hasDecreased <- waitForDecreased counter -- If we get 'NoConnections', that means the file -- descriptor exhaustion is outside of our control. -- We flag it so that 'settingsOnException' doesn't get -- called until the exhaustion issue is resolved. when (hasDecreased == NoConnections) $ setFdExhaustion fdRef -- Fork a new worker thread for this connection maker, and ask for a -- function to unmask (i.e., allow async exceptions to be thrown). fork :: Settings -> IO (Connection, Transport) -> SockAddr -> Application -> Counter -> InternalInfo -> IO () fork set mkConn addr app counter ii = do -- Count the connection here rather than in the thread below. The -- accept loop does not wait for that thread to be scheduled, so -- counting there leaves a window in which the connection is accepted -- and not counted, and 'gracefulShutdown' waits on this counter. increase counter settingsFork set $ \unmask -> runConnection unmask `E.finally` decrease counter where runConnection unmask = do tid <- myThreadId labelThread tid "Warp just forked" -- Call the user-supplied on exception code if any -- exceptions are thrown. -- -- Intentionally using Control.Exception.handle, since we want to -- catch all exceptions and avoid them from propagating, even -- async exceptions. See: -- https://github.com/yesodweb/wai/issues/850 E.handle (onConnectionException set addr) $ -- Run the connection maker to get a new connection, and ensure -- that the connection is closed. If the mkConn call throws an -- exception, we will leak the connection. If the mkConn call is -- vulnerable to attacks (e.g., Slowloris), we do nothing to -- protect the server. It is therefore vital that mkConn is well -- vetted. -- -- We grab the connection before registering timeouts since the -- timeouts will be useless during connection creation, due to the -- fact that async exceptions are still masked. E.bracket mkConn cleanUp (serve unmask) cleanUp (conn, _) = connClose conn `E.finally` do writeBuffer <- readIORef $ connWriteBuffer conn bufFree writeBuffer -- We need to register a timeout handler for this thread, and -- cancel that handler as soon as we exit. serve unmask (conn, transport) = T.withHandleKillThread (timeoutManager ii) (return ()) $ \th -> do -- We now have fully registered a connection close handler in -- the case of all exceptions, so it is safe to once again -- allow async exceptions. unmask . -- Call the user-supplied code for connection open and -- close events E.bracket (onOpen addr) (onClose addr) $ \goingon -> -- Actually serve this connection. bracket with closeConn -- above ensures the connection is closed. when goingon $ serveConnection conn ii th addr transport set app onOpen adr = settingsOnOpen set adr onClose adr _ = settingsOnClose set adr serveConnection :: Connection -> InternalInfo -> T.Handle -> SockAddr -> Transport -> Settings -> Application -> IO () serveConnection conn ii th origAddr transport settings app = do -- fixme: Upgrading to HTTP/2 should be supported. tid <- myThreadId (h2, bs) <- if isHTTP2 transport then return (True, "") else do bs0 <- recv4 "" if "PRI " `S.isPrefixOf` bs0 then return (True, bs0) else return (False, bs0) let appsInProgress = connAppsInProgress conn app' req rsp = E.bracket_ (atomically $ modifyTVar' appsInProgress $ (+ 1)) (atomically $ modifyTVar' appsInProgress $ \i -> (i - 1)) $ app req rsp if settingsHTTP2Enabled settings && h2 then do labelThread tid ("Warp HTTP/2 " ++ show origAddr) http2 settings ii conn transport app' origAddr th bs else do labelThread tid ("Warp HTTP/1.1 " ++ show origAddr) http1 settings ii conn transport app' origAddr th bs where recv4 bs0 = do bs1 <- connRecv conn if S.null bs1 then return bs0 else do -- In the case where bs0 is "", (<>) is called unnecessarily. -- But we adopt this logic for simplicity. let bs2 = bs0 <> bs1 if S.length bs2 >= 4 then return bs2 else recv4 bs2 -- | Set flag FileCloseOnExec flag on a socket (on Unix) -- -- Copied from: https://github.com/mzero/plush/blob/master/src/Plush/Server/Warp.hs -- -- @since 3.2.17 setSocketCloseOnExec :: Socket -> IO () #if WINDOWS setSocketCloseOnExec _ = return () #else setSocketCloseOnExec socket = do #if MIN_VERSION_network(3,0,0) fd <- fdSocket socket #else let fd = fdSocket socket #endif F.setFileCloseOnExec $ fromIntegral fd #endif gracefulShutdown :: Settings -> Counter -> IO () gracefulShutdown set counter = do setShuttingDown case settingsGracefulShutdownTimeout set of Nothing -> waitForZero counter (Just seconds) -> void (timeout (seconds * microsPerSecond) (waitForZero counter)) where microsPerSecond = 1000000 setShuttingDown = case settingsServerState set of Nothing -> pure () Just ServerState{serverShuttingDown} -> writeShuttingDown serverShuttingDown True data FdExhaustion = NoFdIssue | FdExhausted deriving (Eq, Show) initFdExhaustionRef :: IO (IORef FdExhaustion) initFdExhaustionRef = newIORef NoFdIssue -- [FD_EXHAUSTION] -- No need for "atomic" variants, since this is only used in a tight loop in -- 'acceptConnection'. resetFdExhaustion :: IORef FdExhaustion -> IO () resetFdExhaustion = flip writeIORef NoFdIssue setFdExhaustion :: IORef FdExhaustion -> IO () setFdExhaustion = flip writeIORef FdExhausted -- [FD_EXHAUSTION]