nats-queue 0.1.0.1 → 0.1.1.0
raw patch · 3 files changed
+72/−29 lines, 3 filesdep +network-uridep ~network
Dependencies added: network-uri
Dependency ranges changed: network
Files
- Network/Nats.hs +28/−26
- Network/Nats/Json.hs +39/−0
- nats-queue.cabal +5/−3
Network/Nats.hs view
@@ -55,6 +55,7 @@ -- command is synchronous, it waits until the server responds with +OK. The commands 'publish' -- and 'unsubscribe' are asynchronous, no confirmation from server is required. Nats+ , NatsSID , connect -- * Exceptions , NatsException@@ -81,6 +82,7 @@ import qualified Data.Foldable as FOLD import Control.Exception (bracket, bracketOnError, throwIO, catch, IOException, AsyncException, Exception) import System.Random (randomRIO)+import Data.IORef import qualified Data.Map.Strict as Map import qualified Data.ByteString.Lazy.Char8 as BL@@ -169,8 +171,8 @@ MVar () -- Empty mvar that gets full in the moment we connect ) , natsThreadId :: MVar ThreadId- , natsNextSid :: MVar NatsSID- , natsSubMap :: MVar (Map.Map NatsSID NatsSubscription)+ , natsNextSid :: IORef NatsSID+ , natsSubMap :: IORef (Map.Map NatsSID NatsSubscription) } @@ -220,32 +222,32 @@ _makeClntMsg (NatsClntConnect info) = "CONNECT " : (BL.toChunks $ AE.encode info) -- | Decode NATS server message; result is message + payload (payload is 'undefined' in NatsSvrMsg)-decodeMessage :: BS.ByteString -> Maybe (NatsSvrMessage, Int)+decodeMessage :: BS.ByteString -> Maybe (NatsSvrMessage, Maybe Int) decodeMessage line = decodeMessage_ mid mpayload where (mid, mpayload) = (BS.takeWhile (\x -> x/=' ' && x/='\r') line, BS.drop 1 $ BS.dropWhile (\x -> x/=' ' && x/='\r') line) - decodeMessage_ :: BS.ByteString -> BS.ByteString -> Maybe (NatsSvrMessage, Int)- decodeMessage_ "PING" _ = Just (NatsSvrPing, 0)- decodeMessage_ "PONG" _ = Just (NatsSvrPong, 0)- decodeMessage_ "+OK" _ = Just (NatsSvrOK, 0)- decodeMessage_ "-ERR" msg = Just (NatsSvrError (decodeUtf8 msg), 0)+ decodeMessage_ :: BS.ByteString -> BS.ByteString -> Maybe (NatsSvrMessage, Maybe Int)+ decodeMessage_ "PING" _ = Just (NatsSvrPing, Nothing)+ decodeMessage_ "PONG" _ = Just (NatsSvrPong, Nothing)+ decodeMessage_ "+OK" _ = Just (NatsSvrOK, Nothing)+ decodeMessage_ "-ERR" msg = Just (NatsSvrError (decodeUtf8 msg), Nothing) decodeMessage_ "INFO" msg = do info <- AE.decode $ BL.fromChunks [msg]- return $ (NatsSvrInfo info, 0)+ return $ (NatsSvrInfo info, Nothing) decodeMessage_ "MSG" msg = do let fields = BS.split ' ' msg case (map BS.unpack fields) of- [subj, sid, len] -> return (NatsSvrMsg subj (read sid) undefined Nothing, read len)- [subj, sid, reply, len] -> return (NatsSvrMsg subj (read sid) undefined (Just $ reply), read $ len)+ [subj, sid, len] -> return (NatsSvrMsg subj (read sid) undefined Nothing, Just $ read len)+ [subj, sid, reply, len] -> return (NatsSvrMsg subj (read sid) undefined (Just $ reply), Just $ read len) _ -> fail "" decodeMessage_ _ _ = Nothing -- | Returns next sid and updates MVar newNatsSid :: Nats -> IO NatsSID-newNatsSid nats = modifyMVar (natsNextSid nats) $ \sid -> return (sid + 1, sid)+newNatsSid nats = atomicModifyIORef' (natsNextSid nats) $ \sid -> (sid + 1, sid) -- | Generates a new INBOX name for request/response communication newInbox :: IO String@@ -338,15 +340,15 @@ authenticate nats handle = do info <- BS.hGetLine handle case (decodeMessage info) of- Just (NatsSvrInfo (NatsServerInfo {natsSvrAuthRequired=True}), 0) -> do+ Just (NatsSvrInfo (NatsServerInfo {natsSvrAuthRequired=True}), Nothing) -> do BL.hPut handle $ makeClntMsg (NatsClntConnect $ natsConnOptions nats) BS.hPut handle "\r\n" response <- BS.hGetLine handle case (decodeMessage response) of- Just (NatsSvrOK, 0) -> return ()- Just (NatsSvrError err, 0)-> throwIO $ NatsException $ "Authentication error: " ++ (show err)+ Just (NatsSvrOK, Nothing) -> return ()+ Just (NatsSvrError err, Nothing)-> throwIO $ NatsException $ "Authentication error: " ++ (show err) _ -> throwIO $ NatsException $ "Incorrect server response"- Just (NatsSvrInfo _, 0) -> return ()+ Just (NatsSvrInfo _, Nothing) -> return () _ -> throwIO $ NatsException "Incorrect input from server" -- | Open and authenticate a connection@@ -396,7 +398,7 @@ connectionHandler nats = do (handle, _, _, _) <- readMVar (natsRuntime nats) -- Subscribe channels that are supposed to be subscribed- subscriptions <- readMVar (natsSubMap nats)+ subscriptions <- readIORef (natsSubMap nats) FOLD.forM_ subscriptions $ \(NatsSubscription subject queue _ sid) -> sendMessage nats True (NatsClntSubscribe subject sid queue) Nothing -- Perform the job@@ -422,7 +424,7 @@ sendMessage nats True (NatsClntConnect $ natsConnOptions nats) Nothing handleMessage (NatsSvrInfo _) = return () handleMessage (NatsSvrMsg {..}) = do- msubscription <- Map.lookup msgSid <$> readMVar (natsSubMap nats)+ msubscription <- Map.lookup msgSid <$> readIORef (natsSubMap nats) case msubscription of Just subscription -> (subCallback subscription) msgSid msgSubject (BL.fromChunks [msgText]) msgReply -- SID not found in map, force unsubscribe@@ -430,9 +432,9 @@ in do line <- BS.hGetLine handle case (decodeMessage line) of- Just (msg, 0) -> do+ Just (msg, Nothing) -> do handleMessage msg- Just (msg@(NatsSvrMsg {}), paylen) -> do+ Just (msg@(NatsSvrMsg {}), Just paylen) -> do payload <- BS.hGet handle paylen _ <- BS.hGet handle 2 -- CRLF handleMessage msg{msgText=payload}@@ -461,8 +463,8 @@ csig <- newEmptyMVar mruntime <- newMVar (undefined, undefined, False, csig) mthreadid <- newEmptyMVar - nextsid <- newMVar 1- submap <- newMVar Map.empty+ nextsid <- newIORef 1+ submap <- newIORef Map.empty let opts = defaultConnectionOptions{natsConnUser=T.pack user, natsConnPass=T.pack password} let nats = Nats{ natsConnOptions=opts @@ -492,12 +494,12 @@ sendMessage nats True (NatsClntSubscribe ssubject sid squeue) $ Just $ \err -> do case err of Just _ -> return ()- Nothing -> modifyMVarMasked_ (natsSubMap nats) - (return . Map.insert sid (NatsSubscription{subSubject=ssubject, subQueue=squeue, subCallback=cb, subSid=sid})) + Nothing -> atomicModifyIORef' (natsSubMap nats) $ \ioref ->+ (Map.insert sid (NatsSubscription{subSubject=ssubject, subQueue=squeue, subCallback=cb, subSid=sid}) ioref, ()) putMVar mvar err merr <- takeMVar mvar case merr of- Just err -> error $ show err+ Just err -> throwIO $ NatsException $ T.unpack err Nothing -> return $ sid -- | Unsubscribe from a channel@@ -506,7 +508,7 @@ -> IO () unsubscribe nats sid = do -- Remove from internal tables- modifyMVarMasked_ (natsSubMap nats) (return . Map.delete sid)+ atomicModifyIORef' (natsSubMap nats) $ \ioref -> (Map.delete sid ioref, ()) -- Unsubscribe from server, ignore errors sendMessage nats False (NatsClntUnsubscribe sid) Nothing `catch` ((\_ -> return ()) :: IOException -> IO ())
+ Network/Nats/Json.hs view
@@ -0,0 +1,39 @@+{-# LANGUAGE RecordWildCards,PatternGuards,Rank2Types #-}+++module Network.Nats.Json (+ subscribe+ , publish+) where++import Network.Nats (Nats, NatsSID)+import qualified Network.Nats as N+import qualified Data.Aeson as AE++-- | Publish a message+publish :: AE.ToJSON a =>+ Nats + -> String -- ^ Subject+ -> a -- ^ Data+ -> IO ()+publish nats subject body = N.publish nats subject (AE.encode body)++-- | Subscribe to a channel, optionally specifying queue group +-- If the JSON cannot be properly parsed, the message is ignored+subscribe :: AE.FromJSON a =>+ Nats + -> String -- ^ Subject+ -> (Maybe String) -- ^ Queue+ -> (NatsSID+ -> String+ -> a+ -> Maybe String+ -> IO ()+ )+ -- ^ Callback+ -> IO NatsSID -- ^ SID of subscription+subscribe nats subject queue jcallback = N.subscribe nats subject queue cb+ where+ cb sid subj msg repl+ | Just body <- AE.decode msg = jcallback sid subj body repl+ | True = return () -- Ignore when there is an error decoding
nats-queue.cabal view
@@ -1,5 +1,5 @@ name: nats-queue-version: 0.1.0.1+version: 0.1.1.0 synopsis: Haskell API for NATS messaging system description: This library is a Haskell driver for NATS <http://nats.io>. @@ -22,8 +22,8 @@ location: https://github.com/ondrap/nats-queue.git library- exposed-modules: Network.Nats- build-depends: base==4.*, network, dequeue, random, + exposed-modules: Network.Nats, Network.Nats.Json+ build-depends: base==4.*, network>=2.6, dequeue, random, network-uri>=2.6, containers >= 0.5, bytestring, text, aeson >= 0.7 ghc-options: -Wall extensions: TemplateHaskell@@ -32,3 +32,5 @@ GeneralizedNewtypeDeriving PatternGuards DeriveDataTypeable+ Rank2Types+