nqe 0.3.0.0 → 0.4.0
raw patch · 7 files changed
+235/−48 lines, 7 filesdep +mtlPVP ok
version bump matches the API change (PVP)
Dependencies added: mtl
API changes (from Hackage documentation)
+ Control.Concurrent.NQE: boundedPublisher :: (MonadIO m, Mailbox mbox, Mailbox events) => Publisher mbox TBQueue msg -> events msg -> m ()
+ Control.Concurrent.NQE: publisher :: (MonadIO m, Mailbox mbox, Mailbox events, Mailbox ch, Eq (ch msg)) => Publisher mbox ch msg -> events msg -> m ()
+ Control.Concurrent.NQE: type Publisher mbox ch msg = mbox (ControlMsg ch msg)
+ Control.Concurrent.NQE: withBoundedPubSub :: (MonadUnliftIO m, Mailbox mbox) => Int -> Publisher mbox TBQueue msg -> (TBQueue msg -> m a) -> m a
+ Control.Concurrent.NQE: withPubSub :: (MonadUnliftIO m, Mailbox mbox) => Publisher mbox TQueue msg -> (TQueue msg -> m a) -> m a
- Control.Concurrent.NQE: AddChild :: (IO ()) -> (Reply (Async ())) -> SupervisorMessage
+ Control.Concurrent.NQE: AddChild :: (n ()) -> (Reply (Async ())) -> SupervisorMessage n
- Control.Concurrent.NQE: RemoveChild :: (Async ()) -> SupervisorMessage
+ Control.Concurrent.NQE: RemoveChild :: (Async ()) -> SupervisorMessage n
- Control.Concurrent.NQE: StopSupervisor :: SupervisorMessage
+ Control.Concurrent.NQE: StopSupervisor :: SupervisorMessage n
- Control.Concurrent.NQE: addChild :: (MonadIO m, Mailbox mbox) => mbox SupervisorMessage -> IO () -> m (Async ())
+ Control.Concurrent.NQE: addChild :: (MonadUnliftIO n, MonadIO m, Mailbox mbox) => mbox (SupervisorMessage n) -> n () -> m (Async ())
- Control.Concurrent.NQE: data SupervisorMessage
+ Control.Concurrent.NQE: data SupervisorMessage n
- Control.Concurrent.NQE: dispatch :: (MonadIO m, Mailbox mbox) => [(msg -> Maybe a, a -> IO b)] -> mbox msg -> m b
+ Control.Concurrent.NQE: dispatch :: (MonadIO m, Mailbox mbox) => [(msg -> Maybe a, a -> m b)] -> mbox msg -> m b
- Control.Concurrent.NQE: removeChild :: (MonadIO m, Mailbox mbox) => mbox SupervisorMessage -> Async () -> m ()
+ Control.Concurrent.NQE: removeChild :: (MonadUnliftIO n, MonadIO m, Mailbox mbox) => mbox (SupervisorMessage n) -> Async () -> m ()
- Control.Concurrent.NQE: stopSupervisor :: (MonadIO m, Mailbox mbox) => mbox SupervisorMessage -> m ()
+ Control.Concurrent.NQE: stopSupervisor :: (MonadUnliftIO n, MonadIO m, Mailbox mbox) => mbox (SupervisorMessage n) -> m ()
- Control.Concurrent.NQE: supervisor :: (MonadUnliftIO m, Mailbox mbox) => Strategy -> mbox SupervisorMessage -> [m ()] -> m ()
+ Control.Concurrent.NQE: supervisor :: (MonadUnliftIO m, Mailbox mbox) => Strategy -> mbox (SupervisorMessage m) -> [m ()] -> m ()
Files
- CHANGELOG.md +16/−0
- nqe.cabal +12/−4
- src/Control/Concurrent/NQE.hs +2/−0
- src/Control/Concurrent/NQE/Process.hs +6/−6
- src/Control/Concurrent/NQE/PubSub.hs +131/−0
- src/Control/Concurrent/NQE/Supervisor.hs +36/−33
- test/Spec.hs +32/−5
+ CHANGELOG.md view
@@ -0,0 +1,16 @@+# Changelog+All notable changes to this project will be documented in this file.++The format is based on [Keep a Changelog](http://keepachangelog.com/en/1.0.0/)+and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0.html).++## [Unreleased]+### Added+- Changelog and semantic versions.+- Raw TCP actors.+- Move to `package.yaml` and `hpack`.+- Type-safe asynchronous messages.+- Supervisors for `MonadUnliftIO` actions.+- Test suite.+- PubSub actor.+- Support for bounded PubSub subscribers.
nqe.cabal view
@@ -1,11 +1,11 @@--- This file has been generated from package.yaml by hpack version 0.20.0.+-- This file has been generated from package.yaml by hpack version 0.28.2. -- -- see: https://github.com/sol/hpack ----- hash: a9ba15118eebb5f3b7c7996d90bc459e31f1d09bc3773d436c88a55e4fb6b3a6+-- hash: 1fd3035db301b81bce304fc8b3046396a8ed2fe1f7305463ac9bdac15f4029f7 name: nqe-version: 0.3.0.0+version: 0.4.0 synopsis: Concurrency library in the style of Erlang/OTP description: Minimalistic actor library inspired by Erlang/OTP with support for supervisor hierarchies and asynchronous messages, as well as abstractions for synchronous communication and easy management of TCP connections. category: Control@@ -17,8 +17,8 @@ license-file: UNLICENSE build-type: Simple cabal-version: >= 1.10- extra-source-files:+ CHANGELOG.md README.md source-repository head@@ -34,13 +34,19 @@ , conduit , conduit-extra , containers+ , exceptions+ , hspec+ , mtl , stm+ , stm-conduit+ , text , unliftio exposed-modules: Control.Concurrent.NQE other-modules: Control.Concurrent.NQE.Network Control.Concurrent.NQE.Process+ Control.Concurrent.NQE.PubSub Control.Concurrent.NQE.Supervisor Paths_nqe default-language: Haskell2010@@ -56,8 +62,10 @@ , bytestring , conduit , conduit-extra+ , containers , exceptions , hspec+ , mtl , nqe , stm , stm-conduit
src/Control/Concurrent/NQE.hs view
@@ -2,8 +2,10 @@ ( module Control.Concurrent.NQE.Process , module Control.Concurrent.NQE.Network , module Control.Concurrent.NQE.Supervisor+ , module Control.Concurrent.NQE.PubSub ) where import Control.Concurrent.NQE.Network import Control.Concurrent.NQE.Process+import Control.Concurrent.NQE.PubSub import Control.Concurrent.NQE.Supervisor
src/Control/Concurrent/NQE/Process.hs view
@@ -39,10 +39,10 @@ requeueMsg msg (Inbox mbox) = msg `requeueMsg` mbox mailboxEmpty :: (MonadIO m, Mailbox mbox) => mbox msg -> m Bool-mailboxEmpty = liftIO . atomically . mailboxEmptySTM+mailboxEmpty = atomically . mailboxEmptySTM send :: (MonadIO m, Mailbox mbox) => msg -> mbox msg -> m ()-send msg = liftIO . atomically . sendSTM msg+send msg = atomically . sendSTM msg requeue :: (Mailbox mbox) => [msg] -> mbox msg -> STM () requeue xs mbox = mapM_ (`requeueMsg` mbox) xs@@ -72,16 +72,16 @@ -> mbox msg -> m b query f mbox = do- box <- liftIO $ atomically newEmptyTMVar+ box <- atomically newEmptyTMVar f (putTMVar box) `send` mbox- liftIO . atomically $ takeTMVar box+ atomically (takeTMVar box) dispatch :: (MonadIO m, Mailbox mbox)- => [(msg -> Maybe a, a -> IO b)] -- ^ action to dispatch+ => [(msg -> Maybe a, a -> m b)] -- ^ action to dispatch -> mbox msg -- ^ mailbox to read from -> m b-dispatch hs = liftIO . join . atomically . extractMsg hs+dispatch hs = join . atomically . extractMsg hs dispatchSTM :: (Mailbox mbox) => [msg -> Maybe a] -> mbox msg -> STM a dispatchSTM = extractMsg . map (\x -> (x, id))
+ src/Control/Concurrent/NQE/PubSub.hs view
@@ -0,0 +1,131 @@+{-# LANGUAGE FlexibleContexts #-}+module Control.Concurrent.NQE.PubSub+ ( Publisher+ , publisher+ , boundedPublisher+ , withPubSub+ , withBoundedPubSub+ ) where++import Control.Applicative+import Control.Concurrent.NQE.Process+import Control.Monad.Reader+import Data.List+import UnliftIO++data ControlMsg ch msg+ = Subscribe (ch msg)+ | Unsubscribe (ch msg)++data Incoming ch msg+ = Control (ControlMsg ch msg)+ | Event msg++type Publisher mbox ch msg = mbox (ControlMsg ch msg)++withPubSub ::+ (MonadUnliftIO m, Mailbox mbox)+ => Publisher mbox TQueue msg+ -> (TQueue msg -> m a)+ -> m a+withPubSub pub f = bracket subscribe unsubscribe action+ where+ subscribe = do+ mbox <- newTQueueIO+ Subscribe mbox `send` pub+ return mbox+ unsubscribe mbox = Unsubscribe mbox `send` pub+ action mbox = f mbox++withBoundedPubSub ::+ (MonadUnliftIO m, Mailbox mbox)+ => Int+ -> Publisher mbox TBQueue msg+ -> (TBQueue msg -> m a)+ -> m a+withBoundedPubSub bound pub f = bracket subscribe unsubscribe action+ where+ subscribe = do+ mbox <- newTBQueueIO bound+ Subscribe mbox `send` pub+ return mbox+ unsubscribe mbox = Unsubscribe mbox `send` pub+ action mbox = f mbox++publisher ::+ ( MonadIO m+ , Mailbox mbox+ , Mailbox events+ , Mailbox ch+ , Eq (ch msg)+ )+ => Publisher mbox ch msg+ -> events msg+ -> m ()+publisher pub events = do+ box <- newTVarIO []+ runReaderT go box+ where+ go =+ forever $ do+ incoming <-+ atomically $+ Control <$> receiveSTM pub <|> Event <$> receiveSTM events+ process incoming++boundedPublisher ::+ (MonadIO m, Mailbox mbox, Mailbox events)+ => Publisher mbox TBQueue msg+ -> events msg+ -> m ()+boundedPublisher pub events = do+ box <- newTVarIO []+ runReaderT go box+ where+ go =+ forever $ do+ incoming <-+ atomically $+ Control <$> receiveSTM pub <|> Event <$> receiveSTM events+ processBound incoming++processBound ::+ (MonadIO m, MonadReader (TVar [TBQueue msg]) m)+ => Incoming TBQueue msg+ -> m ()+processBound (Control (Subscribe mbox)) = do+ box <- ask+ atomically $ do+ subscribers <- readTVar box+ when (mbox `notElem` subscribers) $ writeTVar box (mbox : subscribers)++processBound (Control (Unsubscribe mbox)) = do+ box <- ask+ atomically (modifyTVar box (delete mbox))++processBound (Event event) =+ ask >>= \box ->+ atomically $+ readTVar box >>= \subs ->+ forM_ subs $ \sub ->+ isFullTBQueue sub >>= \full ->+ when (not full) (event `sendSTM` sub)++process ::+ (Eq (ch msg), Mailbox ch, MonadIO m, MonadReader (TVar [ch msg]) m)+ => Incoming ch msg+ -> m ()+process (Control (Subscribe mbox)) = do+ box <- ask+ atomically $ do+ subscribers <- readTVar box+ when (mbox `notElem` subscribers) $+ writeTVar box (mbox : subscribers)++process (Control (Unsubscribe mbox)) = do+ box <- ask+ atomically (modifyTVar box (delete mbox))++process (Event event) = do+ box <- ask+ readTVarIO box >>= mapM_ (send event)
src/Control/Concurrent/NQE/Supervisor.hs view
@@ -17,8 +17,9 @@ import Control.Monad.STM (catchSTM) import UnliftIO -data SupervisorMessage- = AddChild (IO ())+data SupervisorMessage n+ = MonadUnliftIO n =>+ AddChild (n ()) (Reply (Async ())) | RemoveChild (Async ()) | StopSupervisor@@ -32,11 +33,11 @@ supervisor :: (MonadUnliftIO m, Mailbox mbox) => Strategy- -> mbox SupervisorMessage+ -> mbox (SupervisorMessage m) -> [m ()] -> m () supervisor strat mbox children = do- state <- liftIO $ newTVarIO []+ state <- newTVarIO [] finally (go state) (down state) where go state = do@@ -44,14 +45,15 @@ loop state loop state = do e <-- liftIO . atomically $+ atomically $ Right <$> receiveSTM mbox <|> Left <$> waitForChild state- again <- case e of- Right m -> processMessage state m- Left x -> processDead state strat x+ again <-+ case e of+ Right m -> processMessage state m+ Left x -> processDead state strat x when again $ loop state down state = do- as <- liftIO . atomically $ readTVar state+ as <- atomically (readTVar state) mapM_ cancel as waitForChild :: TVar [Async ()] -> STM (Async (), Either SomeException ())@@ -62,37 +64,37 @@ processMessage :: (MonadUnliftIO m) => TVar [Async ()]- -> SupervisorMessage+ -> SupervisorMessage m -> m Bool processMessage state (AddChild ch r) = do- a <- async $ liftIO ch- liftIO . atomically $ do+ a <- async ch+ atomically $ do modifyTVar' state (a:) r a return True processMessage state (RemoveChild a) = do- liftIO . atomically $ modifyTVar' state (filter (/= a))+ atomically (modifyTVar' state (filter (/= a))) cancel a return True processMessage state StopSupervisor = do- as <- liftIO $ readTVarIO state+ as <- readTVarIO state forM_ as (stopChild state) return False processDead ::- (MonadUnliftIO m)+ (MonadIO m) => TVar [Async ()] -> Strategy -> (Async (), Either SomeException ()) -> m Bool processDead state IgnoreAll (a, _) = do- liftIO . atomically $ modifyTVar' state (filter (/= a))+ atomically (modifyTVar' state (filter (/= a))) return True processDead state KillAll (a, e) = do- as <- liftIO . atomically $ do+ as <- atomically $ do modifyTVar' state (filter (/= a)) readTVar state mapM_ (stopChild state) as@@ -101,11 +103,11 @@ Right () -> return False processDead state IgnoreGraceful (a, Right ()) = do- liftIO . atomically $ modifyTVar' state (filter (/= a))+ atomically (modifyTVar' state (filter (/= a))) return True processDead state IgnoreGraceful (a, Left e) = do- as <- liftIO . atomically $ do+ as <- atomically $ do modifyTVar' state (filter (/= a)) readTVar state mapM_ (stopChild state) as@@ -113,14 +115,14 @@ processDead state (Notify notif) (a, e) = do x <-- liftIO . atomically $ do+ atomically $ do modifyTVar' state (filter (/= a)) catchSTM (notif (a, e) >> return Nothing) $ \x -> return $ Just (x :: SomeException) case x of Nothing -> return True Just ex -> do- as <- liftIO $ readTVarIO state+ as <- readTVarIO state forM_ as (stopChild state) throwIO ex @@ -131,34 +133,35 @@ -> m (Async ()) startChild state run = do a <- async run- liftIO . atomically $ modifyTVar' state (a:)+ atomically (modifyTVar' state (a:)) return a -stopChild ::- (MonadUnliftIO m) => TVar [Async ()] -> Async () -> m ()+stopChild :: MonadIO m => TVar [Async ()] -> Async () -> m () stopChild state a = do isChild <-- liftIO . atomically $ do+ atomically $ do cur <- readTVar state let new = filter (/= a) cur writeTVar state new- return $ cur /= new- when isChild $ cancel a+ return (cur /= new)+ when isChild (cancel a) addChild ::- (MonadIO m, Mailbox mbox)- => mbox SupervisorMessage- -> IO ()+ (MonadUnliftIO n, MonadIO m, Mailbox mbox)+ => mbox (SupervisorMessage n)+ -> n () -> m (Async ()) addChild mbox action = AddChild action `query` mbox removeChild ::- (MonadIO m, Mailbox mbox)- => mbox SupervisorMessage+ (MonadUnliftIO n, MonadIO m, Mailbox mbox)+ => mbox (SupervisorMessage n) -> Async () -> m () removeChild mbox child = RemoveChild child `send` mbox stopSupervisor ::- (MonadIO m, Mailbox mbox) => mbox SupervisorMessage -> m ()+ (MonadUnliftIO n, MonadIO m, Mailbox mbox)+ => mbox (SupervisorMessage n)+ -> m () stopSupervisor mbox = StopSupervisor `send` mbox
test/Spec.hs view
@@ -1,13 +1,14 @@ {-# LANGUAGE LambdaCase #-} {-# LANGUAGE OverloadedStrings #-} {-# LANGUAGE ScopedTypeVariables #-}+import Conduit import Control.Concurrent hiding (yield) import Control.Concurrent.NQE+import Control.Concurrent.STM (retry) import Control.Exception import Control.Monad import Control.Monad.Catch import Data.ByteString (ByteString)-import Conduit import Data.Conduit.Text (decode, encode, utf8) import qualified Data.Conduit.Text as CT import Data.Conduit.TMChan@@ -89,12 +90,12 @@ main :: IO () main = hspec $ do- describe "two communicating processes" $ do+ describe "two communicating processes" $ it "exchange ping/pong messages" $ do mbox <- newTQueueIO g <- withAsync (pong mbox) $ const $ query Ping mbox g `shouldBe` Pong- describe "network process" $ do+ describe "network process" $ it "responds to a ping" $ do (source1, sink1, source2, sink2) <- conduits msg <-@@ -105,7 +106,7 @@ n <- timeout 0xbeef (threadDelay 0xdeadbeef) n `shouldBe` Nothing it "timeout action succeeds" $ do- n <- timeout 0xdeadbeef (return 0xbeef)+ n <- timeout 0xdeadbeef (return (0xbeef :: Integer)) n `shouldBe` Just 0xbeef describe "supervisor" $ do let p1 m = forever $ receive m >>= \r -> atomically $ r ()@@ -162,8 +163,34 @@ case fromException x of Just TestError1 -> True Just TestError2 -> True- _ -> False+ _ -> False snd t1 `shouldSatisfy` er snd t2 `shouldSatisfy` er stopSupervisor sup wait g `shouldReturn` ()+ describe "pubsub" $ do+ it "sends messages to all subscribers" $ do+ let msgs = words "hello world"+ pub <- newTQueueIO+ events <- newTQueueIO+ withAsync (publisher pub events) $ \_ ->+ withPubSub pub $ \sub1 ->+ withPubSub pub $ \sub2 -> do+ mapM_ (`send` events) msgs+ sub1msgs <- replicateM 2 (receive sub1)+ sub2msgs <- replicateM 2 (receive sub2)+ sub1msgs `shouldBe` msgs+ sub2msgs `shouldBe` msgs+ it "drops messages when bounded queue full" $ do+ let msgs = words "hello world drop"+ pub <- newTQueueIO+ events <- newTQueueIO+ withAsync (boundedPublisher pub events) $ \_ ->+ withBoundedPubSub 2 pub $ \sub -> do+ mapM_ (`send` events) msgs+ atomically $+ isFullTBQueue sub >>= \full -> when (not full) retry+ msgs <- replicateM 2 (receive sub)+ "meh" `send` events+ msg <- receive sub+ msgs <> [msg] `shouldBe` (words "hello world meh")