packages feed

haskell-pgmq-0.1.0.0: tests/Test/Integration/PGMQ/Simple.hs

{-# LANGUAGE QuasiQuotes #-}

module Test.Integration.PGMQ.Simple
  ( pgmqSimpleTests )
where

import Control.Exception (bracket)
-- import Data.Aeson qualified as Aeson
import Data.Maybe (isJust, isNothing, fromJust)
import Data.Text qualified as T
import Database.PostgreSQL.Simple qualified as PSQL
import Database.PostgreSQL.Simple.Newtypes qualified as PSQL (Aeson(..))
import Database.PostgreSQL.Simple.SqlQQ (sql)
import Database.PostgreSQL.Simple.Types qualified as PSQL (QualifiedIdentifier(..))
import Database.PGMQ.Simple qualified as PGMQ
import Database.PGMQ.Types qualified as PGMQ
import Test.Hspec
import Test.Integration.Utils (getPSQLEnvConnectInfo, randomQueueName)
import Test.RandomStrings (randomASCII, randomString, onlyAlphaNum)


data TestEnv =
  TestEnv {
    conn  :: PSQL.Connection
  , queue :: PGMQ.Queue
  }
    
-- NOTE These tests expect a local pgmq server runnign on port 5432.

testQueuePrefix :: PGMQ.Queue
testQueuePrefix = "test_pgmq"

setUpConn :: IO TestEnv
setUpConn = do
  connInfo <- getPSQLEnvConnectInfo
  conn <- PSQL.connect connInfo
  queue <- randomQueueName testQueuePrefix
  return $ TestEnv { conn, queue }

dropConn :: TestEnv -> IO ()
dropConn (TestEnv { conn }) = do
  PSQL.close conn

withConn :: (TestEnv -> IO ()) -> IO ()
withConn = bracket setUpConn dropConn

withPGMQ :: (TestEnv -> IO ()) -> IO ()
withPGMQ f = withConn $ \testEnv -> bracket (setUpPGMQ testEnv) (tearDownPGMQ testEnv) (\_ -> f testEnv)
  where
    setUpPGMQ (TestEnv { conn, queue }) = do
      PGMQ.initialize conn
      PGMQ.createQueue conn queue

    tearDownPGMQ (TestEnv { conn = _conn, queue = _queue }) _ = do
      -- PGMQ.dropQueue conn queue
      pure ()

      
pgmqSimpleTests :: Spec
pgmqSimpleTests = parallel $ around withPGMQ $ describe "PGMQ Simple" $ do
  -- it "can get metrics for non-existing queue" $ \(TestEnv { conn, queue }) -> do
  --   -- first of all, this should also work for non-existing queues
  --   metrics <- PGMQ.getMetrics conn queue
  --   metrics `shouldSatisfy` isNothing

  it "can get metrics for an empty queue" $ \(TestEnv { conn, queue }) -> do
    metrics <- PGMQ.getMetrics conn queue
    metrics `shouldSatisfy` isJust
    (PGMQ.queueLength <$> metrics) `shouldBe` Just 0

  it "listQueues properly returns our queue" $ \(TestEnv { conn, queue }) -> do
    queues <- PGMQ.listQueues conn
    ((\(PGMQ.QueueInfo { queueName }) -> queueName) <$> queues) `shouldContain` [queue]
    
  it "can archive messages properly and read them" $ \(TestEnv { conn, queue }) -> do
    let message = "hello" :: String
    msgId <- PGMQ.sendMessage conn queue message 0
    mMsg <- PGMQ.readMessage conn queue 0 :: IO (Maybe (PGMQ.Message String))
    mMsg `shouldSatisfy` isJust
    let (PGMQ.Message { msgId = msgId', message = message', archivedAt }) = fromJust mMsg
    message `shouldBe` message'
    msgId `shouldBe` msgId'
    archivedAt `shouldSatisfy` isNothing
    PGMQ.archiveMessage conn queue msgId
    mMsgArchive <- PGMQ.readMessageFromArchive conn queue msgId
    mMsgArchive `shouldSatisfy` isJust
    let (PGMQ.Message { msgId = msgArchiveId, message = messageArchive, archivedAt = archivedAt' }) = fromJust mMsgArchive
    msgId `shouldBe` msgArchiveId
    message `shouldBe` messageArchive
    archivedAt' `shouldSatisfy` isJust

  it "send message properly returns msg id" $ \(TestEnv { conn, queue }) -> do
    let queueId = PSQL.QualifiedIdentifier (Just "pgmq") $ T.pack ("q_" <> queue)
    let iter = [1..20] :: [Int]  -- number of steps
    mapM_ (\_i -> do
              -- Generate random strings and make sure that the
              -- message ids we get from sendMessage match our data
              message <- randomString (onlyAlphaNum randomASCII) 20
              msgId <- PGMQ.sendMessage conn queue message 0
              [PSQL.Only (PSQL.Aeson msg)] <- PSQL.query conn [sql| SELECT message FROM ? WHERE msg_id = ? |]
                                                              (queueId, msgId)
              msg `shouldBe` message
              ) iter

  it "send messages properly return msg ids" $ \(TestEnv { conn, queue }) -> do
    let queueId = PSQL.QualifiedIdentifier (Just "pgmq") $ T.pack ("q_" <> queue)
    let iter = [1..20] :: [Int]  -- number of steps
    mapM_ (\_i -> do
              -- Generate random strings and make sure that the
              -- message ids we get from sendMessage match our data
              message1 <- randomString (onlyAlphaNum randomASCII) 20
              message2 <- randomString (onlyAlphaNum randomASCII) 20
              [msgId1, msgId2] <- PGMQ.sendMessages conn queue [message1, message2] 0
              [PSQL.Only (PSQL.Aeson msg1)] <- PSQL.query conn [sql| SELECT message FROM ? WHERE msg_id = ? |]
                                                               (queueId, msgId1)
              [PSQL.Only (PSQL.Aeson msg2)] <- PSQL.query conn [sql| SELECT message FROM ? WHERE msg_id = ? |]
                                                               (queueId, msgId2)
              msg1 `shouldBe` message1
              msg2 `shouldBe` message2
              ) iter

  it "can count messages properly, including vt (https://github.com/tembo-io/pgmq/issues/301)" $ \(TestEnv { conn, queue }) -> do
    let message = "hello vt test" :: String
    msgId <- PGMQ.sendMessage conn queue message 1
    -- immediately after creating such message, we should get metrics.length = 1 and available length = 0
    mMetrics <- PGMQ.getMetrics conn queue
    let mQLen = PGMQ.queueLength <$> mMetrics
    mQLen `shouldBe` (Just 1)
    aLen <- PGMQ.queueAvailableLength conn queue
    aLen `shouldBe` 0

    -- Now reset the vt
    PGMQ.setMessageVt conn queue msgId 0
    mMetrics' <- PGMQ.getMetrics conn queue
    let mQLen' = PGMQ.queueLength <$> mMetrics'
    mQLen' `shouldBe` (Just 1)
    aLen' <- PGMQ.queueAvailableLength conn queue
    aLen' `shouldBe` 1