packages feed

hadoop-rpc 0.1.1.1 → 1.0.0.0

raw patch · 6 files changed

+180/−159 lines, 6 filesdep +randomdep +uuidPVP ok

version bump matches the API change (PVP)

Dependencies added: random, uuid

API changes (from Hackage documentation)

- Data.Hadoop.Protobuf.Headers: RpcRequest :: Required 1 (Value Text) -> Optional 2 (Value ByteString) -> Required 3 (Value Text) -> Required 4 (Value Word64) -> RpcRequest
- Data.Hadoop.Protobuf.Headers: data RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Constructor C1_0RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Datatype D1RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Decode RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Encode RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Generic RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Selector S1_0_0RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Selector S1_0_1RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Selector S1_0_2RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Selector S1_0_3RpcRequest
- Data.Hadoop.Protobuf.Headers: instance Show RpcRequest
- Data.Hadoop.Protobuf.Headers: reqBytes :: RpcRequest -> Optional 2 (Value ByteString)
- Network.Hadoop.Rpc: initConnectionV7 :: HadoopConfig -> Protocol -> Socket -> IO Connection
- Network.Hadoop.Socket: SockAddrInet :: PortNumber -> HostAddress -> SockAddr
- Network.Hadoop.Socket: SockAddrInet6 :: PortNumber -> FlowInfo -> HostAddress6 -> ScopeID -> SockAddr
- Network.Hadoop.Socket: SockAddrUnix :: String -> SockAddr
- Network.Hadoop.Socket: data SockAddr :: *
- Network.Hadoop.Socket: newSocket :: AddrInfo -> IO Socket
- Network.Hadoop.Socket: runTcp :: (MonadMask m, MonadIO m) => Maybe SocksProxy -> Endpoint -> (Socket -> m a) -> m a
+ Data.Hadoop.Protobuf.Headers: ErrorApplication :: Error
+ Data.Hadoop.Protobuf.Headers: ErrorCode :: Int -> Error
+ Data.Hadoop.Protobuf.Headers: ErrorNoSuchMethod :: Error
+ Data.Hadoop.Protobuf.Headers: ErrorNoSuchProtocol :: Error
+ Data.Hadoop.Protobuf.Headers: ErrorRpcServer :: Error
+ Data.Hadoop.Protobuf.Headers: ErrorRpcVersionMismatch :: Error
+ Data.Hadoop.Protobuf.Headers: ErrorSerializingResponse :: Error
+ Data.Hadoop.Protobuf.Headers: FatalCode :: Int -> Error
+ Data.Hadoop.Protobuf.Headers: FatalDeserializingRequest :: Error
+ Data.Hadoop.Protobuf.Headers: FatalInvalidRpcHeader :: Error
+ Data.Hadoop.Protobuf.Headers: FatalUnauthorized :: Error
+ Data.Hadoop.Protobuf.Headers: FatalUnknown :: Error
+ Data.Hadoop.Protobuf.Headers: FatalUnsupportedSerialization :: Error
+ Data.Hadoop.Protobuf.Headers: FatalVersionMismatch :: Error
+ Data.Hadoop.Protobuf.Headers: RequestHeader :: Required 1 (Value Text) -> Required 2 (Value Text) -> Required 3 (Value Word64) -> RequestHeader
+ Data.Hadoop.Protobuf.Headers: data Error
+ Data.Hadoop.Protobuf.Headers: data RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_0Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_0RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_10Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_11Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_12Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_13Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_1Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_2Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_3Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_4Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_5Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_6Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_7Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_8Error
+ Data.Hadoop.Protobuf.Headers: instance Constructor C1_9Error
+ Data.Hadoop.Protobuf.Headers: instance Datatype D1Error
+ Data.Hadoop.Protobuf.Headers: instance Datatype D1RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Decode RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Encode RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Enum Error
+ Data.Hadoop.Protobuf.Headers: instance Generic Error
+ Data.Hadoop.Protobuf.Headers: instance Generic RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_0RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_1RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_2RequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_3RpcRequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_3RpcResponseHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_4RpcRequestHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_4RpcResponseHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_5RpcResponseHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_6RpcResponseHeader
+ Data.Hadoop.Protobuf.Headers: instance Selector S1_0_7RpcResponseHeader
+ Data.Hadoop.Protobuf.Headers: instance Show Error
+ Data.Hadoop.Protobuf.Headers: instance Show RequestHeader
+ Data.Hadoop.Protobuf.Headers: reqClientId :: RpcRequestHeader -> Required 4 (Value ByteString)
+ Data.Hadoop.Protobuf.Headers: reqRetryCount :: RpcRequestHeader -> Optional 5 (Value (Signed Int32))
+ Data.Hadoop.Protobuf.Headers: rspClientId :: RpcResponseHeader -> Optional 7 (Value ByteString)
+ Data.Hadoop.Protobuf.Headers: rspErrorDetail :: RpcResponseHeader -> Optional 6 (Enumeration Error)
+ Data.Hadoop.Protobuf.Headers: rspErrorMsg :: RpcResponseHeader -> Optional 5 (Value Text)
+ Data.Hadoop.Protobuf.Headers: rspExceptionClassName :: RpcResponseHeader -> Optional 4 (Value Text)
+ Data.Hadoop.Protobuf.Headers: rspRetryCount :: RpcResponseHeader -> Optional 8 (Value Int32)
+ Network.Hadoop.Rpc: initConnectionV9 :: HadoopConfig -> Protocol -> Socket -> IO Connection
+ Network.Hadoop.Socket: bracketSocket :: (MonadMask m, MonadIO m) => Maybe SocksProxy -> Endpoint -> (Socket -> m a) -> m a
- Data.Hadoop.Protobuf.Headers: RpcRequestHeader :: Optional 1 (Enumeration RpcKind) -> Optional 2 (Enumeration RpcOperation) -> Required 3 (Value Word32) -> RpcRequestHeader
+ Data.Hadoop.Protobuf.Headers: RpcRequestHeader :: Optional 1 (Enumeration RpcKind) -> Optional 2 (Enumeration RpcOperation) -> Required 3 (Value (Signed Int32)) -> Required 4 (Value ByteString) -> Optional 5 (Value (Signed Int32)) -> RpcRequestHeader
- Data.Hadoop.Protobuf.Headers: RpcResponseHeader :: Required 1 (Value Word32) -> Required 2 (Enumeration RpcStatus) -> Optional 3 (Value Word32) -> RpcResponseHeader
+ Data.Hadoop.Protobuf.Headers: RpcResponseHeader :: Required 1 (Value Word32) -> Required 2 (Enumeration RpcStatus) -> Optional 3 (Value Word32) -> Optional 4 (Value Text) -> Optional 5 (Value Text) -> Optional 6 (Enumeration Error) -> Optional 7 (Value ByteString) -> Optional 8 (Value Int32) -> RpcResponseHeader
- Data.Hadoop.Protobuf.Headers: reqCallId :: RpcRequestHeader -> Required 3 (Value Word32)
+ Data.Hadoop.Protobuf.Headers: reqCallId :: RpcRequestHeader -> Required 3 (Value (Signed Int32))
- Data.Hadoop.Protobuf.Headers: reqMethodName :: RpcRequest -> Required 1 (Value Text)
+ Data.Hadoop.Protobuf.Headers: reqMethodName :: RequestHeader -> Required 1 (Value Text)
- Data.Hadoop.Protobuf.Headers: reqProtocolName :: RpcRequest -> Required 3 (Value Text)
+ Data.Hadoop.Protobuf.Headers: reqProtocolName :: RequestHeader -> Required 2 (Value Text)
- Data.Hadoop.Protobuf.Headers: reqProtocolVersion :: RpcRequest -> Required 4 (Value Word64)
+ Data.Hadoop.Protobuf.Headers: reqProtocolVersion :: RequestHeader -> Required 3 (Value Word64)
- Network.Hadoop.Socket: closeSocket :: Socket -> IO ()
+ Network.Hadoop.Socket: closeSocket :: MonadIO m => Socket -> m ()
- Network.Hadoop.Socket: connectSocket :: Endpoint -> IO (Socket, SockAddr)
+ Network.Hadoop.Socket: connectSocket :: (MonadMask m, MonadIO m) => Maybe SocksProxy -> Endpoint -> m Socket

