diff --git a/app/Main.hs b/app/Main.hs
--- a/app/Main.hs
+++ b/app/Main.hs
@@ -221,7 +221,6 @@
                 runHttp2Client frameConn _encoderBufsize _decoderBufsize conf (throwTo parentThread) ignoreFallbackHandler
 
     withConn $ \conn -> do
-      linkAsyncs conn
       _addCredit (_incomingFlowControl conn) _initialWindowKick
       _ <- forkIO $ forever $ do
               updated <- _updateWindow $ _incomingFlowControl conn
@@ -262,7 +261,7 @@
     tlsParams = TLS.ClientParams {
           TLS.clientWantSessionResume    = Nothing
         , TLS.clientUseMaxFragmentLength = Nothing
-        , TLS.clientServerIdentification = ("127.0.0.1", "")
+        , TLS.clientServerIdentification = (_host, ByteString.pack $ show _port)
         , TLS.clientUseServerNameIndication = True
         , TLS.clientShared               = def
         , TLS.clientHooks                = def { TLS.onServerCertificate = \_ _ _ _ -> return []
diff --git a/http2-client.cabal b/http2-client.cabal
--- a/http2-client.cabal
+++ b/http2-client.cabal
@@ -1,5 +1,5 @@
 name:                http2-client
-version:             0.6.0.0
+version:             0.7.0.0
 synopsis:            A native HTTP2 client library.
 description:         Please read the README.md at the homepage.
 homepage:            https://github.com/lucasdicioccio/http2-client
@@ -22,14 +22,15 @@
   other-modules:       Network.HTTP2.Client.Channels
                      , Network.HTTP2.Client.Dispatch
   build-depends:       base >= 4.7 && < 5
-                     , async
-                     , bytestring
-                     , connection
-                     , containers
-                     , http2
-                     , network
-                     , time
-                     , tls
+                     , async >= 2.1 && < 3
+                     , bytestring >= 0.10 && < 1
+                     , containers >= 0.5 && < 1
+                     , deepseq >= 1.4 && < 2
+                     , http2 >= 1.6 && < 2
+                     , network >= 2.6 && < 3
+                     , stm >= 2.4 && < 3
+                     , time >= 1.8 && < 2
+                     , tls >= 1.4 && < 2
   default-language:    Haskell2010
 
 executable http2-client-exe
diff --git a/src/Network/HTTP2/Client.hs b/src/Network/HTTP2/Client.hs
--- a/src/Network/HTTP2/Client.hs
+++ b/src/Network/HTTP2/Client.hs
@@ -46,7 +46,7 @@
 
 import           Control.Concurrent.Async (Async, async, race, withAsync, link)
 import           Control.Exception (bracket, throwIO, SomeException, catch)
-import           Control.Concurrent.MVar (newMVar, takeMVar, putMVar)
+import           Control.Concurrent.MVar (newEmptyMVar, newMVar, putMVar, takeMVar, tryPutMVar)
 import           Control.Concurrent (threadDelay)
 import           Control.Monad (forever, void, when, forM_)
 import           Data.ByteString (ByteString)
@@ -177,7 +177,8 @@
   , _asyncs         :: !Http2ClientAsyncs
   -- ^ Asynchronous operations threads.
   , _close           :: IO ()
-  -- ^ Closes the network connection.
+  -- ^ Immediately stop processing incoming frames and closes the network
+  -- connection.
   }
 
 data InitHttp2Client = InitHttp2Client {
@@ -189,6 +190,9 @@
   , _initOutgoingFlowControl :: OutgoingFlowControl
   , _initPaylodSplitter      :: IO PayloadSplitter
   , _initClose               :: IO ()
+  -- ^ Immediately closes the connection.
+  , _initStop                :: IO Bool
+  -- ^ Stops receiving frames.
   }
 
 -- | Set of Async threads running an Http2Client.
@@ -308,8 +312,9 @@
 --
 -- This function is slightly safer than 'startHttp2Client' because it uses
 -- 'Control.Concurrent.Async.withAsync' instead of
--- 'Control.Concurrent.Async.async' however, this with-pattern takes the
--- control of the thread and can be annoying at times.
+-- 'Control.Concurrent.Async.async'; plus this function calls 'linkAsyncs' to
+-- make sure that a network error kills the controlling thread. However, this
+-- with-pattern takes the control of the thread and can be annoying at times.
 runHttp2Client
   :: Http2FrameConnection
   -- ^ A frame connection.
@@ -332,17 +337,20 @@
     withAsync incomingLoop $ \aIncoming -> do
         settsIO <- _initSettings initClient initSettings
         withAsync settsIO $ \aSettings -> do
-            mainHandler $ Http2Client {
+            let client = Http2Client {
               _settings            = _initSettings initClient
             , _ping                = _initPing initClient
             , _goaway              = _initGoaway initClient
-            , _close               = _initClose initClient
+            , _close               =
+                _initStop initClient >> _initClose initClient
             , _startStream         = _initStartStream initClient
             , _incomingFlowControl = _initIncomingFlowControl initClient
             , _outgoingFlowControl = _initOutgoingFlowControl initClient
             , _payloadSplitter     = _initPaylodSplitter initClient
             , _asyncs              = Http2ClientAsyncs aSettings aIncoming
             }
+            linkAsyncs client
+            mainHandler client
 
 -- | Starts a new Http2Client around a frame connection.
 --
@@ -372,7 +380,8 @@
         _settings            = _initSettings initClient
       , _ping                = _initPing initClient
       , _goaway              = _initGoaway initClient
-      , _close               = _initClose initClient
+      , _close               =
+          _initStop initClient >> _initClose initClient
       , _startStream         = _initStartStream initClient
       , _incomingFlowControl = _initIncomingFlowControl initClient
       , _outgoingFlowControl = _initOutgoingFlowControl initClient
@@ -404,7 +413,7 @@
     (_initOutgoingFlowControl,windowUpdatesChan) <- newOutgoingFlowControl dispatchControl 0
 
     dispatchHPACK <- newDispatchHPACKIO decoderBufSize
-    let incomingLoop = dispatchLoop conn dispatch dispatchControl windowUpdatesChan _initIncomingFlowControl dispatchHPACK
+    (incomingLoop,endIncomingLoop) <- dispatchLoop conn dispatch dispatchControl windowUpdatesChan _initIncomingFlowControl dispatchHPACK
 
     {- Setup for client-initiated streams. -}
     conccurentStreams <- newIORef 0
@@ -454,7 +463,10 @@
 
     let _initPaylodSplitter = settingsPayloadSplitter <$> readSettings dispatchControl
 
+    let _initStop = endIncomingLoop
+
     let _initClose = closeConnection conn
+
     return (incomingLoop, InitHttp2Client{..})
 
 initializeStream
@@ -518,39 +530,43 @@
   -> Chan (FrameHeader, FramePayload)
   -> IncomingFlowControl
   -> DispatchHPACK
-  -> IO ()
+  -> IO (IO (), IO Bool)
 dispatchLoop conn d dc windowUpdatesChan inFlowControl dh = do
     let getNextFrame = next conn
-    delayException . forever $ do
-        frame <- getNextFrame
-        dispatchFramesStep frame d
-        whenFrame (hasStreamId 0) frame $ \got ->
-            dispatchControlFramesStep windowUpdatesChan got dc
-        whenFrame (hasTypeId [FrameData]) frame $ \got ->
-            creditDataFramesStep d inFlowControl got
-        whenFrame (hasTypeId [FrameWindowUpdate]) frame $ \got -> do
-            updateWindowsStep d got
-        whenFrame (hasTypeId [FramePushPromise, FrameHeaders]) frame $ \got -> do
-            let hpackLoop (FinishedWithHeaders curFh sId mkNewHdrs) = do
-                    newHdrs <- mkNewHdrs
-                    chan <- fmap _streamStateEvents <$> lookupStreamState d sId
-                    let msg = StreamHeadersEvent curFh newHdrs
-                    maybe (return ()) (flip writeChan msg) chan
-                hpackLoop (FinishedWithPushPromise curFh parentSid newSid mkNewHdrs) = do
-                    newHdrs <- mkNewHdrs
-                    chan <- fmap _streamStateEvents <$> lookupStreamState d parentSid
-                    let msg = StreamPushPromiseEvent curFh newSid newHdrs
-                    maybe (return ()) (flip writeChan msg) chan
-                hpackLoop (WaitContinuation act)        =
-                    getNextFrame >>= act >>= hpackLoop
-                hpackLoop (FailedHeaders curFh sId err)        = do
-                    chan <- fmap _streamStateEvents <$> lookupStreamState d sId
-                    let msg = StreamErrorEvent curFh err
-                    maybe (return ()) (flip writeChan msg) chan
-            hpackLoop (dispatchHPACKFramesStep got dh)
-        whenFrame (hasTypeId [FrameRSTStream]) frame $ \got -> do
-            handleRSTStep d got
-        finalizeFramesStep frame d
+    let go = delayException . forever $ do
+            frame <- getNextFrame
+            dispatchFramesStep frame d
+            whenFrame (hasStreamId 0) frame $ \got ->
+                dispatchControlFramesStep windowUpdatesChan got dc
+            whenFrame (hasTypeId [FrameData]) frame $ \got ->
+                creditDataFramesStep d inFlowControl got
+            whenFrame (hasTypeId [FrameWindowUpdate]) frame $ \got -> do
+                updateWindowsStep d got
+            whenFrame (hasTypeId [FramePushPromise, FrameHeaders]) frame $ \got -> do
+                let hpackLoop (FinishedWithHeaders curFh sId mkNewHdrs) = do
+                        newHdrs <- mkNewHdrs
+                        chan <- fmap _streamStateEvents <$> lookupStreamState d sId
+                        let msg = StreamHeadersEvent curFh newHdrs
+                        maybe (return ()) (flip writeChan msg) chan
+                    hpackLoop (FinishedWithPushPromise curFh parentSid newSid mkNewHdrs) = do
+                        newHdrs <- mkNewHdrs
+                        chan <- fmap _streamStateEvents <$> lookupStreamState d parentSid
+                        let msg = StreamPushPromiseEvent curFh newSid newHdrs
+                        maybe (return ()) (flip writeChan msg) chan
+                    hpackLoop (WaitContinuation act)        =
+                        getNextFrame >>= act >>= hpackLoop
+                    hpackLoop (FailedHeaders curFh sId err)        = do
+                        chan <- fmap _streamStateEvents <$> lookupStreamState d sId
+                        let msg = StreamErrorEvent curFh err
+                        maybe (return ()) (flip writeChan msg) chan
+                hpackLoop (dispatchHPACKFramesStep got dh)
+            whenFrame (hasTypeId [FrameRSTStream]) frame $ \got -> do
+                handleRSTStep d got
+            finalizeFramesStep frame d
+    end <- newEmptyMVar
+    let run = void $ race go (takeMVar end)
+    let stop = tryPutMVar end ()
+    return (run, stop)
 
 handleRSTStep
   :: Dispatch
diff --git a/src/Network/HTTP2/Client/Channels.hs b/src/Network/HTTP2/Client/Channels.hs
--- a/src/Network/HTTP2/Client/Channels.hs
+++ b/src/Network/HTTP2/Client/Channels.hs
@@ -1,89 +1,28 @@
 
 module Network.HTTP2.Client.Channels (
    FramesChan
- , HeadersChan
- , PushPromisesChan
- , waitHeaders
- , waitHeadersWithStreamId
- , waitFrame
- , waitFrameWithStreamId
- , waitPushPromiseWithParentStreamId
- , waitFrameWithTypeId
- , waitFrameWithTypeIdForStreamId
- , isPingReply
- , isSettingsReply
  , hasStreamId
  , hasTypeId
  , whenFrame
  , whenFrameElse
+ -- re-exports
  , module Control.Concurrent.Chan
  ) where
 
-import           Control.Concurrent.Chan (Chan, readChan, newChan, dupChan, writeChan)
+import           Control.Concurrent.Chan (Chan, readChan, newChan, writeChan)
 import           Control.Exception (Exception, throwIO)
-import           Data.ByteString (ByteString)
-import           Network.HPACK as HPACK
-import           Network.HTTP2 as HTTP2
+import           Network.HTTP2 (StreamId, FrameHeader, FramePayload, FrameTypeId, framePayloadToFrameTypeId, streamId)
 
 type FramesChan e = Chan (FrameHeader, Either e FramePayload)
 
-type HeadersChanContent = (FrameHeader, StreamId, Either ErrorCode HeaderList)
-
-type HeadersChan = Chan HeadersChanContent
-
-type PushPromisesChanContent = (StreamId, StreamId, HeaderList)
-
-type PushPromisesChan = Chan PushPromisesChanContent
-
-waitFrameWithStreamId
-  :: Exception e
-  => StreamId
-  -> FramesChan e
-  -> IO (FrameHeader, FramePayload)
-waitFrameWithStreamId sid = waitFrame (\h _ -> streamId h == sid)
-
-waitFrameWithTypeId
-  :: (Exception e)
-  => [FrameTypeId]
-  -> FramesChan e
-  -> IO (FrameHeader, FramePayload)
-waitFrameWithTypeId tids = waitFrame (\_ p -> HTTP2.framePayloadToFrameTypeId p `elem` tids)
-
-waitFrameWithTypeIdForStreamId
-  :: (Exception e)
-  => StreamId
-  -> [FrameTypeId]
-  -> FramesChan e
-  -> IO (FrameHeader, FramePayload)
-waitFrameWithTypeIdForStreamId sid tids =
-    waitFrame (\h p -> streamId h == sid && HTTP2.framePayloadToFrameTypeId p `elem` tids)
-
-waitFrame
-  :: Exception e
-  => (FrameHeader -> FramePayload -> Bool)
-  -> FramesChan e
-  -> IO (FrameHeader, FramePayload)
-waitFrame test chan =
-    loop
-  where
-    loop = do
-        (fHead, fPayload) <- readChan chan
-        dat <- either throwIO pure fPayload
-        if test fHead dat
-        then return (fHead, dat)
-        else loop
-
 whenFrame
   :: Exception e
   => (FrameHeader -> FramePayload -> Bool)
   -> (FrameHeader, Either e FramePayload)
   -> ((FrameHeader, FramePayload) -> IO ())
   -> IO ()
-whenFrame test (fHead, fPayload) handle = do
-    dat <- either throwIO pure fPayload
-    if test fHead dat
-    then handle (fHead, dat)
-    else pure ()
+whenFrame test frame handle = do
+    whenFrameElse test frame handle (const $ pure ())
 
 whenFrameElse
   :: Exception e
@@ -102,45 +41,4 @@
 hasStreamId sid h _ = streamId h == sid
 
 hasTypeId :: [FrameTypeId] -> FrameHeader -> FramePayload -> Bool
-hasTypeId tids _ p = HTTP2.framePayloadToFrameTypeId p `elem` tids
-
-isPingReply :: ByteString -> FrameHeader -> FramePayload -> Bool
-isPingReply datSent _ (PingFrame datRcv) = datSent == datRcv
-isPingReply _       _ _                  = False
-
-isSettingsReply :: FrameHeader -> FramePayload -> Bool
-isSettingsReply fh (SettingsFrame _) = HTTP2.testAck (flags fh)
-isSettingsReply _ _                  = False
-
-waitHeadersWithStreamId
-  :: StreamId
-  -> HeadersChan
-  -> IO HeadersChanContent
-waitHeadersWithStreamId sid =
-    waitHeaders (\_ s _ -> s == sid)
-
-waitHeaders
-  :: (FrameHeader -> StreamId -> Either ErrorCode HeaderList -> Bool)
-  -> HeadersChan
-  -> IO HeadersChanContent
-waitHeaders test chan =
-    loop
-  where
-    loop = do
-        tuple@(fH, sId, hdrs) <- readChan chan
-        if test fH sId hdrs
-        then return tuple
-        else loop
-
-waitPushPromiseWithParentStreamId
-  :: StreamId
-  -> PushPromisesChan
-  -> IO PushPromisesChanContent
-waitPushPromiseWithParentStreamId sid chan =
-    loop
-  where
-    loop = do
-        tuple@(parentSid,_,_) <- readChan chan
-        if parentSid == sid
-        then return tuple
-        else loop
+hasTypeId tids _ p = framePayloadToFrameTypeId p `elem` tids
diff --git a/src/Network/HTTP2/Client/FrameConnection.hs b/src/Network/HTTP2/Client/FrameConnection.hs
--- a/src/Network/HTTP2/Client/FrameConnection.hs
+++ b/src/Network/HTTP2/Client/FrameConnection.hs
@@ -15,9 +15,10 @@
     , closeConnection
     ) where
 
+import           Control.DeepSeq (deepseq)
 import           Control.Exception (bracket)
 import           Control.Concurrent.MVar (newMVar, takeMVar, putMVar)
-import           Control.Monad (void)
+import           Control.Monad ((>=>), void)
 import           Network.HTTP2 (FrameHeader(..), FrameFlags, FramePayload, HTTP2Error, encodeInfo, decodeFramePayload)
 import qualified Network.HTTP2 as HTTP2
 import           Network.Socket (HostName, PortNumber)
@@ -68,15 +69,11 @@
 next :: Http2FrameConnection -> IO (FrameHeader, Either HTTP2Error FramePayload)
 next = _nextHeaderAndFrame . _serverStream
 
--- | Creates a new 'Http2FrameConnection' to a given host for a frame-to-frame communication.
-newHttp2FrameConnection :: HostName
-                        -> PortNumber
-                        -> Maybe TLS.ClientParams
-                        -> IO Http2FrameConnection
-newHttp2FrameConnection host port params = do
-    -- Spawns an HTTP2 connection.
-    http2conn <- newRawHttp2Connection host port params
-
+-- | Adds framing around a 'RawHttp2Connection'.
+frameHttp2RawConnection
+  :: RawHttp2Connection
+  -> IO Http2FrameConnection
+frameHttp2RawConnection http2conn = do
     -- Prepare a local mutex, this mutex should never escape the
     -- function's scope. Else it might lead to bugs (e.g.,
     -- https://ro-che.info/articles/2014-07-30-bracket ) 
@@ -87,13 +84,15 @@
 
     -- Define handlers.
     let makeClientStream streamID = 
-            let putFrame modifyFF frame = do
+            let putFrame modifyFF frame =
                     let info = encodeInfo modifyFF streamID
-                    _sendRaw http2conn $
-                        HTTP2.encodeFrame info frame
+                    in HTTP2.encodeFrame info frame
                 putFrames f = writeProtect . void $ do
                     xs <- f
-                    traverse (uncurry putFrame) xs
+                    let ys = fmap (uncurry putFrame) xs
+                    -- Force evaluation of frames serialization whilst
+                    -- write-protected to avoid out-of-order errrors.
+                    deepseq ys (_sendRaw http2conn ys)
              in Http2FrameClientStream putFrames streamID
 
         nextServerFrameChunk = Http2ServerStream $ do
@@ -108,3 +107,11 @@
         gtfo = _close http2conn
 
     return $ Http2FrameConnection makeClientStream nextServerFrameChunk gtfo
+
+-- | Creates a new 'Http2FrameConnection' to a given host for a frame-to-frame communication.
+newHttp2FrameConnection :: HostName
+                        -> PortNumber
+                        -> Maybe TLS.ClientParams
+                        -> IO Http2FrameConnection
+newHttp2FrameConnection host port params = do
+    frameHttp2RawConnection =<< newRawHttp2Connection host port params
diff --git a/src/Network/HTTP2/Client/Helpers.hs b/src/Network/HTTP2/Client/Helpers.hs
--- a/src/Network/HTTP2/Client/Helpers.hs
+++ b/src/Network/HTTP2/Client/Helpers.hs
@@ -122,9 +122,12 @@
 waitStream stream streamFlowControl ppHandler = do
     ev <- _waitEvent stream
     case ev of
-        StreamHeadersEvent _ hdrs -> do
-            (dfrms,trls) <- waitDataFrames []
-            return (Right hdrs, reverse dfrms, trls)
+        StreamHeadersEvent fH hdrs
+            | HTTP2.testEndStream (HTTP2.flags fH) -> do
+                return (Right hdrs, [], Nothing)
+            | otherwise -> do
+                (dfrms,trls) <- waitDataFrames []
+                return (Right hdrs, reverse dfrms, trls)
         StreamPushPromiseEvent _ ppSid ppHdrs -> do
             _handlePushPromise stream ppSid ppHdrs ppHandler
             waitStream stream streamFlowControl ppHandler
diff --git a/src/Network/HTTP2/Client/RawConnection.hs b/src/Network/HTTP2/Client/RawConnection.hs
--- a/src/Network/HTTP2/Client/RawConnection.hs
+++ b/src/Network/HTTP2/Client/RawConnection.hs
@@ -7,16 +7,22 @@
     , newRawHttp2Connection
     ) where
 
