distributed-process 0.7.8 → 0.7.9
raw patch · 10 files changed
+830/−421 lines, 10 filesdep +tasty-benchPVP: major bump suggested
API removals or changes: PVP suggests a major version bump
Dependencies added: tasty-bench
API changes (from Hackage documentation)
- Control.Distributed.Process: catchExit :: forall a b. (Show a, Serializable a) => Process b -> (ProcessId -> a -> Process b) -> Process b
+ Control.Distributed.Process: catchExit :: (Show a, Serializable a) => Process b -> (ProcessId -> a -> Process b) -> Process b
- Control.Distributed.Process: data () => Closure a
+ Control.Distributed.Process: data Closure a
- Control.Distributed.Process: data () => RemoteTable
+ Control.Distributed.Process: data RemoteTable
- Control.Distributed.Process: data () => Static a
+ Control.Distributed.Process: data Static a
- Control.Distributed.Process: expect :: forall a. Serializable a => Process a
+ Control.Distributed.Process: expect :: Serializable a => Process a
- Control.Distributed.Process: expectTimeout :: forall a. Serializable a => Int -> Process (Maybe a)
+ Control.Distributed.Process: expectTimeout :: Serializable a => Int -> Process (Maybe a)
- Control.Distributed.Process: handleMessage :: forall m a b. (Monad m, Serializable a) => Message -> (a -> m b) -> m (Maybe b)
+ Control.Distributed.Process: handleMessage :: (Monad m, Serializable a) => Message -> (a -> m b) -> m (Maybe b)
- Control.Distributed.Process: handleMessageIf :: forall m a b. (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m b) -> m (Maybe b)
+ Control.Distributed.Process: handleMessageIf :: (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m b) -> m (Maybe b)
- Control.Distributed.Process: handleMessageIf_ :: forall m a. (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m ()) -> m ()
+ Control.Distributed.Process: handleMessageIf_ :: (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m ()) -> m ()
- Control.Distributed.Process: handleMessage_ :: forall m a. (Monad m, Serializable a) => Message -> (a -> m ()) -> m ()
+ Control.Distributed.Process: handleMessage_ :: (Monad m, Serializable a) => Message -> (a -> m ()) -> m ()
- Control.Distributed.Process: mask :: ((forall a. Process a -> Process a) -> Process b) -> Process b
+ Control.Distributed.Process: mask :: ((forall a. () => Process a -> Process a) -> Process b) -> Process b
- Control.Distributed.Process: match :: forall a b. Serializable a => (a -> Process b) -> Match b
+ Control.Distributed.Process: match :: Serializable a => (a -> Process b) -> Match b
- Control.Distributed.Process: matchAny :: forall b. (Message -> Process b) -> Match b
+ Control.Distributed.Process: matchAny :: (Message -> Process b) -> Match b
- Control.Distributed.Process: matchAnyIf :: forall a b. Serializable a => (a -> Bool) -> (Message -> Process b) -> Match b
+ Control.Distributed.Process: matchAnyIf :: Serializable a => (a -> Bool) -> (Message -> Process b) -> Match b
- Control.Distributed.Process: matchIf :: forall a b. Serializable a => (a -> Bool) -> (a -> Process b) -> Match b
+ Control.Distributed.Process: matchIf :: Serializable a => (a -> Bool) -> (a -> Process b) -> Match b
- Control.Distributed.Process: monitorPort :: forall a. Serializable a => SendPort a -> Process MonitorRef
+ Control.Distributed.Process: monitorPort :: Serializable a => SendPort a -> Process MonitorRef
- Control.Distributed.Process: spawnChannel :: forall a. Serializable a => Static (SerializableDict a) -> NodeId -> Closure (ReceivePort a -> Process ()) -> Process (SendPort a)
+ Control.Distributed.Process: spawnChannel :: Serializable a => Static (SerializableDict a) -> NodeId -> Closure (ReceivePort a -> Process ()) -> Process (SendPort a)
- Control.Distributed.Process: unwrapMessage :: forall m a. (Monad m, Serializable a) => Message -> m (Maybe a)
+ Control.Distributed.Process: unwrapMessage :: (Monad m, Serializable a) => Message -> m (Maybe a)
- Control.Distributed.Process.Closure: [SerializableDict] :: Serializable a => SerializableDict a
+ Control.Distributed.Process.Closure: [SerializableDict] :: forall a. Serializable a => SerializableDict a
- Control.Distributed.Process.Closure: bindCP :: forall a b. (Typeable a, Typeable b) => Closure (Process a) -> CP a b -> Closure (Process b)
+ Control.Distributed.Process.Closure: bindCP :: (Typeable a, Typeable b) => Closure (Process a) -> CP a b -> Closure (Process b)
- Control.Distributed.Process.Closure: call' :: forall a. Serializable a => NodeId -> Closure (Process a) -> Process a
+ Control.Distributed.Process.Closure: call' :: Serializable a => NodeId -> Closure (Process a) -> Process a
- Control.Distributed.Process.Closure: cpSend :: forall a. Typeable a => Static (SerializableDict a) -> ProcessId -> CP a ()
+ Control.Distributed.Process.Closure: cpSend :: Typeable a => Static (SerializableDict a) -> ProcessId -> CP a ()
- Control.Distributed.Process.Closure: mkClosureValSingle :: forall a b. (Serializable a, Typeable b, MkTDict b) => String -> (a -> b) -> (a -> Closure b, RemoteRegister)
+ Control.Distributed.Process.Closure: mkClosureValSingle :: (Serializable a, Typeable b, MkTDict b) => String -> (a -> b) -> (a -> Closure b, RemoteRegister)
- Control.Distributed.Process.Closure: returnCP :: forall a. Serializable a => Static (SerializableDict a) -> a -> Closure (Process a)
+ Control.Distributed.Process.Closure: returnCP :: Serializable a => Static (SerializableDict a) -> a -> Closure (Process a)
- Control.Distributed.Process.Closure: type CP a b = Closure (a -> Process b)
+ Control.Distributed.Process.Closure: type CP a b = Closure a -> Process b
- Control.Distributed.Process.Debug: withFlags :: forall a. TraceFlags -> Process a -> Process (Either SomeException a)
+ Control.Distributed.Process.Debug: withFlags :: TraceFlags -> Process a -> Process (Either SomeException a)
- Control.Distributed.Process.Debug: withTracer :: forall a. (MxEvent -> Process ()) -> Process a -> Process (Either SomeException a)
+ Control.Distributed.Process.Debug: withTracer :: (MxEvent -> Process ()) -> Process a -> Process (Either SomeException a)
- Control.Distributed.Process.Internal.CQueue: dequeue :: forall m a. CQueue m -> BlockSpec -> [MatchOn m a] -> IO (Maybe a)
+ Control.Distributed.Process.Internal.CQueue: dequeue :: CQueue m -> BlockSpec -> [MatchOn m a] -> IO (Maybe a)
- Control.Distributed.Process.Internal.Closure.BuiltIn: bindCP :: forall a b. (Typeable a, Typeable b) => Closure (Process a) -> CP a b -> Closure (Process b)
+ Control.Distributed.Process.Internal.Closure.BuiltIn: bindCP :: (Typeable a, Typeable b) => Closure (Process a) -> CP a b -> Closure (Process b)
- Control.Distributed.Process.Internal.Closure.BuiltIn: cpSend :: forall a. Typeable a => Static (SerializableDict a) -> ProcessId -> CP a ()
+ Control.Distributed.Process.Internal.Closure.BuiltIn: cpSend :: Typeable a => Static (SerializableDict a) -> ProcessId -> CP a ()
- Control.Distributed.Process.Internal.Closure.BuiltIn: returnCP :: forall a. Serializable a => Static (SerializableDict a) -> a -> Closure (Process a)
+ Control.Distributed.Process.Internal.Closure.BuiltIn: returnCP :: Serializable a => Static (SerializableDict a) -> a -> Closure (Process a)
- Control.Distributed.Process.Internal.Closure.BuiltIn: type CP a b = Closure (a -> Process b)
+ Control.Distributed.Process.Internal.Closure.BuiltIn: type CP a b = Closure a -> Process b
- Control.Distributed.Process.Internal.Closure.Explicit: call' :: forall a. Serializable a => NodeId -> Closure (Process a) -> Process a
+ Control.Distributed.Process.Internal.Closure.Explicit: call' :: Serializable a => NodeId -> Closure (Process a) -> Process a
- Control.Distributed.Process.Internal.Closure.Explicit: mkClosureValSingle :: forall a b. (Serializable a, Typeable b, MkTDict b) => String -> (a -> b) -> (a -> Closure b, RemoteRegister)
+ Control.Distributed.Process.Internal.Closure.Explicit: mkClosureValSingle :: (Serializable a, Typeable b, MkTDict b) => String -> (a -> b) -> (a -> Closure b, RemoteRegister)
- Control.Distributed.Process.Internal.Primitives: catchExit :: forall a b. (Show a, Serializable a) => Process b -> (ProcessId -> a -> Process b) -> Process b
+ Control.Distributed.Process.Internal.Primitives: catchExit :: (Show a, Serializable a) => Process b -> (ProcessId -> a -> Process b) -> Process b
- Control.Distributed.Process.Internal.Primitives: expect :: forall a. Serializable a => Process a
+ Control.Distributed.Process.Internal.Primitives: expect :: Serializable a => Process a
- Control.Distributed.Process.Internal.Primitives: expectTimeout :: forall a. Serializable a => Int -> Process (Maybe a)
+ Control.Distributed.Process.Internal.Primitives: expectTimeout :: Serializable a => Int -> Process (Maybe a)
- Control.Distributed.Process.Internal.Primitives: handleMessage :: forall m a b. (Monad m, Serializable a) => Message -> (a -> m b) -> m (Maybe b)
+ Control.Distributed.Process.Internal.Primitives: handleMessage :: (Monad m, Serializable a) => Message -> (a -> m b) -> m (Maybe b)
- Control.Distributed.Process.Internal.Primitives: handleMessageIf :: forall m a b. (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m b) -> m (Maybe b)
+ Control.Distributed.Process.Internal.Primitives: handleMessageIf :: (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m b) -> m (Maybe b)
- Control.Distributed.Process.Internal.Primitives: handleMessageIf_ :: forall m a. (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m ()) -> m ()
+ Control.Distributed.Process.Internal.Primitives: handleMessageIf_ :: (Monad m, Serializable a) => Message -> (a -> Bool) -> (a -> m ()) -> m ()
- Control.Distributed.Process.Internal.Primitives: handleMessage_ :: forall m a. (Monad m, Serializable a) => Message -> (a -> m ()) -> m ()
+ Control.Distributed.Process.Internal.Primitives: handleMessage_ :: (Monad m, Serializable a) => Message -> (a -> m ()) -> m ()
- Control.Distributed.Process.Internal.Primitives: mask :: ((forall a. Process a -> Process a) -> Process b) -> Process b
+ Control.Distributed.Process.Internal.Primitives: mask :: ((forall a. () => Process a -> Process a) -> Process b) -> Process b
- Control.Distributed.Process.Internal.Primitives: match :: forall a b. Serializable a => (a -> Process b) -> Match b
+ Control.Distributed.Process.Internal.Primitives: match :: Serializable a => (a -> Process b) -> Match b
- Control.Distributed.Process.Internal.Primitives: matchAny :: forall b. (Message -> Process b) -> Match b
+ Control.Distributed.Process.Internal.Primitives: matchAny :: (Message -> Process b) -> Match b
- Control.Distributed.Process.Internal.Primitives: matchAnyIf :: forall a b. Serializable a => (a -> Bool) -> (Message -> Process b) -> Match b
+ Control.Distributed.Process.Internal.Primitives: matchAnyIf :: Serializable a => (a -> Bool) -> (Message -> Process b) -> Match b
- Control.Distributed.Process.Internal.Primitives: matchIf :: forall a b. Serializable a => (a -> Bool) -> (a -> Process b) -> Match b
+ Control.Distributed.Process.Internal.Primitives: matchIf :: Serializable a => (a -> Bool) -> (a -> Process b) -> Match b
- Control.Distributed.Process.Internal.Primitives: monitorPort :: forall a. Serializable a => SendPort a -> Process MonitorRef
+ Control.Distributed.Process.Internal.Primitives: monitorPort :: Serializable a => SendPort a -> Process MonitorRef
- Control.Distributed.Process.Internal.Primitives: unwrapMessage :: forall m a. (Monad m, Serializable a) => Message -> m (Maybe a)
+ Control.Distributed.Process.Internal.Primitives: unwrapMessage :: (Monad m, Serializable a) => Message -> m (Maybe a)
- Control.Distributed.Process.Internal.Spawn: spawnChannel :: forall a. Serializable a => Static (SerializableDict a) -> NodeId -> Closure (ReceivePort a -> Process ()) -> Process (SendPort a)
+ Control.Distributed.Process.Internal.Spawn: spawnChannel :: Serializable a => Static (SerializableDict a) -> NodeId -> Closure (ReceivePort a -> Process ()) -> Process (SendPort a)
- Control.Distributed.Process.Management: mxDeactivate :: forall s. String -> MxAgent s MxAction
+ Control.Distributed.Process.Management: mxDeactivate :: String -> MxAgent s MxAction
- Control.Distributed.Process.Management: mxReady :: forall s. MxAgent s MxAction
+ Control.Distributed.Process.Management: mxReady :: MxAgent s MxAction
- Control.Distributed.Process.Management: mxReceive :: forall s. MxAgent s MxAction
+ Control.Distributed.Process.Management: mxReceive :: MxAgent s MxAction
- Control.Distributed.Process.Management: mxReceiveChan :: forall s. MxAgent s MxAction
+ Control.Distributed.Process.Management: mxReceiveChan :: MxAgent s MxAction
- Control.Distributed.Process.Management: mxSkip :: forall s. MxAgent s MxAction
+ Control.Distributed.Process.Management: mxSkip :: MxAgent s MxAction
- Control.Distributed.Process.Management: type MxSink s = Message -> MxAgent s (Maybe MxAction)
+ Control.Distributed.Process.Management: type MxSink s = Message -> MxAgent s Maybe MxAction
- Control.Distributed.Process.Management.Internal.Agent: type AgentConfig = (Tracer, Weak (CQueue Message), (((TChan Message, TChan Message) -> Process ()) -> IO ProcessId))
+ Control.Distributed.Process.Management.Internal.Agent: type AgentConfig = (Tracer, Weak CQueue Message, (TChan Message, TChan Message) -> Process () -> IO ProcessId)
- Control.Distributed.Process.Management.Internal.Types: type Fork = (Process () -> IO ProcessId)
+ Control.Distributed.Process.Management.Internal.Types: type Fork = Process () -> IO ProcessId
- Control.Distributed.Process.Management.Internal.Types: type MxSink s = Message -> MxAgent s (Maybe MxAction)
+ Control.Distributed.Process.Management.Internal.Types: type MxSink s = Message -> MxAgent s Maybe MxAction
- Control.Distributed.Process.Serializable: [SerializableDict] :: Serializable a => SerializableDict a
+ Control.Distributed.Process.Serializable: [SerializableDict] :: forall a. Serializable a => SerializableDict a
- Control.Distributed.Process.Serializable: [TypeableDict] :: Typeable a => TypeableDict a
+ Control.Distributed.Process.Serializable: [TypeableDict] :: forall a. Typeable a => TypeableDict a
- Control.Distributed.Process.Serializable: data () => Fingerprint
+ Control.Distributed.Process.Serializable: data Fingerprint
Files
- ChangeLog +11/−3
- benchmarks/Channels.hs +0/−43
- benchmarks/Latency.hs +0/−39
- benchmarks/Main.hs +753/−0
- benchmarks/ProcessRing.hs +0/−117
- benchmarks/Spawns.hs +0/−49
- benchmarks/Throughput.hs +0/−74
- distributed-process.cabal +12/−51
- src/Control/Distributed/Process/Internal/CQueue.hs +47/−43
- src/Control/Distributed/Process/Internal/Primitives.hs +7/−2
ChangeLog view
@@ -1,3 +1,11 @@+2026-09-21 Laurent P. René de Cotret <laurent.decotret@outlook.com> 0.7.9++* `receiveTimeout`, `expectTimeout` and `receiveChanTimeout` now arm the+ system timer only once they are about to block. This results in up to a 9x speedup+ in cases where a process receives lots of messages (in which case, a timeout isn't necessary). (#490)+* Reworked benchmarks, which can now be run using `cabal bench distributed-process`. (#487)+* Added support for `containers-0.8`.+ 2025-02-04 Laurent P. René de Cotret <laurent.decotret@outlook.com> 0.7.8 * Added documentation on the unit of measurement for timeout durations (#340)@@ -190,7 +198,7 @@ * Numerous memory leaks plugged * Relax upper bound on dependency on 'network' * New primitive 'matchAny'-* Remove 'whereisRemote' (see comment of 'whereisRemoteAsync') +* Remove 'whereisRemote' (see comment of 'whereisRemoteAsync') 2012-08-16 Edsko de Vries <edsko@well-typed.com> 0.3.1 @@ -240,9 +248,9 @@ 2012-07-11 Edsko de Vries <edsko@well-typed.com> 0.2.1 * Complete redesign of the underlying implementation of static values and-closures. +closures. -* Add support for 'spawnChannel' +* Add support for 'spawnChannel' 2012-07-09 Edsko de Vries <edsko@well-typed.com> 0.2.0.1
− benchmarks/Channels.hs
@@ -1,43 +0,0 @@--- | Like Latency, but creating lots of channels-import System.Environment-import Control.Monad-import Control.Applicative-import Control.Distributed.Process-import Control.Distributed.Process.Node-import Network.Transport.TCP (createTransport, defaultTCPAddr, defaultTCPParameters)-import Data.Binary (encode, decode)-import qualified Data.ByteString.Lazy as BSL--pingServer :: Process ()-pingServer = forever $ do- them <- expect- sendChan them ()- -- TODO: should this be automatic?- reconnectPort them--pingClient :: Int -> ProcessId -> Process ()-pingClient n them = do- replicateM_ n $ do- (sc, rc) <- newChan :: Process (SendPort (), ReceivePort ())- send them sc- receiveChan rc- liftIO . putStrLn $ "Did " ++ show n ++ " pings"--initialProcess :: String -> Process ()-initialProcess "SERVER" = do- us <- getSelfPid- liftIO $ BSL.writeFile "pingServer.pid" (encode us)- pingServer-initialProcess "CLIENT" = do- n <- liftIO $ getLine- them <- liftIO $ decode <$> BSL.readFile "pingServer.pid"- pingClient (read n) them--main :: IO ()-main = do- [role, host, port] <- getArgs- trans <- createTransport (defaultTCPAddr host port) defaultTCPParameters- case trans of- Right transport -> do node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role- Left other -> error $ show other
− benchmarks/Latency.hs
@@ -1,39 +0,0 @@-import System.Environment-import Control.Monad-import Control.Applicative-import Control.Distributed.Process-import Control.Distributed.Process.Node-import Network.Transport.TCP (createTransport, defaultTCPAddr, defaultTCPParameters)-import Data.Binary (encode, decode)-import qualified Data.ByteString.Lazy as BSL--pingServer :: Process ()-pingServer = forever $ do- them <- expect- send them ()--pingClient :: Int -> ProcessId -> Process ()-pingClient n them = do- us <- getSelfPid- replicateM_ n $ send them us >> (expect :: Process ())- liftIO . putStrLn $ "Did " ++ show n ++ " pings"--initialProcess :: String -> Process ()-initialProcess "SERVER" = do- us <- getSelfPid- liftIO $ BSL.writeFile "pingServer.pid" (encode us)- pingServer-initialProcess "CLIENT" = do- n <- liftIO $ getLine- them <- liftIO $ decode <$> BSL.readFile "pingServer.pid"- pingClient (read n) them--main :: IO ()-main = do- [role, host, port] <- getArgs- trans <- createTransport- (defaultTCPAddr host port) defaultTCPParameters- case trans of- Right transport -> do node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role- Left other -> error $ show other
+ benchmarks/Main.hs view
@@ -0,0 +1,753 @@+{-# LANGUAGE BangPatterns #-}+{-# LANGUAGE DeriveGeneric #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE TemplateHaskell #-}+{-# OPTIONS_GHC -Wno-unused-top-binds #-}++module Main (main) where++import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)+import Control.Concurrent.STM+ ( TQueue,+ atomically,+ newTQueueIO,+ readTQueue,+ writeTQueue,+ )+import Control.Distributed.Process+ ( Handler (Handler),+ MonitorRef,+ NodeId,+ Process,+ ProcessId,+ ProcessMonitorNotification (ProcessMonitorNotification),+ ReceivePort,+ SendPort,+ WhereIsReply (WhereIsReply),+ call,+ callLocal,+ catchExit,+ catches,+ catchesExit,+ delegate,+ die,+ exit,+ expect,+ expectTimeout,+ forward,+ getLocalNodeStats,+ getNodeStats,+ getProcessInfo,+ getSelfNode,+ getSelfPid,+ handleMessage,+ kill,+ liftIO,+ link,+ match,+ matchAny,+ matchChan,+ matchIf,+ matchMessage,+ matchSTM,+ matchUnknown,+ mergePortsBiased,+ mergePortsRR,+ monitor,+ monitorNode,+ monitorPort,+ newChan,+ nsend,+ nsendRemote,+ proxy,+ receiveChan,+ receiveChanTimeout,+ receiveTimeout,+ receiveWait,+ register,+ relay,+ reregister,+ send,+ sendChan,+ spawn,+ spawnChannel,+ spawnChannelLocal,+ spawnLocal,+ spawnMonitor,+ uforward,+ unlink,+ unmonitor,+ unregister,+ unsafeSend,+ unwrapMessage,+ usend,+ whereis,+ whereisRemoteAsync,+ withMonitor_,+ wrapMessage,+ )+import Control.Distributed.Process.Closure+ ( functionTDict,+ mkClosure,+ remotable,+ sdictUnit,+ )+import Control.Distributed.Process.Node+ ( LocalNode (..),+ closeLocalNode,+ forkProcess,+ initRemoteTable,+ newLocalNode,+ runProcess,+ )+import Control.Distributed.Process.Serializable (Serializable)+import qualified Control.Exception as E+import Control.Monad (forever, replicateM, replicateM_, void, when)+import qualified Control.Monad.Catch as Catch+import Data.Binary (Binary)+import qualified Data.ByteString.Char8 as BS+import GHC.Generics (Generic)+import qualified Network.Transport as NT+import Network.Transport.TCP+ ( createTransport,+ defaultTCPAddr,+ defaultTCPParameters,+ )+import Test.Tasty.Bench+ ( Benchmark,+ bench,+ bgroup,+ defaultMain,+ whnfIO,+ )++-- A top-level splice only brings names into scope for later declaration+-- groups, so these must precede 'main' and every use of 'mkClosure'.++remoteSignal :: ProcessId -> Process ()+remoteSignal them = send them ()++remoteChanEcho :: ProcessId -> ReceivePort () -> Process ()+remoteChanEcho them rp = receiveChan rp >> send them ()++remoteAnswer :: () -> Process Int+remoteAnswer () = return 42++remotable ['remoteSignal, 'remoteChanEcho, 'remoteAnswer]++main :: IO ()+main = do+ let rtable = __remoteTable initRemoteTable+ transport <-+ either E.throwIO return+ =<< createTransport (defaultTCPAddr "127.0.0.1" "0") defaultTCPParameters+ ( E.bracket (newLocalNode transport rtable) closeLocalNode $ \node1 ->+ E.bracket (newLocalNode transport rtable) closeLocalNode $ \node2 -> do+ fx <- setup node1 node2+ defaultMain (benchmarks fx)+ )+ `E.finally` NT.closeTransport transport++benchmarks :: Fixture -> [Benchmark]+benchmarks fx =+ [ bgroup+ "local"+ [ baseline fx,+ messaging fx,+ channels fx,+ receiving fx,+ timeouts fx,+ messages fx,+ processes fx,+ monitoring fx,+ registry fx,+ exceptions fx,+ introspection fx,+ ring fx+ ],+ remote fx+ ]++-- | Cost of the harness alone, and so the floor below which the other numbers+-- say nothing.+baseline :: Fixture -> Benchmark+baseline fx =+ bgroup+ "baseline"+ [ oneBench fx "runner round trip (no work)" (return ()),+ repsBench fx "empty loop" 1000 (return ())+ ]++messaging :: Fixture -> Benchmark+messaging fx =+ bgroup+ "messaging"+ [ repsBench fx "send/expect" 1000 $ do+ us <- getSelfPid+ send echo us+ expect :: Process (),+ repsBench fx "usend/expect" 1000 $ do+ us <- getSelfPid+ usend echo us+ expect :: Process (),+ repsBench fx "unsafeSend/expect" 1000 $ do+ us <- getSelfPid+ unsafeSend echo us+ expect :: Process (),+ repsBench fx "nsend/expect" 100 $ do+ us <- getSelfPid+ nsend echoName us+ expect :: Process (),+ bgroup+ "throughput/bytestring"+ [ oneBench fx (show sz ++ "B") $+ sendThrough (fxCounter fx) 1000 (BS.replicate sz 'x')+ | sz <- [8, 1024, 65536]+ ],+ bgroup+ "throughput/list-of-int"+ [ oneBench fx (show n ++ "elems") $+ sendThrough (fxCounter fx) 1000 [1 .. n]+ | n <- [1, 100 :: Int]+ ]+ ]+ where+ echo = fxEcho fx++channels :: Fixture -> Benchmark+channels fx =+ bgroup+ "channels"+ [ repsBench fx "newChan" 100 $+ void (newChan :: Process (SendPort (), ReceivePort ())),+ repsBench fx "sendChan/receiveChan" 1000 $ do+ sendChan sp ()+ receiveChan rp,+ repsBench fx "receiveChanTimeout (empty)" 1000 $+ void (receiveChanTimeout 0 rp),+ repsBench fx "newChan + roundtrip via echo server" 100 $ do+ (sp', rp') <- newChan+ send (fxEcho fx) sp'+ receiveChan rp',+ repsBench fx "spawnChannelLocal" 100 $ do+ us <- getSelfPid+ sp' <- spawnChannelLocal $ \rp' ->+ (receiveChan rp' :: Process ()) >> send us ()+ sendChan sp' ()+ expect :: Process (),+ repsBench fx "mergePortsBiased" 100 (mergeBench mergePortsBiased),+ repsBench fx "mergePortsRR" 100 (mergeBench mergePortsRR)+ ]+ where+ (sp, rp) = fxChan fx++ mergeBench merge = do+ (sps, rps) <-+ unzip+ <$> replicateM 4 (newChan :: Process (SendPort (), ReceivePort ()))+ merged <- merge rps+ mapM_ (`sendChan` ()) sps+ replicateM_ 4 (receiveChan merged)++-- | Each of these puts one message in the runner's own mailbox and takes it+-- out again, so the spread between them is the cost of the 'Match'.+receiving :: Fixture -> Benchmark+receiving fx =+ bgroup+ "receiving"+ [ repsBench fx "expect" 1000 $ do+ selfSend ()+ expect :: Process (),+ repsBench fx "receiveWait (first of 1 match)" 1000 $ do+ selfSend ()+ receiveWait [match (\() -> return ())],+ repsBench fx "receiveWait (last of 6 matches)" 1000 $ do+ selfSend ()+ receiveWait+ [ match (\(_ :: Int) -> return ()),+ match (\(_ :: Bool) -> return ()),+ match (\(_ :: Char) -> return ()),+ match (\(_ :: BS.ByteString) -> return ()),+ match (\(_ :: Ping) -> return ()),+ match (\() -> return ())+ ],+ repsBench fx "matchIf" 1000 $ do+ selfSend (1 :: Int)+ receiveWait [matchIf (> (0 :: Int)) (\_ -> return ())],+ repsBench fx "matchAny" 1000 $ do+ selfSend ()+ receiveWait [matchAny (\_ -> return ())],+ repsBench fx "matchUnknown" 1000 $ do+ selfSend ()+ receiveWait [match (\(_ :: Int) -> return ()), matchUnknown (return ())],+ repsBench fx "matchMessage" 1000 $ do+ selfSend ()+ void (receiveWait [matchMessage return]),+ repsBench fx "matchChan" 1000 $ do+ sendChan sp ()+ receiveWait [matchChan rp return],+ repsBench fx "matchSTM" 1000 $ do+ liftIO (atomically (writeTQueue q ()))+ receiveWait [matchSTM (readTQueue q) return],+ repsBench fx "receiveTimeout (empty mailbox)" 1000 $+ void (receiveTimeout 0 [match (\() -> return ())]),+ repsBench fx "expectTimeout (hit)" 1000 $ do+ selfSend ()+ void (expectTimeout 0 :: Process (Maybe ()))+ ]+ where+ (sp, rp) = fxChan fx+ q = fxQueue fx++ selfSend :: (Serializable a) => a -> Process ()+ selfSend x = getSelfPid >>= \us -> unsafeSend us x++timeouts :: Fixture -> Benchmark+timeouts fx =+ bgroup+ "timeouts"+ [ repsBench fx "expectTimeout 0 (message waiting)" 1000 $ do+ leaveInMailbox+ void (expectTimeout 0 :: Process (Maybe Int)),+ repsBench fx "expectTimeout 1s (message waiting)" 1000 $ do+ leaveInMailbox+ void (expectTimeout 1000000 :: Process (Maybe Int)),+ repsBench fx "receiveTimeout 0 (matchSTM ready)" 1000 $ do+ liftIO (atomically (writeTQueue q ()))+ void (receiveTimeout 0 [matchSTM (readTQueue q) return]),+ repsBench fx "receiveTimeout 1s (matchSTM ready)" 1000 $ do+ liftIO (atomically (writeTQueue q ()))+ void (receiveTimeout 1000000 [matchSTM (readTQueue q) return]),+ repsBench fx "receiveChanTimeout 0 (value waiting)" 1000 $ do+ sendChan sp ()+ barrier+ void (receiveChanTimeout 0 rp),+ repsBench fx "receiveChanTimeout 1s (value waiting)" 1000 $ do+ sendChan sp ()+ barrier+ void (receiveChanTimeout 1000000 rp)+ ]+ where+ (sp, rp) = fxChan fx+ q = fxQueue fx++ barrier :: Process ()+ barrier = do+ us <- getSelfPid+ unsafeSend us ()+ expect :: Process ()++ leaveInMailbox :: Process ()+ leaveInMailbox = do+ us <- getSelfPid+ unsafeSend us (1 :: Int)+ barrier++messages :: Fixture -> Benchmark+messages fx =+ bgroup+ "messages"+ [ repsBench fx "unwrapMessage (hit)" 1000 $+ void (unwrapMessage intMessage :: Process (Maybe Int)),+ repsBench fx "unwrapMessage (miss)" 1000 $+ void (unwrapMessage intMessage :: Process (Maybe Bool)),+ repsBench fx "handleMessage (hit)" 1000 $+ void (handleMessage intMessage (\(_ :: Int) -> return ())),+ repsBench fx "handleMessage (miss)" 1000 $+ void (handleMessage intMessage (\(_ :: Bool) -> return ())),+ repsBench fx "wrapMessage + unwrapMessage" 1000 $+ void (unwrapMessage (wrapMessage (42 :: Int)) :: Process (Maybe Int)),+ -- A 'ProcessId' is sent rather than @()@ so that 'echoServer' recognises+ -- the forwarded message and replies.+ repsBench fx "forward" 1000 $ do+ us <- getSelfPid+ unsafeSend us us+ receiveWait [matchAny (`forward` fxEcho fx)]+ expect :: Process (),+ repsBench fx "uforward" 1000 $ do+ us <- getSelfPid+ unsafeSend us us+ receiveWait [matchAny (`uforward` fxEcho fx)]+ expect :: Process (),+ repsBench fx "relay" 1000 $ do+ send (fxRelay fx) ()+ expect :: Process (),+ repsBench fx "delegate" 1000 $ do+ send (fxDelegate fx) ()+ expect :: Process (),+ repsBench fx "proxy" 1000 $ do+ send (fxProxy fx) ()+ expect :: Process ()+ ]+ where+ intMessage = wrapMessage (42 :: Int)++processes :: Fixture -> Benchmark+processes fx =+ bgroup+ "processes"+ [ repsBench fx "spawnLocal (sequential)" 100 $ do+ us <- getSelfPid+ _ <- spawnLocal (send us ())+ expect :: Process (),+ oneBench fx "spawnLocal (pipelined)" $ do+ us <- getSelfPid+ replicateM_ 100 (spawnLocal (send us ()))+ replicateM_ 100 (expect :: Process ()),+ repsBench fx "callLocal" 100 $+ callLocal (return ()),+ repsBench fx "getSelfPid" 1000 $+ void getSelfPid,+ repsBench fx "getSelfNode" 1000 $+ void getSelfNode+ ]++monitoring :: Fixture -> Benchmark+monitoring fx =+ bgroup+ "monitoring"+ [ repsBench fx "monitor/unmonitor" 100 $+ monitor echo >>= unmonitor,+ repsBench fx "withMonitor_" 100 $+ withMonitor_ echo (return ()),+ repsBench fx "link/unlink" 100 $+ link echo >> unlink echo,+ repsBench fx "monitorNode/unmonitor" 100 $+ (getSelfNode >>= monitorNode) >>= unmonitor,+ repsBench fx "monitorPort/unmonitor" 100 $+ monitorPort (fst (fxChan fx)) >>= unmonitor,+ repsBench fx "notification (normal exit)" 100 $ do+ pid <- spawnLocal (expect :: Process ())+ ref <- monitor pid+ send pid ()+ awaitDown ref,+ repsBench fx "notification (kill)" 100 $ do+ pid <- spawnLocal (expect :: Process ())+ ref <- monitor pid+ kill pid "benchmark"+ awaitDown ref,+ repsBench fx "notification (die)" 100 $ do+ pid <- spawnLocal (die "benchmark")+ ref <- monitor pid+ awaitDown ref,+ repsBench fx "exit caught by catchExit" 100 $ do+ pid <-+ spawnLocal $+ catchExit (expect :: Process ()) (\_ (_ :: String) -> return ())+ ref <- monitor pid+ exit pid "benchmark"+ awaitDown ref,+ repsBench fx "exit caught by catchesExit" 100 $ do+ pid <-+ spawnLocal $+ catchesExit+ (expect :: Process ())+ [\_ m -> handleMessage m (\(_ :: String) -> return ())]+ ref <- monitor pid+ exit pid "benchmark"+ awaitDown ref+ ]+ where+ echo = fxEcho fx++registry :: Fixture -> Benchmark+registry fx =+ bgroup+ "registry"+ [ repsBench fx "whereis (hit)" 100 $+ void (whereis echoName),+ repsBench fx "whereis (miss)" 100 $+ void (whereis "benchmarks.absent"),+ repsBench fx "register/unregister" 100 $ do+ register "benchmarks.tmp" (fxEcho fx)+ unregister "benchmarks.tmp",+ repsBench fx "reregister" 100 $+ reregister echoName (fxEcho fx)+ ]++exceptions :: Fixture -> Benchmark+exceptions fx =+ bgroup+ "exceptions"+ [ repsBench fx "catch (not thrown)" 1000 $+ Catch.catch (return ()) (\(_ :: E.SomeException) -> return ()),+ repsBench fx "catch (thrown)" 1000 $+ Catch.catch (Catch.throwM Boom) (\Boom -> return ()),+ repsBench fx "try" 1000 $+ void (Catch.try (return ()) :: Process (Either E.SomeException ())),+ repsBench fx "catches (distributed-process Handler)" 1000 $+ catches+ (return ())+ [ Handler (\(_ :: E.ArithException) -> return ()),+ Handler (\(_ :: E.SomeException) -> return ())+ ],+ repsBench fx "bracket" 1000 $+ Catch.bracket (return ()) (\_ -> return ()) (\_ -> return ()),+ repsBench fx "finally" 1000 $+ Catch.finally (return ()) (return ()),+ repsBench fx "onException" 1000 $+ Catch.onException (return ()) (return ()),+ repsBench fx "mask_" 1000 $+ Catch.mask_ (return ())+ ]++introspection :: Fixture -> Benchmark+introspection fx =+ bgroup+ "introspection"+ [ repsBench fx "getProcessInfo" 100 $+ void (getProcessInfo (fxEcho fx)),+ repsBench fx "getLocalNodeStats" 100 $+ void getLocalNodeStats,+ repsBench fx "getNodeStats" 100 $+ void (getSelfNode >>= getNodeStats)+ ]++-- | 100 laps around each of the rings built by 'setup'.+ring :: Fixture -> Benchmark+ring fx =+ bgroup+ "ring"+ [ oneBench fx nm $ do+ replicateM_ 100 (send entry (Ping 0))+ replicateM_ 100 (void (expect :: Process Ping))+ | (nm, entry) <- fxRings fx+ ]++remote :: Fixture -> Benchmark+remote fx =+ bgroup+ "remote"+ [ repsBench fx "send/expect" 100 $ do+ us <- getSelfPid+ send echo us+ expect :: Process (),+ repsBench fx "usend/expect" 100 $ do+ us <- getSelfPid+ usend echo us+ expect :: Process (),+ repsBench fx "newChan + sendChan/receiveChan" 100 $ do+ (sp, rp) <- newChan+ send echo sp+ receiveChan rp,+ bgroup+ "throughput/bytestring"+ [ oneBench fx (show sz ++ "B") $+ sendThrough (fxRemoteCounter fx) 100 (BS.replicate sz 'x')+ | sz <- [8, 1024, 65536]+ ],+ repsBench fx "nsendRemote/expect" 100 $ do+ us <- getSelfPid+ nsendRemote nid echoName us+ expect :: Process (),+ repsBench fx "whereisRemoteAsync" 100 $ do+ whereisRemoteAsync nid echoName+ receiveWait+ [ matchIf+ (\(WhereIsReply n _) -> n == echoName)+ (\_ -> return ())+ ],+ repsBench fx "spawn" 100 $ do+ us <- getSelfPid+ _ <- spawn nid ($(mkClosure 'remoteSignal) us)+ expect :: Process (),+ repsBench fx "spawnMonitor + notification" 100 $ do+ us <- getSelfPid+ (_, ref) <- spawnMonitor nid ($(mkClosure 'remoteSignal) us)+ expect :: Process ()+ awaitDown ref,+ repsBench fx "spawnChannel" 100 $ do+ us <- getSelfPid+ sp <- spawnChannel sdictUnit nid ($(mkClosure 'remoteChanEcho) us)+ sendChan sp ()+ expect :: Process (),+ repsBench fx "call" 100 $+ void+ ( call+ $(functionTDict 'remoteAnswer)+ nid+ ($(mkClosure 'remoteAnswer) ())+ ),+ repsBench fx "getNodeStats" 100 $+ void (getNodeStats nid),+ repsBench fx "getProcessInfo" 100 $+ void (getProcessInfo echo)+ ]+ where+ echo = fxRemoteEcho fx+ nid = fxRemoteNodeId fx++-- | tasty-bench already repeats the body of the benchmark, but the benchmark+-- fixture adds a baseline amount of time which drowns some of the faster benchmarks.+--+-- Therefore, we amortize the fixture overhead by looping.+repsBench :: Fixture -> String -> Int -> Process () -> Benchmark+repsBench fx name reps act =+ bench (name ++ " (x" ++ show reps ++ ")") $+ whnfIO (fxRun fx (replicateM_ reps act))++oneBench :: Fixture -> String -> Process () -> Benchmark+oneBench fx name act = bench name $ whnfIO (fxRun fx act)++data Fixture = Fixture+ { fxRun :: Process () -> IO (),+ fxRemoteNodeId :: NodeId,+ fxEcho :: ProcessId,+ fxCounter :: ProcessId,+ fxRemoteEcho :: ProcessId,+ fxRemoteCounter :: ProcessId,+ fxChan :: (SendPort (), ReceivePort ()),+ fxQueue :: TQueue (),+ fxRelay :: ProcessId,+ fxDelegate :: ProcessId,+ fxProxy :: ProcessId,+ fxRings :: [(String, ProcessId)]+ }++echoName :: String+echoName = "benchmarks.echo"++setup :: LocalNode -> LocalNode -> IO Fixture+setup node1 node2 = do+ run <- newRunner node1+ queue <- newTQueueIO+ echo <- forkProcess node1 echoServer+ counter <- forkProcess node1 counterServer+ remoteEcho <- forkProcess node2 echoServer+ remoteCount <- forkProcess node2 counterServer+ -- 'register' acts on the caller's node.+ runProcess node1 (register echoName echo)+ runProcess node2 (register echoName remoteEcho)+ -- 'relay', 'delegate' and 'proxy' never return, so they cannot be spawned+ -- per iteration. They, the rings and the shared channel all have to be+ -- rooted at the runner, since that is the process each iteration runs on.+ var <- newEmptyMVar+ run $ do+ self <- getSelfPid+ chan <- newChan+ rly <- spawnLocal (relay self)+ dlg <- spawnLocal (delegate self (const True))+ prx <- spawnLocal (proxy self (\() -> return True))+ rings <-+ mapM+ (\(nm, mode) -> (,) nm <$> makeRing mode 10 self)+ [ ("send", RelaySend),+ ("unsafeSend", RelayUnsafeSend),+ ("forward", RelayForward)+ ]+ liftIO $ putMVar var (chan, rly, dlg, prx, rings)+ (chan, rly, dlg, prx, rings) <- takeMVar var+ return+ Fixture+ { fxRun = run,+ fxRemoteNodeId = localNodeId node2,+ fxEcho = echo,+ fxCounter = counter,+ fxRemoteEcho = remoteEcho,+ fxRemoteCounter = remoteCount,+ fxChan = chan,+ fxQueue = queue,+ fxRelay = rly,+ fxDelegate = dlg,+ fxProxy = prx,+ fxRings = rings+ }++-- | Runs actions on one long-lived process. Using 'runProcess' instead would+-- fold a 'forkProcess' into every measurement and give each iteration a fresh+-- 'ProcessId', defeating the connection caching real applications rely on.+newRunner :: LocalNode -> IO (Process () -> IO ())+newRunner node = do+ reqVar <- newEmptyMVar+ respVar <- newEmptyMVar+ _ <- forkProcess node $ forever $ do+ act <- liftIO (takeMVar reqVar)+ r <- Catch.try act+ drainMailbox+ liftIO $ putMVar respVar (r :: Either E.SomeException ())+ return $ \act -> do+ putMVar reqVar act+ takeMVar respVar >>= either E.throwIO return++-- | Keeps a benchmark from perturbing later ones through the runner's mailbox.+drainMailbox :: Process ()+drainMailbox = do+ r <- receiveTimeout 0 [matchAny (\_ -> return ())]+ case r of+ Nothing -> return ()+ Just () -> drainMailbox++awaitDown :: MonitorRef -> Process ()+awaitDown ref =+ receiveWait+ [ matchIf+ (\(ProcessMonitorNotification ref' _ _) -> ref' == ref)+ (\_ -> return ())+ ]++-- | Pipelined throughput: @n@ one-way sends, then one round trip to confirm+-- they all arrived.+sendThrough :: (Serializable a) => ProcessId -> Int -> a -> Process ()+sendThrough srv n payload = do+ us <- getSelfPid+ replicateM_ n (send srv payload)+ send srv (Report us)+ n' <- expect+ when (n' /= n) $+ die ("expected " ++ show n ++ " messages, server saw " ++ show n')++-- | The trailing 'matchAny' stops the mailbox growing if a benchmark sends+-- something unexpected; a growing mailbox is rescanned on every 'receiveWait'+-- and would skew every benchmark that follows.+echoServer :: Process ()+echoServer =+ forever $+ receiveWait+ [ match $ \(them :: ProcessId) -> send them (),+ match $ \(them, n :: Int) -> send them n,+ match $ \(them, bs :: BS.ByteString) -> send them bs,+ match $ \(sp :: SendPort ()) -> sendChan sp (),+ matchAny $ \_ -> return ()+ ]++-- | Counts one-way messages, and on 'Report' replies with the number seen+-- since the last report.+counterServer :: Process ()+counterServer = go 0+ where+ go :: Int -> Process ()+ go !n =+ receiveWait+ [ match $ \(Report them) -> send them n >> go 0,+ matchAny $ \_ -> go (n + 1)+ ]++data RelayMode = RelaySend | RelayUnsafeSend | RelayForward++relayLoop :: RelayMode -> ProcessId -> Process ()+relayLoop mode next = forever $ case mode of+ RelaySend -> expect >>= \m -> send next (m :: Ping)+ RelayUnsafeSend -> expect >>= \m -> unsafeSend next (m :: Ping)+ RelayForward -> receiveWait [matchAny (`forward` next)]++-- | Ring of @n@ relays whose last member relays to @target@; returns the entry+-- point.+makeRing :: RelayMode -> Int -> ProcessId -> Process ProcessId+makeRing mode n target+ | n <= 0 = return target+ | otherwise = makeRing mode (n - 1) =<< spawnLocal (relayLoop mode target)++newtype Ping = Ping Int+ deriving (Generic)++instance Binary Ping++newtype Report = Report ProcessId+ deriving (Generic)++instance Binary Report++data Boom = Boom+ deriving (Show)++instance E.Exception Boom
− benchmarks/ProcessRing.hs
@@ -1,117 +0,0 @@-{- ProcessRing benchmarks.--To run the benchmarks, select a value for the ring size (sz) and-the number of times to send a message around the ring---}--{-# LANGUAGE BangPatterns #-}-{-# LANGUAGE ScopedTypeVariables #-}--import Control.Monad-import Control.Distributed.Process hiding (catch)-import Control.Distributed.Process.Node-import Control.Exception (catch, SomeException)-import Network.Transport.TCP (createTransport, defaultTCPAddr, defaultTCPParameters)-import System.Environment-import System.Console.GetOpt--data Options = Options- { optRingSize :: Int- , optIterations :: Int- , optForward :: Bool- , optParallel :: Bool- , optUnsafe :: Bool- } deriving Show--initialProcess :: Options -> Process ()-initialProcess op =- let ringSz = optRingSize op- msgCnt = optIterations op- fwd = optForward op- unsafe = optUnsafe op- msg = ("foobar", "baz")- in do- self <- getSelfPid- ring <- makeRing fwd unsafe ringSz self- forM_ [1..msgCnt] (\_ -> send ring msg)- collect msgCnt- where relay fsend pid = do- msg <- expect :: Process (String, String)- fsend pid msg- relay fsend pid-- forward' pid =- receiveWait [ matchAny (\m -> forward m pid) ] >> forward' pid-- makeRing :: Bool -> Bool -> Int -> ProcessId -> Process ProcessId- makeRing !f !u !n !pid- | n == 0 = go f u pid- | otherwise = go f u pid >>= makeRing f u (n - 1)-- go :: Bool -> Bool -> ProcessId -> Process ProcessId- go False False next = spawnLocal $ relay send next- go False True next = spawnLocal $ relay unsafeSend next- go True _ next = spawnLocal $ forward' next-- collect :: Int -> Process ()- collect !n- | n == 0 = return ()- | otherwise = do- receiveWait [- matchIf (\(a, b) -> a == "foobar" && b == "baz")- (\_ -> return ())- , matchAny (\_ -> error "unexpected input!")- ]- collect (n - 1)--defaultOptions :: Options-defaultOptions = Options- { optRingSize = 10- , optIterations = 100- , optForward = False- , optParallel = False- , optUnsafe = False- }--options :: [OptDescr (Options -> Options)]-options =- [ Option ['s'] ["ring-size"] (OptArg optSz "SIZE") "# of processes in ring"- , Option ['i'] ["iterations"] (OptArg optMsgCnt "ITER") "# of times to send"- , Option ['f'] ["forward"]- (NoArg (\opts -> opts { optForward = True }))- "use `forward' instead of send - default = False"- , Option ['u'] ["unsafe-send"]- (NoArg (\opts -> opts { optUnsafe = True }))- "use 'unsafeSend' (ignored with -f) - default = False"- , Option ['p'] ["parallel"]- (NoArg (\opts -> opts { optParallel = True }))- "send in parallel and consume sequentially - default = False"- ]--optMsgCnt :: Maybe String -> Options -> Options-optMsgCnt Nothing opts = opts-optMsgCnt (Just c) opts = opts { optIterations = ((read c) :: Int) }--optSz :: Maybe String -> Options -> Options-optSz Nothing opts = opts-optSz (Just s) opts = opts { optRingSize = ((read s) :: Int) }--parseArgv :: [String] -> IO (Options, [String])-parseArgv argv = do- pn <- getProgName- case getOpt Permute options argv of- (o,n,[] ) -> return (foldl (flip id) defaultOptions o, n)- (_,_,errs) -> ioError (userError (concat errs ++ usageInfo (header pn) options))- where header pn' = "Usage: " ++ pn' ++ " [OPTION...]"--main :: IO ()-main = do- argv <- getArgs- (opt, _) <- parseArgv argv- putStrLn $ "options: " ++ (show opt)- Right transport <- createTransport- (defaultTCPAddr "127.0.0.1" "8090" ) defaultTCPParameters- node <- newLocalNode transport initRemoteTable- catch (void $ runProcess node $ initialProcess opt)- (\(e :: SomeException) -> putStrLn $ "ERROR: " ++ (show e))
− benchmarks/Spawns.hs
@@ -1,49 +0,0 @@-{-# LANGUAGE BangPatterns #-}---- | Like Throughput, but send every ping from a different process--- (i.e., require a lightweight connection per ping)-import System.Environment-import Control.Monad-import Control.Applicative-import Control.Distributed.Process-import Control.Distributed.Process.Node-import Network.Transport.TCP (createTransport, defaultTCPAddr, defaultTCPParameters)-import Data.Binary (encode, decode)-import qualified Data.ByteString.Lazy as BSL--counter :: Process ()-counter = go 0- where- go :: Int -> Process ()- go !n = do- b <- expect- case b of- Nothing -> go (n + 1)- Just them -> send them n >> go 0--count :: Int -> ProcessId -> Process ()-count n them = do- us <- getSelfPid- replicateM_ n . spawnLocal $ send them (Nothing :: Maybe ProcessId)- send them (Just us)- n' <- expect- liftIO $ print (n == n')--initialProcess :: String -> Process ()-initialProcess "SERVER" = do- us <- getSelfPid- liftIO $ BSL.writeFile "counter.pid" (encode us)- counter-initialProcess "CLIENT" = do- n <- liftIO $ getLine- them <- liftIO $ decode <$> BSL.readFile "counter.pid"- count (read n) them--main :: IO ()-main = do- [role, host, port] <- getArgs- trans <- createTransport (defaultTCPAddr host port) defaultTCPParameters- case trans of- Right transport -> do node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role- Left other -> error $ show other
− benchmarks/Throughput.hs
@@ -1,74 +0,0 @@-{-# LANGUAGE BangPatterns #-}-{-# LANGUAGE DeriveDataTypeable #-}--import System.Environment-import Control.Monad-import Control.Applicative-import Control.Distributed.Process-import Control.Distributed.Process.Node-import Network.Transport.TCP (createTransport, defaultTCPParameters, defaultTCPAddr)-import Data.Binary-import qualified Data.ByteString.Lazy as BSL-import Data.Typeable--data SizedList a = SizedList { size :: Int , elems :: [a] }- deriving (Typeable)--instance Binary a => Binary (SizedList a) where- put (SizedList sz xs) = put sz >> mapM_ put xs- get = do- sz <- get- xs <- getMany sz- return (SizedList sz xs)---- Copied from Data.Binary-getMany :: Binary a => Int -> Get [a]-getMany = go []- where- go xs 0 = return $! reverse xs- go xs i = do x <- get- x `seq` go (x:xs) (i-1)-{-# INLINE getMany #-}--nats :: Int -> SizedList Int-nats = \n -> SizedList n (aux n)- where- aux 0 = []- aux n = n : aux (n - 1)--counter :: Process ()-counter = go 0- where- go :: Int -> Process ()- go !n =- receiveWait- [ match $ \xs -> go (n + size (xs :: SizedList Int))- , match $ \them -> send them n >> go 0- ]--count :: (Int, Int) -> ProcessId -> Process ()-count (packets, sz) them = do- us <- getSelfPid- replicateM_ packets $ send them (nats sz)- send them us- n' <- expect- liftIO $ print (packets * sz, n' == packets * sz)--initialProcess :: String -> Process ()-initialProcess "SERVER" = do- us <- getSelfPid- liftIO $ BSL.writeFile "counter.pid" (encode us)- counter-initialProcess "CLIENT" = do- n <- liftIO getLine- them <- liftIO $ decode <$> BSL.readFile "counter.pid"- count (read n) them--main :: IO ()-main = do- [role, host, port] <- getArgs- trans <- createTransport (defaultTCPAddr host port) defaultTCPParameters- case trans of- Right transport -> do node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role- Left other -> error $ show other
distributed-process.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.0 Name: distributed-process-Version: 0.7.8+Version: 0.7.9 Build-Type: Simple License: BSD-3-Clause License-File: LICENSE@@ -21,7 +21,7 @@ You will probably also want to install a Cloud Haskell backend such as distributed-process-simplelocalnet.-tested-with: GHC==8.10.7 GHC==9.0.2 GHC==9.2.8 GHC==9.4.5 GHC==9.6.4 GHC==9.8.2 GHC==9.10.1 GHC==9.12.1+tested-with: GHC==8.10.7 GHC==9.0.2 GHC==9.2.8 GHC==9.4.8 GHC==9.6.7 GHC==9.8.4 GHC==9.10.3 GHC==9.12.2 GHC==9.14.1 Category: Control extra-doc-files: ChangeLog @@ -60,7 +60,7 @@ rank1dynamic >= 0.1 && < 0.5, syb >= 0.3 && < 0.8, exceptions >= 0.10,- containers >= 0.6 && < 0.8,+ containers >= 0.6 && < 0.9, deepseq >= 1.4 && < 1.7, time >= 1.9 Exposed-modules: Control.Distributed.Process@@ -113,64 +113,25 @@ UndecidableInstances if flag(th) other-extensions: TemplateHaskell- Build-Depends: template-haskell >= 2.6 && <2.24+ Build-Depends: template-haskell >= 2.6 && <2.25 Exposed-modules: Control.Distributed.Process.Internal.Closure.TH CPP-Options: -DTemplateHaskellSupport -- Tests are in distributed-process-test package, for convenience. -benchmark distributed-process-throughput+benchmark distributed-process-benchmarks import: warnings Type: exitcode-stdio-1.0+ Main-Is: Main.hs+ HS-Source-Dirs: benchmarks Build-Depends: base >= 4.14 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && <= 0.9,+ binary >= 0.8 && < 0.10, bytestring >= 0.10 && < 0.13,- binary >= 0.8 && < 0.10- Main-Is: benchmarks/Throughput.hs- default-language: Haskell2010--benchmark distributed-process-latency- import: warnings- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.14 && < 5, distributed-process,- network-transport-tcp >= 0.3 && <= 0.9,- bytestring >= 0.10 && < 0.13,- binary >= 0.8 && < 0.10- Main-Is: benchmarks/Latency.hs- default-language: Haskell2010--benchmark distributed-process-channels- import: warnings- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.14 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && <= 0.9,- bytestring >= 0.10 && < 0.13,- binary >= 0.8 && < 0.10- Main-Is: benchmarks/Channels.hs- default-language: Haskell2010--benchmark distributed-process-spawns- import: warnings- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.14 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && <= 0.9,- bytestring >= 0.10 && < 0.13,- binary >= 0.8 && < 0.10- Main-Is: benchmarks/Spawns.hs- default-language: Haskell2010--benchmark distributed-process-ring- import: warnings- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.14 && < 5,- distributed-process,+ exceptions >= 0.10,+ network-transport >= 0.4.1.0 && < 0.6, network-transport-tcp >= 0.3 && <= 0.9,- bytestring >= 0.10 && < 0.13,- binary >= 0.8 && < 0.10- Main-Is: benchmarks/ProcessRing.hs+ stm >= 2.4 && < 2.6,+ tasty-bench >= 0.3.4 && < 0.6 default-language: Haskell2010 ghc-options: -threaded -O2 -rtsopts
src/Control/Distributed/Process/Internal/CQueue.hs view
@@ -43,7 +43,7 @@ ( StrictList(..) , append )-import Data.Maybe (fromJust)+import Control.Monad (join) import GHC.MVar (MVar(MVar)) import GHC.IO (IO(IO), unIO) import GHC.Exts (mkWeak#)@@ -126,16 +126,21 @@ -> [MatchOn m a] -- ^ List of matches -> IO (Maybe a) -- ^ 'Nothing' only on timeout dequeue (CQueue arrived incoming size) blockSpec matchons = mask_ $ decrementJust $- case blockSpec of- Timeout n -> timeout n $ fmap fromJust run- _other ->- case chunks of- [Right ports] -> -- channels only, this is easy:- case blockSpec of- NonBlocking -> atomically $ waitChans ports (return Nothing)- _ -> atomically $ waitChans ports retry- -- no onException needed- _other -> run+ case chunks of+ [Right ports] -> -- channels only, this is easy:+ case blockSpec of+ NonBlocking -> atomically $ waitChans ports (return Nothing)+ Blocking -> atomically $ waitChans ports retry+ -- no onException needed+ Timeout n -> do+ -- Arming the timer is not cheap, and get can get+ -- much higher throughput in cases where the mailbox+ -- is not empty by first checking if we even need a timeout+ r <- atomically $ waitChans ports (return Nothing)+ case r of+ Just _ -> return r+ Nothing -> join <$> timeout n (atomically $ waitChans ports retry)+ _other -> run where -- Decrement counter is smth is returned from the queue, -- this is safe to use as method is called under a mask@@ -155,7 +160,13 @@ Nothing -> return xs Just x -> grabNew (Snoc xs x) arr' <- grabNew arr- goCheck chunks arr'+ checked <- goCheck chunks arr'+ case checked of+ Left r -> return r+ Right old -> case blockSpec of+ NonBlocking -> returnOld old Nothing+ Blocking -> goWait old+ Timeout n -> join <$> timeout n (goWait old) -- Yields the value of the first succesful STM transaction as -- @Just (Left v)@. If all transactions fail, yields the value of the second@@ -169,20 +180,20 @@ -- mailbox. For channel matches, we do a non-blocking check at -- this point. --- -- Yields @Just (Left a)@ when a channel is matched, @Just (Right a)@- -- when a message is matched and @Nothing@ when there are no messages and we- -- aren't blocking.- --+ -- Yields @Left (Just (Left a))@ when a channel is matched and+ -- @Left (Just (Right a))@ when a message is matched. When nothing+ -- matched it yields @Right old@: the messages to hold on to, for the+ -- caller to decide whether to wait for more. goCheck :: MatchChunks m a -> StrictList m -- messages to check, in this order- -> IO (Maybe (Either a a))+ -> IO (Either (Maybe (Either a a)) (StrictList m)) - goCheck [] old = goWait old+ goCheck [] old = return (Right old) goCheck (Right ports : rest) old = do r <- atomically $ waitChans ports (return Nothing) -- does not block case r of- Just _ -> returnOld old r+ Just _ -> Left <$> returnOld old r Nothing -> goCheck rest old goCheck (Left matches : rest) old = do@@ -192,7 +203,7 @@ -- of passing around restore and setting up exception handlers is -- high. So just don't use expensive matchIfs! case checkArrived matches old of- (old', Just r) -> returnOld old' (Just (Right r))+ (old', Just r) -> Left <$> returnOld old' (Just (Right r)) (old', Nothing) -> goCheck rest old' -- use the result list, which is now left-biased @@ -207,12 +218,8 @@ mkSTM (Right ports : rest) = foldr orElse (mkSTM rest) (map (fmap Right) ports) - waitIncoming :: IO (Maybe (Either m a))- waitIncoming = case blockSpec of- NonBlocking -> atomically $ fmap Just stm `orElse` return Nothing- _ -> atomically $ fmap Just stm- where- stm = mkSTM chunks+ waitIncoming :: IO (Either m a)+ waitIncoming = atomically (mkSTM chunks) -- -- The initial pass didn't find a message, so now we go into blocking@@ -223,23 +230,20 @@ -- goWait :: StrictList m -> IO (Maybe (Either a a)) goWait old = do- r <- waitIncoming `onException` putMVar arrived old- case r of- -- Nothing => non-blocking and no message- Nothing -> returnOld old Nothing- Just e -> case e of- --- -- Left => message arrived in the process mailbox. We now have to- -- run through the MatchChunks checking each one, because we might- -- have a situation where the first chunk fails to match and the- -- second chunk is a channel match and there *is* a message in the- -- channel. In that case the channel wins.- --- Left m -> goCheck1 chunks m old- --- -- Right => message arrived on a channel first- --- Right a -> returnOld old (Just (Left a))+ e <- waitIncoming `onException` putMVar arrived old+ case e of+ --+ -- Left => message arrived in the process mailbox. We now have to+ -- run through the MatchChunks checking each one, because we might+ -- have a situation where the first chunk fails to match and the+ -- second chunk is a channel match and there *is* a message in the+ -- channel. In that case the channel wins.+ --+ Left m -> goCheck1 chunks m old+ --+ -- Right => message arrived on a channel first+ --+ Right a -> returnOld old (Just (Left a)) -- -- A message arrived in the process inbox; check the MatchChunks for
src/Control/Distributed/Process/Internal/Primitives.hs view
@@ -363,8 +363,13 @@ receiveChanTimeout :: Serializable a => Int -> ReceivePort a -> Process (Maybe a) receiveChanTimeout 0 ch = liftIO . atomically $ (Just <$> receiveSTM ch) `orElse` return Nothing-receiveChanTimeout n ch = liftIO . timeout n . atomically $- receiveSTM ch+receiveChanTimeout n ch = liftIO $ do+ -- Checking if the mailbox has a message /before/ arming,+ -- because arming a timeout can be expensive+ r <- atomically $ (Just <$> receiveSTM ch) `orElse` return Nothing+ case r of+ Just _ -> return r+ Nothing -> timeout n . atomically $ receiveSTM ch -- | Merge a list of typed channels. --