Files

hadoop-rpc.cabal view
@@ -1,5 +1,5 @@ name:          hadoop-rpc-version:       0.1.1.1+version:       1.0.0.0  synopsis:   Use the Hadoop RPC interface from Haskell.@@ -7,9 +7,9 @@ description:   Use the Hadoop RPC interface from Haskell.   .-  Currently we only support v7 of the RPC protocol (< CDH5).+  This version only supports v9 of the RPC protocol (CDH 5.x and above).   .-  Support for v9 (>= CDH5) is coming soon.+  Use hadoop-rpc-0.x.x.x if you need v7 support.   .   /The API is still evolving and is highly likely to change in the future./ @@ -53,12 +53,14 @@     , hashable             >= 1.2.1     , network              >= 2.5     , protobuf             >= 0.2.0.4+    , random               >= 1.0.1.1     , socks                >= 0.5     , stm                  >= 2.4     , text                 >= 1.1     , transformers         >= 0.4     , unix                 >= 2.7     , unordered-containers >= 0.2+    , uuid                 >= 1.3.4     , vector               >= 0.10     , xmlhtml              >= 0.2 
src/Data/Hadoop/Protobuf/Headers.hs view
@@ -4,6 +4,7 @@ module Data.Hadoop.Protobuf.Headers where  import Data.ByteString (ByteString)+import Data.Int (Int32) import Data.ProtocolBuffers import Data.ProtocolBuffers.Orphans () import Data.Text (Text)@@ -52,30 +53,43 @@ data RpcRequestHeader = RpcRequestHeader     { reqKind       :: Optional 1 (Enumeration RpcKind)     , reqOp         :: Optional 2 (Enumeration RpcOperation)-    , reqCallId     :: Required 3 (Value Word32)     -- ^ A sequence number that is sent back in the response--    -- Fields below don't apply until v9-    --, reqClientId   :: Required 4 (Value ByteString) -- ^ Globally unique client ID-    --, reqRetryCount :: Optional 5 (Value Int32)      -- ^ Retry count, 1 means this is the first retry+    , reqCallId     :: Required 3 (Value (Signed Int32)) -- ^ Sequence number that is sent back in response+    , reqClientId   :: Required 4 (Value ByteString)     -- ^ Globally unique client ID+    , reqRetryCount :: Optional 5 (Value (Signed Int32)) -- ^ Retry count, 1 means this is the first retry     } deriving (Generic, Show)  instance Encode RpcRequestHeader instance Decode RpcRequestHeader --- | This message is used for Protobuf RPC Engine.--- The message is used to marshal a RPC request from RPC client to the--- RPC server. The response to the RPC call (including errors) are handled--- as part of the standard RPC response.-data RpcRequest = RpcRequest-    { reqMethodName      :: Required 1 (Value Text)       -- ^ Name of the RPC method-    , reqBytes           :: Optional 2 (Value ByteString) -- ^ Bytes corresponding to the client protobuf request-    , reqProtocolName    :: Required 3 (Value Text)       -- ^ Protocol name of class declaring the called method-    , reqProtocolVersion :: Required 4 (Value Word64)     -- ^ Protocol version of class declaring the called method-    } deriving (Generic, Show)+-- | This message is the header for the Protobuf Rpc Engine+-- when sending a RPC request from  RPC client to the RPC server.+-- The actual request (serialized as protobuf) follows this request.+--+-- No special header is needed for the Rpc Response for Protobuf Rpc Engine.+-- The normal RPC response header (see RpcHeader.proto) are sufficient.+data RequestHeader = RequestHeader {+  -- | Name of the RPC method+    reqMethodName :: Required 1 (Value Text) -instance Encode RpcRequest-instance Decode RpcRequest+  -- | RPCs for a particular interface (ie protocol) are done using an IPC+  -- connection that is setup using rpcProxy.  The rpcProxy has a declared+  -- protocol name that is sent from client to server at connection time.+  --+  -- Each Rpc call also sends a protocol name (reqProtocolName). This name is+  -- usually the same as the connection protocol name, but not always. For example,+  -- meta protocols, such as ProtocolInfoProto, reuse the connection but need to+  -- indicate that the actual protocol is different (i.e. the protocol is+  -- ProtocolInfoProto) since they reuse the connection; in this case the+  -- protocol name is set to ProtocolInfoProto.+  , reqProtocolName :: Required 2 (Value Text) +  -- | Protocol version of class declaring the called method.+  , reqProtocolVersion :: Required 3 (Value Word64)+} deriving (Generic, Show)++instance Encode RequestHeader+instance Decode RequestHeader+ ------------------------------------------------------------------------  -- | Success or failure. The reponse header's error detail, exception@@ -88,52 +102,39 @@ -- | Note that RPC response header is also used when connection setup fails. -- (i.e. the response looks like an RPC response with a fake callId) ----- For v7:---  - If successfull then the Respose follows after this header---      - length (4 byte int), followed by the response---  - If error or fatal - the exception info follow---      - length (4 byte int) Class name of exception - UTF-8 string---      - length (4 byte int) Stacktrace - UTF-8 string---      - if the strings are null then the length is -1--- --  In case of Fatal error then the respose contains the Serverside's IPC version.  data RpcResponseHeader = RpcResponseHeader     { rspCallId             :: Required 1 (Value Word32)          -- ^ Call ID used in request     , rspStatus             :: Required 2 (Enumeration RpcStatus)     , rspServerIpcVersion   :: Optional 3 (Value Word32)          -- ^ v7: Sent if fatal v9: Sent if success or fail--    -- Fields below don't apply until v9-    --, rspExceptionClassName :: Optional 4 (Value Text)            -- ^ If the request fails-    --, rspErrorMsg           :: Optional 5 (Value Text)            -- ^ If the request fails, often contains stack trace-    --, rspErrorDetail        :: Optional 6 (Enumeration Error)  -- ^ In case of error-    --, rspClientId           :: Optional 7 (Value ByteString)      -- ^ Globally unique client ID-    --, rspRetryCount         :: Optional 8 (Value Int32)+    , rspExceptionClassName :: Optional 4 (Value Text)            -- ^ If the request fails+    , rspErrorMsg           :: Optional 5 (Value Text)            -- ^ If the request fails, often contains stack trace+    , rspErrorDetail        :: Optional 6 (Enumeration Error)     -- ^ In case of error+    , rspClientId           :: Optional 7 (Value ByteString)      -- ^ Globally unique client ID+    , rspRetryCount         :: Optional 8 (Value Int32)     } deriving (Generic, Show)  instance Encode RpcResponseHeader instance Decode RpcResponseHeader -{---- Error doesn't apply until v9- -- | Describes why an RPC error occurred. data Error = ErrorApplication         -- ^ RPC failed - RPC app threw exception-              | ErrorNoSuchMethod        -- ^ RPC error - no such method-              | ErrorNoSuchProtocol      -- ^ RPC error - no such protocol-              | ErrorRpcServer           -- ^ RPC error on server side-              | ErrorSerializingResponse -- ^ Error serializing response-              | ErrorRpcVersionMismatch  -- ^ RPC protocol version mismatch-              | ErrorCode Int            -- ^ RPC error that we don't know about+           | ErrorNoSuchMethod        -- ^ RPC error - no such method+           | ErrorNoSuchProtocol      -- ^ RPC error - no such protocol+           | ErrorRpcServer           -- ^ RPC error on server side+           | ErrorSerializingResponse -- ^ Error serializing response+           | ErrorRpcVersionMismatch  -- ^ RPC protocol version mismatch+           | ErrorCode Int            -- ^ RPC error that we don't know about -                -- starts at 10-              | FatalUnknown                  -- ^ Unknown fatal error-              | FatalUnsupportedSerialization -- ^ IPC layer serilization type invalid-              | FatalInvalidRpcHeader         -- ^ Fields of RPC header are invalid-              | FatalDeserializingRequest     -- ^ Could not deserialize RPC request-              | FatalVersionMismatch          -- ^ IPC layer version mismatch-              | FatalUnauthorized             -- ^ Auth failed-              | FatalCode Int                 -- ^ Fatal error that we don't know about+             -- starts at 10+           | FatalUnknown                  -- ^ Unknown fatal error+           | FatalUnsupportedSerialization -- ^ IPC layer serilization type invalid+           | FatalInvalidRpcHeader         -- ^ Fields of RPC header are invalid+           | FatalDeserializingRequest     -- ^ Could not deserialize RPC request+           | FatalVersionMismatch          -- ^ IPC layer version mismatch+           | FatalUnauthorized             -- ^ Auth failed+           | FatalCode Int                 -- ^ Fatal error that we don't know about   deriving (Generic, Show)  instance Enum Error where@@ -167,4 +168,3 @@       FatalVersionMismatch          -> 14       FatalUnauthorized             -> 15       FatalCode n                   -> n--}
src/Network/Hadoop/Hdfs.hs view
@@ -103,10 +103,10 @@     runHdfs' config hdfs  runHdfs' :: HadoopConfig -> Hdfs a -> IO a-runHdfs' config@HadoopConfig{..} hdfs = S.runTcp hcProxy nameNode session+runHdfs' config@HadoopConfig{..} hdfs = S.bracketSocket hcProxy nameNode session   where     session socket = do-        conn <- initConnectionV7 config hdfsProtocol socket+        conn <- initConnectionV9 config hdfsProtocol socket         unHdfs hdfs conn      nameNode = case hcNameNodes of
src/Network/Hadoop/Read.hs view
@@ -109,7 +109,7 @@      runBlock proxy endpoint offset len0 extended token = do         let len = fromMaybe len0 . getField . ebNumBytes $ extended-        S.runTcp proxy endpoint $ readBlock offset len extended token+        S.bracketSocket proxy endpoint $ readBlock offset len extended token      readBlock offset len extended token sock = go 0 offset len acc0       where
src/Network/Hadoop/Rpc.hs view
@@ -10,32 +10,34 @@     , RawRequest     , RawResponse -    , initConnectionV7+    , initConnectionV9     , invokeAsync     , invoke     ) where -import           Control.Applicative ((<$>), (<*>))+import           Control.Applicative ((<$>)) import           Control.Concurrent (ThreadId, forkIO, newEmptyMVar, putMVar, takeMVar) import           Control.Concurrent.STM import           Control.Exception (SomeException(..), throwIO, handle) import           Control.Monad (forever, when) import           Data.ByteString (ByteString)-import qualified Data.ByteString.Char8 as B+import qualified Data.ByteString.Lazy as L import qualified Data.HashMap.Strict as H import           Data.Hashable (Hashable) import           Data.Maybe (fromMaybe, isNothing)+import           Data.Monoid ((<>)) import           Data.Monoid (mempty) import           Data.Text (Text) import qualified Data.Text as T-import qualified Data.Text.Encoding as T+import qualified Data.UUID as UUID+import           System.Random (randomIO)  import           Data.ProtocolBuffers import           Data.ProtocolBuffers.Orphans () import           Data.Serialize.Get import           Data.Serialize.Put -import           Data.Hadoop.Protobuf.Headers+import qualified Data.Hadoop.Protobuf.Headers as P import           Data.Hadoop.Types import qualified Network.Hadoop.Stream as S import           Network.Socket (Socket)@@ -43,10 +45,10 @@ ------------------------------------------------------------------------  data Connection = Connection-    { cnVersion    :: !Int-    , cnConfig     :: !HadoopConfig-    , cnProtocol   :: !Protocol-    , invokeRaw    :: !(Method -> RawRequest -> (RawResponse -> IO ()) -> IO ())+    { cnVersion  :: !Int+    , cnConfig   :: !HadoopConfig+    , cnProtocol :: !Protocol+    , invokeRaw  :: !(Method -> RawRequest -> (RawResponse -> IO ()) -> IO ())     }  data Protocol = Protocol@@ -60,6 +62,8 @@  type CallId = Int +newtype ClientId = ClientId { unClientId :: ByteString }+ ------------------------------------------------------------------------  -- hadoop-2.1.0-beta is on version 9@@ -67,25 +71,25 @@  data ConnectionState = ConnectionState     { csStream        :: !S.Stream+    , csClientId      :: !ClientId     , csCallId        :: !(TVar CallId)     , csRecvCallbacks :: !(TVar (H.HashMap CallId (RawResponse -> IO ())))     , csSendQueue     :: !(TQueue (Method, RawRequest, RawResponse -> IO ()))     , csFatalError    :: !(TVar (Maybe SomeException))     } -initConnectionV7 :: HadoopConfig -> Protocol -> Socket -> IO Connection-initConnectionV7 config@HadoopConfig{..} protocol sock = do-    csStream <- S.mkSocketStream sock+initConnectionV9 :: HadoopConfig -> Protocol -> Socket -> IO Connection+initConnectionV9 config@HadoopConfig{..} protocol sock = do+    csStream   <- S.mkSocketStream sock+    csClientId <- mkClientId      S.runPut csStream $ do         putByteString "hrpc"-        putWord8 7  -- version-        putWord8 80 -- auth method (80 = simple, 81 = kerberos/gssapi, 82 = token/digest-md5)-        putWord8 0  -- ipc serialization type (0 = protobuf)--        let bs = runPut (encodeMessage context)-        putWord32be (fromIntegral (B.length bs))-        putByteString bs+        putWord8 9  -- version+        putWord8 0 -- rpc service class (0 = default/protobuf, 1 = built-in, 2 = writable, 3 = protobuf+        putWord8 0 -- auth protocol (0 = none, -33/0xDF = sasl)+        putMessage $ delimitedBytesL (rpcRequestHeaderProto csClientId (-3))+                  <> delimitedBytesL (contextProto protocol hcUser)      csCallId        <- newTVarIO 0     csRecvCallbacks <- newTVarIO H.empty@@ -122,25 +126,36 @@             modifyTVar' csCallId succ             modifyTVar' csRecvCallbacks (H.insert callId k) -            return $ runPut $ encodeLengthPrefixedMessage (requestHeaderProto callId)-                           >> encodeLengthPrefixedMessage (requestProto method requestBytes)+            return $ delimitedBytesL (rpcRequestHeaderProto csClientId callId)+                  <> delimitedBytesL (requestHeaderProto protocol method)+                  <> L.fromStrict requestBytes -        S.runPut csStream $ do-            putWord32be (fromIntegral (B.length bs))-            putByteString bs+        S.runPut csStream (putMessage bs)      forkRecv :: ConnectionState -> IO ThreadId     forkRecv cs@ConnectionState{..} = forkIO $ handle (onSocketError cs) $ forever $ do-        hdr <- S.maybeGet csStream decodeLengthPrefixedMessage-        case hdr of-          Nothing     -> throwIO ConnectionClosed-          Just rspHdr -> do+        mget <- S.maybeGet csStream $ do+            n <- fromIntegral <$> getWord32be+            -- TODO Would be nice if we didn't have to isolate here+            -- TODO and could stream instead. We could stream if we+            -- TODO were able to read the varint length prefix+            -- TODO ourselves and keep track of how many bytes were+            -- TODO remaining instead of calling `getRemaining`.+            isolate n $ do+                hdr <- decodeLengthPrefixedMessage+                msg <- case getField (P.rspStatus hdr) of+                    P.Success -> Right <$> getRemaining+                    _         -> return . Left . SomeException $ rspError hdr+                return (hdr, msg)++        case mget of+          Nothing -> throwIO ConnectionClosed+          Just (hdr, msg) -> do             onResponse <- fromMaybe (return $ return ())-                      <$> lookupDelete csRecvCallbacks (fromIntegral $ getField $ rspCallId rspHdr)-            case getField (rspStatus rspHdr) of-              Success -> S.runGet csStream getResponse >>= onResponse . Right-              _       -> S.runGet csStream getError    >>= onResponse . Left . SomeException+                      <$> lookupDelete csRecvCallbacks (fromIntegral $ getField $ P.rspCallId hdr) +            onResponse msg+     onSocketError :: ConnectionState -> SomeException -> IO ()     onSocketError ConnectionState{..} ex = do         ks <- atomically $ do@@ -154,26 +169,8 @@     ignore :: SomeException -> IO ()     ignore _ = return () -    context = IpcConnectionContext-        { ctxProtocol = putField (Just (prName protocol))-        , ctxUserInfo = putField (Just UserInformation-            { effectiveUser = putField (Just hcUser)-            , realUser      = mempty-            })-        }--    requestHeaderProto callId = RpcRequestHeader-        { reqKind       = putField (Just ProtocolBuffer)-        , reqOp         = putField (Just FinalPacket)-        , reqCallId     = putField (fromIntegral callId)-        }--    requestProto method bytes = RpcRequest-        { reqMethodName      = putField method-        , reqBytes           = putField (Just bytes)-        , reqProtocolName    = putField (prName protocol)-        , reqProtocolVersion = putField (fromIntegral (prVersion protocol))-        }+mkClientId :: IO ClientId+mkClientId = ClientId . L.toStrict . UUID.toByteString <$> randomIO  unfoldM :: Monad m => m (Maybe a) -> m [a] unfoldM f = go []@@ -192,18 +189,45 @@  ------------------------------------------------------------------------ -getResponse :: Get ByteString-getResponse = do-    n <- fromIntegral <$> getWord32be-    getByteString n+contextProto :: Protocol -> User -> P.IpcConnectionContext+contextProto protocol user = P.IpcConnectionContext+    { P.ctxProtocol = putField (Just (prName protocol))+    , P.ctxUserInfo = putField (Just P.UserInformation+        { P.effectiveUser = putField (Just user)+        , P.realUser      = mempty+        })+    } -getError :: Get RemoteError-getError = RemoteError <$> getText <*> getText-  where-    getText = do-        n <- fromIntegral <$> getWord32be-        T.decodeUtf8 <$> getByteString n+rpcRequestHeaderProto :: ClientId -> CallId -> P.RpcRequestHeader+rpcRequestHeaderProto clientId callId = P.RpcRequestHeader+    { P.reqKind       = putField (Just P.ProtocolBuffer)+    , P.reqOp         = putField (Just P.FinalPacket)+    , P.reqCallId     = putField (fromIntegral callId)+    , P.reqClientId   = putField (unClientId clientId)+    , P.reqRetryCount = putField (Just (-1))+    } +requestHeaderProto :: Protocol -> Method -> P.RequestHeader+requestHeaderProto protocol method = P.RequestHeader+    { P.reqMethodName      = putField method+    , P.reqProtocolName    = putField (prName protocol)+    , P.reqProtocolVersion = putField (fromIntegral (prVersion protocol))+    }++rspError :: P.RpcResponseHeader -> RemoteError+rspError rsp = RemoteError (fromMaybe "unknown error" $ getField $ P.rspExceptionClassName rsp)+                           (fromMaybe "unknown error" $ getField $ P.rspErrorMsg rsp)++putMessage :: L.ByteString -> Put+putMessage body = do+    putWord32be (fromIntegral (L.length body))+    putLazyByteString body++getRemaining :: Get ByteString+getRemaining = do+    n <- remaining+    getByteString n+ ------------------------------------------------------------------------  invoke :: (Decode b, Encode a) => Connection -> Text -> a -> IO b@@ -216,16 +240,19 @@       Right x -> return x  invokeAsync :: (Decode b, Encode a) => Connection -> Text -> a -> (Either SomeException b -> IO ()) -> IO ()-invokeAsync Connection{..} method arg k = invokeRaw method (encodeBytes arg) k'+invokeAsync Connection{..} method arg k = invokeRaw method (delimitedBytes arg) k'   where     k' (Left err) = k (Left err)-    k' (Right bs) = k (decodeBytes bs)+    k' (Right bs) = k (fromDelimitedBytes bs) -encodeBytes :: Encode a => a -> ByteString-encodeBytes = runPut . encodeMessage+delimitedBytes :: Encode a => a -> ByteString+delimitedBytes = runPut . encodeLengthPrefixedMessage -decodeBytes :: Decode a => ByteString -> Either SomeException a-decodeBytes bs = case runGetState decodeMessage bs 0 of+delimitedBytesL :: Encode a => a -> L.ByteString+delimitedBytesL = L.fromStrict . delimitedBytes++fromDelimitedBytes :: Decode a => ByteString -> Either SomeException a+fromDelimitedBytes bs = case runGetState decodeLengthPrefixedMessage bs 0 of     Left err      -> decodeError (T.pack err)     Right (x, "") -> Right x     Right (_, _)  -> decodeError "decoded response but did not consume enough bytes"
src/Network/Hadoop/Socket.hs view
@@ -1,15 +1,10 @@ module Network.Hadoop.Socket     ( S.Socket-    , S.SockAddr(..)--    , runTcp-+    , bracketSocket     , connectSocket-    , newSocket     , closeSocket     ) where -import           Control.Applicative ((<$>)) import           Control.Monad.Catch (MonadMask, bracket, bracketOnError) import           Control.Monad.IO.Class (MonadIO(..)) import           Data.Hadoop.Types@@ -20,45 +15,42 @@  ------------------------------------------------------------------------ -runTcp :: (MonadMask m, MonadIO m) => Maybe SocksProxy -> Endpoint -> (S.Socket -> m a) -> m a-runTcp Nothing      = runTcp'-runTcp (Just proxy) = runSocks proxy--runTcp' :: (MonadMask m, MonadIO m) => Endpoint -> (S.Socket -> m a) -> m a-runTcp' endpoint = bracket-    (liftIO $ fst <$> connectSocket endpoint)-    (liftIO . closeSocket)+bracketSocket :: (MonadMask m, MonadIO m) => Maybe SocksProxy -> Endpoint -> (S.Socket -> m a) -> m a+bracketSocket proxy endpoint = bracket (connectSocket proxy endpoint) closeSocket -runSocks :: (MonadMask m, MonadIO m) => SocksProxy -> Endpoint -> (S.Socket -> m a) -> m a-runSocks proxy endpoint = bracket-    (liftIO $ socksConnectWith proxyConf host port)-    (liftIO . closeSocket)-  where-    proxyConf = defaultSocksConf (T.unpack $ epHost proxy)-                                 (fromIntegral $ epPort proxy)+connectSocket :: (MonadMask m, MonadIO m) => Maybe SocksProxy -> Endpoint -> m S.Socket+connectSocket Nothing      = connectDirect+connectSocket (Just proxy) = connectSocks proxy -    host = T.unpack $ epHost endpoint-    port = PortNumber $ fromIntegral $ epPort endpoint+closeSocket :: MonadIO m => S.Socket -> m ()+closeSocket = liftIO . S.sClose  ------------------------------------------------------------------------ -connectSocket :: Endpoint -> IO (S.Socket, S.SockAddr)-connectSocket endpoint = do-    (addr:_) <- S.getAddrInfo (Just hints) (Just host) (Just port)+connectDirect :: (MonadMask m, MonadIO m) => Endpoint -> m S.Socket+connectDirect endpoint = do+    (addr:_) <- liftIO $ S.getAddrInfo (Just hints) (Just host) (Just port)     bracketOnError (newSocket addr) closeSocket $ \sock -> do-       let sockAddr = S.addrAddress addr-       S.connect sock sockAddr-       return (sock, sockAddr)+       liftIO $ S.connect sock (S.addrAddress addr)+       return sock   where     host  = T.unpack (epHost endpoint)     port  = show (epPort endpoint)     hints = S.defaultHints { S.addrFlags = [S.AI_ADDRCONFIG]                            , S.addrSocketType = S.Stream } -newSocket :: S.AddrInfo -> IO S.Socket-newSocket addr = S.socket (S.addrFamily addr)-                          (S.addrSocketType addr)-                          (S.addrProtocol addr)+newSocket :: MonadIO m => S.AddrInfo -> m S.Socket+newSocket addr = liftIO $ S.socket (S.addrFamily addr)+                                   (S.addrSocketType addr)+                                   (S.addrProtocol addr) -closeSocket :: S.Socket -> IO ()-closeSocket = S.sClose+------------------------------------------------------------------------++connectSocks :: (MonadMask m, MonadIO m) => SocksProxy -> Endpoint -> m S.Socket+connectSocks proxy endpoint = liftIO (socksConnectWith proxyConf host port)+  where+    proxyConf = defaultSocksConf (T.unpack $ epHost proxy)+                                 (fromIntegral $ epPort proxy)++    host = T.unpack $ epHost endpoint+    port = PortNumber $ fromIntegral $ epPort endpoint