extensible-effects-concurrent 0.1.0.1 → 0.1.1.0
raw patch · 10 files changed
+602/−495 lines, 10 files
Files
- ChangeLog.md +6/−2
- extensible-effects-concurrent.cabal +5/−5
- src/Control/Eff/Concurrent/Dispatcher.hs +206/−0
- src/Control/Eff/Concurrent/Examples.hs +75/−0
- src/Control/Eff/Concurrent/GenServer.hs +211/−0
- src/Control/Eff/Concurrent/MessagePassing.hs +99/−0
- src/Control/Eff/Processes.hs +0/−72
- src/Control/Eff/Processes/Examples.hs +0/−68
- src/Control/Eff/Processes/STM.hs +0/−198
- src/Control/Eff/Processes/Server.hs +0/−150
ChangeLog.md view
@@ -1,11 +1,15 @@ # Changelog for extensible-effects-concurrent -## Unreleased Changes+## 0.1.1.0 +* Substantial API reoganization+* Rename/Move modules++## 0.1.0.1+ * Stack/Cabal/Github Cosmetics * Travis build job ## 0.1.0.0 * Initial Version-
extensible-effects-concurrent.cabal view
@@ -1,5 +1,5 @@ name: extensible-effects-concurrent-version: 0.1.0.1+version: 0.1.1.0 description: Please see the README on GitHub at <https://github.com/sheyll/extensible-effects-concurrent#readme> synopsis: Message passing concurrency as extensible-effect homepage: https://github.com/sheyll/extensible-effects-concurrent#readme@@ -41,12 +41,12 @@ tagged exposed-modules: Control.Eff.Interactive,- Control.Eff.Processes,- Control.Eff.Processes.STM,- Control.Eff.Processes.Server+ Control.Eff.Concurrent.GenServer,+ Control.Eff.Concurrent.MessagePassing,+ Control.Eff.Concurrent.Dispatcher other-modules: Paths_extensible_effects_concurrent,- Control.Eff.Processes.Examples+ Control.Eff.Concurrent.Examples other-extensions: ConstraintKinds, DeriveFoldable,
+ src/Control/Eff/Concurrent/Dispatcher.hs view
@@ -0,0 +1,206 @@+{-# LANGUAGE GeneralizedNewtypeDeriving #-}+{-# LANGUAGE ConstraintKinds #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE TypeFamilies #-}+{-# LANGUAGE KindSignatures #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE DeriveFunctor #-}+{-# LANGUAGE StandaloneDeriving #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE TypeOperators #-}+{-# LANGUAGE GADTs #-}+module Control.Eff.Concurrent.Dispatcher+ ( ProcIO+ , HasProcesses+ , Scheduler+ -- , nextPid, processTable+ , ProcessInfo+ , processId+ -- , messageQ+ , runProcesses+ , getProcessInfo+ , spawn+ , dispatchMessages+ -- , withMessageQueue+ )+where++import GHC.Stack+import Data.Kind+import Control.Concurrent as Concurrent+import Control.Concurrent.STM as STM+import Control.Eff+import Control.Eff.Lift+import Control.Eff.Concurrent.MessagePassing+import Control.Eff.Reader.Strict as Reader+import Control.Lens+import qualified Control.Monad.State as Mtl+import Data.Dynamic+import Data.Typeable ( typeRep )+import Data.Map ( Map )+import qualified Data.Map as Map++-- * MessagePassing Scheduling++data ProcessInfo =+ ProcessInfo { _processId :: ProcessId+ , _messageQ :: STM.TQueue Dynamic+ }++makeLenses ''ProcessInfo++instance Show ProcessInfo where+ show p =+ "ProcessInfo: " ++ show (p ^. processId)++data Scheduler =+ Scheduler { _nextPid :: ProcessId+ , _processTable :: Map ProcessId ProcessInfo+ }+ deriving Show+makeLenses ''Scheduler++type family Members (es :: [* -> *]) (r :: [* -> *]) :: Constraint where+ Members '[] r = ()+ Members (e ': es) r = (Member e r, Members es r)++type HasProcesses r = ( HasCallStack+ , SetMember Lift (Lift IO) r+ , Member (Reader (STM.TVar Scheduler)) r)++type ProcIO = '[MessagePassing, Process, Reader (STM.TVar Scheduler), Lift IO]++runProcesses+ :: (SetMember Lift (Lift IO) r, HasCallStack)+ => Eff (Reader (STM.TVar Scheduler) ': r) ()+ -> Eff r ()+runProcesses e = do+ v <- lift (newTVarIO (Scheduler 1 Map.empty))+ runReader e v++getProcessInfo :: HasProcesses r => ProcessId -> Eff r (Maybe ProcessInfo)+getProcessInfo pid = do+ p <- getScheduler+ return (p ^. processTable . at pid)++-- ** MessagePassing execution++spawn :: HasProcesses r => Eff ProcIO () -> Eff r ProcessId+spawn mfa = do+ processes <- ask+ pidVar <- lift newEmptyTMVarIO+ _threadId <- lift+ (Concurrent.forkIO+ (runLift+ (runReader+ (dispatchMessages+ (do+ pid <- self+ lift (atomically (STM.putTMVar pidVar pid))+ mfa+ )+ )+ processes+ )+ )+ )+ lift (atomically (STM.takeTMVar pidVar))++dispatchMessages+ :: forall r a+ . (HasProcesses r, HasCallStack)+ => Eff (MessagePassing ': Process ': r) a+ -> Eff r a+dispatchMessages processAction = withMessageQueue+ (\pinfo ->+ handle_relay return (goProc pinfo) (handle_relay return (go pinfo) processAction))+ where+ go+ :: forall v+ . HasCallStack+ => ProcessInfo+ -> MessagePassing v+ -> (v -> Eff (Process ': r) a)+ -> Eff (Process ': r) a+ go _ (SendMessage toPid reqIn) k = do+ psVar <- getSchedulerTVar+ lift+ (atomically+ (do+ p <- readTVar psVar+ let mto = p ^. processTable . at toPid+ case mto of+ Just toProc ->+ let dReq = toDyn reqIn+ in do+ writeTQueue (toProc ^. messageQ) dReq+ return True+ Nothing -> return False+ )+ )+ >>= k+ go (ProcessInfo selfPidInt channel) (ReceiveMessage onMsg) k = do+ mDynMsg <- lift (atomically (readTQueue channel))+ case fromDynamic mDynMsg of+ Just req -> let result = onMsg req in k (Just result)+ nix -> do+ lift+ (putStrLn+ ( show selfPidInt+ ++ " got unexpected msg: "+ ++ show mDynMsg+ ++ " expected: "+ ++ show (typeRep nix)+ )+ )+ k Nothing+ goProc+ :: forall v+ . HasCallStack+ => ProcessInfo+ -> Process v+ -> (v -> Eff r a)+ -> Eff r a+ goProc (ProcessInfo selfPidInt _) SelfPid k = k selfPidInt++withMessageQueue :: HasProcesses r => (ProcessInfo -> Eff r a) -> Eff r a+withMessageQueue m = do+ pinfo <- createQueue+ res <- m pinfo+ destroyQueue pinfo+ return res+ where+ createQueue = overScheduler+ (do+ pid <- nextPid <<+= 1+ channel <- Mtl.lift newTQueue+ let pinfo = ProcessInfo pid channel+ processTable . at pid .= Just pinfo+ return pinfo+ )+ destroyQueue pinfo =+ overScheduler (processTable . at (pinfo ^. processId) .= Nothing)+++overScheduler :: HasProcesses r => Mtl.StateT Scheduler STM.STM a -> Eff r a+overScheduler stAction = do+ psVar <- ask+ lift+ (STM.atomically+ (do+ ps <- STM.readTVar psVar+ (result, psModified) <- Mtl.runStateT stAction ps+ STM.writeTVar psVar psModified+ return result+ )+ )++getSchedulerTVar :: HasProcesses r => Eff r (TVar Scheduler)+getSchedulerTVar = ask++getScheduler :: HasProcesses r => Eff r Scheduler+getScheduler = do+ processesVar <- getSchedulerTVar+ lift (atomically (readTVar processesVar))
+ src/Control/Eff/Concurrent/Examples.hs view
@@ -0,0 +1,75 @@+{-# LANGUAGE GeneralizedNewtypeDeriving #-}+{-# LANGUAGE ConstraintKinds #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE TypeFamilies #-}+{-# LANGUAGE KindSignatures #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE DeriveFunctor #-}+{-# LANGUAGE StandaloneDeriving #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE TypeOperators #-}+{-# LANGUAGE TypeApplications #-}+{-# LANGUAGE GADTs #-}+module Control.Eff.Concurrent.Examples where++import GHC.Stack+import Control.Eff+import Control.Eff.Lift+import Control.Monad+import Data.Dynamic++import Control.Eff.Concurrent.MessagePassing+import Control.Eff.Concurrent.GenServer+import Control.Eff.Concurrent.Dispatcher++data TestApi+ deriving Typeable++data instance Api TestApi x where+ SayHello :: String -> Api TestApi ('Synchronous Bool)+ Shout :: String -> Api TestApi 'Asynchronous+ deriving (Typeable)++deriving instance Show (Api TestApi x)++runExample :: IO ()+runExample = runLift $ runProcesses $ dispatchMessages example++example+ :: ( HasCallStack+ , HasProcesses r+ , Member MessagePassing r+ , Member Process r+ , SetMember Lift (Lift IO) r)+ => Eff r ()+example = do+ me <- self+ lift (putStrLn ("I am " ++ show me))+ server <- asServer @TestApi <$> spawn testServerLoop+ lift (putStrLn ("Started server " ++ show server))+ let go = do+ x <- lift (putStr ">>> " >> getLine)+ res <- call server (SayHello x)+ lift (putStrLn ("Result: " ++ show res))+ case x of+ ('q':_) -> return ()+ _ -> go+ go++testServerLoop+ :: forall r. (HasCallStack, Member MessagePassing r, Member Process r, SetMember Lift (Lift IO) r)+ => Eff r ()+testServerLoop = forever $ serve_ $ ApiHandler handleCast handleCall+ where+ handleCast :: Api TestApi 'Asynchronous -> Eff r ()+ handleCast (Shout x) = do+ me <- self+ lift (putStrLn (show me ++ " Shouting: " ++ x))+ handleCall :: Api TestApi ('Synchronous x) -> (x -> Eff r Bool) -> Eff r ()+ handleCall (SayHello x) reply = do+ me <- self+ lift (putStrLn (show me ++ " Got Hello: " ++ x))+ void (reply (length x > 3))
+ src/Control/Eff/Concurrent/GenServer.hs view
@@ -0,0 +1,211 @@+{-# LANGUAGE GeneralizedNewtypeDeriving #-}+{-# LANGUAGE ConstraintKinds #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE TypeFamilies #-}+{-# LANGUAGE KindSignatures #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE DeriveFunctor #-}+{-# LANGUAGE StandaloneDeriving #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE TypeOperators #-}+{-# LANGUAGE TypeApplications #-}+{-# LANGUAGE GADTs #-}++-- | Type safe /server/ API processes++module Control.Eff.Concurrent.GenServer+ ( Api+ , Synchronicity(..)+ , ApiHandler(ApiHandler)+ , Server(..)+ , fromServer+ , proxyAsServer+ , asServer+ , cast+ , cast_+ , call+ , serve+ , serve_+ , unhandledCallError+ , unhandledCastError+ )+where++import GHC.Stack+import Data.Kind+import Control.Eff+import Control.Lens+import Control.Monad+import Data.Dynamic+import Data.Typeable ( typeRep )+import Data.Proxy++import Control.Eff.Concurrent.MessagePassing++-- | This data family defines an API implemented by a server.+-- The first parameter is the API /index/ and the second parameter+-- (the @* -> *@)+data family Api o :: Synchronicity -> *++data Synchronicity =+ Synchronous Type | Asynchronous+ deriving (Typeable)++newtype Server o = Server { _fromServer :: ProcessId }+ deriving (Eq,Ord,Typeable)++instance Read (Server o) where+ readsPrec _ ('[':'#':rest1) =+ case reads (dropWhile (/= '⇒') rest1) of+ [(c, ']':rest2)] -> [(Server c, rest2)]+ _ -> []+ readsPrec _ _ = []++instance Typeable o => Show (Server o) where+ show s@(Server c) =+ "[#" ++ show (typeRep s) ++ "⇒" ++ show c ++ "]"++makeLenses ''Server++proxyAsServer :: proxy api -> ProcessId -> Server api+proxyAsServer = const Server++asServer :: ProcessId -> Server api+asServer = Server++data Request p where+ Call :: forall p x . (Typeable p, Typeable x, Typeable (Api p ('Synchronous x)))+ => ProcessId -> Api p ('Synchronous x) -> Request p+ Cast :: forall p . (Typeable p, Typeable (Api p 'Asynchronous))+ => Api p 'Asynchronous -> Request p+ deriving Typeable++data Response p x where+ Response :: (Typeable p, Typeable x) => Proxy p -> x -> Response p x+ deriving Typeable++cast+ :: forall r o+ . ( HasCallStack+ , Member MessagePassing r+ , Typeable o+ , Typeable (Api o 'Asynchronous)+ )+ => Server o+ -> Api o 'Asynchronous+ -> Eff r Bool+cast (Server pid) callMsg = sendMessage pid (Cast callMsg)++cast_+ :: forall r o+ . ( HasCallStack+ , Member MessagePassing r+ , Typeable o+ , Typeable (Api o 'Asynchronous)+ )+ => Server o+ -> Api o 'Asynchronous+ -> Eff r ()+cast_ = ((.) . (.)) void cast++call+ :: forall result o r+ . ( Member MessagePassing r+ , Member Process r+ , Typeable o+ , Typeable (Api o ('Synchronous result))+ , Typeable result+ , HasCallStack+ )+ => Server o+ -> Api o ('Synchronous result)+ -> Eff r (Maybe result)+call (Server pidInt) req = do+ fromPid <- self+ let requestMessage = Call fromPid req+ wasSent <- sendMessage pidInt requestMessage+ if wasSent+ then+ let extractResult :: Response o result -> result+ extractResult (Response _pxResult result) = result+ in do mResp <- receiveMessage (Proxy @(Response o result))+ return (extractResult <$> mResp)+ else fail "Could not send request message " >> return Nothing++data ApiHandler p r e where+ ApiHandler ::+ { _handleCast+ :: (Typeable p, Typeable (Api p 'Asynchronous), HasCallStack)+ => Api p 'Asynchronous -> Eff r e+ , _handleCall+ :: forall x . (Typeable p, Typeable (Api p ('Synchronous x)), Typeable x, HasCallStack)+ => Api p ('Synchronous x) -> (x -> Eff r Bool) -> Eff r e+ } -> ApiHandler p r e++serve_+ :: forall r p+ . (Typeable p, Member MessagePassing r, Member Process r, HasCallStack)+ => ApiHandler p r ()+ -> Eff r ()+serve_ = void . serve++serve+ :: forall r p e+ . (Typeable p, Member MessagePassing r, Member Process r, HasCallStack)+ => ApiHandler p r e+ -> Eff r (Maybe e)+serve (ApiHandler handleCast handleCall) = do+ mReq <- receiveMessage (Proxy @(Request p))+ mapM receiveCallReq mReq+ where+ receiveCallReq :: Request p -> Eff r e+ receiveCallReq (Cast request ) = handleCast request+ receiveCallReq (Call fromPid request) = handleCall request+ (sendReply request)+ where+ sendReply :: Typeable x => Api p ('Synchronous x) -> x -> Eff r Bool+ sendReply _ reply = sendMessage fromPid (Response (Proxy :: Proxy p) reply)++unhandledCallError+ :: ( Show (Api p ('Synchronous x))+ , Typeable p+ , Typeable (Api p ('Synchronous x))+ , Typeable x+ , HasCallStack+ , Member Process r+ )+ => Api p ('Synchronous x)+ -> (x -> Eff r Bool)+ -> Eff r e+unhandledCallError api _ = do+ me <- self+ fail+ ( show me+ ++ " Unhandled call: ("+ ++ show api+ ++ " :: "+ ++ show (typeRep api)+ ++ ")"+ )++unhandledCastError+ :: ( Show (Api p 'Asynchronous)+ , Typeable p+ , Typeable (Api p 'Asynchronous)+ , HasCallStack+ , Member Process r+ )+ => Api p 'Asynchronous+ -> Eff r e+unhandledCastError api = do+ me <- self+ fail+ ( show me+ ++ " Unhandled cast: ("+ ++ show api+ ++ " :: "+ ++ show (typeRep api)+ ++ ")"+ )
+ src/Control/Eff/Concurrent/MessagePassing.hs view
@@ -0,0 +1,99 @@+{-# LANGUAGE GeneralizedNewtypeDeriving #-}+{-# LANGUAGE ConstraintKinds #-}+{-# LANGUAGE FlexibleContexts #-}+{-# LANGUAGE TypeFamilies #-}+{-# LANGUAGE KindSignatures #-}+{-# LANGUAGE DataKinds #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE DeriveFunctor #-}+{-# LANGUAGE StandaloneDeriving #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE TypeOperators #-}+{-# LANGUAGE GADTs #-}+module Control.Eff.Concurrent.MessagePassing+ ( ProcessId(..)+ , fromProcessId+ , Process(..)+ , self+ , MessagePassing(..)+ , sendMessage+ , kill+ , receiveMessage+ )+where++import GHC.Stack+import Control.Eff+import Control.Lens+import Data.Dynamic+import Data.Proxy+import Data.Void+import Text.Printf+++-- * Process Effects++data Process b where+ SelfPid :: Process ProcessId+ -- LinkProcesses :: ProcessId -> ProcessId -> Process ()++self :: Member Process r => Eff r ProcessId+self = send SelfPid++newtype ProcessId = ProcessId { _fromProcessId :: Int }+ deriving (Eq,Ord,Typeable,Bounded,Num, Enum, Integral, Real)++instance Read ProcessId where+ readsPrec _ ('<':'0':'.':rest1) =+ case reads rest1 of+ [(c, '.':'0':'>':rest2)] -> [(ProcessId c, rest2)]+ _ -> []+ readsPrec _ _ = []++instance Show ProcessId where+ show (ProcessId c) =+ printf "<0.%d.0>" c++makeLenses ''ProcessId+++-- * MessagePassing Effect++data MessagePassing b where+ SendMessage :: Typeable m+ => ProcessId+ -> Message m+ -> MessagePassing Bool+ ReceiveMessage+ :: forall e m . (Typeable m, Typeable (Message m))+ => (Message m -> e)+ -> MessagePassing (Maybe e)++data Message m where+ Shutdown :: Message Void+ Message :: m -> Message m+ deriving Typeable++sendMessage+ :: forall o r+ . (HasCallStack, Member MessagePassing r, Typeable o)+ => ProcessId+ -> o+ -> Eff r Bool+sendMessage pid message = send (SendMessage pid (Message message))++kill :: (HasCallStack, Member MessagePassing r)+ => ProcessId -> Eff r Bool+kill pid = send (SendMessage pid Shutdown)++receiveMessage+ :: forall o r . (HasCallStack, Member MessagePassing r, Member Process r, Typeable o)+ => Proxy o -> Eff r (Maybe o)+receiveMessage _ = do+ me <- self+ mRes <- send (ReceiveMessage id)+ case mRes of+ Just Shutdown -> fail (show me ++ " SHUTDOWN")+ Just (Message m) -> return (Just m)+ Nothing -> return Nothing
− src/Control/Eff/Processes.hs
@@ -1,72 +0,0 @@-{-# LANGUAGE GeneralizedNewtypeDeriving #-}-{-# LANGUAGE ConstraintKinds #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE TypeFamilies #-}-{-# LANGUAGE KindSignatures #-}-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE RankNTypes #-}-{-# LANGUAGE TemplateHaskell #-}-{-# LANGUAGE DeriveFunctor #-}-{-# LANGUAGE StandaloneDeriving #-}-{-# LANGUAGE ScopedTypeVariables #-}-{-# LANGUAGE TypeOperators #-}-{-# LANGUAGE GADTs #-}-module Control.Eff.Processes- ( ProcessId(..)- , self- , sendMessage- , receiveMessage- , Process(..)- , fromProcessId- )-where--import Control.Eff-import Control.Lens-import Data.Dynamic-import Data.Proxy-import Text.Printf---- * Process Types--newtype ProcessId = ProcessId { _fromProcessId :: Int }- deriving (Eq,Ord,Typeable,Bounded,Num, Enum, Integral, Real)--instance Read ProcessId where- readsPrec _ ('<':'0':'.':rest1) =- case reads rest1 of- [(c, '.':'0':'>':rest2)] -> [(ProcessId c, rest2)]- _ -> []- readsPrec _ _ = []--instance Show ProcessId where- show (ProcessId c) =- printf "<0.%d.0>" c--makeLenses ''ProcessId---- * Process Effect--data Process b where- SendMessage :: Typeable m- => ProcessId- -> m- -> Process Bool- ReceiveMessage- :: forall e m . (Typeable m)- => (m -> e)- -> Process (Maybe e)- SelfPid :: Process ProcessId---- * Process Effects--self :: Member Process r => Eff r ProcessId-self = send SelfPid--sendMessage- :: forall o r . (Member Process r, Typeable o) => ProcessId -> o -> Eff r Bool-sendMessage pid message = send (SendMessage pid message)--receiveMessage- :: forall o r . (Member Process r, Typeable o) => Proxy o -> Eff r (Maybe o)-receiveMessage _ = send (ReceiveMessage id)
− src/Control/Eff/Processes/Examples.hs
@@ -1,68 +0,0 @@-{-# LANGUAGE GeneralizedNewtypeDeriving #-}-{-# LANGUAGE ConstraintKinds #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE TypeFamilies #-}-{-# LANGUAGE KindSignatures #-}-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE RankNTypes #-}-{-# LANGUAGE TemplateHaskell #-}-{-# LANGUAGE DeriveFunctor #-}-{-# LANGUAGE StandaloneDeriving #-}-{-# LANGUAGE ScopedTypeVariables #-}-{-# LANGUAGE TypeOperators #-}-{-# LANGUAGE TypeApplications #-}-{-# LANGUAGE GADTs #-}-module Control.Eff.Processes.Examples where--import GHC.Stack-import Control.Eff-import Control.Eff.Lift-import Control.Monad-import Data.Dynamic-import Data.Void--import Control.Eff.Processes-import Control.Eff.Processes.Server-import Control.Eff.Processes.STM--data TestServer- deriving Typeable--data instance Api TestServer x where- SayHello :: String -> Api TestServer Bool- Shout :: String -> Api TestServer Void- deriving Typeable--example :: (HasCallStack, HasScheduler r, Member Process r, SetMember Lift (Lift IO) r) => Eff r ()-example = do- me <- self- lift (putStrLn ("I am " ++ show me))- server <- asServer @TestServer <$> spawn testa- lift (putStrLn ("Started server " ++ show server))- let go = do- x <- lift (putStr ">>> " >> getLine)- res <- call server (SayHello x)- lift (putStrLn ("Result: " ++ show res))- pid2 <- asServer @TestServer <$>- spawn (do Just m <- serve (\(Shout m) k -> (m, k undefined))- cast_ server (Shout ("Who: " ++ m)))- x2 <- lift (putStr "<<< " >> getLine)- cast_ pid2 (Shout x2)- case x2 of- ('q':_) -> return ()- _ -> go- go--testa :: (HasCallStack, Member Process r, SetMember Lift (Lift IO) r) => Eff r ()-testa =- forever $- do me <- self- Just msg <- serve (testaRespond me)- lift (putStrLn msg)- return ()--testaRespond :: ProcessId -> Api TestServer x -> (x -> z) -> (String, z)-testaRespond me (SayHello x) k =- (show me ++ " Got Hello: " ++ x, k (length x > 3))-testaRespond me (Shout x) k =- (show me ++ " Shouting: " ++ x, k undefined )
− src/Control/Eff/Processes/STM.hs
@@ -1,198 +0,0 @@-{-# LANGUAGE GeneralizedNewtypeDeriving #-}-{-# LANGUAGE ConstraintKinds #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE TypeFamilies #-}-{-# LANGUAGE KindSignatures #-}-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE RankNTypes #-}-{-# LANGUAGE TemplateHaskell #-}-{-# LANGUAGE DeriveFunctor #-}-{-# LANGUAGE StandaloneDeriving #-}-{-# LANGUAGE ScopedTypeVariables #-}-{-# LANGUAGE TypeOperators #-}-{-# LANGUAGE GADTs #-}-module Control.Eff.Processes.STM- ( ProcessIOEffects- , HasScheduler- , Scheduler- -- , nextPid, processTable- , ProcessInfo- , processId- -- , messageQ- , runScheduler- , getProcessInfo- , spawn- , dispatchMessages- -- , withMessageQueue- )-where--import GHC.Stack-import Data.Kind-import Control.Concurrent as Concurrent-import Control.Concurrent.STM as STM-import Control.Eff-import Control.Eff.Lift-import Control.Eff.Processes-import Control.Eff.Reader.Strict as Reader-import Control.Lens-import qualified Control.Monad.State as Mtl-import Data.Dynamic-import Data.Typeable ( typeRep )-import Data.Map ( Map )-import qualified Data.Map as Map---- * Process Scheduling--data ProcessInfo =- ProcessInfo { _processId :: ProcessId- , _messageQ :: STM.TQueue Dynamic- }--makeLenses ''ProcessInfo--instance Show ProcessInfo where- show p =- "Process: " ++ show (p ^. processId)--data Scheduler =- Scheduler { _nextPid :: ProcessId- , _processTable :: Map ProcessId ProcessInfo- }- deriving Show-makeLenses ''Scheduler--type family Members (es :: [* -> *]) (r :: [* -> *]) :: Constraint where- Members '[] r = ()- Members (e ': es) r = (Member e r, Members es r)--type HasScheduler r = ( HasCallStack- , SetMember Lift (Lift IO) r- , Member (Reader (STM.TVar Scheduler)) r)--type ProcessIOEffects = '[Process, Reader (STM.TVar Scheduler), Lift IO]--runScheduler- :: (SetMember Lift (Lift IO) r, HasCallStack)- => Eff (Reader (STM.TVar Scheduler) ': r) ()- -> Eff r ()-runScheduler e = do- v <- lift (newTVarIO (Scheduler 1 Map.empty))- runReader e v--getProcessInfo :: HasScheduler r => ProcessId -> Eff r (Maybe ProcessInfo)-getProcessInfo pid = do- p <- getScheduler- return (p ^. processTable . at pid)---- ** Process execution--spawn :: HasScheduler r => Eff ProcessIOEffects () -> Eff r ProcessId-spawn mfa = do- processes <- ask- pidVar <- lift newEmptyTMVarIO- _threadId <- lift- (Concurrent.forkIO- (runLift- (runReader- (dispatchMessages- (do- pid <- self- lift (atomically (STM.putTMVar pidVar pid))- mfa- )- )- processes- )- )- )- lift (atomically (STM.takeTMVar pidVar))--dispatchMessages- :: forall r a- . (HasScheduler r, HasCallStack)- => Eff (Process ': r) a- -> Eff r a-dispatchMessages processAction = withMessageQueue- (\pinfo -> handle_relay return (go pinfo) processAction)- where- go- :: forall v- . HasCallStack- => ProcessInfo- -> Process v- -> (v -> Eff r a)- -> Eff r a- go (ProcessInfo selfPidInt _) SelfPid k = k selfPidInt- go _ (SendMessage toPid reqIn) k = do- psVar <- getSchedulerTVar- lift- (atomically- (do- p <- readTVar psVar- let mto = p ^. processTable . at toPid- case mto of- Just toProc ->- let dReq = toDyn reqIn- in do- writeTQueue (toProc ^. messageQ) dReq- return True- Nothing -> return False- )- )- >>= k- go (ProcessInfo selfPidInt channel) (ReceiveMessage onMsg) k = do- mDynMsg <- lift (atomically (readTQueue channel))- case fromDynamic mDynMsg of- Just req -> let result = onMsg req in k (Just result)- nix -> do- lift- (putStrLn- ( show selfPidInt- ++ " got unexpected msg: "- ++ show mDynMsg- ++ " expected: "- ++ show (typeRep nix)- )- )- k Nothing--withMessageQueue :: HasScheduler r => (ProcessInfo -> Eff r a) -> Eff r a-withMessageQueue m = do- pinfo <- createQueue- res <- m pinfo- destroyQueue pinfo- return res- where- createQueue = overScheduler- (do- pid <- nextPid <<+= 1- channel <- Mtl.lift newTQueue- let pinfo = ProcessInfo pid channel- processTable . at pid .= Just pinfo- return pinfo- )- destroyQueue pinfo =- overScheduler (processTable . at (pinfo ^. processId) .= Nothing)---overScheduler :: HasScheduler r => Mtl.StateT Scheduler STM.STM a -> Eff r a-overScheduler stAction = do- psVar <- ask- lift- (STM.atomically- (do- ps <- STM.readTVar psVar- (result, psModified) <- Mtl.runStateT stAction ps- STM.writeTVar psVar psModified- return result- )- )--getSchedulerTVar :: HasScheduler r => Eff r (TVar Scheduler)-getSchedulerTVar = ask--getScheduler :: HasScheduler r => Eff r Scheduler-getScheduler = do- processesVar <- getSchedulerTVar- lift (atomically (readTVar processesVar))
− src/Control/Eff/Processes/Server.hs
@@ -1,150 +0,0 @@-{-# LANGUAGE GeneralizedNewtypeDeriving #-}-{-# LANGUAGE ConstraintKinds #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE TypeFamilies #-}-{-# LANGUAGE KindSignatures #-}-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE RankNTypes #-}-{-# LANGUAGE TemplateHaskell #-}-{-# LANGUAGE DeriveFunctor #-}-{-# LANGUAGE StandaloneDeriving #-}-{-# LANGUAGE ScopedTypeVariables #-}-{-# LANGUAGE TypeOperators #-}-{-# LANGUAGE GADTs #-}---- | Type safe /server/ API processes--module Control.Eff.Processes.Server- ( Api- , Server(..)- , fromServer- , proxyAsServer- , asServer- , cast- , cast_- , call- , serve- )-where--import GHC.Stack-import Data.Kind-import Control.Eff-import Control.Lens-import Control.Monad-import Data.Dynamic-import Data.Typeable ( typeRep )-import Data.Proxy--import Control.Eff.Processes---- | This data family defines an API implemented by a server.--- The first parameter is the API /index/ and the second parameter--- (the @* -> *@)-data family Api o :: * -> *--newtype Server o = Server { _fromServer :: ProcessId }- deriving (Eq,Ord,Typeable)--instance Read (Server o) where- readsPrec _ ('[':'#':rest1) =- case reads (dropWhile (/= '⇒') rest1) of- [(c, ']':rest2)] -> [(Server c, rest2)]- _ -> []- readsPrec _ _ = []--instance Typeable o => Show (Server o) where- show s@(Server c) =- "[#" ++ show (typeRep s) ++ "⇒" ++ show c ++ "]"--makeLenses ''Server--proxyAsServer :: proxy api -> ProcessId -> Server api-proxyAsServer = const Server--asServer :: ProcessId -> Server api-asServer = Server--data Request p where- Call :: forall p x . (Typeable p, Typeable x, Typeable (Api p x)) => Maybe ProcessId -> Api p x -> Request p- Cast :: forall p x . (Typeable p, Typeable x, Typeable (Api p x)) => Api p x -> Request p- deriving Typeable--data Response p x where- Response :: (Typeable p, Typeable x) => Proxy p -> x -> Response p x- deriving Typeable--cast- :: forall r o result- . ( HasCallStack- , Member Process r- , Typeable o- , Typeable result- , Typeable (Api o result)- )- => Server o- -> Api o result- -> Eff r Bool-cast (Server pid) callMsg = sendMessage pid (Cast callMsg)--cast_- :: forall r o result- . ( HasCallStack- , Member Process r- , Typeable o- , Typeable result- , Typeable (Api o result)- )- => Server o- -> Api o result- -> Eff r ()-cast_ = ((.) . (.)) void cast--call- :: forall result o r- . ( Member Process r- , Typeable o- , Typeable (Api o result)- , Typeable result- , HasCallStack- )- => Server o- -> Api o result- -> Eff r (Maybe result)-call (Server pidInt) req = do- fromPid <- self- let requestMessage = Call (Just fromPid) req- wasSent <- sendMessage pidInt requestMessage- if wasSent- then- let responseMessage :: Response o result -> result- responseMessage (Response _pxReply reply) = reply- in send (ReceiveMessage responseMessage)- else fail "Could not send request message " >> return Nothing--serve- :: forall r p e- . (Member Process r, Typeable p, HasCallStack)- => ( forall x- . (Typeable (Api p x), Typeable x)- => Api p x- -> (x -> Maybe (Process Bool))- -> (e, Maybe (Process Bool))- )- -> Eff r (Maybe e)-serve handle = do- mrequest <- send (ReceiveMessage receiveCallReq)- case mrequest of- Just (result, mReplyCast) -> do- mapM_ (void . send) mReplyCast- return (Just result)- _ -> fail "sdf4446" >> return Nothing- where- receiveCallReq :: Request p -> (e, Maybe (Process Bool))- receiveCallReq (Cast request ) = handle request (const Nothing)- receiveCallReq (Call mFromPid request) = handle request (packReply request)- where- packReply :: Typeable x => Api p x -> x -> Maybe (Process Bool)- packReply _ reply = do- fromPid <- mFromPid- return (SendMessage fromPid (Response (Proxy :: Proxy p) reply))