amqp 0.8.0 → 0.8.1
raw patch · 10 files changed
+507/−1 lines, 10 files
Files
- amqp.cabal +11/−1
- test/BasicPublishSpec.hs +62/−0
- test/BasicRejectSpec.hs +107/−0
- test/ChannelSpec.hs +44/−0
- test/ConnectionSpec.hs +32/−0
- test/ExchangeDeclareSpec.hs +57/−0
- test/ExchangeDeleteSpec.hs +41/−0
- test/QueueDeclareSpec.hs +63/−0
- test/QueueDeleteSpec.hs +40/−0
- test/QueuePurgeSpec.hs +50/−0
amqp.cabal view
@@ -1,5 +1,5 @@ Name: amqp -Version: 0.8.0 +Version: 0.8.1 Synopsis: Client library for AMQP servers (currently only RabbitMQ) Description: Client library for AMQP servers (currently only RabbitMQ) . @@ -42,6 +42,16 @@ ., test main-is: Runner.hs + other-modules: + BasicPublishSpec + BasicRejectSpec + ChannelSpec + ConnectionSpec + ExchangeDeclareSpec + ExchangeDeleteSpec + QueueDeclareSpec + QueueDeleteSpec + QueuePurgeSpec build-depends: base >= 4 && < 5, binary >= 0.7, containers>=0.2, bytestring>=0.9, network>=2.2.3.1, data-binary-ieee754>=0.4.2.1, text>=0.11.2, split>=0.2, clock >= 0.4.0.1 , hspec >= 1.3
+ test/BasicPublishSpec.hs view
@@ -0,0 +1,62 @@+{-# OPTIONS -XOverloadedStrings #-} + +module BasicPublishSpec (main, spec) where + +import Test.Hspec +import Network.AMQP + +import Data.ByteString.Lazy.Char8 as BL + +import Control.Concurrent (threadDelay) + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "publishMsg" $ do + context "with a routing key" $ do + it "publishes a message" $ do + let q = "haskell-amqp.queues.publish-over-default-exchange1" + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + (_, n1, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + n1 `shouldBe` 0 + + -- publishes using default exchange + publishMsg ch "" q + (newMsg {msgBody = (BL.pack "hello")}) + threadDelay (1000 * 100) + + (_, n2, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + n2 `shouldBe` 1 + + n3 <- deleteQueue ch q + n3 `shouldBe` 1 + closeConnection conn + + context "with a blank routing key" $ do + it "publishes a message" $ do + let q = "haskell-amqp.queues.publish-over-fanout1" + e = "haskell-amqp.fanout.d.na" + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + _ <- declareExchange ch (newExchange {exchangeName = e, + exchangeType = "fanout", + exchangeDurable = True}) + + (_, _, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + _ <- purgeQueue ch q + bindQueue ch q e "" + + publishMsg ch e "" + (newMsg {msgBody = (BL.pack "hello")}) + threadDelay (1000 * 100) + + (_, n, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + n `shouldBe` 1 + + _ <- deleteQueue ch q + closeConnection conn
+ test/BasicRejectSpec.hs view
@@ -0,0 +1,107 @@+{-# OPTIONS -XOverloadedStrings #-} + +module BasicRejectSpec (main, spec) where + +import Test.Hspec +import Network.AMQP + +import Data.ByteString.Lazy.Char8 as L8 +import Control.Concurrent (threadDelay) +import Control.Exception (bracket) + +main :: IO () +main = hspec spec + +withTestConnection :: (Channel -> IO c) -> IO c +withTestConnection job = do + bracket (openConnection "127.0.0.1" "/" "guest" "guest") closeConnection $ \conn -> do + ch <- openChannel conn + job ch + +spec :: Spec +spec = do + describe "rejectMsg" $ do + context "requeue = True" $ do + it "requeues a message" $ withTestConnection $ \ch -> do + let q = "haskell-amqp.basic.reject.with-requeue-true" + + (_, n1, _) <- declareQueue ch $ newQueue {queueName = q, queueDurable = False} + n1 `shouldBe` 0 + + -- publishes using default exchange + publishMsg ch "" q $ newMsg {msgBody = (L8.pack "hello")} + threadDelay (1000 * 100) + + (_, n2, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + n2 `shouldBe` 1 + + Just (_msg, env) <- getMsg ch Ack q + rejectMsg (envChannel env) (envDeliveryTag env) True + threadDelay (1000 * 100) + + n3 <- deleteQueue ch q + n3 `shouldBe` 1 + + context "requeue = False" $ do + it "rejects a message" $ withTestConnection $ \ch -> do + let q = "haskell-amqp.basic.reject.with-requeue-false" + + (_, n1, _) <- declareQueue ch $ newQueue {queueName = q, queueDurable = False} + n1 `shouldBe` 0 + + -- publishes using default exchange + publishMsg ch "" q $ newMsg {msgBody = (L8.pack "hello")} + threadDelay (1000 * 100) + + (_, n2, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + n2 `shouldBe` 1 + + Just (_msg, env) <- getMsg ch Ack q + rejectMsg (envChannel env) (envDeliveryTag env) False + threadDelay (1000 * 100) + + n3 <- deleteQueue ch q + n3 `shouldBe` 0 + + describe "rejectEnv" $ do + context "requeue = True" $ do + it "requeues a message" $ withTestConnection $ \ch -> do + let q = "haskell-amqp.basic.reject.with-requeue-true" + + (_, n1, _) <- declareQueue ch $ newQueue {queueName = q, queueDurable = False} + n1 `shouldBe` 0 + + -- publishes using default exchange + publishMsg ch "" q $ newMsg {msgBody = (L8.pack "hello")} + threadDelay (1000 * 100) + + (_, n2, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + n2 `shouldBe` 1 + + Just (_msg, env) <- getMsg ch Ack q + rejectEnv env True + threadDelay (1000 * 100) + + n3 <- deleteQueue ch q + n3 `shouldBe` 1 + + context "requeue = False" $ do + it "rejects a message" $ withTestConnection $ \ch -> do + let q = "haskell-amqp.basic.reject.with-requeue-false" + + (_, n1, _) <- declareQueue ch $ newQueue {queueName = q, queueDurable = False} + n1 `shouldBe` 0 + + -- publishes using default exchange + publishMsg ch "" q $ newMsg {msgBody = (L8.pack "hello")} + threadDelay (1000 * 100) + + (_, n2, _) <- declareQueue ch (newQueue {queueName = q, queueDurable = False}) + n2 `shouldBe` 1 + + Just (_msg, env) <- getMsg ch Ack q + rejectEnv env False + threadDelay (1000 * 100) + + n3 <- deleteQueue ch q + n3 `shouldBe` 0
+ test/ChannelSpec.hs view
@@ -0,0 +1,44 @@+{-# OPTIONS -XOverloadedStrings #-} + +module ChannelSpec (main, spec) where + +import Test.Hspec +import Network.AMQP +import Network.AMQP.Internal (channelID) + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "openChannel" $ do + context "with automatically allocated channel id" $ do + it "opens a new channel with unique id" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch1 <- openChannel conn + ch2 <- openChannel conn + ch3 <- openChannel conn + + channelID ch1 `shouldBe` 1 + channelID ch2 `shouldBe` 2 + channelID ch3 `shouldBe` 3 + + closeConnection conn + + describe "closeChannel" $ do + context "with an open channel" $ do + it "closes the channel" $ do + pending + + describe "qos" $ do + context "with prefetchCount = 5" $ do + it "sets prefetch count" $ do + -- we won't demonstrate how basic.qos works in concert + -- with acks here, it's more of a basic.consume functionality + -- aspect + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + qos ch 0 5 + + closeConnection conn
+ test/ConnectionSpec.hs view
@@ -0,0 +1,32 @@+{-# OPTIONS -XOverloadedStrings #-} + +module ConnectionSpec (main, spec) where + +import Test.Hspec +import Network.AMQP + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "openConnection" $ do + context "with default vhost and default admin credentials" $ do + it "connects successfully" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + closeConnection conn + + context "with custom vhost and valid credentials" $ do + it "connects successfully" $ do + -- see ./bin/ci/before_build.sh + conn <- openConnection "127.0.0.1" "haskell_amqp_testbed" + "haskell_amqp" + "haskell_amqp_password" + closeConnection conn + + context "with custom vhost and valid credentials" $ do + it "raises an exception" $ do + let ex = ConnectionClosedException "Handshake failed. Please check the RabbitMQ logs for more information" + (openConnection "127.0.0.1" "haskell_amqp_testbed" + "NxqrbLaNiN3TAenNu:r9Pq]XwABuRs" + "RyRxfVDyrKjC8yhJ6htCp}P>FnJxfc") `shouldThrow` (== ex)
+ test/ExchangeDeclareSpec.hs view
@@ -0,0 +1,57 @@+{-# OPTIONS -XOverloadedStrings #-} + +module ExchangeDeclareSpec (main, spec) where + +import Test.Hspec +import Network.AMQP + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "declareExchange" $ do + context "client-named, fanout, durable, non-autodelete" $ do + it "declares the exchange" $ do + let eName = "haskell-amqp.fanout.d.na" + + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + _ <- declareExchange ch (newExchange {exchangeName = eName, + exchangeType = "fanout", + exchangeDurable = True}) + + + _ <- declareExchange ch (newExchange {exchangeName = eName, + exchangePassive = True}) + + closeConnection conn + + context "client-named, topic, non-durable, non-autodelete" $ do + it "declares the exchange" $ do + let eName = "haskell-amqp.topic.nd.na" + + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + _ <- declareExchange ch (newExchange {exchangeName = eName, + exchangeType = "topic", + exchangeDurable = False}) + + + _ <- declareExchange ch (newExchange {exchangeName = eName, + exchangePassive = True}) + + closeConnection conn + + context "passive declaration when the exchange DOES NOT exist" $ do + it "throws an exception" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + let x = "haskell-amqp.exchanges.Xiz2mQozyYcFrQgGmN8r" + ex = ChannelClosedException "NOT_FOUND - no exchange 'haskell-amqp.exchanges.Xiz2mQozyYcFrQgGmN8r' in vhost '/'" + (declareExchange ch $ newExchange {exchangeName = x, exchangePassive = True}) `shouldThrow` (== ex) + + closeConnection conn
+ test/ExchangeDeleteSpec.hs view
@@ -0,0 +1,41 @@+{-# OPTIONS -XOverloadedStrings #-} + +module ExchangeDeleteSpec (main, spec) where + +import Test.Hspec +import Network.AMQP + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "deleteExchange" $ do + context "when exchange exists" $ do + it "deletes the exchange" $ do + let eName = "haskell-amqp.exchanges.to-be-deleted" + + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + _ <- declareExchange ch (newExchange {exchangeName = eName, + exchangeType = "topic", + exchangeDurable = False}) + + _ <- deleteExchange ch eName + + let ex = ChannelClosedException "NOT_FOUND - no exchange 'haskell-amqp.exchanges.to-be-deleted' in vhost '/'" + (declareExchange ch $ newExchange {exchangeName = eName, exchangePassive = True}) `shouldThrow` (== ex) + + closeConnection conn + + context "when exchange DOES NOT exist" $ do + it "throws an exception" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + let q = "haskell-amqp.exchanges.GmN8rozyXiz2mQYcFrQg" + ex = ChannelClosedException "NOT_FOUND - no exchange 'haskell-amqp.exchanges.GmN8rozyXiz2mQYcFrQg' in vhost '/'" + (declareExchange ch $ newExchange {exchangeName = q, exchangePassive = True}) `shouldThrow` (== ex) + + closeConnection conn
+ test/QueueDeclareSpec.hs view
@@ -0,0 +1,63 @@+{-# OPTIONS -XOverloadedStrings #-} + +module QueueDeclareSpec (main, spec) where + +import Test.Hspec +import Network.AMQP +import Data.Text (isPrefixOf) + + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "declareQueue" $ do + context "client named, durable, non-autodelete, non-exclusive" $ do + it "declares the queue" $ do + let qName = "haskell-amqp.client-named.d.na.ne" + + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + (_, _, _) <- declareQueue ch (newQueue {queueName = qName, + queueDurable = True, + queueExclusive = False, + queueAutoDelete = False}) + + -- ensure the queue was declared + (_, _, _) <- declareQueue ch (newQueue {queueName = qName, queuePassive = True}) + closeConnection conn + + + context "server- named, non-durable, non-autodelete, exclusive" $ do + it "declares the queue, providing access to the server-generated name" $ do + let qName = "" + + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + (q, _, _) <- declareQueue ch (newQueue {queueName = qName, + queueDurable = False, + queueExclusive = True, + queueAutoDelete = False}) + + (_, cn, mn) <- declareQueue ch (newQueue {queueName = q, queuePassive = True}) + + (isPrefixOf "amq.gen" q) `shouldBe` True + -- consumer count, undelivered message count + cn `shouldBe` 0 + mn `shouldBe` 0 + + closeConnection conn + + context "passive declaration when the queue DOES NOT exist" $ do + it "throws an exception" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + let x = "haskell-amqp.queues.mQozyYcFrQgGmN8Xiz2r" + ex = ChannelClosedException "NOT_FOUND - no queue 'haskell-amqp.queues.mQozyYcFrQgGmN8Xiz2r' in vhost '/'" + (declareQueue ch $ newQueue {queueName = x, queuePassive = True}) `shouldThrow` (== ex) + + closeConnection conn
+ test/QueueDeleteSpec.hs view
@@ -0,0 +1,40 @@+{-# OPTIONS -XOverloadedStrings #-} + +module QueueDeleteSpec (main, spec) where + +import Test.Hspec +import Network.AMQP + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "deleteQueue" $ do + context "when queue exists" $ do + it "deletes the queue" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + (q, _, _) <- declareQueue ch (newQueue {queueName = "haskell-amqp.queues.to-be-deleted", + queueExclusive = True}) + + + n <- deleteQueue ch q + n `shouldBe` 0 + + let ex = (ChannelClosedException "NOT_FOUND - no queue 'haskell-amqp.queues.to-be-deleted' in vhost '/'") + (declareQueue ch $ newQueue {queueName = q, queuePassive = True}) `shouldThrow` (== ex) + + closeConnection conn + + context "when queue DOES NOT exist" $ do + it "throws an exception" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + let q = "haskell-amqp.queues.GmN8rozyXiz2mQYcFrQg" + ex = ChannelClosedException "NOT_FOUND - no queue 'haskell-amqp.queues.GmN8rozyXiz2mQYcFrQg' in vhost '/'" + (declareQueue ch $ newQueue {queueName = q, queuePassive = True}) `shouldThrow` (== ex) + + closeConnection conn
+ test/QueuePurgeSpec.hs view
@@ -0,0 +1,50 @@+{-# OPTIONS -XOverloadedStrings #-} + +module QueuePurgeSpec (main, spec) where + +import Test.Hspec +import Network.AMQP + +import Data.ByteString.Lazy.Char8 as BL hiding (putStrLn) +import Control.Concurrent (threadDelay) + +main :: IO () +main = hspec spec + +spec :: Spec +spec = do + describe "purgeQueue" $ do + context "when queue exists" $ do + it "empties the queue" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + (q, _, _) <- declareQueue ch (newQueue {queueName = "", + queueDurable = True, + queueExclusive = False, + queueAutoDelete = False}) + + publishMsg ch "" q + newMsg {msgBody = (BL.pack "payload")} + + threadDelay (1000 * 100) + (_, n, _) <- declareQueue ch (newQueue {queueName = q, + queuePassive = True}) + n `shouldBe` 1 + _ <- purgeQueue ch q + + threadDelay (1000 * 100) + (_, n2, _) <- declareQueue ch (newQueue {queueName = q, + queuePassive = True}) + n2 `shouldBe` 0 + closeConnection conn + + context "when queue DOES NOT exist" $ do + it "empties the queue" $ do + conn <- openConnection "127.0.0.1" "/" "guest" "guest" + ch <- openChannel conn + + let ex = ChannelClosedException "NOT_FOUND - no queue 'haskell-amqp.queues.avjqmyG{CHrc66MRyzYVA+PwrMVARJ' in vhost '/'" + (purgeQueue ch "haskell-amqp.queues.avjqmyG{CHrc66MRyzYVA+PwrMVARJ") `shouldThrow` (== ex) + + closeConnection conn