network-messagepack-rpc (empty) → 0.1.0.0
raw patch · 5 files changed
+374/−0 lines, 5 filesdep +basedep +bytestringdep +data-msgpacksetup-changed
Dependencies added: base, bytestring, data-msgpack, safe-exceptions, text, unordered-containers
Files
- LICENSE +29/−0
- Setup.hs +2/−0
- network-messagepack-rpc.cabal +30/−0
- src/Data/MessagePack/RPC.hs +129/−0
- src/Network/MessagePack/RPC/Client.hs +184/−0
+ LICENSE view
@@ -0,0 +1,29 @@+Copyright (c) 2018, IIJ Innovation Institute Inc.+All rights reserved.++Redistribution and use in source and binary forms, with or without+modification, are permitted provided that the following conditions+are met:++ * Redistributions of source code must retain the above copyright+ notice, this list of conditions and the following disclaimer.+ * Redistributions in binary form must reproduce the above copyright+ notice, this list of conditions and the following disclaimer in+ the documentation and/or other materials provided with the+ distribution.+ * Neither the name of the copyright holders nor the names of its+ contributors may be used to endorse or promote products derived+ from this software without specific prior written permission.++THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS+"AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT+LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS+FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE+COPYRIGHT OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT,+INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING,+BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES;+LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER+CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT+LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN+ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE+POSSIBILITY OF SUCH DAMAGE.
+ Setup.hs view
@@ -0,0 +1,2 @@+import Distribution.Simple+main = defaultMain
+ network-messagepack-rpc.cabal view
@@ -0,0 +1,30 @@+name: network-messagepack-rpc+version: 0.1.0.0+synopsis: MessagePack RPC+description: [MessagePack RPC](https://github.com/msgpack-rpc/msgpack-rpc/blob/master/spec.md) library based on the "data-msgpack" package.+homepage: https://github.com/iij-ii/network-messagepack-rpc+license: BSD3+license-file: LICENSE+author: Yuji Yamamoto and Kazu Yamamoto+maintainer: yuji-yamamoto@iij.ad.jp, kazu@iij.ad.jp+category: Data Network+build-type: Simple+cabal-version: >=1.10++library+ hs-source-dirs: src+ ghc-options: -Wall+ exposed-modules: Data.MessagePack.RPC+ , Network.MessagePack.RPC.Client+ build-depends:+ base >= 4.7 && < 5+ , bytestring+ , data-msgpack+ , safe-exceptions+ , text+ , unordered-containers+ default-language: Haskell2010++source-repository head+ type: git+ location: https://github.com/iij-ii/network-messagepack-rpc
+ src/Data/MessagePack/RPC.hs view
@@ -0,0 +1,129 @@+-- | Types in [MessagePack RPC](https://github.com/msgpack-rpc/msgpack-rpc/blob/master/spec.md)++module Data.MessagePack.RPC (+ MessageId+ , MethodName+ , Message(..)+ ) where++import Data.MessagePack (MessagePack (..), Object (..))+import Data.List (intercalate)+import qualified Data.Text as T+import Data.Word (Word64)++-- | Message ID.+type MessageId = Word64++-- | Method name.+type MethodName = T.Text++-- | Message type of MessagePack PRC.+-- Use 'toObject' and 'fromObject' for conversion.+data Message =+ -- | Request+ RequestMessage MessageId MethodName [Object]+ -- | Response. 'Left' is an error. 'Right' is a result.+ | ResponseMessage MessageId (Either Object Object)+ -- | Notification.+ | NotificationMessage MethodName [Object]+ deriving Eq++instance MessagePack Message where+ toObject (RequestMessage mid methodName args) =+ ObjectArray+ [ ObjectWord 0+ , ObjectWord mid+ , ObjectStr methodName+ , ObjectArray args+ ]++ toObject (ResponseMessage mid (Right result)) =+ ObjectArray+ [ ObjectWord 1+ , ObjectWord mid+ , ObjectNil+ , result+ ]++ toObject (ResponseMessage mid (Left err)) =+ ObjectArray+ [ ObjectWord 1+ , ObjectWord mid+ , err+ , ObjectNil+ ]++ toObject (NotificationMessage methodName params) =+ ObjectArray+ [ ObjectWord 2+ , ObjectStr methodName+ , ObjectArray params+ ]++ fromObject+ ( ObjectArray+ [ ObjectWord 0+ , ObjectWord mid+ , ObjectStr methodName+ , ObjectArray args+ ]+ ) =+ return $ RequestMessage mid methodName args++ fromObject+ ( ObjectArray+ [ ObjectWord 1+ , ObjectWord mid+ , ObjectNil+ , result+ ]+ ) =+ return $ ResponseMessage mid (Right result)+ fromObject+ ( ObjectArray+ [ ObjectWord 1+ , ObjectWord mid+ , err+ , ObjectNil+ ]+ ) =+ return $ ResponseMessage mid (Left err)++ fromObject+ ( ObjectArray+ [ ObjectWord 2+ , ObjectStr methodName+ , ObjectArray params+ ]+ ) =+ return $ NotificationMessage methodName params++ fromObject other =+ fail $ "Unexpected object:" ++ show other++instance Show Message where+ show (RequestMessage mid method objs) =+ "request(" ++ show mid ++ ") " ++ T.unpack method ++ " " ++ showObjs objs+ show (ResponseMessage mid (Left obj)) =+ "response error(" ++ show mid ++ ") " ++ showObj obj+ show (ResponseMessage mid (Right obj)) =+ "response(" ++ show mid ++ ") " ++ showObj obj+ show (NotificationMessage method objs) =+ "notification " ++ T.unpack method ++ " " ++ showObjs objs++showObjs :: [Object] -> String+showObjs objs = "[" ++ intercalate "," (map showObj objs) ++ "]"++showObj :: Object -> String+showObj (ObjectWord w) = "+" ++ show w+showObj (ObjectInt n) = show n+showObj ObjectNil = "nil"+showObj (ObjectBool b) = show b+showObj (ObjectStr s) = "\"" ++ T.unpack s ++ "\""+showObj (ObjectArray v) = "[" ++ intercalate "," (map showObj v) ++ "]"+showObj (ObjectMap m) = "{" ++ intercalate "," (map showPair m) ++ "}"+ where showPair (x, y) = "(" ++ showObj x ++ "," ++ showObj y ++ ")"+showObj (ObjectBin _ ) = error "ObjectBin"+showObj (ObjectExt _ _ ) = error "ObjectExt"+showObj (ObjectFloat _) = error "ObjectFloat"+showObj (ObjectDouble _) = error "ObjectDouble"
+ src/Network/MessagePack/RPC/Client.hs view
@@ -0,0 +1,184 @@+{-# LANGUAGE OverloadedStrings #-}++-- | Backend-free MessagePack RPC Client.+module Network.MessagePack.RPC.Client+ (+ -- * Config+ Config(..)+ , NotificationHandler+ , RequestHandler+ , Logger+ , Formatter+ , defaultConfig+ -- * Backend+ , Backend(..)+ -- * Client+ , Client+ , withClient+ -- * Call and reply+ , Result+ , call+ , reply+ ) where++import Control.Concurrent (forkIO, killThread)+import Control.Concurrent.MVar (MVar)+import qualified Control.Concurrent.MVar as MVar+import qualified Control.Exception.Safe as E+import Control.Monad (forever, void)+import qualified Data.ByteString as B+import qualified Data.ByteString.Lazy as BL+import Data.HashMap.Strict (HashMap)+import qualified Data.HashMap.Strict as HM+import Data.IORef (IORef)+import qualified Data.IORef as IORef+import qualified Data.MessagePack as MsgPack+import Data.Monoid ((<>))+import System.Timeout (timeout)++import Data.MessagePack.RPC++-- | A client data type for MessagePack RPC.+data Client = Client {+ clientSessionState :: !SessionState+ , clientBackend :: !Backend+ , clientLog :: Logger+ , clientFormat :: Formatter+ }++data SessionState = SessionState {+ lastMessageId :: IORef MessageId+ , dispatchTable :: IORef (HashMap MessageId (MVar Result))+ }++-- | Result type of a RPC call.+-- Described as "error" and "result" of "Response Message"+-- in [the spec of MessagePack RPC](https://github.com/msgpack-rpc/msgpack-rpc/blob/master/spec.md#response-message).+type Result = Either MsgPack.Object MsgPack.Object++-- | Notification handler. The 3rd argument is response objects.+type NotificationHandler = Client -> MethodName -> [MsgPack.Object] -> IO ()++-- | Notification handler. The 2nd argument is message id to be used+-- for replying. The 3rd argument is response objects.+type RequestHandler = Client -> MessageId -> MethodName -> [MsgPack.Object] -> IO ()++-- | Logger type. Should print out the message passed as a first argument somewhere.+type Logger = String -> IO ()++-- | Convert 'Message' into a @String@ to print out by 'Logger'+type Formatter = Message -> String++-- | Configuration for MessagePack RPC.+data Config = Config {+ notificationHandler :: NotificationHandler+ , requestHandler :: RequestHandler+ , logger :: Logger+ , exceptionHandlers :: [E.Handler IO ()]+ -- ^ Handles an exception thrown from the receiver thread,+ -- which is the only thread to receive 'Message's via 'Backend'.+ -- Until exiting from the block of 'withClient', the receiver thread+ -- indefinitely waits for frames via 'backendRecv'.+ , formatter :: Formatter+ }++-- | The default configuration.+-- 'formatter' is 'show'.+-- Others do nothing.+defaultConfig :: Config+defaultConfig = Config+ { notificationHandler = \_ _ _ -> return ()+ , requestHandler = \_ _ _ _ -> return ()+ , logger = \_ -> return ()+ , exceptionHandlers = [E.Handler $ \(E.SomeException _) -> return ()]+ , formatter = show+ }++-- | Backend IO functions.+data Backend = Backend {+ backendSend :: B.ByteString -> IO () -- ^ Sending+ , backendRecv :: IO B.ByteString -- ^ Receiving+ , backendClose :: IO () -- ^ Closing+ }++-- TODO: Returns any exception+-- | Calling RPC.+call :: Client -> MethodName -> [MsgPack.Object] -> IO Result+call client funName args = do+ rrsp <- E.bracket register unregister sendAndRecv+ case rrsp of+ Nothing -> return $ Left MsgPack.ObjectNil+ Just rsp -> return rsp+ where+ sendAndRecv (requestId, rspVar) = do+ let request = RequestMessage requestId funName args+ backendSend (clientBackend client) $ BL.toStrict $ MsgPack.pack request+ clientLog client $ "sent: " <> clientFormat client request+ timeout 3000000 $ MVar.takeMVar rspVar+ register = do+ requestId <- getNewMessageId ss+ rspVar <- MVar.newEmptyMVar+ IORef.atomicModifyIORef' (dispatchTable ss)+ $ \tbl -> (HM.insert requestId rspVar tbl, ())+ return (requestId, rspVar)+ unregister (requestId, _) = IORef.atomicModifyIORef' (dispatchTable ss)+ $ \tbl -> (HM.delete requestId tbl, ())+ ss = clientSessionState client++-- | Replying RPC. This should be used in 'RequestHandler'.+reply :: Client -> MessageId -> Result -> IO ()+reply client mid result = do+ let response = ResponseMessage mid result+ let p = BL.toStrict $ MsgPack.pack response+ backendSend (clientBackend client) p+ clientLog client $ "sent: " <> clientFormat client response++getNewMessageId :: SessionState -> IO MessageId+getNewMessageId ss =+ IORef.atomicModifyIORef (lastMessageId ss) $ \cur -> (cur + 1, cur)++receiverThread :: Client -> Config -> IO ()+receiverThread client config =+ (`E.catches` exceptionHandlers config) $ forever $ do+ response <- MsgPack.unpack . BL.fromStrict =<< backendRecv+ (clientBackend client)+ clientLog client $ "received: " <> clientFormat client response+ case response of+ ResponseMessage mid result -> do+ tbl <- IORef.readIORef $ dispatchTable ss+ case HM.lookup mid tbl of+ Just rspVar -> MVar.putMVar rspVar result+ Nothing ->+ clientLog client+ $ "ERROR: No MVar assinged with request ID "+ ++ show mid+ ++ "."+ NotificationMessage methodName params ->+ void . forkIO $ notificationHandler config+ client+ methodName+ params+ RequestMessage mid methodName params ->+ void . forkIO $ requestHandler config+ client+ mid+ methodName+ params+ where ss = clientSessionState client++initSessionState :: IO SessionState+initSessionState =+ SessionState <$> IORef.newIORef 0 <*> IORef.newIORef HM.empty++-- | Executing the action in the 3rd argument with a 'Client'.+withClient :: Config -> Backend -> (Client -> IO a) -> IO a+withClient config backend action = do+ ss <- initSessionState+ let client = Client ss backend (logger config) (formatter config)+ tid <- forkIO $ receiverThread client config+ takeAction client `E.finally` killThread tid+ where+ takeAction client = do+ returned <- action client+ backendClose backend+ return returned