distributed-process 0.6.6 → 0.7.9
raw patch · 26 files changed
Files
- ChangeLog +58/−3
- Setup.hs +0/−2
- benchmarks/Channels.hs +0/−41
- benchmarks/Latency.hs +0/−36
- benchmarks/Main.hs +753/−0
- benchmarks/ProcessRing.hs +0/−116
- benchmarks/Spawns.hs +0/−47
- benchmarks/Throughput.hs +0/−72
- distributed-process.cabal +98/−116
- src/Control/Distributed/Process.hs +2/−4
- src/Control/Distributed/Process/Internal/CQueue.hs +48/−50
- src/Control/Distributed/Process/Internal/Closure/Explicit.hs +3/−1
- src/Control/Distributed/Process/Internal/Closure/TH.hs +27/−1
- src/Control/Distributed/Process/Internal/Primitives.hs +146/−73
- src/Control/Distributed/Process/Internal/StrictMVar.hs +0/−19
- src/Control/Distributed/Process/Internal/Types.hs +9/−5
- src/Control/Distributed/Process/Internal/WeakTQueue.hs +0/−4
- src/Control/Distributed/Process/Management.hs +38/−81
- src/Control/Distributed/Process/Management/Internal/Table.hs +0/−222
- src/Control/Distributed/Process/Management/Internal/Trace/Primitives.hs +8/−2
- src/Control/Distributed/Process/Management/Internal/Trace/Tracer.hs +0/−8
- src/Control/Distributed/Process/Management/Internal/Trace/Types.hs +0/−10
- src/Control/Distributed/Process/Management/Internal/Types.hs +9/−19
- src/Control/Distributed/Process/Node.hs +37/−29
- src/Control/Distributed/Process/Serializable.hs +5/−15
- src/Control/Distributed/Process/UnsafePrimitives.hs +83/−33
ChangeLog view
@@ -1,3 +1,58 @@+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)+* Added upper bound on `template-haskell` to prevent future breakage.+* Addressed some compilation warnings (#467)++2024-09-03 Laurent P. René de Cotret <laurent.decotret@outlook.com> 0.7.7++* Bumped dependency bounds to support GHC 8.10.7 - GHC 9.10.1+* Updated links to point to Distributed Haskell monorepo++2024-04-03 David Simmons-Duffin <dsd@caltech.edu> 0.7.6++* Bumped hashable upper bound.++2018-06-12 David Simmons-Duffin <dsd@caltech.edu> 0.7.5++* Bumped dependencies to build with ghc-9.8+* Turn Serializable into a type synonym+* Remove tests for ghc-8.8.* and earlier++2018-06-12 Facundo Domínguez <facundo.dominguez@tweag.io> 0.7.4++* Added support for exceptions >= 0.10++2017-08-31 Facundo Domínguez <facundo.dominguez@tweag.io> 0.7.3++* Drop support for ghc-7.8.* and earlier.++2017-08-31 Facundo Domínguez <facundo.dominguez@tweag.io> 0.7.2++* Fixed build errors with ghc-8.2.1.++2017-08-22 Facundo Domínguez <facundo.dominguez@tweag.io> 0.7.1++* Relax upper bounds in dependencies to build with ghc-8.2.1.++2017-08-21 Facundo Domínguez <facundo.dominguez@tweag.io> 0.7.0++* Change type of message sent by `say` from a 3-tuple to a proper+type (`SayMessage`) with a proper `UTCTime`. (#291)+* Expose the MonitorRef in the withMonitor call.+* Have unmonitor remove the monitor message in the inbox. (#268)+* Remove Mx Data Tables. This API isn't used, is easy to replace with various+other packages. (#276)+* Add Ord instance to SpawnRef.+ 2016-10-13 Facundo Domínguez <facundo.dominguez@tweag.io> 0.6.6 * Remove monitors from remote nodes when a process dies. (#295)@@ -143,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 @@ -193,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
− Setup.hs
@@ -1,2 +0,0 @@-import Distribution.Simple-main = defaultMain
− benchmarks/Channels.hs
@@ -1,41 +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, 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- Right transport <- createTransport host port defaultTCPParameters- node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role
− benchmarks/Latency.hs
@@ -1,36 +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, 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- Right transport <- createTransport host port defaultTCPParameters- node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role
+ 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,116 +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, 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 "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,47 +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, 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- Right transport <- createTransport host port defaultTCPParameters- node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role
− benchmarks/Throughput.hs
@@ -1,72 +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)-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- Right transport <- createTransport host port defaultTCPParameters- node <- newLocalNode transport initRemoteTable- runProcess node $ initialProcess role
distributed-process.cabal view
@@ -1,155 +1,137 @@+cabal-version: 3.0 Name: distributed-process-Version: 0.6.6-Cabal-Version: >=1.8+Version: 0.7.9 Build-Type: Simple-License: BSD3+License: BSD-3-Clause License-File: LICENSE Copyright: Well-Typed LLP, Tweag I/O Limited Author: Duncan Coutts, Nicolas Wu, Edsko de Vries-Maintainer: Facundo Domínguez <facundo.dominguez@tweag.io>+maintainer: The Distributed Haskell team Stability: experimental-Homepage: http://haskell-distributed.github.com/+Homepage: https://haskell-distributed.github.io/ Bug-Reports: https://github.com/haskell-distributed/distributed-process/issues Synopsis: Cloud Haskell: Erlang-style concurrency in Haskell Description: This is an implementation of Cloud Haskell, as described in /Towards Haskell in the Cloud/ by Jeff Epstein, Andrew Black, and Simon Peyton Jones- (<http://research.microsoft.com/en-us/um/people/simonpj/papers/parallel/>),+ (<https://simon.peytonjones.org/haskell-cloud/>), although some of the details are different. The precise message passing semantics are based on /A unified semantics for future Erlang/ by Hans Svensson, Lars-Åke Fredlund and Clara Benac Earle. You will probably also want to install a Cloud Haskell backend such as distributed-process-simplelocalnet.-Tested-With: GHC==7.2.2 GHC==7.4.1 GHC==7.4.2 GHC==7.6.2+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-source-files: ChangeLog+extra-doc-files: ChangeLog -Source-Repository head+common warnings+ ghc-options: -Wall+ -Wcompat+ -Widentities+ -Wincomplete-uni-patterns+ -Wincomplete-record-updates+ -Wredundant-constraints+ -fhide-source-paths+ -Wpartial-fields+ -Wunused-packages++source-repository head Type: git Location: https://github.com/haskell-distributed/distributed-process- SubDir: distributed-process+ SubDir: packages/distributed-process flag th description: Build with Template Haskell support default: True -flag old-locale- description: If false then depend on time >= 1.5.- .- If true then depend on time < 1.5 together with old-locale.- default: False- Library- Build-Depends: base >= 4.4 && < 5,- binary >= 0.6.3 && < 0.9,- hashable >= 1.2.0.5 && < 1.3,- network-transport >= 0.4.1.0 && < 0.5,- stm >= 2.4 && < 2.5,- transformers >= 0.2 && < 0.6,+ import: warnings+ Build-Depends: base >= 4.14 && < 5,+ binary >= 0.8 && < 0.10,+ hashable >= 1.2.0.5 && < 1.6,+ network-transport >= 0.4.1.0 && < 0.6,+ stm >= 2.4 && < 2.6, mtl >= 2.0 && < 2.4, data-accessor >= 0.2 && < 0.3,- bytestring >= 0.9 && < 0.11,- random >= 1.0 && < 1.2,+ bytestring >= 0.10 && < 0.13,+ random >= 1.0 && < 1.4, distributed-static >= 0.2 && < 0.4,- rank1dynamic >= 0.1 && < 0.4,- syb >= 0.3 && < 0.7,- exceptions >= 0.5- Exposed-modules: Control.Distributed.Process,- Control.Distributed.Process.Closure,- Control.Distributed.Process.Debug,- Control.Distributed.Process.Internal.BiMultiMap,- Control.Distributed.Process.Internal.Closure.BuiltIn,- Control.Distributed.Process.Internal.Closure.Explicit,- Control.Distributed.Process.Internal.CQueue,- Control.Distributed.Process.Internal.Messaging,- Control.Distributed.Process.Internal.Primitives,- Control.Distributed.Process.Internal.Spawn,- Control.Distributed.Process.Internal.StrictContainerAccessors,- Control.Distributed.Process.Internal.StrictList,- Control.Distributed.Process.Internal.StrictMVar,- Control.Distributed.Process.Internal.Types,- Control.Distributed.Process.Internal.WeakTQueue,- Control.Distributed.Process.Management,- Control.Distributed.Process.Node,- Control.Distributed.Process.Serializable,+ rank1dynamic >= 0.1 && < 0.5,+ syb >= 0.3 && < 0.8,+ exceptions >= 0.10,+ containers >= 0.6 && < 0.9,+ deepseq >= 1.4 && < 1.7,+ time >= 1.9+ Exposed-modules: Control.Distributed.Process+ Control.Distributed.Process.Closure+ Control.Distributed.Process.Debug+ Control.Distributed.Process.Internal.BiMultiMap+ Control.Distributed.Process.Internal.Closure.BuiltIn+ Control.Distributed.Process.Internal.Closure.Explicit+ Control.Distributed.Process.Internal.CQueue+ Control.Distributed.Process.Internal.Messaging+ Control.Distributed.Process.Internal.Primitives+ Control.Distributed.Process.Internal.Spawn+ Control.Distributed.Process.Internal.StrictContainerAccessors+ Control.Distributed.Process.Internal.StrictList+ Control.Distributed.Process.Internal.StrictMVar+ Control.Distributed.Process.Internal.Types+ Control.Distributed.Process.Internal.WeakTQueue+ Control.Distributed.Process.Management+ Control.Distributed.Process.Node+ Control.Distributed.Process.Serializable Control.Distributed.Process.UnsafePrimitives- Control.Distributed.Process.Management.Internal.Agent,- Control.Distributed.Process.Management.Internal.Bus,- Control.Distributed.Process.Management.Internal.Table,- Control.Distributed.Process.Management.Internal.Types,- Control.Distributed.Process.Management.Internal.Trace.Primitives,- Control.Distributed.Process.Management.Internal.Trace.Remote,- Control.Distributed.Process.Management.Internal.Trace.Types,+ Control.Distributed.Process.Management.Internal.Agent+ Control.Distributed.Process.Management.Internal.Bus+ Control.Distributed.Process.Management.Internal.Types+ Control.Distributed.Process.Management.Internal.Trace.Primitives+ Control.Distributed.Process.Management.Internal.Trace.Remote+ Control.Distributed.Process.Management.Internal.Trace.Types Control.Distributed.Process.Management.Internal.Trace.Tracer- ghc-options: -Wall+ default-language: Haskell2010 HS-Source-Dirs: src- if impl(ghc <= 7.4.2)- Build-Depends: containers >= 0.4 && < 0.5,- deepseq == 1.3.0.0- else- Build-Depends: containers >= 0.4 && < 0.6,- deepseq >= 1.3.0.1 && < 1.6- if flag(old-locale)- Build-Depends: time < 1.5, old-locale >= 1.0 && <1.1- else- Build-Depends: time >= 1.5+ other-extensions: BangPatterns+ CPP+ DeriveDataTypeable+ DeriveFunctor+ DeriveGeneric+ ExistentialQuantification+ FlexibleInstances+ GADTs+ GeneralizedNewtypeDeriving+ KindSignatures+ MagicHash+ PatternGuards+ RankNTypes+ RecordWildCards+ ScopedTypeVariables+ StandaloneDeriving+ TypeFamilies+ TypeSynonymInstances+ UnboxedTuples+ UndecidableInstances if flag(th)- if impl(ghc <= 7.4.2)- Build-Depends: template-haskell >= 2.7 && < 2.8- else- Build-Depends: template-haskell >= 2.6 && < 2.12+ other-extensions: TemplateHaskell+ 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- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.4 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && < 0.6,- bytestring >= 0.9 && < 0.11,- binary >= 0.6.3 && < 0.9- Main-Is: benchmarks/Throughput.hs- ghc-options: -Wall--benchmark distributed-process-latency- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.4 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && < 0.6,- bytestring >= 0.9 && < 0.11,- binary >= 0.6.3 && < 0.9- Main-Is: benchmarks/Latency.hs- ghc-options: -Wall--benchmark distributed-process-channels- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.4 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && < 0.6,- bytestring >= 0.9 && < 0.11,- binary >= 0.6.3 && < 0.9- Main-Is: benchmarks/Channels.hs- ghc-options: -Wall--benchmark distributed-process-spawns- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.4 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && < 0.6,- bytestring >= 0.9 && < 0.11,- binary >= 0.6.3 && < 0.9- Main-Is: benchmarks/Spawns.hs- ghc-options: -Wall--benchmark distributed-process-ring- Type: exitcode-stdio-1.0- Build-Depends: base >= 4.4 && < 5,- distributed-process,- network-transport-tcp >= 0.3 && < 0.6,- bytestring >= 0.9 && < 0.11,- binary >= 0.6.3 && < 0.9- Main-Is: benchmarks/ProcessRing.hs- ghc-options: -Wall -threaded -O2 -rtsopts+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,+ binary >= 0.8 && < 0.10,+ bytestring >= 0.10 && < 0.13,+ distributed-process,+ exceptions >= 0.10,+ network-transport >= 0.4.1.0 && < 0.6,+ network-transport-tcp >= 0.3 && <= 0.9,+ stm >= 2.4 && < 2.6,+ tasty-bench >= 0.3.4 && < 0.6+ default-language: Haskell2010+ ghc-options: -threaded -O2 -rtsopts
src/Control/Distributed/Process.hs view
@@ -102,6 +102,7 @@ , monitorPort , unmonitor , withMonitor+ , withMonitor_ , MonitorRef -- opaque , ProcessLinkException(..) , NodeLinkException(..)@@ -269,6 +270,7 @@ , monitorPort , unmonitor , withMonitor+ , withMonitor_ -- Logging , say -- Registry@@ -320,11 +322,7 @@ ) import qualified Control.Monad.Catch as Catch -#if MIN_VERSION_base(4,6,0) import Prelude-#else-import Prelude hiding (catch)-#endif import qualified Control.Exception as Exception (onException) import Data.Accessor ((^.)) import Data.Foldable (forM_)
src/Control/Distributed/Process/Internal/CQueue.hs view
@@ -31,7 +31,6 @@ , orElse , retry )-import Control.Applicative ((<$>), (<*>)) import Control.Exception (mask_, onException) import System.Timeout (timeout) import Control.Distributed.Process.Internal.StrictMVar@@ -44,8 +43,7 @@ ( StrictList(..) , append )-import Data.Maybe (fromJust)-import Data.Traversable (traverse)+import Control.Monad (join) import GHC.MVar (MVar(MVar)) import GHC.IO (IO(IO), unIO) import GHC.Exts (mkWeak#)@@ -74,7 +72,7 @@ data BlockSpec = NonBlocking | Blocking- | Timeout Int+ | Timeout Int -- ^ Timeout in microseconds -- Match operations --@@ -128,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@@ -157,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@@ -171,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@@ -194,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 @@ -209,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@@ -225,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@@ -290,11 +292,7 @@ -- | Weak reference to a CQueue mkWeakCQueue :: CQueue a -> IO () -> IO (Weak (CQueue a)) mkWeakCQueue m@(CQueue (StrictMVar (MVar m#)) _ _) f = IO $ \s ->-#if MIN_VERSION_base(4,9,0) case mkWeak# m# m (unIO f) s of (# s1, w #) -> (# s1, Weak w #)-#else- case mkWeak# m# m f s of (# s1, w #) -> (# s1, Weak w #)-#endif queueSize :: CQueue a -> IO Int queueSize (CQueue _ _ size) = readTVarIO size
src/Control/Distributed/Process/Internal/Closure/Explicit.hs view
@@ -7,6 +7,7 @@ , KindSignatures , GADTs , EmptyDataDecls+ , TypeOperators , DeriveDataTypeable #-} module Control.Distributed.Process.Internal.Closure.Explicit (@@ -29,6 +30,7 @@ import Data.Rank1Typeable import Data.Binary(encode,put,get,Binary) import qualified Data.ByteString.Lazy as B+import Data.Kind (Type) -- | A RemoteRegister is a trasformer on a RemoteTable to register additional static values. type RemoteRegister = RemoteTable -> RemoteTable@@ -118,7 +120,7 @@ -- This generic uncurry courtesy Andrea Vezzosi data HTrue data HFalse-data Fun :: * -> * -> * -> * where+data Fun :: Type -> Type -> Type -> Type where Done :: Fun EndOfTuple r r Moar :: Fun xs f r -> Fun (x,xs) (x -> f) r
src/Control/Distributed/Process/Internal/Closure/TH.hs view
@@ -12,7 +12,6 @@ ) where import Prelude hiding (succ, any)-import Control.Applicative ((<$>)) import Language.Haskell.TH ( -- Q monad and operations Q@@ -28,6 +27,9 @@ , Exp , Type(AppT, ForallT, VarT, ArrowT) , Info(VarI)+#if MIN_VERSION_template_haskell(2,17,0)+ , Specificity+#endif , TyVarBndr(PlainTV, KindedTV) , Pred #if MIN_VERSION_template_haskell(2,10,0)@@ -203,7 +205,11 @@ , concat [register, registerSDict, registerTDict] ) where+#if MIN_VERSION_template_haskell(2,17,0)+ makeStatic :: [TyVarBndr Specificity] -> Type -> Q ([Dec], [Q Exp])+#else makeStatic :: [TyVarBndr] -> Type -> Q ([Dec], [Q Exp])+#endif makeStatic typVars typ = do static <- generateStatic origName typVars typ let dyn = case typVars of@@ -222,7 +228,12 @@ ) -- | Turn a polymorphic type into a monomorphic type using ANY and co+#if MIN_VERSION_template_haskell(2,17,0)+monomorphize :: [TyVarBndr Specificity] -> Type -> Q Type+#else monomorphize :: [TyVarBndr] -> Type -> Q Type+#endif+ monomorphize tvs = let subst = zip (map tyVarBndrName tvs) anys in everywhereM (mkM (applySubst subst))@@ -247,7 +258,11 @@ applySubst s t = gmapM (mkM (applySubst s)) t -- | Generate a static value+#if MIN_VERSION_template_haskell(2,17,0)+generateStatic :: Name -> [TyVarBndr Specificity] -> Type -> Q [Dec]+#else generateStatic :: Name -> [TyVarBndr] -> Type -> Q [Dec]+#endif generateStatic n xs typ = do staticTyp <- [t| Static |] sequence@@ -259,7 +274,11 @@ , sfnD (staticName n) [| staticLabel $(showFQN n) |] ] where+#if MIN_VERSION_template_haskell(2,17,0)+ typeable :: TyVarBndr Specificity -> Q Pred+#else typeable :: TyVarBndr -> Q Pred+#endif typeable tv = #if MIN_VERSION_template_haskell(2,10,0) conT (mkName "Typeable") `appT` varT (tyVarBndrName tv)@@ -315,9 +334,16 @@ sfnD n e = funD n [clause [] (normalB e) []] -- | The name of a type variable binding occurrence+#if MIN_VERSION_template_haskell(2,17,0)+tyVarBndrName :: TyVarBndr Specificity -> Name+tyVarBndrName (PlainTV n _) = n+tyVarBndrName (KindedTV n _ _) = n+#else tyVarBndrName :: TyVarBndr -> Name tyVarBndrName (PlainTV n) = n tyVarBndrName (KindedTV n _) = n+#endif+ -- | Fully qualified name; that is, the name and the _current_ module --
src/Control/Distributed/Process/Internal/Primitives.hs view
@@ -76,8 +76,11 @@ , unlink , monitor , unmonitor+ , unmonitorAsync , withMonitor+ , withMonitor_ -- * Logging+ , SayMessage(..) , say -- * Registry , register@@ -121,20 +124,13 @@ , sendCtrlMsg ) where -#if ! MIN_VERSION_base(4,6,0)-import Prelude hiding (catch)-#endif--import Data.Binary (decode)-import Data.Time.Clock (getCurrentTime)+import Data.Binary (Binary(..), Put, Get, decode)+import Data.Time.Clock (getCurrentTime, UTCTime(..))+import Data.Time.Calendar (Day(..)) import Data.Time.Format (formatTime)-#if MIN_VERSION_time(1,5,0) import Data.Time.Format (defaultTimeLocale)-#else-import System.Locale (defaultTimeLocale)-#endif import System.Timeout (timeout)-import Control.Monad (when)+import Control.Monad (when, void) import Control.Monad.Reader (ask) import Control.Monad.IO.Class (liftIO) import Control.Monad.Catch@@ -184,6 +180,7 @@ , SpawnRef(..) , ProcessSignal(..) , NodeMonitorNotification(..)+ , ProcessMonitorNotification(..) , monitorCounter , spawnCounter , SendPort(..)@@ -248,18 +245,16 @@ let us = processId proc node = processNode proc nodeId = localNodeId node- destNode = (processNodeId them) in do- case destNode == nodeId of- True -> sendLocal them msg- False -> liftIO $ sendMessage (processNode proc)- (ProcessIdentifier (processId proc))- (ProcessIdentifier them)- NoImplicitReconnect- msg- -- We do not fire the trace event until after the sending is complete;- -- In the remote case, 'sendMessage' can block in the networking stack.+ destNode = (processNodeId them) liftIO $ traceEvent (localEventBus node) (MxSent them us (createUnencodedMessage msg))+ if destNode == nodeId+ then sendLocal them msg+ else liftIO $ sendMessage (processNode proc)+ (ProcessIdentifier (processId proc))+ (ProcessIdentifier them)+ NoImplicitReconnect+ msg -- | /Unsafe/ variant of 'send'. This function makes /no/ attempt to serialize -- and (in the case when the destination process resides on the same local@@ -279,9 +274,12 @@ -- usend :: Serializable a => ProcessId -> a -> Process () usend them msg = do- here <- getSelfNode+ proc <- ask let there = processNodeId them- if here == there+ let (us, node) = (processId proc, processNode proc)+ let msg' = wrapMessage msg+ liftIO $ traceEvent (localEventBus node) (MxSent them us msg')+ if localNodeId (processNode proc) == there then sendLocal them msg else sendCtrlMsg (Just there) $ UnreliableSend (processLocalId them) (createMessage msg)@@ -301,7 +299,13 @@ -- Channels -- -------------------------------------------------------------------------------- --- | Create a new typed channel+-- | Create a new typed channel, bound to the calling Process+--+-- Note that the channel is bound to the lifecycle of the process that evaluates+-- this function, such that when it dies/exits, the channel will no longer+-- function, but will remain accessible. Thus reading from the ReceivePort will+-- fail silently thereafter, blocking indefinitely (unless a timeout is used).+-- newChan :: Serializable a => Process (SendPort a, ReceivePort a) newChan = do proc <- ask@@ -329,13 +333,16 @@ sendChan :: Serializable a => SendPort a -> a -> Process () sendChan (SendPort cid) msg = do proc <- ask- let node = localNodeId (processNode proc)- destNode = processNodeId (sendPortProcessId cid) in do- case destNode == node of+ let node = processNode proc+ pid = processId proc+ us = localNodeId node+ them = processNodeId (sendPortProcessId cid)+ liftIO $ traceEvent (localEventBus node) (MxSentToPort pid cid $ wrapMessage msg)+ case them == us of True -> sendChanLocal cid msg False -> do- liftIO $ sendBinary (processNode proc)- (ProcessIdentifier (processId proc))+ liftIO $ sendBinary node+ (ProcessIdentifier pid) (SendPortIdentifier cid) NoImplicitReconnect msg@@ -356,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. --@@ -396,15 +408,20 @@ receiveWait :: [Match b] -> Process b receiveWait ms = do queue <- processQueue <$> ask- Just proc <- liftIO $ dequeue queue Blocking (map unMatch ms)- proc+ mProc <- liftIO $ dequeue queue Blocking (map unMatch ms)+ case mProc of+ Just proc' -> proc'+ Nothing -> die $ "System Invariant Violation: CQueue.hs returned `Nothing` "+ ++ "in the absence of a timeout value." -- | Like 'receiveWait' but with a timeout. -- -- If the timeout is zero do a non-blocking check for matching messages. A -- non-zero timeout is applied only when waiting for incoming messages (that is, -- /after/ we have checked the messages that are already in the mailbox).-receiveTimeout :: Int -> [Match b] -> Process (Maybe b)+receiveTimeout :: Int -- ^ Timeout in microseconds+ -> [Match b] + -> Process (Maybe b) receiveTimeout t ms = do queue <- processQueue <$> ask let blockSpec = if t == 0 then NonBlocking else Timeout t@@ -475,18 +492,15 @@ let node = processNode proc us = processId proc nid = localNodeId node- destNode = (processNodeId them) in do- case destNode == nid of- True -> sendCtrlMsg Nothing (LocalSend them msg)- False -> liftIO $ sendPayload (processNode proc)- (ProcessIdentifier (processId proc))- (ProcessIdentifier them)- NoImplicitReconnect- (messageToPayload msg)- -- We do not fire the trace event until after the sending is complete;- -- In the remote case, 'sendMessage' can block in the networking stack.- liftIO $ traceEvent (localEventBus node)- (MxSent them us msg)+ destNode = (processNodeId them)+ liftIO $ traceEvent (localEventBus node) (MxSent them us msg)+ if destNode == nid+ then sendCtrlMsg Nothing (LocalSend them msg)+ else liftIO $ sendPayload (processNode proc)+ (ProcessIdentifier (processId proc))+ (ProcessIdentifier them)+ NoImplicitReconnect+ (messageToPayload msg) -- | Forward a raw 'Message' to the given 'ProcessId'. --@@ -499,15 +513,11 @@ let node = processNode proc us = processId proc nid = localNodeId node- destNode = (processNodeId them) in do- case destNode == nid of- True -> sendCtrlMsg Nothing (LocalSend them msg)- False -> sendCtrlMsg (Just destNode) $ UnreliableSend (processLocalId them)- msg- -- We do not fire the trace event until after the sending is complete;- -- In the remote case, 'sendCtrlMsg' can block in the networking stack.- liftIO $ traceEvent (localEventBus node)- (MxSent them us msg)+ destNode = (processNodeId them)+ liftIO $ traceEvent (localEventBus node) (MxSent them us msg)+ if destNode == nid+ then sendCtrlMsg Nothing (LocalSend them msg)+ else sendCtrlMsg (Just destNode) $ UnreliableSend (processLocalId them) msg -- | Wrap a 'Serializable' value in a 'Message'. Note that 'Message's are -- 'Serializable' - like the datum they contain - but also note, deserialising@@ -885,13 +895,20 @@ -- @withMonitor@ returns, there might still be unreceived monitor -- messages in the queue. ---withMonitor :: ProcessId -> Process a -> Process a-withMonitor pid code = bracket (monitor pid) unmonitor (\_ -> code)- -- unmonitor blocks waiting for the response, so there's a possibility- -- that an exception might interrupt withMonitor before the unmonitor- -- has completed. I think that's better than making the unmonitor- -- uninterruptible.+withMonitor :: ProcessId -> (MonitorRef -> Process a) -> Process a+withMonitor pid = bracket (monitor pid) unmonitor +-- | Establishes temporary monitoring of another process.+--+-- @withMonitor_ pid code@ sets up monitoring of @pid@ for the duration+-- of @code@. Note: although monitoring is no longer active when+-- @withMonitor_@ returns, there might still be unreceived monitor+-- messages in the queue.+--+-- Since 0.6.1+withMonitor_ :: ProcessId -> Process a -> Process a+withMonitor_ p = withMonitor p . const+ -- | Remove a link -- -- This is synchronous in the sense that once it returns you are guaranteed@@ -928,12 +945,21 @@ -- | Remove a monitor -- -- This has the same synchronous/asynchronous nature as 'unlink'.+--+-- ProcessMonitorNotification messages for the given MonitorRef are removed from+-- the mailbox. unmonitor :: MonitorRef -> Process () unmonitor ref = do unmonitorAsync ref receiveWait [ matchIf (\(DidUnmonitor ref') -> ref' == ref) (\_ -> return ()) ]+ -- Discard the notification if any. With the current NC implementation at most+ -- one notification is in the mailbox for any given ref.+ void $ receiveTimeout 0+ [ matchIf (\(ProcessMonitorNotification ref' _ _) -> ref' == ref)+ (const $ return ())+ ] -------------------------------------------------------------------------------- -- Exception handling --@@ -947,7 +973,7 @@ -- | Lift 'Control.Exception.try' try :: Exception e => Process a -> Process (Either e a) try = Catch.try-{-# DEPRECATED try "Use Control.Monad.Catch.mask_ instead" #-}+{-# DEPRECATED try "Use Control.Monad.Catch.try instead" #-} -- | Lift 'Control.Exception.mask' mask :: ((forall a. Process a -> Process a) -> Process b) -> Process b@@ -1001,7 +1027,9 @@ -------------------------------------------------------------------------------- -- | Like 'expect' but with a timeout-expectTimeout :: forall a. Serializable a => Int -> Process (Maybe a)+expectTimeout :: forall a. Serializable a + => Int -- ^ Timeout in microseconds+ -> Process (Maybe a) expectTimeout n = receiveTimeout n [match return] -- | Asynchronous version of 'spawn'@@ -1059,17 +1087,52 @@ -- Logging -- -------------------------------------------------------------------------------- +data SayMessage = SayMessage { sayTime :: UTCTime+ , sayProcess :: ProcessId+ , sayMessage :: String }+ deriving (Typeable)++-- There is sadly no Show UTCTime instance+instance Show SayMessage where+ showsPrec p msg =+ showParen (p >= 11)+ $ showString "SayMessage "+ . showString (formatTime defaultTimeLocale "%c" (sayTime msg))+ . showChar ' '+ . showsPrec 11 (sayProcess msg) . showChar ' '+ . showsPrec 11 (sayMessage msg) . showChar ' '++instance Binary SayMessage where+ put s = do+ putUTCTime (sayTime s)+ put (sayProcess s)+ put (sayMessage s)+ get = SayMessage <$> getUTCTime <*> get <*> get++-- Sadly there is no Binary UTCTime instance+putUTCTime :: UTCTime -> Put+putUTCTime (UTCTime (ModifiedJulianDay day) tod) = do+ put day+ put (toRational tod)++getUTCTime :: Get UTCTime+getUTCTime = do+ day <- get+ tod <- get+ return $! UTCTime (ModifiedJulianDay day)+ (fromRational tod)+ -- | Log a string ----- @say message@ sends a message (time, pid of the current process, message)--- to the process registered as 'logger'. By default, this process simply--- sends the string to 'stderr'. Individual Cloud Haskell backends might--- replace this with a different logger process, however.+-- @say message@ sends a message of type 'SayMessage' with the current time and+-- 'ProcessId' of the current process to the process registered as @logger@. By+-- default, this process simply sends the string to @stderr@. Individual Cloud+-- Haskell backends might replace this with a different logger process, however. say :: String -> Process () say string = do now <- liftIO getCurrentTime us <- getSelfPid- nsend "logger" (formatTime defaultTimeLocale "%c" now, us, string)+ nsend "logger" (SayMessage now us string) -------------------------------------------------------------------------------- -- Registry --@@ -1165,8 +1228,12 @@ -- | Named send to a process in the local registry (asynchronous) nsend :: Serializable a => String -> a -> Process ()-nsend label msg =- sendCtrlMsg Nothing (NamedSend label (createUnencodedMessage msg))+nsend label msg = do+ proc <- ask+ let msg' = createUnencodedMessage msg+ liftIO $ traceEvent (localEventBus (processNode proc))+ (MxSentToName label (processId proc) msg')+ sendCtrlMsg Nothing (NamedSend label msg') -- | Named send to a process in the local registry (asynchronous). -- This function makes /no/ attempt to serialize and (in the case when the@@ -1178,9 +1245,15 @@ -- | Named send to a process in a remote registry (asynchronous) nsendRemote :: Serializable a => NodeId -> String -> a -> Process () nsendRemote nid label msg = do- here <- getSelfNode- if here == nid then nsend label msg- else sendCtrlMsg (Just nid) (NamedSend label (createMessage msg))+ proc <- ask+ let us = processId proc+ let node = processNode proc+ if localNodeId node == nid+ then nsend label msg+ else let lbl = label ++ "@" ++ show nid in do+ liftIO $ traceEvent (localEventBus node)+ (MxSentToName lbl us (wrapMessage msg))+ sendCtrlMsg (Just nid) (NamedSend label (createMessage msg)) -- | Named send to a process in a remote registry (asynchronous) -- This function makes /no/ attempt to serialize and (in the case when the
src/Control/Distributed/Process/Internal/StrictMVar.hs view
@@ -14,13 +14,8 @@ , mkWeakMVar ) where -import Control.Applicative ((<$>)) import Control.Monad ((>=>))-#if MIN_VERSION_base(4,6,0)-import Control.Exception (evaluate)-#else import Control.Exception (evaluate, mask_, onException)-#endif import qualified Control.Concurrent.MVar as MVar ( MVar , newEmptyMVar@@ -31,9 +26,6 @@ , withMVar , modifyMVar_ , modifyMVar-#if MIN_VERSION_base(4,6,0)- , modifyMVarMasked-#endif ) import GHC.MVar (MVar(MVar)) import GHC.IO (IO(IO), unIO)@@ -71,23 +63,12 @@ modifyMVarMasked :: StrictMVar a -> (a -> IO (a, b)) -> IO b modifyMVarMasked (StrictMVar v) f =-#if MIN_VERSION_base(4,6,0)- MVar.modifyMVarMasked v (f >=> evaluateFst)-#else mask_ $ do a <- MVar.takeMVar v (a',b) <- (f a >>= evaluate) `onException` MVar.putMVar v a MVar.putMVar v a' return b-#endif- where- evaluateFst :: (a, b) -> IO (a, b)- evaluateFst (x, y) = evaluate x >> return (x, y) mkWeakMVar :: StrictMVar a -> IO () -> IO (Weak (StrictMVar a)) mkWeakMVar q@(StrictMVar (MVar m#)) f = IO $ \s ->-#if MIN_VERSION_base(4,9,0) case mkWeak# m# q (unIO f) s of (# s', w #) -> (# s', Weak w #)-#else- case mkWeak# m# q f s of (# s', w #) -> (# s', Weak w #)-#endif
src/Control/Distributed/Process/Internal/Types.hs view
@@ -353,6 +353,7 @@ deriving ( Applicative , Functor , Monad+ , MonadFail , MonadFix , MonadIO , MonadReader LocalProcess@@ -366,6 +367,13 @@ lproc <- ask liftIO $ catch (runLocalProcess lproc p) (runLocalProcess lproc . h) instance MonadMask Process where+ generalBracket acquire release inner = do+ lproc <- ask+ liftIO $+ generalBracket (runLocalProcess lproc acquire)+ (\a e -> runLocalProcess lproc $ release a e)+ (runLocalProcess lproc . inner)+ mask p = do lproc <- ask liftIO $ mask $ \restore ->@@ -463,11 +471,7 @@ deriving (Typeable) instance NFData Message where-#if MIN_VERSION_bytestring(0,10,0) rnf (EncodedMessage _ e) = rnf e `seq` ()-#else- rnf (EncodedMessage _ e) = BSL.length e `seq` ()-#endif rnf (UnencodedMessage _ a) = e `seq` () where e = BSL.length (encode a) @@ -629,7 +633,7 @@ -- | 'SpawnRef' are used to return pids of spawned processes newtype SpawnRef = SpawnRef Int32- deriving (Show, Binary, Typeable, Eq)+ deriving (Show, Binary, Typeable, Eq, Ord) -- | (Asynchronius) reply from 'spawn' data DidSpawn = DidSpawn SpawnRef ProcessId
src/Control/Distributed/Process/Internal/WeakTQueue.hs view
@@ -100,8 +100,4 @@ mkWeakTQueue :: TQueue a -> IO () -> IO (Weak (TQueue a)) mkWeakTQueue q@(TQueue _read (TVar write#)) f = IO $ \s ->-#if MIN_VERSION_base(4,9,0) case mkWeak# write# q (unIO f) s of (# s', w #) -> (# s', Weak w #)-#else- case mkWeak# write# q f s of (# s', w #) -> (# s', Weak w #)-#endif
src/Control/Distributed/Process/Management.hs view
@@ -76,6 +76,40 @@ -- -- * Whether messages will be taken from the mailbox first, or the event bus. --+-- Since the event bus uses STM broadcast channels to communicate with agents,+-- no message written to the bus successfully can be lost.+--+-- Agents can also receive messages via their mailboxes - these are subject to+-- the same guarantees as all inter-process message sending.+--+-- Messages dispatched on an STM broadcast channel (i.e., management event bus)+-- are guaranteed to be delivered with the same FIFO ordering guarantees that+-- exist between two communicating processes, such that communication from the+-- node controller's threads (i.e., MxEvent's) will never be re-ordered, but+-- messages dispatched to the event bus by other processes (including, but not+-- limited to agents) are only guaranteed to be ordered between one sender and+-- one receiver.+--+-- No guarantee exists for the ordering in which messages sent to an agent's+-- mailbox will be delivered, vs messages dispatched via the event bus.+--+-- Because of the above, there are no ordering guarantees for messages sent+-- between agents, or for processes to agents, except for those that apply to+-- messages sent between regular processes, since agents are+-- implemented as such.+--+-- The event bus is serial and single threaded. Anything that is published by+-- the node controller will be seen in FIFO order. There are no ordering+-- guarantees pertaining to entries published to the event bus by other+-- processes or agents.+--+-- It should not be possible to see, for example, an @MxReceived@ before the+-- corresponding @MxSent@ event, since the places where we issue the @MxSent@+-- write directly to the event bus (using STM) in the calling (green) thread,+-- before dispatching instructions to the node controller to perform the+-- necessary routing to deliver the message to a process (or registered name,+-- or typed channel) locally or remotely.+-- -- [Management Data API] -- -- Both management agents and clients of the API have access to a variety of@@ -83,28 +117,8 @@ -- system information. Agents maintain their own internal state privately (via a -- state transformer - see 'mxGetLocal' et al), however it is possible for -- agents to share additional data with each other (and the outside world)--- using /data tables/.------ Each agent is assigned its own data table, which acts as a shared map, where--- the keys are @String@s and the values are @Serializable@ datum of whatever--- type the agent or its clients stores.------ Because an agent's /data table/ stores its values in raw 'Message' format,--- it works effectively as an /un-typed dictionary/, into which data of varying--- types can be fed and later retrieved. The upside of this is that different--- keys can be mapped to various types without any additional work on the part--- of the developer. The downside is that the code reading these values must--- know in advance what type(s) to expect, and the API provides no additional--- support for handling that.------ Publishing is accomplished using the 'mxPublish' and 'mxSet' APIs, whilst--- querying and deletion are handled by 'mxGet', 'mxClear', 'mxPurgeTable' and--- 'mxDropTable' respectively.------ When a management agent terminates, their tables are left in memory despite--- termination, such that an agent may resume its role (by restarting) or have--- its 'MxAgentId' taken over by another subsequent agent, leaving the data--- originally captured in place.+-- using whatever mechanism the user wishes, e.g., acidstate, or shared memory+-- primitives. -- -- [Defining Agents] --@@ -164,8 +178,6 @@ -- > monitorNames = getSelfPid >>= nsend "name-monitor" -- > monitorNames2 = getSelfPid >>= mxNotify ----- For some real-world examples, see the distributed-process-platform package.--- -- [Performance, Stablity and Scalability] -- -- /Management Agents/ offer numerous advantages over regular processes:@@ -194,10 +206,6 @@ -- the event bus /and/ their own mailboxes, plus searching through the set of -- event sinks (for each agent) to determine the right handler for the event. ----- Each management agent requires not only its own @Process@ (in which the agent--- code is run), but also a peer process that provides its /data table/. These--- data tables also have to be coordinated and manaaged on each agent's behalf.--- -- [Architecture Overview] -- -- The architecture of the management event bus is internal and subject to@@ -264,13 +272,6 @@ , mxGetLocal , mxUpdateLocal , liftMX- -- * Mx Data API- , mxPublish- , mxSet- , mxGet- , mxClear- , mxPurgeTable- , mxDropTable ) where import Control.Applicative@@ -281,10 +282,7 @@ , TChan ) import Control.Distributed.Process.Internal.Primitives- ( newChan- , nsend- , receiveWait- , matchChan+ ( receiveWait , matchAny , matchSTM , unwrapMessage@@ -302,14 +300,12 @@ , unsafeCreateUnencodedMessage ) import Control.Distributed.Process.Management.Internal.Bus (publishEvent)-import qualified Control.Distributed.Process.Management.Internal.Table as Table import Control.Distributed.Process.Management.Internal.Types ( MxAgentId(..) , MxAgent(..) , MxAction(..) , ChannelSelector(..) , MxAgentState(..)- , MxAgentStart(..) , MxSink , MxEvent(..) )@@ -335,42 +331,6 @@ bus <- localEventBus . processNode <$> ask liftIO $ publishEvent bus $ unsafeCreateUnencodedMessage msg --- | Publish an arbitrary @Message@ as a property in the management database.------ For publishing @Serializable@ data, use 'mxSet' instead.----mxPublish :: MxAgentId -> String -> Message -> Process ()-mxPublish a k v = Table.set k v (Table.MxForAgent a)---- | Sets an arbitrary @Serializable@ datum against a key in the management--- database. Note that /no attempt is made to force the argument/, therefore--- it is very important that you do not pass unevaluated thunks that might--- crash some other, arbitrary process (or management agent!) that obtains--- and attempts to force the value later on.----mxSet :: Serializable a => MxAgentId -> String -> a -> Process ()-mxSet mxId key msg = do- Table.set key (unsafeCreateUnencodedMessage msg) (Table.MxForAgent mxId)---- | Fetches a property from the management database for the given key.--- If the property is not set, or does not match the expected type when--- typechecked (at runtime), returns @Nothing@.-mxGet :: Serializable a => MxAgentId -> String -> Process (Maybe a)-mxGet = Table.fetch . Table.MxForAgent---- | Clears a property from the management database using the given key.--- If the key does not exist in the database, this is a noop.-mxClear :: MxAgentId -> String -> Process ()-mxClear mxId key = Table.clear key (Table.MxForAgent mxId)---- | Purges a table in the management database of all its stored properties.-mxPurgeTable :: MxAgentId -> Process ()-mxPurgeTable = Table.purge . Table.MxForAgent---- | Deletes a table from the management database.-mxDropTable :: MxAgentId -> Process ()-mxDropTable = Table.delete . Table.MxForAgent- -------------------------------------------------------------------------------- -- API for writing user defined management extensions (i.e., agents) -- --------------------------------------------------------------------------------@@ -490,10 +450,7 @@ return pid where start (sendTChan, recvTChan) = do- (sp, rp) <- newChan- nsend Table.mxTableCoordinator (MxAgentStart sp mxId)- tablePid <- receiveWait [ matchChan rp (\(p :: ProcessId) -> return p) ]- let nState = MxAgentState mxId sendTChan tablePid initState+ let nState = MxAgentState mxId sendTChan initState runAgent dtor handlers InputChan recvTChan nState runAgent :: MxAgent s ()
− src/Control/Distributed/Process/Management/Internal/Table.hs
@@ -1,222 +0,0 @@-{-# LANGUAGE DeriveDataTypeable #-}-{-# LANGUAGE ScopedTypeVariables #-}-{-# LANGUAGE RankNTypes #-}-{-# LANGUAGE DeriveGeneric #-}-{-# LANGUAGE RecordWildCards #-}--module Control.Distributed.Process.Management.Internal.Table- ( MxTableRequest(..)- , MxTableId(..)- , mxTableCoordinator- , startTableCoordinator- , delete- , purge- , clear- , set- , get- , fetch- ) where--import Control.Distributed.Process.Internal.Primitives- ( receiveWait- , receiveChan- , match- , matchAny- , matchIf- , matchChan- , send- , nsend- , sendChan- , getSelfPid- , link- , monitor- , unwrapMessage- , newChan- , withMonitor- )-import Control.Distributed.Process.Internal.Types- ( Process- , ProcessId- , ProcessMonitorNotification(..)- , SendPort- , ReceivePort- , Message- , unsafeCreateUnencodedMessage- )-import Control.Distributed.Process.Management.Internal.Types- ( MxTableId(..)- , MxAgentId(..)- , MxAgentStart(..)- , Fork)-import Control.Distributed.Process.Serializable (Serializable)-import Control.Monad.IO.Class (liftIO)-import Data.Accessor (Accessor, accessor, (^=), (^:))-import Data.Binary (Binary)-import Data.Map (Map)-import qualified Data.Map as Map-import Data.Typeable (Typeable)--import GHC.Generics---- An extremely lightweight shared Map implementation, for use--- by /management agents/ and their cohorts. Each agent is assigned--- a table, into which any serializable @Message@ can be inserted.--- Data are inserted, removed and searched for via their key, which--- is a string. Tables can be purged, values can be set, fetched or--- cleared/removed.-----data MxTableRequest =- Delete- | Purge- | Clear !String- | Set !String !Message- | Get !String !(SendPort (Maybe Message)) -- see [note: un-typed send port]- deriving (Typeable, Generic)-instance Binary MxTableRequest where--data MxTableState = MxTableState { _name :: !String- , _entries :: !(Map String Message)- }--type MxTables = Map MxAgentId ProcessId--mxTableCoordinator :: String-mxTableCoordinator = "mx.table.coordinator"--delete :: MxTableId -> Process ()-delete = sendReq Delete--purge :: MxTableId -> Process ()-purge = sendReq Purge--clear :: String -> MxTableId -> Process ()-clear k = sendReq (Clear k)--set :: String -> Message -> MxTableId -> Process ()-set k v = sendReq (Set k v)--fetch :: forall a. (Serializable a)- => MxTableId- -> String- -> Process (Maybe a)-fetch (MxForPid pid) key = get pid key-fetch mxId@(MxForAgent _) key = do- (sp, rp) <- newChan :: Process (SendPort (Maybe Message),- ReceivePort (Maybe Message))- sendReq (Get key sp) mxId- receiveChan rp >>= maybe (return Nothing)- (unwrapMessage :: Message -> Process (Maybe a))---- [note: un-typed send port]--- Here, fetch uses a typed channel over a raw Message to obtain--- its result, so type checking is deferred until receipt and will--- be handled in the caller's thread. This is necessary because--- the server portion of the code knows nothing about the types--- involved, nor should it, since these tables can be used to--- store arbitrary serializable data.--get :: forall a. (Serializable a)- => ProcessId- -> String- -> Process (Maybe a)-get pid key = do- safeFetch pid key >>= maybe (return Nothing)- (unwrapMessage :: Message -> Process (Maybe a))--safeFetch :: ProcessId -> String -> Process (Maybe Message)-safeFetch pid key = do- (sp, rp) <- newChan- send pid $ Get key sp- withMonitor pid $ do- receiveWait [- matchChan rp return- , matchIf (\(ProcessMonitorNotification _ pid' _) -> pid' == pid)- (\_ -> return $ Just (unsafeCreateUnencodedMessage ()))- ]--sendReq :: MxTableRequest -> MxTableId -> Process ()-sendReq req tid = (resolve tid) req--resolve :: Serializable a => MxTableId -> (a -> Process ())-resolve (MxForAgent agent) = \msg -> nsend mxTableCoordinator (agent, msg)-resolve (MxForPid pid) = \msg -> send pid msg--startTableCoordinator :: Fork -> Process ()-startTableCoordinator fork = run Map.empty- where- run :: MxTables -> Process ()- run tables =- receiveWait [- -- note that this state change can race with MxAgentStart requests- match (\(ProcessMonitorNotification _ pid _) -> do- return $ Map.filter (/= pid) tables)- , match (\(MxAgentStart ch agent) -> do- lookupAgent tables agent >>= \(p, t) -> do- sendChan ch p >> return t)- , match (\req@(agent, tReq :: MxTableRequest) -> do- case tReq of- Get k sp -> do- lookupAgent tables agent >>= \(p, t) -> do- safeFetch p k >>= sendChan sp >> return t- _ -> do- handleRequest tables req)- , matchAny (\_ -> return tables) -- unrecognised messages are dropped- ] >>= run-- handleRequest :: MxTables- -> (MxAgentId, MxTableRequest)- -> Process MxTables- handleRequest tables' (agent, req) = do- lookupAgent tables' agent >>= \(p, t) -> send p req >> return t-- lookupAgent :: MxTables -> MxAgentId -> Process (ProcessId, MxTables)- lookupAgent tables' agentId' = do- case Map.lookup agentId' tables' of- Nothing -> launchNew agentId' tables'- Just p -> return (p, tables')-- launchNew :: MxAgentId- -> MxTables- -> Process (ProcessId, MxTables)- launchNew mxId tblMap = do- let initState = MxTableState { _name = (agentId mxId)- , _entries = Map.empty- }- (pid, _) <- spawnSup $ tableHandler initState- return $ (pid, mxId `seq` pid `seq` Map.insert mxId pid tblMap)-- spawnSup proc = do- us <- getSelfPid- -- we need to use that passed in "fork", in order to- -- break an import cycle with Node.hs courtesy of the- -- management agent, API and tracing modules- them <- liftIO $ fork $ link us >> proc- ref <- monitor them- return (them, ref)--tableHandler :: MxTableState -> Process ()-tableHandler state = do- ns <- receiveWait [- match (handleTableRequest state)- , matchAny (\_ -> return (Just state))- ]- case ns of- Nothing -> return ()- Just s' -> tableHandler s'- where- handleTableRequest _ Delete = return Nothing- handleTableRequest st Purge = return $ Just $ (entries ^= Map.empty) $ st- handleTableRequest st (Clear k) = return $ Just $ (entries ^: (k `seq` Map.delete k)) $ st- handleTableRequest st (Set k v) = return $ Just $ (entries ^: (k `seq` v `seq` Map.insert k v)) st- handleTableRequest st (Get k c) = getEntry k c st >> return (Just st)--getEntry :: String- -> SendPort (Maybe Message)- -> MxTableState- -> Process ()-getEntry k m MxTableState{..} = do- sendChan m =<< return (Map.lookup k _entries)--entries :: Accessor MxTableState (Map String Message)-entries = accessor _entries (\ls st -> st { _entries = ls })
src/Control/Distributed/Process/Management/Internal/Trace/Primitives.hs view
@@ -33,6 +33,7 @@ ( whereis , newChan , receiveChan+ , die ) import Control.Distributed.Process.Management.Internal.Trace.Types ( TraceArg(..)@@ -168,6 +169,11 @@ withLocalTracer $ \t -> liftIO $ Tracer.getCurrentTraceClient t sp currentTracer <- receiveChan rp case currentTracer of- Nothing -> do { (Just p') <- whereis "tracer.initial"; act p' }+ Nothing -> do mTP <- whereis "tracer.initial"+ -- NB: this should NOT ever happen, but forcing pattern matches+ -- is not considered cool in later versions of MonadFail+ case mTP of+ Just p' -> act p'+ Nothing -> die $ "System Invariant Violation: Tracer Process "+ ++ "Name Not Found (whereis tracer.initial)" (Just p) -> act p-
src/Control/Distributed/Process/Management/Internal/Trace/Tracer.hs view
@@ -77,10 +77,6 @@ import Debug.Trace (traceEventIO) import Prelude -#if ! MIN_VERSION_base(4,6,0)-import Prelude hiding (catch)-#endif- import System.Environment (getEnv) import System.IO ( Handle@@ -91,11 +87,7 @@ , hPutStrLn , hSetBuffering )-#if MIN_VERSION_time(1,5,0) import Data.Time.Format (defaultTimeLocale)-#else-import System.Locale (defaultTimeLocale)-#endif import System.Mem.Weak ( Weak )
src/Control/Distributed/Process/Management/Internal/Trace/Types.hs view
@@ -43,7 +43,6 @@ import Data.Binary import Data.List (intersperse) import Data.Set (Set)-import qualified Data.Set as Set (fromList) import Data.Typeable import GHC.Generics @@ -158,12 +157,3 @@ getCurrentTraceClient :: MxEventBus -> SendPort (Maybe ProcessId) -> IO () getCurrentTraceClient t s = publishEvent t (unsafeCreateUnencodedMessage s)--class Traceable a where- uod :: [a] -> TraceSubject--instance Traceable ProcessId where- uod = TraceProcs . Set.fromList--instance Traceable String where- uod = TraceNames . Set.fromList
src/Control/Distributed/Process/Management/Internal/Types.hs view
@@ -4,27 +4,24 @@ {-# LANGUAGE DeriveGeneric #-} module Control.Distributed.Process.Management.Internal.Types ( MxAgentId(..)- , MxTableId(..) , MxAgentState(..) , MxAgent(..) , MxAction(..) , ChannelSelector(..)- , MxAgentStart(..) , Fork , MxSink , MxEvent(..) , Addressable(..) ) where -import Control.Applicative (Applicative) import Control.Concurrent.STM ( TChan ) import Control.Distributed.Process.Internal.Types ( Process , ProcessId+ , SendPortId , Message- , SendPort , DiedReason , NodeId )@@ -33,6 +30,7 @@ ( MonadState , StateT )+import Control.Monad.Fix (MonadFix) import Data.Binary import Data.Typeable (Typeable) import GHC.Generics@@ -59,8 +57,14 @@ -- ^ fired whenever a node /dies/ (i.e., the connection is broken/disconnected) | MxSent ProcessId ProcessId Message -- ^ fired whenever a message is sent from a local process+ | MxSentToName String ProcessId Message+ -- ^ fired whenever a named send occurs+ | MxSentToPort ProcessId SendPortId Message+ -- ^ fired whenever a sendChan occurs | MxReceived ProcessId Message -- ^ fired whenever a message is received by a local process+ | MxReceivedPort SendPortId Message+ -- ^ fired whenever a message is received via a typed channel | MxConnected ConnectionId EndPointAddress -- ^ fired when a network-transport connection is first established | MxDisconnected ConnectionId EndPointAddress@@ -99,17 +103,10 @@ newtype MxAgentId = MxAgentId { agentId :: String } deriving (Typeable, Binary, Eq, Ord) -data MxTableId =- MxForAgent !MxAgentId- | MxForPid !ProcessId- deriving (Typeable, Generic)-instance Binary MxTableId where- data MxAgentState s = MxAgentState { mxAgentId :: !MxAgentId , mxBus :: !(TChan Message)- , mxSharedTable :: !ProcessId , mxLocalState :: !s } @@ -122,18 +119,11 @@ } deriving ( Functor , Monad , MonadIO+ , MonadFix , ST.MonadState (MxAgentState s) , Typeable , Applicative )--data MxAgentStart = MxAgentStart- {- mxAgentTableChan :: SendPort ProcessId- , mxAgentIdStart :: MxAgentId- }- deriving (Typeable, Generic)-instance Binary MxAgentStart where data ChannelSelector = InputChan | Mailbox
src/Control/Distributed/Process/Node.hs view
@@ -30,13 +30,15 @@ ( empty , toList , fromList- , filter+ , partition , partitionWithKey , elems , size , filterWithKey , foldlWithKey )+import Data.Time.Format (formatTime)+import Data.Time.Format (defaultTimeLocale) import Data.Set (Set) import qualified Data.Set as Set ( empty@@ -48,7 +50,6 @@ , union ) import Data.Foldable (forM_)-import Data.List (foldl') import Data.Maybe (isJust, fromJust, isNothing, catMaybes) import Data.Typeable (Typeable) import Control.Category ((>>>))@@ -178,9 +179,8 @@ import Control.Distributed.Process.Management.Internal.Agent ( mxAgentController )-import qualified Control.Distributed.Process.Management.Internal.Table as Table- ( mxTableCoordinator- , startTableCoordinator+import Control.Distributed.Process.Management.Internal.Types+ ( MxEvent(..) ) import qualified Control.Distributed.Process.Management.Internal.Trace.Remote as Trace ( remoteTable@@ -194,9 +194,6 @@ , traceLogFmt , enableTrace )-import Control.Distributed.Process.Management.Internal.Types- ( MxEvent(..)- ) import Control.Distributed.Process.Serializable (Serializable) import Control.Distributed.Process.Internal.Messaging ( sendBinary@@ -209,6 +206,7 @@ , match , sendChan , unwrapMessage+ , SayMessage(..) ) import Control.Distributed.Process.Internal.Types (SendPort, Tracer(..)) import qualified Control.Distributed.Process.Internal.Closure.BuiltIn as BuiltIn (remoteTable)@@ -316,8 +314,6 @@ -- before /that/ process has started - this is a totally harmless race -- however, so we deliberably ignore it startDefaultTracer node- tableCoordinatorPid <- fork $ Table.startTableCoordinator fork- runProcess node $ register Table.mxTableCoordinator tableCoordinatorPid logger <- forkProcess node loop runProcess node $ do register "logger" logger@@ -326,12 +322,11 @@ -- process which uses 'send' or other primitives which are traced. register "trace.logger" logger where- fork = forkProcess node- loop = do receiveWait- [ match $ \((time, pid, string) ::(String, ProcessId, String)) -> do- liftIO . hPutStrLn stderr $ time ++ " " ++ show pid ++ ": " ++ string+ [ match $ \(SayMessage time pid string) -> do+ let time' = formatTime defaultTimeLocale "%c" time+ liftIO . hPutStrLn stderr $ time' ++ " " ++ show pid ++ ": " ++ string loop , match $ \((time, string) :: (String, String)) -> do -- this is a 'trace' message from the local node tracer@@ -487,7 +482,7 @@ data IncomingTarget = Uninit | ToProc ProcessId (Weak (CQueue Message))- | ToChan TypedChannel+ | ToChan SendPortId TypedChannel | ToNode data ConnectionState = ConnectionState {@@ -542,13 +537,17 @@ enqueue queue msg -- 'enqueue' is strict trace node (MxReceived pid msg) go st- Just (_, ToChan (TypedChannel chan')) -> do+ Just (_, ToChan chId (TypedChannel chan')) -> do mChan <- deRefWeak chan' -- If mChan is Nothing, the process has given up the read end of -- the channel and we simply ignore the incoming message- forM_ mChan $ \chan -> atomically $- -- We make sure the message is fully decoded when it is enqueued- writeTQueue chan $! decode (BSL.fromChunks payload)+ forM_ mChan $ \chan -> do+ msg' <- atomically $ do+ msg <- return $! decode (BSL.fromChunks payload)+ -- We make sure the message is fully decoded when it is enqueued+ writeTQueue chan msg+ return msg+ trace node $ MxReceivedPort chId $ unsafeCreateUnencodedMessage msg' go st Just (_, ToNode) -> do let ctrlMsg = decode . BSL.fromChunks $ payload@@ -575,7 +574,7 @@ mChannel <- withMVar (processState proc) $ return . (^. typedChannelWithId lcid) case mChannel of Just channel ->- go (incomingAt cid ^= Just (src, ToChan channel) $ st)+ go (incomingAt cid ^= Just (src, ToChan chId channel) $ st) Nothing -> invalidRequest cid st $ "incoming attempt to connect to unknown channel of"@@ -904,7 +903,11 @@ modify' $ (links ^= unaffectedLinks') . (monitors ^= unaffectedMons') - modify' $ registeredHere ^: Map.filter (\pid -> not $ ident `impliesDeathOf` ProcessIdentifier pid)+ -- we now consider all labels for this identifier unregistered+ let toDrop pid = not $ ident `impliesDeathOf` ProcessIdentifier pid+ (keepNames, dropNames) <- Map.partition toDrop <$> gets (^. registeredHere)+ mapM_ (\(p, l) -> liftIO $ trace node (MxUnRegistered l p)) (Map.toList dropNames)+ modify' $ registeredHere ^= keepNames remaining <- fmap Map.toList (gets (^. registeredOnNodes)) >>= mapM (\(pid,nidlist) ->@@ -956,7 +959,11 @@ do modify' $ registeredHereFor label ^= mPid updateRemote node currentVal mPid case mPid of- (Just p) -> liftIO $ trace node (MxRegistered p label)+ (Just p) -> do+ if reregistration+ then liftIO $ trace node (MxUnRegistered (fromJust currentVal) label)+ else return ()+ liftIO $ trace node (MxRegistered p label) Nothing -> liftIO $ trace node (MxUnRegistered (fromJust currentVal) label) newVal <- gets (^. registeredHereFor label) ncSendToProcess from $ unsafeCreateUnencodedMessage $@@ -1037,11 +1044,12 @@ -- If ch is Nothing, the process has given up the read end of -- the channel and we simply ignore the incoming message - this ch <- deRefWeak chan'- forM_ ch $ \chan -> deliverChan msg chan- where deliverChan :: forall a . Message -> TQueue a -> IO ()- deliverChan (UnencodedMessage _ raw) chan' =+ forM_ ch $ \chan -> deliverChan node from msg chan+ where deliverChan :: forall a . LocalNode -> SendPortId -> Message -> TQueue a -> IO ()+ deliverChan n p (UnencodedMessage _ raw) chan' = do atomically $ writeTQueue chan' ((unsafeCoerce raw) :: a)- deliverChan (EncodedMessage _ _) _ =+ trace n (MxReceivedPort p $ unsafeCreateUnencodedMessage raw)+ deliverChan _ _ (EncodedMessage _ _) _ = -- this will not happen unless someone screws with Primitives.hs error "invalid local channel delivery" @@ -1071,7 +1079,7 @@ case mProc of Nothing -> dispatch (isLocal node (ProcessIdentifier from)) from (ProcessInfoNone DiedUnknownId)- Just proc -> do+ Just proc -> do itsLinks <- Set.map fst . BiMultiMap.lookupBy1st them <$> gets (^. links) itsMons <- BiMultiMap.lookupBy1st them <$> gets (^. monitors)@@ -1085,8 +1093,8 @@ infoNode = (processNodeId pid) , infoRegisteredNames = reg , infoMessageQueueLength = size- , infoMonitors = Set.toList itsMons- , infoLinks = Set.toList itsLinks+ , infoMonitors = Set.toList itsMons+ , infoLinks = Set.toList itsLinks } where dispatch :: (Serializable a) => Bool
src/Control/Distributed/Process/Serializable.hs view
@@ -1,8 +1,9 @@ {-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE ConstraintKinds #-} {-# LANGUAGE UndecidableInstances #-} {-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE ScopedTypeVariables #-} {-# LANGUAGE GADTs #-}-{-# LANGUAGE CPP #-} module Control.Distributed.Process.Serializable ( Serializable , encodeFingerprint@@ -17,13 +18,7 @@ import Data.Binary (Binary) -#if MIN_VERSION_base(4,7,0)-import Data.Typeable (Typeable)-import Data.Typeable.Internal (TypeRep(TypeRep), typeOf)-#else-import Data.Typeable (Typeable(..))-import Data.Typeable.Internal (TypeRep(TypeRep))-#endif+import Data.Typeable (Typeable, typeRepFingerprint, typeOf) import Numeric (showHex) import Control.Exception (throw)@@ -46,8 +41,7 @@ deriving (Typeable) -- | Objects that can be sent across the network-class (Binary a, Typeable a) => Serializable a-instance (Binary a, Typeable a) => Serializable a+type Serializable a = (Binary a, Typeable a) -- | Encode type representation as a bytestring encodeFingerprint :: Fingerprint -> ByteString@@ -71,11 +65,7 @@ -- | The fingerprint of the typeRep of the argument fingerprint :: Typeable a => a -> Fingerprint-#if MIN_VERSION_base(4,8,0)-fingerprint a = let TypeRep fp _ _ _ = typeOf a in fp-#else-fingerprint a = let TypeRep fp _ _ = typeOf a in fp-#endif+fingerprint = typeRepFingerprint . typeOf -- | Show fingerprint (for debugging purposes) showFingerprint :: Fingerprint -> ShowS
src/Control/Distributed/Process/UnsafePrimitives.hs view
@@ -32,8 +32,25 @@ -- the /normal/ strategy). -- -- Use of the functions in this module can potentially change the runtime--- behaviour of your application. You have been warned!+-- behaviour of your application. In addition, messages passed between Cloud+-- Haskell processes are written to a tracing infrastructure on the local node,+-- to provide improved introspection and debugging facilities for complex actor+-- based systems. This module makes no attempt to force evaluation in these+-- cases either, thus evaluation problems in passed data structures could not+-- only crash your processes, but could also bring down critical internal+-- services on which the node relies to function correctly. --+-- If you wish to repudiate such issues, you are advised to consider the use+-- of NFSerialisable in the distributed-process-extras package, which type+-- class brings NFData into scope along with Serializable, such that we can+-- force evaluation. Intended for use with modules such as this one, this+-- approach guarantees correct evaluatedness in terms of @NFData@. Please note+-- however, that we /cannot/ guarantee that an @NFData@ instance will behave the+-- same way as a @Binary@ one with regards evaluation, so it is still possible+-- to introduce unexpected behaviour by using /unsafe/ primitives in this way.+--+-- You have been warned!+-- -- This module is exported so that you can replace the use of Cloud Haskell's -- /safe/ messaging primitives. If you want to use both variants, then you can -- take advantage of qualified imports, however "Control.Distributed.Process"@@ -55,7 +72,12 @@ , sendBinary , sendCtrlMsg )-+import Control.Distributed.Process.Management.Internal.Types+ ( MxEvent(..)+ )+import Control.Distributed.Process.Management.Internal.Trace.Types+ ( traceEvent+ ) import Control.Distributed.Process.Internal.Types ( ProcessId(..) , NodeId(..)@@ -79,34 +101,48 @@ -- | Named send to a process in the local registry (asynchronous) nsend :: Serializable a => String -> a -> Process ()-nsend label msg =- sendCtrlMsg Nothing (NamedSend label (unsafeCreateUnencodedMessage msg))+nsend label msg = do+ proc <- ask+ let us = processId proc+ let msg' = wrapMessage msg+ -- see [note: tracing]+ liftIO $ traceEvent (localEventBus (processNode proc))+ (MxSentToName label us msg')+ sendCtrlMsg Nothing (NamedSend label msg') -- | Named send to a process in a remote registry (asynchronous) nsendRemote :: Serializable a => NodeId -> String -> a -> Process () nsendRemote nid label msg = do proc <- ask- if localNodeId (processNode proc) == nid+ let us = processId proc+ let node = processNode proc+ if localNodeId node == nid then nsend label msg- else sendCtrlMsg (Just nid) (NamedSend label (createMessage msg))+ else+ let lbl = label ++ "@" ++ show nid in do+ -- see [note: tracing] NB: this is a remote call to another NC...+ liftIO $ traceEvent (localEventBus node)+ (MxSentToName lbl us (wrapMessage msg))+ sendCtrlMsg (Just nid) (NamedSend label (createMessage msg)) -- | Send a message send :: Serializable a => ProcessId -> a -> Process () send them msg = do proc <- ask let node = localNodeId (processNode proc)- destNode = (processNodeId them) in do- case destNode == node of- True -> unsafeSendLocal them msg- False -> liftIO $ sendMessage (processNode proc)- (ProcessIdentifier (processId proc))- (ProcessIdentifier them)- NoImplicitReconnect- msg- where- unsafeSendLocal :: (Serializable a) => ProcessId -> a -> Process ()- unsafeSendLocal pid msg' =- sendCtrlMsg Nothing $ LocalSend pid (unsafeCreateUnencodedMessage msg')+ destNode = (processNodeId them)+ us = (processId proc)+ msg' = wrapMessage msg in do+ -- see [note: tracing]+ liftIO $ traceEvent (localEventBus (processNode proc))+ (MxSent them us msg')+ if destNode == node+ then sendCtrlMsg Nothing $ LocalSend them msg'+ else liftIO $ sendMessage (processNode proc)+ (ProcessIdentifier (processId proc))+ (ProcessIdentifier them)+ NoImplicitReconnect+ msg -- | Send a message unreliably. --@@ -121,30 +157,44 @@ usend them msg = do proc <- ask let there = processNodeId them+ let (us, node) = (processId proc, processNode proc)+ let msg' = wrapMessage msg+ -- see [note: tracing]+ liftIO $ traceEvent (localEventBus node) (MxSent them us msg') if localNodeId (processNode proc) == there- then sendCtrlMsg Nothing $- LocalSend them (unsafeCreateUnencodedMessage msg)+ then sendCtrlMsg Nothing $ LocalSend them msg' else sendCtrlMsg (Just there) $ UnreliableSend (processLocalId them) (createMessage msg) +-- [note: tracing]+-- Note that tracing writes to the local node's control channel, and this+-- module explicitly specifies to its clients that it does unsafe message+-- encoding. The same is true for the messages it puts onto the Management+-- event bus, however we do *not* want unevaluated thunks hitting the event+-- bus control thread. Hence the word /Unsafe/ in this module's name!+--+ -- | Send a message on a typed channel sendChan :: Serializable a => SendPort a -> a -> Process () sendChan (SendPort cid) msg = do proc <- ask- let node = localNodeId (processNode proc)- destNode = processNodeId (sendPortProcessId cid) in do- case destNode == node of- True -> unsafeSendChanLocal cid msg- False -> do- liftIO $ sendBinary (processNode proc)- (ProcessIdentifier (processId proc))- (SendPortIdentifier cid)- NoImplicitReconnect- msg+ let+ node = processNode proc+ pid = processId proc+ us = localNodeId node+ them = processNodeId (sendPortProcessId cid)+ msg' = wrapMessage msg+ liftIO $ traceEvent (localEventBus node) (MxSentToPort pid cid msg')+ if them == us+ then unsafeSendChanLocal cid msg' -- NB: we wrap to P.Message !!!+ else liftIO $ sendBinary node+ (ProcessIdentifier pid)+ (SendPortIdentifier cid)+ NoImplicitReconnect+ msg where- unsafeSendChanLocal :: (Serializable a) => SendPortId -> a -> Process ()- unsafeSendChanLocal spId msg' =- sendCtrlMsg Nothing $ LocalPortSend spId (unsafeCreateUnencodedMessage msg')+ unsafeSendChanLocal :: SendPortId -> Message -> Process ()+ unsafeSendChanLocal p m = sendCtrlMsg Nothing $ LocalPortSend p m -- | Create an unencoded @Message@ for any @Serializable@ type. wrapMessage :: Serializable a => a -> Message