packages feed

Z-MessagePack 0.4.0.1 → 0.4.1.0

raw patch · 6 files changed

+80/−73 lines, 6 filesdep ~Z-Datadep ~Z-IOPVP: major bump suggested

API removals or changes: PVP suggests a major version bump

Dependency ranges changed: Z-Data, Z-IO

API changes (from Hackage documentation)

- Z.Data.MessagePack: infixl 9 <?>
- Z.Data.MessagePack: parseValueChunks :: Monad m => m Bytes -> Bytes -> m (Bytes, Either ParseError Value)
- Z.Data.MessagePack: parseValueChunks' :: Monad m => m Bytes -> Bytes -> m (Either ParseError Value)
- Z.Data.MessagePack.Base: infixl 9 <?>
- Z.Data.MessagePack.Base: parseValueChunks :: Monad m => m Bytes -> Bytes -> m (Bytes, Either ParseError Value)
- Z.Data.MessagePack.Base: parseValueChunks' :: Monad m => m Bytes -> Bytes -> m (Either ParseError Value)
- Z.Data.MessagePack.Value: parseValueChunks :: Monad m => m Bytes -> Bytes -> m (Bytes, Either ParseError Value)
- Z.Data.MessagePack.Value: parseValueChunks' :: Monad m => m Bytes -> Bytes -> m (Either ParseError Value)
+ Z.Data.MessagePack: decodeChunk :: MessagePack a => Bytes -> Result DecodeError a
+ Z.Data.MessagePack.Base: decodeChunk :: MessagePack a => Bytes -> Result DecodeError a
- Z.Data.MessagePack: decodeChunks :: (MessagePack a, Monad m) => m Bytes -> Bytes -> m (Bytes, Either DecodeError a)
+ Z.Data.MessagePack: decodeChunks :: (MessagePack a, Monad m) => ParseChunks m DecodeError a
- Z.Data.MessagePack.Base: decodeChunks :: (MessagePack a, Monad m) => m Bytes -> Bytes -> m (Bytes, Either DecodeError a)
+ Z.Data.MessagePack.Base: decodeChunks :: (MessagePack a, Monad m) => ParseChunks m DecodeError a
- Z.Data.MessagePack.Base: type ParseChunks (m :: Type -> Type) chunk err x = m chunk -> chunk -> m (chunk, Either err x)
+ Z.Data.MessagePack.Base: type ParseChunks (m :: Type -> Type) err x = m Bytes -> Bytes -> m (Bytes, Either err x)
- Z.IO.RPC.MessagePack: [StreamHandler] :: (MessagePack req, MessagePack res) => (SessionCtx a -> req -> IO (Source res)) -> ServerHandler a
+ Z.IO.RPC.MessagePack: [StreamHandler] :: (MessagePack req, MessagePack res) => (SessionCtx a -> IORef Bool -> req -> IO (Source res)) -> ServerHandler a

Files