+import           Control.Monad (forever, when)
+import           Control.Concurrent.Async (Async, async, cancel, pollSTM)
+import           Control.Concurrent.STM (STM, atomically, retry, throwSTM)
+import           Control.Concurrent.STM.TVar (TVar, modifyTVar', newTVarIO, readTVar, writeTVar)
 import           Data.ByteString (ByteString)
-import           Network.Connection (connectTo, initConnectionContext, ConnectionParams(..), TLSSettings(..), connectionPut, connectionGetExact, connectionClose)
+import qualified Data.ByteString as ByteString
+import           Data.ByteString.Lazy (fromChunks)
+import           Data.Monoid ((<>))
 import qualified Network.HTTP2 as HTTP2
-import           Network.Socket (HostName, PortNumber)
+import           Network.Socket hiding (recv)
+import           Network.Socket.ByteString
 import qualified Network.TLS as TLS
 
-
 -- TODO: catch connection errrors
 data RawHttp2Connection = RawHttp2Connection {
-    _sendRaw :: ByteString -> IO ()
+    _sendRaw :: [ByteString] -> IO ()
   -- ^ Function to send raw data to the server.
   , _nextRaw :: Int -> IO ByteString
   -- ^ Function to block reading a datachunk of a given size from the server.
@@ -35,22 +41,88 @@
                       -- overwritten to always return ["h2", "h2-17"].
                       -> IO RawHttp2Connection
 newRawHttp2Connection host port mparams = do
-    -- Connects to SSL.
-    ctx <- initConnectionContext
-    conn <- connectTo ctx connParams
+    -- Connects to TCP.
+    let hints = defaultHints { addrFlags = [AI_NUMERICSERV], addrSocketType = Stream }
+    addr:_ <- getAddrInfo (Just hints) (Just host) (Just $ show port)
+    skt <- socket (addrFamily addr) (addrSocketType addr) (addrProtocol addr)
+    connect skt (addrAddress addr)
 
-    -- Define raw byte-stream handlers.
-    let putRaw dat    = connectionPut conn dat
-    let getRaw amount = connectionGetExact conn amount
-    let doClose       = connectionClose conn
+    -- Prepare structure with abstract API.
+    conn <- maybe (plainTextRaw skt) (tlsRaw skt) mparams
 
     -- Initializes the HTTP2 stream.
-    putRaw HTTP2.connectionPreface
+    _sendRaw conn [HTTP2.connectionPreface]
 
-    return $ RawHttp2Connection putRaw getRaw doClose
+    return conn
+
+plainTextRaw :: Socket -> IO RawHttp2Connection
+plainTextRaw skt = do
+    (b,putRaw) <- startWriteWorker (sendMany skt)
+    (a,getRaw) <- startReadWorker (recv skt)
+    let doClose = cancel a >> cancel b >> close skt
+    return $ RawHttp2Connection (atomically . putRaw) (atomically . getRaw) doClose
+
+tlsRaw :: Socket -> TLS.ClientParams -> IO RawHttp2Connection
+tlsRaw skt params = do
+    -- Connects to SSL
+    tlsContext <- TLS.contextNew skt (modifyParams params)
+    TLS.handshake tlsContext
+
+    (b,putRaw) <- startWriteWorker (TLS.sendData tlsContext . fromChunks)
+    (a,getRaw) <- startReadWorker (const $ TLS.recvData tlsContext)
+    let doClose       = cancel a >> cancel b >> TLS.bye tlsContext >> TLS.contextClose tlsContext
+
+    return $ RawHttp2Connection (atomically . putRaw) (atomically . getRaw) doClose
   where
-    overwriteALPNHook params = (TLS.clientHooks params) {
-        TLS.onSuggestALPN = return $ Just [ "h2", "h2-17" ]
+    modifyParams prms = prms {
+        TLS.clientHooks = (TLS.clientHooks prms) {
+            TLS.onSuggestALPN = return $ Just [ "h2", "h2-17" ]
+          }
       }
-    modifyParams params = params { TLS.clientHooks = overwriteALPNHook params }
-    connParams = ConnectionParams host port (TLSSettings . modifyParams <$> mparams) Nothing
+
+startWriteWorker
+  :: ([ByteString] -> IO ())
+  -> IO (Async (), [ByteString] -> STM ())
+startWriteWorker sendChunks = do
+    outQ <- newTVarIO []
+    let putRaw chunks = modifyTVar' outQ (\xs -> xs ++ chunks)
+    b <- async $ writeWorkerLoop outQ sendChunks
+    return (b, putRaw)
+
+writeWorkerLoop :: TVar [ByteString] -> ([ByteString] -> IO ()) -> IO ()
+writeWorkerLoop outQ sendChunks = forever $ do
+    xs <- atomically $ do
+        chunks <- readTVar outQ
+        when (null chunks) retry
+        writeTVar outQ []
+        return chunks
+    sendChunks xs
+
+startReadWorker
+  :: (Int -> IO ByteString)
+  -> IO (Async (), (Int -> STM ByteString))
+startReadWorker get = do
+    buf <- newTVarIO ""
+    a <- async $ readWorkerLoop buf get
+    return $ (a, getRawWorker a buf)
+
+readWorkerLoop :: TVar ByteString -> (Int -> IO ByteString) -> IO ()
+readWorkerLoop buf next = forever $ do
+    dat <- next 4096
+    atomically $ modifyTVar' buf (\bs -> (bs <> dat))
+
+getRawWorker :: Async () -> TVar ByteString -> Int -> STM ByteString
+getRawWorker a buf amount = do
+    -- Verifies if the STM is alive, if dead, we re-throw the original
+    -- exception.
+    asyncStatus <- pollSTM a
+    case asyncStatus of
+        (Just (Left e)) -> throwSTM e
+        _               -> return ()
+    -- Read data consume, if there's enough, retry otherwise.
+    dat <- readTVar buf
+    if amount > ByteString.length dat
+    then retry
+    else do
+        writeTVar buf (ByteString.drop amount dat)
+        return $ ByteString.take amount dat
