jidmap <- newTVarIO M.empty
sendq <- newSendQueue
nextjid <- newTVarIO (JobId 1)
- void $ async $ sendloop st sendq
- void $ async $ receiveloop external st jidmap sendq
+ sender <- async $ sendloop st sendq
+ receiver <- async $ receiveloop external st jidmap sendq sender
return $ ExternalAsyncRelay $ do
receiveq <- newReceiveQueue
jid <- atomically $ do
(toAsyncWrapped msg, jid)
, externalReceive = atomically (readTBMChan receiveq)
-- This shuts down the whole relay.
- , externalShutdown = shutdown external st sendq
+ , externalShutdown = shutdown external st sendq sender receiver
-- These three TMVars are shared amoung all
-- ExternalStates that use this relay; they're
-- common state about the external process.
newSendQueue :: IO SendQueue
newSendQueue = newTBMChanIO 10
-receiveloop :: External -> ExternalState -> JidMap -> SendQueue -> IO ()
-receiveloop external st jidmap sendq = externalReceive st >>= \case
+receiveloop :: External -> ExternalState -> JidMap -> SendQueue -> Async () -> IO ()
+receiveloop external st jidmap sendq sendthread = externalReceive st >>= \case
Just l -> case parseMessage l :: Maybe AsyncMessage of
Just (AsyncMessage jid msg) ->
M.lookup jid <$> readTVarIO jidmap >>= \case
Just c -> do
atomically $ writeTBMChan c msg
- receiveloop external st jidmap sendq
+ receiveloop external st jidmap sendq sendthread
Nothing -> protoerr "unknown job number"
Nothing -> case parseMessage l :: Maybe ExceptionalMessage of
Just _ -> do
m <- readTVarIO jidmap
forM_ (M.elems m) $ \c ->
atomically $ writeTBMChan c l
- receiveloop external st jidmap sendq
+ receiveloop external st jidmap sendq sendthread
Nothing -> protoerr "unexpected non-async message"
Nothing -> closeandshutdown
where
closeandshutdown
closeandshutdown = do
- shutdown external st sendq True
+ dummy <- async noop
+ shutdown external st sendq sendthread dummy True
m <- atomically $ readTVar jidmap
forM_ (M.elems m) (atomically . closeTBMChan)
where
wrapjid msg jid = AsyncMessage jid $ unwords $ Proto.formatMessage msg
-shutdown :: External -> ExternalState -> SendQueue -> Bool -> IO ()
-shutdown external st sendq b = do
+shutdown :: External -> ExternalState -> SendQueue -> Async () -> Async () -> Bool -> IO ()
+shutdown external st sendq sendthread receivethread b = do
+ -- Receive thread is normally blocked reading from a handle.
+ -- That can block closing the handle, so it needs to be canceled.
+ cancel receivethread
+ -- Cleanly shutdown the send thread as well, allowing it to finish
+ -- writing anything that was buffered.
+ atomically $ closeTBMChan sendq
+ wait sendthread
r <- atomically $ do
r <- tryTakeTMVar (externalAsync external)
putTMVar (externalAsync external)
case r of
Just (ExternalAsync _) -> externalShutdown st b
_ -> noop
- atomically $ closeTBMChan sendq