mqtt-hs 0.2.0 → 1.0.0
raw patch · 8 files changed
+706/−772 lines, 8 filesdep +asyncdep +stmdep +transformersdep ~attoparsecdep ~basedep ~singletons
Dependencies added: async, stm, transformers
Dependency ranges changed: attoparsec, base, singletons, text
Files
- Network/MQTT.hs +74/−464
- Network/MQTT/Encoding.hs +36/−39
- Network/MQTT/Internal.hs +375/−0
- Network/MQTT/Logger.hs +0/−39
- Network/MQTT/Parser.hs +55/−55
- Network/MQTT/Types.hs +150/−154
- README.md +1/−12
- mqtt-hs.cabal +15/−9
Network/MQTT.hs view
@@ -5,36 +5,22 @@ DeriveDataTypeable #-} {-| Module: MQTT-Copyright: Lukas Braun 2014+Copyright: Lukas Braun 2014-2016 License: GPL-3 Maintainer: koomi+mqtt@hackerspace-bamberg.de -A MQTT client library.--A simple example, assuming a broker is running on localhost-(needs -XOverloadedStrings):-->>> import Network.MQTT->>> import Network.MQTT.Logger->>> Just mqtt <- connect defaultConfig { cLogger = warnings stdLogger }->>> let f t payload = putStrLn $ "A message was published to " ++ show t ++ ": " ++ show payload->>> subscribe mqtt NoConfirm "#" f-NoConfirm->>> publish mqtt Handshake False "some random/topic" "Some content!"-A message was published to "some random/topic": "Some content!"+An MQTT client library. -} module Network.MQTT- ( -- * Creating connections- connect- , MQTT+ ( -- * Setup+ run+ , Terminated(..) , disconnect- , reconnect- , onReconnect- , resubscribe- -- * Connection settings- , MQTTConfig+ , Config , defaultConfig- -- ** Field accessors+ , Commands+ , mkCommands+ -- ** Config accessors , cHost , cPort , cClean@@ -43,115 +29,37 @@ , cPassword , cKeepAlive , cClientID- , cConnectTimeout- , cReconnPeriod- , cLogger+ , cLogDebug+ , cPublished+ , cCommands+ , cInputBufferSize -- * Subscribing and publishing , subscribe , unsubscribe , publish- -- * Sending and receiving 'Message's- , send- , addHandler- , removeHandler- , awaitMsg- , awaitMsg' -- * Reexports , module Network.MQTT.Types ) where -import Control.Applicative (pure, (<$>), (<*>), (<$))+import Control.Applicative ((<$>)) import Control.Concurrent-import Control.Exception hiding (handle)-import Control.Monad hiding (sequence_)-import Data.Attoparsec.ByteString (parseOnly)-import Data.Bits ((.&.))-import Data.ByteString (hGet, ByteString)-import qualified Data.ByteString as BS-import Data.Foldable (for_, sequence_, traverse_)-import Data.Maybe (isJust, fromJust)-import Data.Singletons (withSomeSing, SingI(..))-import Data.Singletons.Decide-import Data.Text (Text)-import Data.Traversable (for)-import Data.Typeable (Typeable)+import Control.Concurrent.STM+import Control.Exception (finally)+import Control.Monad (void)+import Data.ByteString (ByteString)+import Data.Maybe (fromJust) import Data.Unique-import Data.Word import Network-import Prelude hiding (sequence_)-import System.IO (Handle, hClose, hIsEOF, hSetBinaryMode)-import System.Timeout (timeout)+import System.IO (hSetBinaryMode) +import Network.MQTT.Internal import Network.MQTT.Types-import Network.MQTT.Parser (mqttBody, mqttHeader)-import Network.MQTT.Encoding-import qualified Network.MQTT.Logger as L --------------------------------------------- Interface------------------------------------------ --- | Abstract type representing a connection to a broker.-data MQTT- = MQTT- { config :: MQTTConfig- , handle :: MVar Handle- , handlers :: MVar [MessageHandler]- , topicHandlers :: MVar [TopicHandler]- , recvThread :: MVar ThreadId- , reconnectHandler :: MVar (IO ())- , keepAliveThread :: MVar ThreadId- , sendSem :: Maybe QSem- }--data TopicHandler- = TopicHandler- { thTopic :: Topic- , thQoS :: QoS- , thHandler :: Topic -> ByteString -> IO ()- }--data MessageHandler where- MessageHandler :: SingI t- => Unique- -> (Message t -> IO ())- -> MessageHandler---- | The various options when establishing a connection.-data MQTTConfig- = MQTTConfig- { cHost :: HostName- -- ^ Hostname of the broker.- , cPort :: PortNumber- -- ^ Port of the broker.- , cClean :: Bool- -- ^ Should the server forget subscriptions and other state on- -- disconnects?- , cWill :: Maybe Will- -- ^ Optional 'Will' message.- , cUsername :: Maybe Text- -- ^ Optional username used for authentication.- , cPassword :: Maybe Text- -- ^ Optional password used for authentication.- , cKeepAlive :: Maybe Int- -- ^ Maximum interval (in seconds) in which a message must be sent.- -- 0 means no limit.- , cClientID :: Text- -- ^ Client ID used by the server to identify clients.- , cConnectTimeout :: Maybe Int- -- ^ Time in seconds after which waiting for a CONNACK is aborted.- -- 'Nothing' means no timeout.- , cReconnPeriod :: Maybe Int- -- ^ Time in seconds to wait between reconnect attempts.- -- 'Nothing' means no reconnects are attempted.- , cLogger :: L.Logger- -- ^ Functions for logging, see 'Network.MQTT.Logger.Logger'.- }---- | Defaults for 'MQTTConfig', connects to a server running on+-- | Defaults for 'Config', connects to a server running on -- localhost.-defaultConfig :: MQTTConfig-defaultConfig = MQTTConfig+defaultConfig :: Commands -> TChan (Message 'PUBLISH) -> Config+defaultConfig commands published = Config { cHost = "localhost" , cPort = 1883 , cClean = True@@ -160,375 +68,77 @@ , cPassword = Nothing , cKeepAlive = Nothing , cClientID = "mqtt-haskell"- , cConnectTimeout = Nothing- , cReconnPeriod = Nothing- , cLogger = L.stdLogger+ , cResendTimeout = 20+ , cLogDebug = const $ return ()+ , cCommands = commands+ , cPublished = published+ , cInputBufferSize = 0x1000 } --- | Establish a connection. This might fail with an 'IOException' or--- return 'Nothing' if the server did not accept the connection.-connect :: MQTTConfig -> IO (Maybe MQTT)-connect conf = do+-- | Connect to the configured broker, write received 'Publish' messages to the+-- 'cPublished' channel and handle commands from the 'cCommands' channel.+--+-- Exceptions are propagated.+run :: Config -> IO Terminated+run conf = do h <- connectTo (cHost conf) (PortNumber $ cPort conf) hSetBinaryMode h True- mqtt <- MQTT conf- <$> newMVar h- <*> newMVar []- <*> newMVar []- <*> newEmptyMVar- <*> newEmptyMVar- <*> newEmptyMVar- <*> for (cKeepAlive conf) (const (newQSem 0))- mCode <- handshake mqtt- if mCode == Just 0- then Just mqtt <$ do forkIO (recvLoop mqtt) >>= putMVar (recvThread mqtt)- forkIO (keepAliveLoop mqtt) >>=- putMVar (keepAliveThread mqtt)- addHandler mqtt (publishHandler mqtt)- else Nothing <$ hClose h---- | Send a 'Message' to the server.-send :: MQTT -> Message t -> IO ()-send mqtt msg = do- logInfo mqtt $ "Sending " ++ show (toMsgType msg)- h <- readMVar (handle mqtt)- writeTo h msg- for_ (sendSem mqtt) signalQSem--handshake :: MQTT -> IO (Maybe Word8)-handshake mqtt = do- let timeout' = maybe (fmap Just) (timeout . (* 1000000))- (cConnectTimeout (config mqtt))- sendConnect mqtt- msg <- timeout' (getMessage mqtt) `catch` \e ->- Nothing <$ logError mqtt (show (e :: MQTTException) ++- " while waiting for CONNACK")- return $ case msg of- Just (SomeMessage (Message _ (MConnAck (ConnAck code)))) -> Just code- _ -> Nothing--sendConnect :: MQTT -> IO ()-sendConnect mqtt = send mqtt connect- where- conf = config mqtt- connect = Message- (Header False NoConfirm False)- (MConnect $ Connect- (cClean conf)- (cWill conf)- (MqttText $ cClientID conf)- (MqttText <$> cUsername conf)- (MqttText <$> cPassword conf)- (maybe 0 fromIntegral $ cKeepAlive conf))+ terminatedVar <- newEmptyTMVarIO+ sendSignal <- newEmptyMVar+ mainLoop conf h (readTMVar terminatedVar) sendSignal+ `finally` atomically (putTMVar terminatedVar ()) --- | Block until a 'Message' of the given type, optionally with the given--- 'MsgID', arrives.+-- | Close the connection after sending a 'Disconnect' message. ----- Note this expects a singleton to guarantee the returned 'Message' is of--- the 'MsgType' that is being waited for. Singleton constructors are the--- 'MsgType' constructors prefixed with a capital @S@, e.g. 'SPUBLISH'.-awaitMsg :: SingI t => MQTT -> SMsgType t -> Maybe MsgID -> IO (Message t)-awaitMsg mqtt _ mMsgID = do- var <- newEmptyMVar- handlerID <- addHandler mqtt (putMVar var)- let wait = do- msg <- readMVar var- if isJust mMsgID- then if mMsgID == getMsgID (body msg)- then removeHandler mqtt handlerID >> return msg- else wait- else removeHandler mqtt handlerID >> return msg- wait---- | A version of 'awaitMsg' that infers the type of the 'Message' that--- is expected.-awaitMsg' :: SingI t => MQTT -> Maybe MsgID -> IO (Message t)-awaitMsg' mqtt mMsgID = awaitMsg mqtt sing mMsgID---- | Register a callback that gets invoked whenever a 'Message' of the--- expected 'MsgType' is received. Returns the ID of the handler which can be--- passed to 'removeHandler'.-addHandler :: SingI t => MQTT -> (Message t -> IO ()) -> IO Unique-addHandler mqtt handler = do- mhID <- newUnique- modifyMVar_ (handlers mqtt) $ \hs ->- return $ MessageHandler mhID handler : hs- return mhID---- | Remove the handler with the given ID.-removeHandler :: MQTT -> Unique -> IO ()-removeHandler mqtt mhID = modifyMVar_ (handlers mqtt) $ \hs ->- return $ filter (\(MessageHandler mhID' _) -> mhID' /= mhID) hs+-- See also: 'Will'+disconnect :: Config -> IO ()+disconnect mqtt = writeCmd mqtt CmdDisconnect --- | Subscribe to a 'Topic' with the given 'QoS' and invoke the callback--- whenever something is published to the 'Topic'. Returns the 'QoS' that--- was granted by the broker (lower or equal to the one requested).+-- | Subscribe to the 'Topic's with the corresponding 'QoS'.+-- Returns the 'QoS' that were granted (lower or equal to the ones requested)+-- in the same order. ----- The 'Topic' may contain--- <http://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html#appendix-a wildcars>.--- The 'Topic' passed to the callback is the fully expanded version where--- the message was actually published.-subscribe :: MQTT -> QoS -> Topic -> (Topic -> ByteString -> IO ())- -> IO QoS-subscribe mqtt qos topic handler = do- qosGranted <- sendSubscribe mqtt qos topic- modifyMVar_ (topicHandlers mqtt) $ \hs ->- return $ TopicHandler topic qosGranted handler : hs- return qosGranted--sendSubscribe :: MQTT -> QoS -> Topic -> IO QoS-sendSubscribe mqtt qos topic = do+-- The 'Topic's may contain+-- <http://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html#appendix-a wildcards>.+subscribe :: Config -> [(Topic, QoS)] -> IO [QoS]+subscribe mqtt topics = do msgID <- fromIntegral . hashUnique <$> newUnique- send mqtt $ Message- (Header False Confirm False)- (MSubscribe $ Subscribe- msgID- [(topic, qos)])- msg <- awaitMsg mqtt SSUBACK (Just msgID)- case msg of- (Message _ (MSubAck (SubAck _ [qosGranted]))) -> return qosGranted- _ -> fail $ "Received invalid message as response to subscribe: "- ++ show (toMsgType msg)+ msg <- sendAwait mqtt+ (Message+ (Header False Confirm False)+ (Subscribe msgID topics))+ SSUBACK+ return $ granted $ body $ msg --- | Unsubscribe from the given 'Topic' and remove any handlers.-unsubscribe :: MQTT -> Topic -> IO ()-unsubscribe mqtt topic = do- modifyMVar_ (topicHandlers mqtt) $ return . filter ((== topic) . thTopic)+-- | Unsubscribe from the given 'Topic's.+unsubscribe :: Config -> [Topic] -> IO ()+unsubscribe mqtt topics = do msgID <- fromIntegral . hashUnique <$> newUnique- send mqtt $ Message- (Header False Confirm False)- (MUnsubscribe $ Unsubscribe msgID [topic])- void $ awaitMsg mqtt SUNSUBACK (Just msgID)+ void $ sendAwait mqtt+ (Message (Header False Confirm False) (Unsubscribe msgID topics))+ SUNSUBACK -- | Publish a message to the given 'Topic' at the requested 'QoS' level.--- The payload can be any sequence of bytes, including none at all. The 'Bool'--- parameter decides if the server should retain the message for future--- subscribers to the topic.+-- The payload can be any sequence of bytes, including none at all.+-- 'True' means the server should retain the message for future subscribers to+-- the topic. -- -- The 'Topic' must not contain -- <http://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html#appendix-a wildcards>.-publish :: MQTT -> QoS -> Bool -> Topic -> ByteString -> IO ()+publish :: Config -> QoS -> Bool -> Topic -> ByteString -> IO () publish mqtt qos retain topic body = do msgID <- if qos > NoConfirm then Just . fromIntegral . hashUnique <$> newUnique else return Nothing- send mqtt $ Message- (Header False qos retain)- (MPublish $ Publish topic msgID body)+ let pub = Message (Header False qos retain) (Publish topic msgID body) case qos of- NoConfirm -> return ()- Confirm -> void $ awaitMsg mqtt SPUBACK msgID+ NoConfirm -> send mqtt pub+ Confirm -> void $ sendAwait mqtt pub SPUBACK Handshake -> do- void $ awaitMsg mqtt SPUBREC msgID- send mqtt $ Message- (Header False Confirm False)- (MPubRel $ SimpleMsg (fromJust msgID))- void $ awaitMsg mqtt SPUBCOMP msgID---- | Close the connection to the server.-disconnect :: MQTT -> IO ()-disconnect mqtt = do- h <- takeMVar $ handle mqtt- writeTo h $- Message- (Header False NoConfirm False)- MDisconnect- readMVar (recvThread mqtt) >>= killThread- readMVar (keepAliveThread mqtt) >>= killThread- hClose h---- | Try creating a new connection with the same config (retrying after the--- specified amount of seconds has passed) and invoke the callback that is--- set with 'onReconnect' once a new connection has been established.------ Does not terminate the old connection.-reconnect :: MQTT -> Int -> IO ()-reconnect mqtt period = do- -- Other threads can't write while the MVar is empty- _ <- takeMVar (handle mqtt)- logInfo mqtt "Reconnecting..."- -- Temporarily create a new MVar for the handshake so other threads- -- don't write before the connection is fully established- handleVar <- newEmptyMVar- go (mqtt { handle = handleVar })- readMVar handleVar >>= putMVar (handle mqtt)- -- forkIO so recvLoop isn't blocked- tryReadMVar (reconnectHandler mqtt) >>= traverse_ (void . forkIO)- logInfo mqtt "Reconnect successfull"- where- -- try reconnecting until it works- go mqtt' = do- let conf = config mqtt- connectTo (cHost conf) (PortNumber $ cPort conf)- >>= putMVar (handle mqtt')- mCode <- handshake mqtt'- unless (mCode == Just 0) $ do- void $ takeMVar (handle mqtt')- threadDelay (period * 10^6)- go mqtt'- `catch`- \e -> do- logWarning mqtt $ "reconnect: " ++ show (e :: IOException)- threadDelay (period * 10^6)- go mqtt'---- | Register a callback that will be invoked when a reconnect has--- happened.-onReconnect :: MQTT -> IO () -> IO ()-onReconnect mqtt io = do- let mvar = reconnectHandler mqtt- empty <- isEmptyMVar mvar- unless empty (void $ takeMVar mvar)- putMVar mvar io---- | Resubscribe to all topics. Returns the new list of granted 'QoS'.-resubscribe :: MQTT -> IO [QoS]-resubscribe mqtt = do- ths <- readMVar (topicHandlers mqtt)- mapM (\th -> sendSubscribe mqtt (thQoS th) (thTopic th)) ths--maybeReconnect :: MQTT -> IO ()-maybeReconnect mqtt = do- catch- (readMVar (handle mqtt) >>= hClose)- (const (pure ()) :: IOException -> IO ())- for_ (cReconnPeriod $ config mqtt) $ reconnect mqtt----------------------------------------------- Logger utility functions--------------------------------------------logInfo :: MQTT -> String -> IO ()-logInfo mqtt = L.logInfo (cLogger (config mqtt))--logWarning :: MQTT -> String -> IO ()-logWarning mqtt = L.logWarning (cLogger (config mqtt))--logError :: MQTT -> String -> IO ()-logError mqtt = L.logError (cLogger (config mqtt))----------------------------------------------- Internal--------------------------------------------recvLoop :: MQTT -> IO ()-recvLoop mqtt = forever $ do- h <- readMVar (handle mqtt)- eof <- hIsEOF h- if eof- then do- logError mqtt "EOF in recvLoop"- maybeReconnect mqtt- else getMessage mqtt >>= dispatchMessage mqtt- `catches`- [ Handler $ \e -> do- logError mqtt $ "recvLoop: Caught " ++ show (e :: IOException)- maybeReconnect mqtt- , Handler $ \e ->- logWarning mqtt $ "recvLoop: Caught " ++ show (e :: MQTTException)- ]--dispatchMessage :: MQTT -> SomeMessage -> IO ()-dispatchMessage mqtt (SomeMessage (msg :: Message t)) =- readMVar (handlers mqtt) >>= mapM_ applyMsg- where- typeSing :: SMsgType t- typeSing = toSMsgType msg-- applyMsg :: MessageHandler -> IO ()- applyMsg (MessageHandler _ (handler :: Message t' -> IO ())) =- case typeSing %~ (sing :: SMsgType t') of- Proved Refl -> void $ forkIO $ handler msg- Disproved _ -> return ()---- | Block on a semaphore that is signaled by 'send'. If a timeout occurs--- while waiting, send a 'PINGREQ' to the server and wait for PINGRESP.--- Ignores errors that occur while writing to the handle, reconnects are--- initiated by 'recvLoop'.------ Returns immediately if no Keep Alive is specified.-keepAliveLoop :: MQTT -> IO ()-keepAliveLoop mqtt =- sequence_ (loop <$> cKeepAlive (config mqtt) <*> sendSem mqtt)- where- loop period sem = forever $ do- rslt <- timeout (period * 1000000) $ waitQSem sem- case rslt of- Nothing -> (do send mqtt $ Message- (Header False NoConfirm False)- MPingReq- void $ awaitMsg mqtt SPINGRESP Nothing)- `catch`- (\e -> logError mqtt $ "keepAliveLoop: " ++ show (e :: IOException))- Just _ -> return ()--publishHandler :: MQTT -> Message PUBLISH -> IO ()-publishHandler mqtt (Message header (MPublish body)) = do- case (qos header, pubMsgID body) of- (Confirm, Just msgid) ->- send mqtt $ Message- (Header False NoConfirm False)- (MPubAck $ SimpleMsg msgid)- (Handshake, Just msgid) -> do- send mqtt $ Message- (Header False NoConfirm False)- (MPubRec $ SimpleMsg msgid)- void $ awaitMsg mqtt SPUBREL (Just msgid)- send mqtt $ Message- (Header False NoConfirm False)- (MPubComp $ SimpleMsg msgid)- _ -> return ()- callbacks <- filter (matches (topic body) . thTopic)- <$> readMVar (topicHandlers mqtt)- for_ callbacks $ \th -> thHandler th (topic body) (payload body)--getMessage :: MQTT -> IO SomeMessage-getMessage mqtt = do- h <- readMVar (handle mqtt)- headerByte <- hGet' h 1- remaining <- getRemaining h 0- rest <- hGet' h remaining- let parseRslt = do- (mType, header) <- parseOnly mqttHeader headerByte- withSomeSing mType $ \sMsgType ->- parseOnly- (SomeMessage . Message header- <$> mqttBody header sMsgType (fromIntegral remaining))- rest- case parseRslt of- Left err -> logError mqtt ("Error while parsing: " ++ err) >>- throw (ParseError err)- Right msg -> msg <$- logInfo mqtt ("Received " ++ show (toMsgType' msg))--getRemaining :: Handle -> Int -> IO Int-getRemaining h n = go n 1- where- go acc fac = do- b <- getByte h- let acc' = acc + (b .&. 127) * fac- if b .&. 128 == 0- then return acc'- else go acc' (fac * 128)--getByte :: Handle -> IO Int-getByte h = fromIntegral . BS.head <$> hGet' h 1--hGet' :: Handle -> Int -> IO BS.ByteString-hGet' h n = do- bs <- hGet h n- if BS.length bs < n- then throw EOF- else return bs---- | Exceptions that may arise while parsing messages. A user should--- never see one of these.-data MQTTException- = EOF- | ParseError String- deriving (Show, Typeable)+ void $ sendAwait mqtt pub SPUBREC+ void $ sendAwait mqtt+ (Message (Header False Confirm False)+ (PubRel (fromJust msgID)))+ SPUBCOMP -instance Exception MQTTException where
Network/MQTT/Encoding.hs view
@@ -1,4 +1,4 @@-{-# Language OverloadedStrings, RecordWildCards, GADTs #-}+{-# Language OverloadedStrings, RecordWildCards, GADTs, DataKinds #-} {-| Module: MQTT.Encoding Copyright: Lukas Braun 2014@@ -71,23 +71,23 @@ -- | Build the 'MessageBody' for any message type. putBody :: MessageBody t -> Builder-putBody (MConnect connect) = putConnect connect-putBody (MConnAck connAck) = putConnAck connAck-putBody (MPublish publish) = putPublish publish-putBody (MPubAck simpleMsg) = putSimple simpleMsg-putBody (MPubRec simpleMsg) = putSimple simpleMsg-putBody (MPubRel simpleMsg) = putSimple simpleMsg-putBody (MPubComp simpleMsg) = putSimple simpleMsg-putBody (MSubscribe subscribe) = putSubscribe subscribe-putBody (MSubAck subAck) = putSubAck subAck-putBody (MUnsubscribe unsubscribe) = putUnsubscribe unsubscribe-putBody (MUnsubAck simpleMsg) = putSimple simpleMsg-putBody MPingReq = mempty-putBody MPingResp = mempty-putBody MDisconnect = mempty+putBody (m@Connect{}) = putConnect m+putBody (m@ConnAck {}) = putConnAck m+putBody (m@Publish{}) = putPublish m+putBody (PubAck m) = putMsgID m+putBody (PubRec m) = putMsgID m+putBody (PubRel m) = putMsgID m+putBody (PubComp m) = putMsgID m+putBody (m@Subscribe{}) = putSubscribe m+putBody (m@SubAck{}) = putSubAck m+putBody (m@Unsubscribe{}) = putUnsubscribe m+putBody (UnsubAck m) = putMsgID m+putBody PingReq = mempty+putBody PingResp = mempty+putBody Disconnect = mempty -putConnect :: Connect -> Builder+putConnect :: MessageBody 'CONNECT -> Builder putConnect Connect{..} = mconcat [ putMqttText "MQIsdp" -- protocol , word8 3 -- version@@ -109,11 +109,11 @@ shiftL (toBit cleanSession) 1 -putConnAck :: ConnAck -> Builder-putConnAck = word8 . returnCode+putConnAck :: MessageBody 'CONNACK -> Builder+putConnAck ConnAck{..} = word8 0 {- reserved -} <> word8 returnCode -putPublish :: Publish -> Builder+putPublish :: MessageBody 'PUBLISH -> Builder putPublish Publish{..} = mconcat [ putTopic topic , maybe mempty putMsgID pubMsgID@@ -121,29 +121,26 @@ ] -putSubscribe :: Subscribe -> Builder+putSubscribe :: MessageBody 'SUBSCRIBE -> Builder putSubscribe Subscribe{..} = mconcat [ putMsgID subscribeMsgID , foldMap (\(txt, qos) -> putTopic txt <> word8 (fromQoS qos)) subTopics ] -putSubAck :: SubAck -> Builder+putSubAck :: MessageBody 'SUBACK -> Builder putSubAck SubAck{..} = mconcat [ putMsgID subAckMsgID , foldMap (word8 . fromQoS) granted ] -putUnsubscribe :: Unsubscribe -> Builder+putUnsubscribe :: MessageBody 'UNSUBSCRIBE -> Builder putUnsubscribe Unsubscribe{..} = mconcat [ putMsgID unsubMsgID , foldMap putTopic unsubTopics ] -putSimple :: SimpleMsg -> Builder-putSimple = putMsgID . msgID - --------------------------------- -- * Utility functions ---------------------------------@@ -174,17 +171,17 @@ -- | Encode the type of a 'MessageBody'. msgType :: (Num a) => MessageBody t -> a-msgType (MConnect _) = 1-msgType (MConnAck _) = 2-msgType (MPublish _) = 3-msgType (MPubAck _) = 4-msgType (MPubRec _) = 5-msgType (MPubRel _) = 6-msgType (MPubComp _) = 7-msgType (MSubscribe _) = 8-msgType (MSubAck _) = 9-msgType (MUnsubscribe _) = 10-msgType (MUnsubAck _) = 11-msgType MPingReq = 12-msgType MPingResp = 13-msgType MDisconnect = 14+msgType (Connect{}) = 1+msgType (ConnAck{}) = 2+msgType (Publish{}) = 3+msgType (PubAck{}) = 4+msgType (PubRec{}) = 5+msgType (PubRel{}) = 6+msgType (PubComp{}) = 7+msgType (Subscribe{}) = 8+msgType (SubAck{}) = 9+msgType (Unsubscribe{}) = 10+msgType (UnsubAck{}) = 11+msgType PingReq = 12+msgType PingResp = 13+msgType Disconnect = 14
+ Network/MQTT/Internal.hs view
@@ -0,0 +1,375 @@+{-# Language ScopedTypeVariables,+ DataKinds,+ GADTs #-}+{-|+Module: Network.MQTT.Internal+Copyright: Lukas Braun 2014-2016+License: GPL-3+Maintainer: koomi+mqtt@hackerspace-bamberg.de++MQTT Internals.++Use with care and expected changes.+-}+module Network.MQTT.Internal+ (+ -- * User interaction+ Config(..)+ , Terminated(..)+ , Command(..)+ , Commands(..)+ , mkCommands+ , send+ , await+ , AwaitMessage(..)+ , stopWaiting+ , sendAwait+ , writeCmd+ -- * Main loop+ , mainLoop+ , WaitTerminate+ , SendSignal+ , MqttState(..)+ , ParseCC+ , Input(..)+ , waitForInput+ , parseBytes+ , handleMessage+ , publishHandler+ -- * Misc+ , secToMicro+ ) where++import Control.Applicative ((<$>))+import Control.Concurrent+import qualified Control.Concurrent.Async as Async+import Control.Concurrent.STM+import Control.Exception (bracketOnError)+import Control.Monad (void, forever, filterM)+import Control.Monad.IO.Class (liftIO, MonadIO)+import Control.Monad.Loops (untilJust)+import Control.Monad.State.Strict (evalStateT, gets, modify, StateT)+import Data.Attoparsec.ByteString (IResult(..) , parse)+import Data.ByteString (ByteString)+import qualified Data.ByteString as BS+import Data.Foldable (for_)+import Data.Maybe (isNothing, fromMaybe)+import Data.Singletons (SingI(..))+import Data.Singletons.Decide+import Data.Text (Text)+import Data.Word (Word16)+import Network+import System.IO (Handle, hLookAhead)+import System.Timeout (timeout)++import Network.MQTT.Types+import Network.MQTT.Parser (message)+import Network.MQTT.Encoding (writeTo)++-----------------------------------------+-- User interaction+-----------------------------------------++-- | Reasons for why the connection was terminated.+data Terminated+ = ParseFailed [String] String+ -- ^ at the context in @['String']@ with the given message.+ | ConnectFailed ConnectError+ | UserRequested+ -- ^ 'disconnect' was called+ deriving Show++-- | Commands from the user for the 'mainLoop'.+data Command+ = CmdDisconnect+ | CmdSend SomeMessage+ | CmdAwait AwaitMessage+ | CmdStopWaiting AwaitMessage++-- | The communication channel used by 'Network.MQTT.publish',+-- 'Network.MQTT.subscribe', etc.+newtype Commands = Cmds { getCmds :: TChan Command }+-- | Create a new 'Commands' channel.+--+-- There should be one channel per MQTT connection. It might be reused by+-- subsequent connections, but never by multiple connections concurrently.+mkCommands :: IO Commands+mkCommands = Cmds <$> newTChanIO++data AwaitMessage where+ AwaitMessage :: SingI t => MVar (Message t) -> Maybe MsgID -> AwaitMessage++instance Eq AwaitMessage where+ AwaitMessage (var :: MVar (Message t)) mMsgID == AwaitMessage (var' :: MVar (Message t')) mMsgID' =+ case (sing :: SMsgType t) %~ (sing :: SMsgType t') of+ Proved Refl -> mMsgID == mMsgID' && var == var'+ Disproved _ -> False++-- | The various options when establishing a connection.+-- See below for available accessors.+data Config+ = Config+ { cHost :: HostName+ -- ^ Hostname of the broker.+ , cPort :: PortNumber+ -- ^ Port of the broker.+ , cClean :: Bool+ -- ^ Should the server forget subscriptions and other state on+ -- disconnects?+ , cWill :: Maybe Will+ -- ^ Optional 'Will' message.+ , cUsername :: Maybe Text+ -- ^ Optional username used for authentication.+ , cPassword :: Maybe Text+ -- ^ Optional password used for authentication.+ , cKeepAlive :: Maybe Word16+ -- ^ Time (in seconds) after which a 'PingReq' is sent to the broker if+ -- no regular message was sent. 'Nothing' means no limit.+ , cClientID :: Text+ -- ^ Client ID used by the server to identify clients.+ , cLogDebug :: String -> IO ()+ -- ^ Function for debug-level logging.+ , cResendTimeout :: Int+ -- ^ Time in seconds after which messages that have not been but+ -- should be acknowledged are retransmitted.+ , cPublished :: TChan (Message 'PUBLISH)+ -- ^ The channel received 'Publish' messages are written to.+ , cCommands :: Commands+ -- ^ The channel used by 'publish', 'subscribe', etc.+ , cInputBufferSize :: Int+ -- ^ Maximum number of bytes read from the network at once.+ }++-- | Tell the 'mainLoop' to send the given 'Message'.+send :: SingI t => Config -> Message t -> IO ()+send mqtt = writeCmd mqtt . CmdSend . SomeMessage++-- | Tell the 'MQTT' instance to place the next 'Message' of correct+-- 'MsgType' and 'MsgID' (if present) into the 'MVar'.+await :: SingI t => Config -> MVar (Message t) -> Maybe MsgID+ -> IO AwaitMessage+await mqtt var mMsgID = do+ writeCmd mqtt $ CmdAwait awaitMsg+ return awaitMsg+ where+ awaitMsg = AwaitMessage var mMsgID++-- | Stop waiting for the described 'Message'.+stopWaiting :: Config -> AwaitMessage -> IO ()+stopWaiting mqtt = writeCmd mqtt . CmdStopWaiting++-- | Execute the common pattern of sending a message and awaiting+-- a response in a safe, non-racy way. The message message is retransmitted+-- if no response has been received after 'cResendTimeout' seconds, with+-- exponential backoff for further retransmissions+--+-- An incoming message is considered a response if it is of the+-- requested type and the 'MsgID's match (if present).+sendAwait :: (SingI t, SingI r)+ => Config -> Message t -> SMsgType r -> IO (Message r)+sendAwait mqtt msg _responseS = do+ var <- newEmptyMVar+ let mMsgID = getMsgID (body msg)+ bracketOnError+ (await mqtt var mMsgID)+ (stopWaiting mqtt)+ (\_ ->+ let wait = do+ received <- readMVar var+ if isNothing mMsgID || mMsgID == getMsgID (body received)+ then return received+ else wait+ keepTrying msg' tout = do+ send mqtt msg'+ let retransmit = do+ cLogDebug mqtt "No response within timeout, retransmitting..."+ keepTrying (setDup msg') (tout * 2)+ timeout tout wait >>= maybe retransmit return+ in keepTrying msg initialTout)+ where+ initialTout = secToMicro $ cResendTimeout mqtt+++-----------------------------------------+-- Main loop+-----------------------------------------++type ParseCC = ByteString -> IResult ByteString SomeMessage++-- | Internal state for the main loop+data MqttState+ = MqttState+ { msParseCC :: ParseCC -- ^ Current parser continuation+ , msUnconsumed :: BS.ByteString -- ^ Not yet parsed input+ , msWaiting :: [AwaitMessage] -- ^ Messages we're waiting for+ }++-- | Input for the main loop+data Input+ = InMsg SomeMessage+ | InErr Terminated+ | InCmd Command++type WaitTerminate = STM ()+type SendSignal = MVar ()++mainLoop :: Config -> Handle -> WaitTerminate -> SendSignal -> IO Terminated+mainLoop mqtt h waitTerminate sendSignal = do+ void $ forkMQTT waitTerminate $ keepAliveLoop mqtt sendSignal+ evalStateT+ (handshake >>= maybe (liftIO (cLogDebug mqtt "Connected") >> go) return)+ (MqttState (parse message) BS.empty [])+ where+ go = do+ input <- waitForInput mqtt h+ case input of+ InErr err -> liftIO $+ return err+ InMsg someMsg -> do+ liftIO $ cLogDebug mqtt $ "Received " ++ show (toMsgType' someMsg)+ handleMessage mqtt waitTerminate someMsg+ go+ InCmd cmd -> case cmd of+ CmdDisconnect -> liftIO $ do+ doSend msgDisconnect+ return UserRequested+ CmdSend (SomeMessage msg) -> do+ doSend msg+ go+ CmdAwait awaitMsg -> do+ modify $ \s -> s { msWaiting = awaitMsg : msWaiting s }+ go+ CmdStopWaiting awaitMsg -> do+ modify $ \s -> s { msWaiting = filter (== awaitMsg) $ msWaiting s }+ go++ handshake :: StateT MqttState IO (Maybe Terminated)+ handshake = do+ doSend msgConnect+ input <- untilJust (getSome mqtt h >>= parseBytes)+ case input of+ InErr err -> return $ Just err+ InMsg someMsg -> return $ case someMsg of+ SomeMessage (Message _ (ConnAck retCode)) ->+ if retCode /= 0+ then Just $ ConnectFailed $ toConnectError retCode+ else Nothing+ _ -> Just $ ConnectFailed InvalidResponse+ InCmd _ -> error "parseBytes returned InCmd, this should not happen."+ where+ msgConnect = Message+ (Header False NoConfirm False)+ (Connect+ (cClean mqtt)+ (cWill mqtt)+ (MqttText $ cClientID mqtt)+ (MqttText <$> cUsername mqtt)+ (MqttText <$> cPassword mqtt)+ (fromMaybe 0 $ cKeepAlive mqtt))++ msgDisconnect = Message (Header False NoConfirm False) Disconnect++ doSend :: (MonadIO io, SingI t) => Message t -> io ()+ doSend msg = liftIO $ do+ cLogDebug mqtt $ "Sending " ++ show (toMsgType msg)+ writeTo h msg+ void $ tryPutMVar sendSignal ()++waitForInput :: Config -> Handle -> StateT MqttState IO Input+waitForInput mqtt h = do+ let cmdChan = getCmds $ cCommands mqtt+ unconsumed <- gets msUnconsumed+ if BS.null unconsumed+ then do+ -- wait until input is available, but don't retrieve it yet to avoid+ -- loosing anything+ input <- liftIO $ Async.race+ (void $ hLookAhead h)+ (void $ atomically $ peekTChan cmdChan)+ -- now we have committed to one source and can actually read it+ case input of+ Left () -> getSome mqtt h >>= parseUntilDone+ Right () -> InCmd <$> liftIO (atomically (readTChan cmdChan))+ else+ parseUntilDone unconsumed+ where+ parseUntilDone bytes = parseBytes bytes >>= maybe (waitForInput mqtt h) return++-- | Parse the given 'ByteString' and update the current 'MqttState'.+--+-- Returns 'Nothing' if more input is needed.+parseBytes :: Monad m => ByteString -> StateT MqttState m (Maybe Input)+parseBytes bytes = do+ parseCC <- gets msParseCC+ case parseCC bytes of+ Fail _unconsumed context err ->+ return $ Just $ InErr $ ParseFailed context err+ Partial cont -> do+ modify $ \s -> s { msParseCC = cont+ , msUnconsumed = BS.empty }+ return Nothing+ Done unconsumed someMsg -> do+ modify $ \s -> s { msParseCC = parse message+ , msUnconsumed = unconsumed }+ return $ Just $ InMsg someMsg++handleMessage :: Config -> WaitTerminate -> SomeMessage -> StateT MqttState IO ()+handleMessage mqtt waitTerminate (SomeMessage msg) =+ case toSMsgType msg %~ SPUBLISH of+ Proved Refl -> liftIO $ void $ forkMQTT waitTerminate $ publishHandler mqtt msg+ Disproved _ -> do+ waiting' <- gets msWaiting >>= liftIO . filterM giveToWaiting+ modify (\s -> s { msWaiting = waiting' })+ where+ giveToWaiting :: AwaitMessage -> IO Bool+ giveToWaiting (AwaitMessage (var :: MVar (Message t')) mMsgID')+ | isNothing mMsgID || mMsgID == mMsgID' =+ case toSMsgType msg %~ (sing :: SMsgType t') of+ Proved Refl -> putMVar var msg >> return False+ Disproved _ -> return True+ | otherwise = return True++ mMsgID = getMsgID (body msg)++keepAliveLoop :: Config -> SendSignal -> IO ()+keepAliveLoop mqtt signal = for_ (cKeepAlive mqtt) $ \tout -> forever $ do+ rslt <- timeout (secToMicro (fromIntegral tout)) (takeMVar signal)+ case rslt of+ Nothing -> void $ sendAwait mqtt+ (Message (Header False NoConfirm False) PingReq)+ SPINGRESP+ Just _ -> return ()++publishHandler :: Config -> Message 'PUBLISH -> IO ()+publishHandler mqtt msg = do+ case (qos (header msg), pubMsgID (body msg)) of+ (Confirm, Just msgid) -> do+ release+ send mqtt $ Message (Header False NoConfirm False) (PubAck msgid)+ (Handshake, Just msgid) -> do+ _ <- sendAwait mqtt+ (Message (Header False NoConfirm False) (PubRec msgid))+ SPUBREL+ release+ send mqtt $ Message (Header False NoConfirm False) (PubComp msgid)+ _ -> release+ where+ release = writeTChanIO (cPublished mqtt) msg++getSome :: MonadIO m => Config -> Handle -> m ByteString+getSome mqtt h = liftIO (BS.hGetSome h (cInputBufferSize mqtt))++-- | Runs the 'IO' action in a seperate thread and cancels it if the 'mainLoop'+-- exits earlier.+forkMQTT :: WaitTerminate -> IO () -> IO (Async.Async ())+forkMQTT waitTerminate action = Async.async $ Async.withAsync action $ \forked ->+ atomically $ waitTerminate `orElse` Async.waitSTM forked++writeTChanIO :: TChan a -> a -> IO ()+writeTChanIO chan = atomically . writeTChan chan++writeCmd :: Config -> Command -> IO ()+writeCmd mqtt = writeTChanIO (getCmds $ cCommands mqtt)++secToMicro :: Int -> Int+secToMicro m = m * 10 ^ (6 :: Int)
− Network/MQTT/Logger.hs
@@ -1,39 +0,0 @@-{-|-Module: MQTT.Logger-Copyright: Lukas Braun 2014-License: GPL-3-Maintainer: koomi+mqtt@hackerspace-bamberg.de--A simple logger abstraction.--}-module Network.MQTT.Logger where--import System.IO---- | Absract logger with three levels of importance.-data Logger- = Logger- { logInfo :: String -> IO ()- , logWarning :: String -> IO ()- , logError :: String -> IO ()- }---- | 'logInfo' prints to stdout, 'logWarning' and 'logError' to stderr--- (with [Warning]/[Error] prefix)-stdLogger :: Logger-stdLogger = Logger- putStrLn- (\msg -> hPutStrLn stderr $ "[Warning] " ++ msg)- (\msg -> hPutStrLn stderr $ "[Error] " ++ msg)---- | Log only warnings and errors, ignoring anything passed to 'logInfo'.-warnings :: Logger -> Logger-warnings l = l { logInfo = ignore }---- | Like 'warnings', but logs only errors.-errors :: Logger -> Logger-errors l = l { logInfo = ignore, logWarning = ignore }---- | Ignore the message.-ignore :: String -> IO ()-ignore _ = return ()
Network/MQTT/Parser.hs view
@@ -1,7 +1,7 @@-{-# Language OverloadedStrings, GADTs #-}+{-# Language OverloadedStrings, GADTs, DataKinds #-} {-| Module: MQTT.Parsers-Copyright: Lukas Braun 2014+Copyright: Lukas Braun 2014-2016 License: GPL-3 Maintainer: koomi+mqtt@hackerspace-bamberg.de @@ -16,8 +16,8 @@ import Data.Attoparsec.ByteString import Data.Bits import qualified Data.ByteString as BS-import Data.Singletons (withSomeSing)-import Data.Text.Encoding (decodeUtf8')+import Data.Text.Encoding (decodeUtf8With)+import Data.Text.Encoding.Error (lenientDecode) import Data.Word import Prelude hiding (takeWhile, take) @@ -31,8 +31,9 @@ message = do (msgType, header) <- mqttHeader remaining <- parseRemaining- withSomeSing msgType $ \sMsgType ->+ msg <- withSomeSingI msgType $ \sMsgType -> SomeMessage . Message header <$> mqttBody header sMsgType remaining+ return msg ---------------------------------@@ -41,7 +42,7 @@ -- | Parser for the fixed header part of a MQTT message. mqttHeader :: Parser (MsgType, MqttHeader)-mqttHeader = do+mqttHeader = ctxt "mqttHeader" $ do byte1 <- anyWord8 qos <- toQoS $ 3 .&. shiftR byte1 1 let retain = testBit byte1 0@@ -68,14 +69,15 @@ -- | Parse the 'remaining length' field that indicates how long the rest of -- the message is. parseRemaining :: Parser Word32-parseRemaining = do+parseRemaining = ctxt "parseRemaining" $ do bytes <- takeWhile (> 0x7f) -- bytes with first bit set when (BS.length bytes > 3) $ fail "'Remaining length' field must not be longer than 4 bytes" stopByte <- anyWord8- return $ snd $ BS.foldr' f (128, fromIntegral stopByte) bytes+ let (factor, acc) = BS.foldl' f (1, 0) bytes+ return $ acc + factor * fromIntegral stopByte where- f byte (factor, acc) =+ f (factor, acc) byte = (factor*128, acc + factor * fromIntegral (0x7f .&. byte)) @@ -86,27 +88,27 @@ -- | «@mqttBody header msgtype remaining@» parses a 'Message' of type -- @msgtype@ that is @remaining@ bytes long. mqttBody :: MqttHeader -> SMsgType t -> Word32 -> Parser (MessageBody t)-mqttBody header msgType remaining =+mqttBody header msgType remaining = ctxt "mqttBody" $ let parser = case msgType of- SCONNECT -> MConnect <$> connect- SCONNACK -> MConnAck <$> connAck- SPUBLISH -> MPublish <$> publish header- SPUBACK -> MPubAck <$> simpleMsg- SPUBREC -> MPubRec <$> simpleMsg- SPUBREL -> MPubRel <$> simpleMsg- SPUBCOMP -> MPubComp <$> simpleMsg- SSUBSCRIBE -> MSubscribe <$> subscribe- SSUBACK -> MSubAck <$> subAck- SUNSUBSCRIBE -> MUnsubscribe <$> unsubscribe- SUNSUBACK -> MUnsubAck <$> simpleMsg- SPINGREQ -> pure MPingReq- SPINGRESP -> pure MPingResp- SDISCONNECT -> pure MDisconnect+ SCONNECT -> connect+ SCONNACK -> connAck+ SPUBLISH -> publish header+ SPUBACK -> PubAck <$> parseMsgID+ SPUBREC -> PubRec <$> parseMsgID+ SPUBREL -> PubRel <$> parseMsgID+ SPUBCOMP -> PubComp <$> parseMsgID+ SSUBSCRIBE -> subscribe+ SSUBACK -> subAck+ SUNSUBSCRIBE -> unsubscribe+ SUNSUBACK -> UnsubAck <$> parseMsgID+ SPINGREQ -> pure PingReq+ SPINGRESP -> pure PingResp+ SDISCONNECT -> pure Disconnect in evalStateT parser remaining -connect :: MessageParser Connect-connect = do+connect :: MessageParser (MessageBody 'CONNECT)+connect = ctxt' "connect" $ do protocol version @@ -123,26 +125,26 @@ mWill <- parseIf willFlag $ Will (testBit flags 5) <$> toQoS (3 .&. shiftR flags 3)- <*> fmap toTopic mqttText- <*> mqttText+ <*> (ctxt' "Will Topic" $ fmap toTopic mqttText)+ <*> (ctxt' "Will Message" mqttText) - username <- parseIf usernameFlag mqttText+ username <- ctxt' "Username" $ parseIf usernameFlag mqttText - password <- parseIf passwordFlag mqttText+ password <- ctxt' "Password" $ parseIf passwordFlag mqttText return $ Connect clean mWill clientID username password keepAlive where- protocol = do+ protocol = ctxt' "protocol" $ do prot <- mqttText when (prot /= "MQIsdp") $ fail $ "Invalid protocol: " ++ show prot - version = do+ version = ctxt' "version" $ do version <- anyWord8' when (version /= 3) $ fail $ "Invalid version: " ++ show version - getClientID = do+ getClientID = ctxt' "getClientID" $ do before <- get clientID <- mqttText after <- get@@ -155,58 +157,50 @@ parseIf :: Applicative f => Bool -> f a -> f (Maybe a) parseIf flag parser = if flag then Just <$> parser else pure Nothing --connAck :: MessageParser ConnAck-connAck = ConnAck <$> anyWord8'+connAck :: MessageParser (MessageBody 'CONNACK)+connAck = ctxt' "connAck" $ anyWord8' {- reserved -} *> (ConnAck <$> anyWord8') -publish :: MqttHeader -> MessageParser Publish-publish header = Publish+publish :: MqttHeader -> MessageParser (MessageBody 'PUBLISH)+publish header = ctxt' "publish" $ Publish <$> getTopic <*> (if qos header > NoConfirm then Just <$> parseMsgID else return Nothing) <*> (get >>= take') -subscribe :: MessageParser Subscribe-subscribe = Subscribe+subscribe :: MessageParser (MessageBody 'SUBSCRIBE)+subscribe = ctxt' "subscribe" $ Subscribe <$> parseMsgID <*> whileM ((0 <) <$> get) ((,) <$> getTopic <*> (anyWord8' >>= toQoS)) -subAck :: MessageParser SubAck-subAck = SubAck+subAck :: MessageParser (MessageBody 'SUBACK)+subAck = ctxt' "subAck" $ SubAck <$> parseMsgID <*> whileM ((0 <) <$> get) (anyWord8' >>= toQoS) -unsubscribe :: MessageParser Unsubscribe-unsubscribe = Unsubscribe+unsubscribe :: MessageParser (MessageBody 'UNSUBSCRIBE)+unsubscribe = ctxt' "unsubscribe" $ Unsubscribe <$> parseMsgID <*> whileM ((0 <) <$> get) getTopic -simpleMsg :: MessageParser SimpleMsg-simpleMsg = SimpleMsg <$> parseMsgID - --------------------------------- -- * Utility functions --------------------------------- -- | Parse a topic name. getTopic :: MessageParser Topic-getTopic = toTopic <$> mqttText+getTopic = ctxt' "getTopic" $ toTopic <$> mqttText -- | Parse a length-prefixed UTF-8 string. mqttText :: MessageParser MqttText-mqttText = do- n <- anyWord16BE- rslt <- decodeUtf8' <$> take' n- case rslt of- Left err -> fail $ "Invalid UTF-8: " ++ show err- Right txt -> return $ MqttText txt+mqttText = ctxt' "mqttText" $+ MqttText . decodeUtf8With lenientDecode <$> (anyWord16BE >>= take') -- | Synonym for 'anyWord16BE'. parseMsgID :: MessageParser Word16-parseMsgID = anyWord16BE+parseMsgID = ctxt' "parseMsgID" anyWord16BE -- | Parse a big-endian 16bit integer. anyWord16BE :: (Num a, Bits a) => MessageParser a@@ -219,6 +213,12 @@ -- the remaining length. anyWord8' :: MessageParser Word8 anyWord8' = parseLength 1 >> lift anyWord8++ctxt :: String -> Parser a -> Parser a+ctxt = flip (<?>)++ctxt' :: String -> MessageParser a -> MessageParser a+ctxt' = mapStateT . ctxt -- | A lifted version of attoparsec's 'take' that also subtracts the -- length.
Network/MQTT/Types.hs view
@@ -1,16 +1,25 @@ {-# Language GeneralizedNewtypeDeriving,- PatternSynonyms,+ DeriveDataTypeable, OverloadedStrings, DataKinds, KindSignatures, GADTs, TypeFamilies, ScopedTypeVariables,+ RankNTypes, TemplateHaskell #-}++-- without -O0 GHC 7.6.3 loops while building, probably related to+-- https://git.haskell.org/ghc.git/commitdiff/c1edbdfd9148ad9f74bfe41e76c524f3e775aaaa+--+-- -fno-warn-unused-binds is used because the Singletons TH magic generates a+-- lot of unused binds and GHC has no way to disable warnings locally+{-# OPTIONS_GHC -O0 -fno-warn-unused-binds #-}+ {-| Module: MQTT.Types-Copyright: Lukas Braun 2014+Copyright: Lukas Braun 2014-2016 License: GPL-3 Maintainer: koomi+mqtt@hackerspace-bamberg.de @@ -21,25 +30,23 @@ Message(..) , SomeMessage(..) , MqttHeader(..)- -- * Message bodies+ , setDup+ -- * Message body , MessageBody(..)- , Connect(..)- , ConnAck(..)- , Publish(..)- , Subscribe(..)- , SubAck(..)- , Unsubscribe(..)- , SimpleMsg(..) -- * Miscellaneous , Will(..) , QoS(..) , MsgID , getMsgID , Topic+ , matches , fromTopic , toTopic- , matches+ , getLevels+ , fromLevels , MqttText(..)+ , ConnectError(..)+ , toConnectError -- * Message types , MsgType(..) , toMsgType@@ -53,6 +60,7 @@ -- about the flow of 'Message's. , toSMsgType , SMsgType+ , withSomeSingI , Sing( SCONNECT , SCONNACK , SPUBLISH@@ -69,11 +77,14 @@ , SDISCONNECT) ) where +import Control.Exception (Exception) import Data.ByteString (ByteString)+import Data.Singletons import Data.Singletons.TH import Data.String (IsString(..)) import Data.Text (Text) import qualified Data.Text as T+import Data.Typeable (Typeable) import Data.Word -- | A MQTT message, indexed by the type of the message ('MsgType').@@ -85,106 +96,75 @@ -- | Any message, hiding the index. data SomeMessage where- SomeMessage :: Message t -> SomeMessage+ SomeMessage :: SingI t => Message t -> SomeMessage -- | Fixed header required in every message. data MqttHeader = Header- { -- msgType :: MsgType -- ^ Type of the message- dup :: Bool -- ^ Has this message been sent before?+ { dup :: Bool -- ^ Has this message been sent before? , qos :: QoS -- ^ Quality of Service-level , retain :: Bool -- ^ Should the broker retain the message for -- future subscribers? } deriving (Eq, Ord, Show) +-- | Set the 'dup' flag to 'True'.+setDup :: Message t -> Message t+setDup (Message h b) = Message h { dup = True } b+ -- | The body of a MQTT message, indexed by the type of the message ('MsgType'). data MessageBody (t :: MsgType) where- MConnect :: Connect -> MessageBody CONNECT- MConnAck :: ConnAck -> MessageBody CONNACK- MPublish :: Publish -> MessageBody PUBLISH- MPubAck :: SimpleMsg -> MessageBody PUBACK- MPubRec :: SimpleMsg -> MessageBody PUBREC- MPubRel :: SimpleMsg -> MessageBody PUBREL- MPubComp :: SimpleMsg -> MessageBody PUBCOMP- MSubscribe :: Subscribe -> MessageBody SUBSCRIBE- MSubAck :: SubAck -> MessageBody SUBACK- MUnsubscribe :: Unsubscribe -> MessageBody UNSUBSCRIBE- MUnsubAck :: SimpleMsg -> MessageBody UNSUBACK- MPingReq :: MessageBody PINGREQ- MPingResp :: MessageBody PINGRESP- MDisconnect :: MessageBody DISCONNECT+ Connect :: { cleanSession :: Bool+ -- ^ Should the server forget subscriptions and other state on+ -- disconnects?+ , will :: Maybe Will+ -- ^ Optional 'Will' message.+ , clientID :: MqttText+ -- ^ Client ID used by the server to identify clients.+ , username :: Maybe MqttText+ -- ^ Optional username used for authentication.+ , password :: Maybe MqttText+ -- ^ Optional password used for authentication.+ , keepAlive :: Word16+ -- ^ Time (in seconds) after which a 'PingReq' is sent to the broker if+ -- no regular message was sent. 0 means no limit.+ } -> MessageBody 'CONNECT+ ConnAck :: { returnCode :: Word8 } -> MessageBody 'CONNACK+ Publish :: { topic :: Topic+ -- ^ The 'Topic' to which the message should be published.+ , pubMsgID :: Maybe MsgID+ -- ^ 'MsgID' of the message if 'QoS' > 'NoConfirm'.+ , payload :: ByteString+ -- ^ The content that will be published.+ } -> MessageBody 'PUBLISH+ PubAck :: { pubAckMsgID :: MsgID } -> MessageBody 'PUBACK+ PubRec :: { pubRecMsgID :: MsgID } -> MessageBody 'PUBREC+ PubRel :: { pubRelMsgID :: MsgID } -> MessageBody 'PUBREL+ PubComp :: { pubCompMsgID :: MsgID } -> MessageBody 'PUBCOMP+ Subscribe :: { subscribeMsgID :: MsgID+ , subTopics :: [(Topic, QoS)]+ -- ^ The 'Topic's and corresponding requested 'QoS'.+ } -> MessageBody 'SUBSCRIBE+ SubAck :: { subAckMsgID :: MsgID+ , granted :: [QoS]+ -- ^ The 'QoS' granted for each 'Topic' in the order they were sent+ -- in the SUBSCRIBE.+ } -> MessageBody 'SUBACK+ Unsubscribe :: { unsubMsgID :: MsgID+ , unsubTopics :: [Topic]+ -- ^ The 'Topic's from which the client should be unsubscribed.+ } -> MessageBody 'UNSUBSCRIBE+ UnsubAck :: { unsubAckMsgID :: MsgID } -> MessageBody 'UNSUBACK+ PingReq :: MessageBody 'PINGREQ+ PingResp :: MessageBody 'PINGRESP+ Disconnect :: MessageBody 'DISCONNECT --- | The fields of a CONNECT message.-data Connect- = Connect- { cleanSession :: Bool- -- ^ Should the server forget subscriptions and other state on- -- disconnects?- , will :: Maybe Will- -- ^ Optional 'Will' message.- , clientID :: MqttText- -- ^ Client ID used by the server to identify clients.- , username :: Maybe MqttText- -- ^ Optional username used for authentication.- , password :: Maybe MqttText- -- ^ Optional password used for authentication.- , keepAlive :: Word16- -- ^ Maximum interval (in seconds) in which a message must be sent.- -- 0 means no limit.- } deriving (Show, Eq) --- | The response to a CONNECT. Anything other than 0 means the broker--- refused the connection--- (<http://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html#connack details>).-newtype ConnAck = ConnAck { returnCode :: Word8 }- deriving (Show, Eq)---- | The fields of a PUBLISH message.-data Publish- = Publish- { topic :: Topic- -- ^ The 'Topic' to which the message should be published.- , pubMsgID :: Maybe MsgID- -- ^ 'MsgID' of the message if 'QoS' > 'NoConfirm'.- , payload :: ByteString- -- ^ The content that will be published.- } deriving (Show, Eq)---- | The fields of a SUBSCRIBE message.-data Subscribe- = Subscribe- { subscribeMsgID :: MsgID- , subTopics :: [(Topic, QoS)]- -- ^ The 'Topic's and corresponding requested 'QoS'.- } deriving (Show, Eq)---- | The fields of a SUBACK message.-data SubAck- = SubAck- { subAckMsgID :: MsgID- , granted :: [QoS]- -- ^ The 'QoS' granted for each 'Topic' in the order they were sent- -- in the SUBSCRIBE.- } deriving (Show, Eq)---- | The fields of a UNSUBSCRIBE message.-data Unsubscribe- = Unsubscribe- { unsubMsgID :: MsgID- , unsubTopics :: [Topic]- -- ^ The 'Topic's from which the client should be unsubscribed.- } deriving (Show, Eq)---- | Any message body that consists only of a 'MsgID'.-newtype SimpleMsg = SimpleMsg { msgID :: MsgID }- deriving (Show, Eq)- -- | The different levels of QoS data QoS- = NoConfirm -- ^ Fire and forget- | Confirm -- ^ Acknowledged delivery (repeat until ack)- | Handshake -- ^ Assured delivery (four-step handshake)+ = NoConfirm -- ^ Fire and forget, message will be published at most once.+ | Confirm -- ^ Acknowledged delivery, message will be published at least once.+ | Handshake -- ^ Assured delivery, message will be published exactly once. deriving (Eq, Ord, Enum, Show) -- | A Will message is published by the broker if a client disconnects@@ -202,42 +182,54 @@ newtype MqttText = MqttText { text :: Text } deriving (Eq, Show, IsString) --- | A topic is a "hierarchical name space that defines a taxonomy of--- information sources for which subscribers can register an interest."------ See--- <http://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html#appendix-a here>--- for more information on topics.-data Topic = Topic { levels :: [Text], orig :: Text }--- levels and orig should always refer to the same topic, this way no text--- has to be copied when converting from/to text- type MsgID = Word16 -- | Get the message ID of any message, if it exists. getMsgID :: MessageBody t -> Maybe MsgID-getMsgID (MConnect _) = Nothing-getMsgID (MConnAck _) = Nothing-getMsgID (MPublish pub) = pubMsgID pub-getMsgID (MPubAck simple) = Just (msgID simple)-getMsgID (MPubRec simple) = Just (msgID simple)-getMsgID (MPubRel simple) = Just (msgID simple)-getMsgID (MPubComp simple) = Just (msgID simple)-getMsgID (MSubscribe sub) = Just (subscribeMsgID sub)-getMsgID (MSubAck subA) = Just (subAckMsgID subA)-getMsgID (MUnsubscribe unsub) = Just (unsubMsgID unsub)-getMsgID (MUnsubAck simple) = Just (msgID simple)-getMsgID MPingReq = Nothing-getMsgID MPingResp = Nothing-getMsgID MDisconnect = Nothing+getMsgID (Connect{}) = Nothing+getMsgID (ConnAck{}) = Nothing+getMsgID (Publish _ mMsgid _) = mMsgid+getMsgID (PubAck msgid) = Just msgid+getMsgID (PubRec msgid) = Just msgid+getMsgID (PubRel msgid) = Just msgid+getMsgID (PubComp msgid) = Just msgid+getMsgID (Subscribe msgid _) = Just msgid+getMsgID (SubAck msgid _) = Just msgid+getMsgID (Unsubscribe msgid _) = Just msgid+getMsgID (UnsubAck msgid) = Just msgid+getMsgID PingReq = Nothing+getMsgID PingResp = Nothing+getMsgID Disconnect = Nothing +-- | A topic is a \"hierarchical name space that defines a taxonomy of+-- information sources for which subscribers can register an interest.\"+-- See the+-- <http://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html#appendix-a specification>+-- for more details.+--+-- A topic can be inspected by using the 'matches' function or after using+-- 'getLevels', e.g.:+--+-- > f1 topic+-- > | topic `matches` "mqtt/hs/example" = putStrLn "example"+-- > | topic `matches` "mqtt/hs/#" = putStrLn "wildcard"+-- >+-- > f2 topic = case getLevels topic of+-- > ["mqtt", "hs", "example"] -> putStrLn "example"+-- > "mqtt" : "hs" : _ -> putStrLn "wildcard"+data Topic = Topic { levels :: [Text], orig :: Text }+-- levels and orig should always refer to the same topic, so no text has to be+-- copied when converting from/to text+ instance Show Topic where show (Topic _ t) = show t instance Eq Topic where Topic _ t1 == Topic _ t2 = t1 == t2 --- | Check if one of the 'Topic's matches the other.+-- | Check if one of the 'Topic's matches the other, taking+-- <http://public.dhe.ibm.com/software/dw/webservices/ws-mqtt/mqtt-v3r1.html#appendix-a wildcards>+-- into consideration. matches :: Topic -> Topic -> Bool matches (Topic t1 _) (Topic t2 _) = go t1 t2 where@@ -254,10 +246,40 @@ fromTopic :: Topic -> MqttText fromTopic = MqttText . orig +-- | Split a topic into its individual levels.+getLevels :: Topic -> [Text]+getLevels = levels++-- | Create a 'Topic' from its individual levels.+fromLevels :: [Text] -> Topic+fromLevels ls = Topic ls (T.intercalate "/" ls)+ instance IsString Topic where fromString str = let txt = T.pack str in Topic (T.split (== '/') txt) txt +-- | Reasons why connecting to a broker might fail.+data ConnectError+ = WrongProtocolVersion+ | IdentifierRejected+ | ServerUnavailable+ | BadLogin+ | Unauthorized+ | UnrecognizedReturnCode+ | InvalidResponse+ deriving (Show, Typeable)++instance Exception ConnectError where++-- | Convert a return code to a 'ConnectError'.+toConnectError :: Word8 -> ConnectError+toConnectError 1 = WrongProtocolVersion+toConnectError 2 = IdentifierRejected+toConnectError 3 = ServerUnavailable+toConnectError 4 = BadLogin+toConnectError 5 = Unauthorized+toConnectError _ = UnrecognizedReturnCode+ -- | The various types of messages. data MsgType = CONNECT@@ -280,43 +302,17 @@ singDecideInstance ''MsgType -- | Determine the 'MsgType' of a 'Message'.-toMsgType :: Message t -> MsgType-toMsgType msg =- case body msg of- MConnect _ -> CONNECT- MConnAck _ -> CONNACK- MPublish _ -> PUBLISH- MPubAck _ -> PUBACK- MPubRec _ -> PUBREC- MPubRel _ -> PUBREL- MPubComp _ -> PUBCOMP- MSubscribe _ -> SUBSCRIBE- MSubAck _ -> SUBACK- MUnsubscribe _ -> UNSUBSCRIBE- MUnsubAck _ -> UNSUBACK- MPingReq -> PINGREQ- MPingResp -> PINGRESP- MDisconnect -> DISCONNECT+toMsgType :: SingI t => Message t -> MsgType+toMsgType = fromSing . toSMsgType -- | Determine the 'MsgType' of a 'SomeMessage'. toMsgType' :: SomeMessage -> MsgType toMsgType' (SomeMessage msg) = toMsgType msg -- | Determine the singleton 'SMsgType' of a 'Message'.-toSMsgType :: Message t -> SMsgType t-toSMsgType msg =- case body msg of- MConnect _ -> SCONNECT- MConnAck _ -> SCONNACK- MPublish _ -> SPUBLISH- MPubAck _ -> SPUBACK- MPubRec _ -> SPUBREC- MPubRel _ -> SPUBREL- MPubComp _ -> SPUBCOMP- MSubscribe _ -> SSUBSCRIBE- MSubAck _ -> SSUBACK- MUnsubscribe _ -> SUNSUBSCRIBE- MUnsubAck _ -> SUNSUBACK- MPingReq -> SPINGREQ- MPingResp -> SPINGRESP- MDisconnect -> SDISCONNECT+toSMsgType :: SingI t => Message t -> SMsgType t+toSMsgType _ = sing++-- | Helper to generate both an implicit and explicit singleton.+withSomeSingI :: MsgType -> (forall t. SingI t => SMsgType t -> r) -> r+withSomeSingI t f = withSomeSing t $ \s -> withSingI s $ f s
README.md view
@@ -3,15 +3,4 @@ A Haskell MQTT client library. -A simple example, assuming a broker is running on localhost (needs -XOverloadedStrings):--```haskell-import Network.MQTT-import Network.MQTT.Logger-Just mqtt <- connect defaultConfig { cLogger = warnings stdLogger }-let f t payload = putStrLn $ "A message was published to " ++ show t ++ ": " ++ show payload-subscribe mqtt NoConfirm "#" f-publish mqtt Handshake False "some random/topic" "Some content!"-```--See the haddock for more documentation.+[Hackage](https://hackage.haskell.org/package/mqtt-hs) and [examples](examples/).
mqtt-hs.cabal view
@@ -2,14 +2,14 @@ -- documentation, see http://haskell.org/cabal/users-guide/ name: mqtt-hs-version: 0.2.0+version: 1.0.0 synopsis: A MQTT client library.-homepage: github.com/k00mi/mqtt-hs+homepage: http://github.com/k00mi/mqtt-hs license: GPL-3 license-file: LICENSE author: Lukas Braun <koomi+mqtt@hackerspace-bamberg.de> maintainer: Lukas Braun <koomi+mqtt@hackerspace-bamberg.de>-copyright: (c) 2014 Lukas Braun+copyright: (c) 2014-2016 Lukas Braun category: Network build-type: Simple extra-source-files: README.md@@ -18,21 +18,27 @@ description: A library to communicate with MQTT brokers. - See the @Network.MQTT@ module for documentation and a small example.+ See the 'Network.MQTT' module for documentation and the project repository+ for some <https://github.com/k00mi/mqtt-hs/blob/master/examples examples>. library exposed-modules: Network.MQTT, Network.MQTT.Parser, Network.MQTT.Encoding,- Network.MQTT.Types, Network.MQTT.Logger- build-depends: base >=4.7 && <4.8,+ Network.MQTT.Types, Network.MQTT.Internal+ build-depends: base >=4.6 && <4.10,+ async >=2.0 && <2.2, mtl >=1.1 && <2.3,- attoparsec >=0.10 && <0.13,+ transformers >=0.2 && <0.6,+ attoparsec >=0.10 && <0.14, bytestring >=0.10.2 && <0.11,- text >=0.11.0.6 && <1.2,+ text >=0.11.0.6 && <1.3, network >=2.0 && <2.7,- singletons >=0.9 && < 1.1,+ singletons >=0.9 && < 2.2,+ stm >=2.4 && <2.5, monad-loops >=0.3 && <0.5 default-language: Haskell2010+ ghc-options: -Wall+ -fno-warn-name-shadowing source-repository head