CHANGELOG.md view
@@ -1,5 +1,9 @@ # Revision history for Z-MessagePack +## 0.4.1.0 -- 2021-04-25++* Adapt to `Z-IO` 0.8. Add stream error handling code.+ ## 0.3.0.1 -- 2021-02-25  * Add `callStream` and `StreamHandler`.
Z-MessagePack.cabal view
@@ -1,6 +1,6 @@ cabal-version:      >=1.10 name:               Z-MessagePack-version:            0.4.0.1+version:            0.4.1.0 synopsis:           MessagePack description:        MessagePack binary serialization format. license:            BSD3@@ -14,6 +14,10 @@ bug-reports:        https://github.com/ZHaskell/z-messagepack/issues extra-source-files: CHANGELOG.md +source-repository head+  type:     git+  location: git://github.com/ZHaskell/z-messagepack.git+ library   exposed-modules:     Z.Data.MessagePack@@ -34,8 +38,8 @@     , tagged                >=0.8   && <0.9     , time                  >=1.9   && <2.0     , unordered-containers  >=0.2   && <0.3-    , Z-Data                >=0.6   && <1.0-    , Z-IO                  >=0.6   && <1.0+    , Z-Data                >=0.8   && <0.9+    , Z-IO                  >=0.8   && <0.9    default-language:   Haskell2010   default-extensions:
Z/Data/MessagePack.hs view
@@ -67,10 +67,10 @@     MessagePack(..), Value(..), defaultSettings, Settings(..), JSON.snakeCase, JSON.trainCase     -- * Encode & Decode   , readMessagePackFile, writeMessagePackFile-  , decode, decode', decodeChunks, encode, encodeChunks+  , decode, decode', decodeChunk, decodeChunks, encode, encodeChunks   , DecodeError, ParseError     -- * parse into MessagePack Value-  , parseValue, parseValue', parseValueChunks, parseValueChunks'+  , parseValue, parseValue'   -- * Generic FromValue, ToValue & EncodeMessagePack   , gToValue, gFromValue, gEncodeMessagePack   -- * Convert 'Value' to Haskell data
Z/Data/MessagePack/Base.hs view
@@ -15,10 +15,10 @@   ( -- * MessagePack Class     MessagePack(..), Value(..), defaultSettings, Settings(..)     -- * Encode & Decode-  , decode, decode', decodeChunks, encode, encodeChunks+  , decode, decode', decodeChunk, decodeChunks, encode, encodeChunks   , DecodeError, P.ParseError, P.ParseChunks     -- * parse into MessagePack Value-  , MV.parseValue, MV.parseValue', MV.parseValueChunks, MV.parseValueChunks'+  , MV.parseValue, MV.parseValue'   -- * Generic FromValue, ToValue & EncodeMessagePack   , gToValue, gFromValue, gEncodeMessagePack   -- * Convert 'Value' to Haskell data@@ -145,17 +145,24 @@         Left cErr -> (bs', Left (Right cErr))         Right r   -> (bs', Right r) +-- | Decode a MessagePack doc chunk.+decodeChunk :: MessagePack a => V.Bytes -> P.Result DecodeError a+{-# INLINE decodeChunk #-}+decodeChunk bs = loop (P.parseChunk MV.value bs)+  where+    loop r = do+        case r of+            P.Success v rest ->+                case convertValue v of+                    Left cErr -> P.Failure (Right cErr) rest+                    Right r'  -> P.Success r' rest+            P.Failure e rest -> P.Failure (Left e) rest+            P.Partial f' -> P.Partial (loop . f')+ -- | Decode MessagePack doc chunks, return trailing bytes.-decodeChunks :: (MessagePack a, Monad m) => m V.Bytes -> V.Bytes -> m (V.Bytes, Either DecodeError a)+decodeChunks :: (MessagePack a, Monad m) => P.ParseChunks m DecodeError a {-# INLINE decodeChunks #-}-decodeChunks mb bs = do-    mr <- P.parseChunks MV.value mb bs-    case mr of-        (bs', Left pErr) -> pure (bs', Left (Left pErr))-        (bs', Right v) ->-            case convertValue v of-                Left cErr -> pure (bs', Left (Right cErr))-                Right r   -> pure (bs', Right r)+decodeChunks = P.parseChunks decodeChunk  -- | Directly encode data to MessagePack bytes. encode :: MessagePack a => a -> V.Bytes
Z/Data/MessagePack/Value.hs view
@@ -11,8 +11,6 @@     -- * parse into MessagePack Value   , parseValue   , parseValue'-  , parseValueChunks-  , parseValueChunks'     -- * Value Parsers   , value   ) where@@ -153,13 +151,3 @@ parseValue' :: V.Bytes -> Either P.ParseError Value {-# INLINE parseValue' #-} parseValue' = P.parse' (value <* P.endOfInput)---- | Increamental parse 'Value' without consuming trailing bytes.-parseValueChunks :: Monad m => m V.Bytes -> V.Bytes -> m (V.Bytes, Either P.ParseError Value)-{-# INLINE parseValueChunks #-}-parseValueChunks = P.parseChunks value---- | Increamental parse 'Value', if there're bytes left, parsing will fail.-parseValueChunks' :: Monad m => m V.Bytes -> V.Bytes -> m (Either P.ParseError Value)-{-# INLINE parseValueChunks' #-}-parseValueChunks' mi inp = snd <$> P.parseChunks (value <* P.endOfInput) mi inp
Z/IO/RPC/MessagePack.hs view
@@ -228,7 +228,7 @@         MP.encodeMessagePack req        -- param     flushBuffer bo     return (sendEOF, sourceFromIO $ do-        res <- pull (sourceParserFromBuffered (do+        res <- readParser (do             tag <- P.anyWord8             -- stream stop             case tag of@@ -245,10 +245,10 @@                         P.fail' $ "wrong response type: " <> T.toText typ                     return (Just (err, v))                 _ -> P.fail' $ "wrong response tag: " <> T.toText tag-            ) bi)+            ) bi          -- we take tcp disconnect as eof too-        case join res of+        case res of             Just (err, v) -> do                 when (err /= MV.Nil) $ throwIO (RPCException err callStack)                 unwrap "EPARSE" (MP.convertValue v)@@ -292,8 +292,10 @@                 => (SessionCtx a -> req -> IO res) -> ServerHandler a     NotifyHandler :: MessagePack req                   => (SessionCtx a -> req -> IO ()) -> ServerHandler a+    -- | 'StreamHandler' will receive an 'IORef' which get updated to 'True'+    -- when client send stream end packet, stream should end up ASAP.     StreamHandler :: (MessagePack req, MessagePack res)-                  => (SessionCtx a -> req -> IO (Source res)) -> ServerHandler a+                  => (SessionCtx a -> IORef Bool -> req -> IO (Source res)) -> ServerHandler a  -- | Simple router using `FlatMap`, lookup name in /O(log(N))/. --@@ -337,15 +339,15 @@     loop ctx bi bo   where     loop ctx bi bo = do-        req <- pull $ sourceParserFromBuffered sourceParser bi+        req <- readParser sourceParser bi         case req of-            Just (Notify (name, v)) -> do+            Notify (name, v) -> do                 case handle name of                     Just (NotifyHandler f) -> do                         f ctx =<< unwrap "EPARSE" (MP.convertValue v)                     _ -> throwOtherError "ENOTFOUND" "notification method not found"                 loop ctx bi bo-            Just (Call (msgid, name, v)) -> do+            Call (msgid, name, v) -> do                 case handle name of                     Just (CallHandler f) -> do                         res <- try (f ctx =<< unwrap "EPARSE" (MP.convertValue v))@@ -370,11 +372,11 @@                             MB.nil                         flushBuffer bo                 loop ctx bi bo-            Just (StreamStart (name, v)) -> do+            StreamStart (name, v) -> do                 eofRef <- newIORef False                 -- fork new thread to get stream end notification                 forkIO $ do-                    pull (sourceParserFromBuffered (do+                    _ <- readParser (do                         tag <- P.anyWord8                         -- stream stop                         when (tag /= 0x91) $@@ -382,49 +384,44 @@                         !typ <- MV.value                         when (typ /= MV.Int 5) $                             P.fail' $ "wrong request type: " <> T.toText typ-                        ) bi)+                        ) bi                     atomicWriteIORef eofRef True                  case handle name of-                    Just (StreamHandler f) -> do-                        src <- f ctx =<< unwrap "EPARSE" (MP.convertValue v)-                        loopSend eofRef src bo-                    _ -> do-                        writeBuilder bo $ do-                            MB.arrayHeader 3-                            MB.int 6                        -- type response-                            MB.str $ "request method: " <> name <> " not found"-                            MB.nil-                        flushBuffer bo+                    Just (StreamHandler f) -> (do+                        src <- f ctx eofRef =<< unwrap "EPARSE" (MP.convertValue v)+                        src (writeItem bo) EOF) `catch` (\ (e :: SomeException) ->+                            writeErrorItem bo $ "error when stream: " <> T.toText e)+                    _ -> writeErrorItem bo $ "request method: " <> name <> " not found"                 loop ctx bi bo -            _ -> return ()-+    writeItem bo = \ mx -> do+        case mx of+            Just x -> do+                writeBuilder bo $ do+                    MB.arrayHeader 3+                    MB.int 6                        -- type stream item+                    MB.nil+                    MP.encodeMessagePack x+                flushBuffer bo+            _ -> do+                writeBuilder bo $ do+                    MB.arrayHeader 1+                    MB.int 7                        -- type stream end+                flushBuffer bo -    loopSend eofRef src bo = do-        eof <- readIORef eofRef-        if eof-        then do-            writeBuilder bo $ do-                MB.arrayHeader 1-                MB.int 7                        -- type response-            flushBuffer bo-        else do-            r <- pull src-            case r of-                Just r' -> do-                    writeBuilder bo $ do-                        MB.arrayHeader 3-                        MB.int 6                        -- type response-                        MB.nil-                        MP.encodeMessagePack r'-                    flushBuffer bo-                _ -> atomicWriteIORef eofRef True-            loopSend eofRef src bo+    writeErrorItem bo msg = do+        writeBuilder bo $ do+            MB.arrayHeader 3+            MB.int 6                                -- type stream item+            MB.str msg+            MB.nil+        flushBuffer bo  -------------------------------------------------------------------------------  sourceParser :: P.Parser (Request MV.Value)+{-# INLINE sourceParser #-} sourceParser = do     tag <- P.anyWord8     case tag of@@ -455,7 +452,6 @@                     _ -> P.fail' $ "wrong msgid: " <> T.toText seq_                 _ -> P.fail' $ "wrong request type: " <> T.toText typ         _ -> P.fail' $ "wrong request tag: " <> T.toText tag-{-# INLINE sourceParser #-}  -- $server-example --@@ -481,8 +477,16 @@ -- >  , ("bar", CallHandler $ \ctx (req :: T.Text) -> do -- >      counter . fromJust <$> readSessionCtx ctx -- >    )--- >  , ("qux", StreamHandler $ \ctx (_ :: ()) -> do--- >     withMVar stdinBuf (pure . sourceFromBuffered)+-- >  , ("qux", StreamHandler $ \ctx eofRef (_ :: ()) -> do+-- >     withMVar stdinBuf (\ stdin -> pure $ \ k _ -> do+--          eof <- readIORef eofRef+--          if eof+--          then k EOF+--          else do+--              r <- readBuffer stdin+--              if V.null r+--              then k EOF+--              else k (Just r)) -- >    ) -- >  ]