packages feed

stakhanov 0.1.0.1 → 0.1.1.0

raw patch · 10 files changed

+343/−45 lines, 10 filesdep ~time

Dependency ranges changed: time

Files

CHANGELOG.md view
@@ -1,13 +1,17 @@ # Revision history for stakhanov -## 0.0.1.0 -- 2026-01-12+## 0.1.1.0 -- 2026-06-21 -* First version.+* Add FIFO Queues API +## 0.1.0.1 -- 2026-05-02++* Updated to use Hasql >= 1.10 && < 1.11.+ ## 0.1.0.0 -- 2026-02-11  * Updated to work with pgmq v1.10.0. The field last_read_at has been added to message record. -## 0.1.0.1 -- 2026-05-02+## 0.0.1.0 -- 2026-01-12 -* Updated to use Hasql >= 1.10 && < 1.11.+* First version.
src/Database/PostgreSQL/Stakhanov.hs view
@@ -40,11 +40,12 @@  , getIsUnlogged   ) where+ import           Control.Monad import           Data.Aeson.Types import           Data.Int import           Data.Maybe-import           Data.Text                                as T hiding (drop)+import qualified Data.Text                                as T hiding (drop) import           Data.Time import qualified Data.Vector                              as V import           Database.PostgreSQL.Stakhanov.Internal
+ src/Database/PostgreSQL/Stakhanov/FIFO.hs view
@@ -0,0 +1,87 @@+-- | [Full PGMQ FIFO documentation](https://pgmq.github.io/pgmq/latest/fifo-queues/#fifo-queues).++module Database.PostgreSQL.Stakhanov.FIFO+  (++  -- * Reading FIFO Messages+    readGrouped+  , readGroupedWithPoll+  , readGroupedRR+  , readGroupedRRWithPoll+  , readGroupedHead++  -- * Utils+  , createFIFOIndex+  , createFIFOIndexesAll++  ) where++import           Database.PostgreSQL.Stakhanov.FIFO.Statements+import           Database.PostgreSQL.Stakhanov.Internal+import           Database.PostgreSQL.Stakhanov.Types+import           Hasql.Connection+import           Hasql.Errors+import           Hasql.Session++-- | Read `Messages` with AWS SQS FIFO-style batch retrieval behavior.+-- Unlike `readGroupedRR` which interleaves fairly across groups,+-- this function attempts to return as many `Messages` as possible from+-- the same message group to maximize throughput for related `Messages`.+readGrouped+  :: Queue+  -> VT+  -> Qty+  -> IO (Either SessionError (Maybe Messages))+readGrouped Queue{..} v q =+  use (unHasqlConn qPGConn) (statement (qName,v,q) readGroupedMessages) >>= pureMap maybeMessages++-- | Same as `readGrouped` but with polling support for real-time processing.+readGroupedWithPoll+  :: Queue              -- ^ The queue to work with+  -> VT                 -- ^ The Visibility Timeout : the time in seconds that message(s) become invisible after reading+  -> Qty                -- ^ The number of messages to read from the queue+  -> Maybe Seconds      -- ^ The max_poll_seconds : the time in seconds to wait for new messages to reach the queue. Defaults to 5+  -> Maybe Milliseconds -- ^ The milliseconds between the internal poll operations. Defaults to 100+  -> IO (Either SessionError (Maybe Messages))+readGroupedWithPoll Queue{..} v q mmp mpi =+  use (unHasqlConn qPGConn) (statement () $ readGroupedMessagesWithPoll qName v q mmp mpi) >>= pureMap maybeMessages++-- | Read `Messages` while respecting FIFO ordering within groups.+readGroupedRR+ :: Queue+ -> VT+ -> Qty+ -> IO (Either SessionError (Maybe Messages))+readGroupedRR Queue{..} v q =+  use (unHasqlConn qPGConn) (statement (qName,v,q) readGroupedRRMessages) >>= pureMap maybeMessages++-- | Same as `readGroupedRR` but with polling support for real-time processing.+readGroupedRRWithPoll+  :: Queue              -- ^ The queue to work with+  -> VT                 -- ^ The Visibility Timeout : the time in seconds that message(s) become invisible after reading+  -> Qty                -- ^ The number of messages to read from the queue+  -> Maybe Seconds      -- ^ The max_poll_seconds : the time in seconds to wait for new messages to reach the queue. Defaults to 5+  -> Maybe Milliseconds -- ^ The milliseconds between the internal poll operations. Defaults to 100+  -> IO (Either SessionError (Maybe Messages))+readGroupedRRWithPoll Queue{..} v q mmp mpi =+  use (unHasqlConn qPGConn) (statement () $ readGroupedRRMessagesWithPoll qName v q mmp mpi) >>= pureMap maybeMessages++-- | Returns exactly one `Message` per FIFO group, up to N groups.+readGroupedHead+  :: Queue -- ^ The queue to work with+  -> VT    -- ^ Visibility timeout in seconds applied to each returned message+  -> Qty   -- ^ Maximum number of groups (and therefore messages) to return+  -> IO (Either SessionError (Maybe Messages))+readGroupedHead Queue{..} v q =+  use (unHasqlConn qPGConn) (statement (qName,v,q) readGroupedHeadMessages) >>= pureMap maybeMessages++-- | Creates a GIN index on the headers column to improve FIFO read performance.+-- Recommended when using FIFO functionality frequently.+createFIFOIndex :: Queue -> IO (Either SessionError ())+createFIFOIndex Queue{..} =+  use (unHasqlConn qPGConn) (statement qName createFIFOIndexQueue)++-- | Creates FIFO indexes on all existing `Queues`.+createFIFOIndexesAll :: Connection -> IO (Either SessionError ())+createFIFOIndexesAll c = use c (statement () createFIFOIndexesAllQueues)+
+ src/Database/PostgreSQL/Stakhanov/FIFO/Statements.hs view
@@ -0,0 +1,66 @@+module Database.PostgreSQL.Stakhanov.FIFO.Statements where++import           Data.Aeson+import           Data.Int+import qualified Data.List                              as L+import qualified Data.Text                              as T+import           Data.Time+import qualified Data.Vector                            as V+import           Database.PostgreSQL.Stakhanov.Internal+import qualified Hasql.DynamicStatements.Snippet        as S+import           Hasql.Statement+import qualified Hasql.TH                               as TH+import           Prelude                                hiding (pi)++readGroupedMessages :: Statement (T.Text,Int32,Int32) (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value))+readGroupedMessages =+  preparable sql readTupleEncoder tupleMessageDecoder+  where+    sql = "select " <> columnsMessage <> " from pgmq.read_grouped($1,$2,$3)"++readGroupedHeadMessages :: Statement (T.Text,Int32,Int32) (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value))+readGroupedHeadMessages =+  preparable sql readTupleEncoder tupleMessageDecoder+  where+    sql = "select "<> columnsMessage <> " from pgmq.read_grouped_head($1,$2,$3)"++readGroupedMessagesWithPoll+  :: T.Text+  -> Int32+  -> Int32+  -> Maybe Int32+  -> Maybe Int32+  -> Statement () (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value))+readGroupedMessagesWithPoll q vt qty mmp mpi =+  let mp = maybe 5 id mmp+      pi = maybe 100 id mpi+      snippet = "select " <> S.sql columnsMessage <> " from pgmq.read_grouped_with_poll(" <>+                mconcat (L.intersperse "," [S.param q, S.param vt, S.param qty, S.param mp, S.param pi]) <> ")"+  in S.toStatement snippet tupleMessageDecoder++readGroupedRRMessages :: Statement (T.Text,Int32,Int32) (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value))+readGroupedRRMessages =+  preparable sql readTupleEncoder tupleMessageDecoder+  where+    sql = "select " <> columnsMessage <> " from pgmq.read_grouped_rr($1,$2,$3)"++readGroupedRRMessagesWithPoll+  :: T.Text+  -> Int32+  -> Int32+  -> Maybe Int32+  -> Maybe Int32+  -> Statement () (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value))+readGroupedRRMessagesWithPoll q vt qty mmp mpi =+  let mp = maybe 5 id mmp+      pi = maybe 100 id mpi+      snippet = "select " <> S.sql columnsMessage <> " from pgmq.read_grouped_rr_with_poll(" <>+                mconcat (L.intersperse "," [S.param q, S.param vt, S.param qty, S.param mp, S.param pi]) <> ")"+  in S.toStatement snippet tupleMessageDecoder++createFIFOIndexQueue :: Statement T.Text ()+createFIFOIndexQueue = [TH.resultlessStatement|select from pgmq.create_fifo_index($1::text)|]++createFIFOIndexesAllQueues :: Statement () ()+createFIFOIndexesAllQueues = [TH.resultlessStatement|select from pgmq.create_fifo_indexes_all()|]+
src/Database/PostgreSQL/Stakhanov/Internal.hs view
@@ -1,5 +1,6 @@ module Database.PostgreSQL.Stakhanov.Internal where +import           Contravariant.Extras.Contrazip      (contrazip3) import           Data.Aeson.Types import           Data.Int import           Data.List                           (intersperse)@@ -9,6 +10,7 @@ import           Data.Vector                         as V import           Database.PostgreSQL.Stakhanov.Types import qualified Hasql.Connection                    as C+import qualified Hasql.Decoders                      as D import qualified Hasql.DynamicStatements.Snippet     as S import qualified Hasql.Encoders                      as E @@ -77,6 +79,25 @@     , message           = e6     , headers           = e7 } +readTupleEncoder :: E.Params (T.Text, Int32, Int32)+readTupleEncoder =+  contrazip3+    (E.param $ E.nonNullable E.text)+    (E.param $ E.nonNullable E.int4)+    (E.param $ E.nonNullable E.int4)++tupleMessageDecoder :: D.Result (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value))+tupleMessageDecoder =+  D.rowVector $+    (,,,,,,) <$>+      D.column (D.nonNullable D.int8) <*>+      D.column (D.nonNullable D.int4) <*>+      D.column (D.nonNullable D.timestamptz) <*>+      D.column (D.nullable D.timestamptz) <*>+      D.column (D.nonNullable D.timestamptz) <*>+      D.column (D.nonNullable D.jsonb) <*>+      D.column (D.nullable D.jsonb)+ maybeHeaders :: Maybe Value -> S.Snippet maybeHeaders (Just v) = "," <> S.encoderAndParam (E.nonNullable E.json) v <> "::jsonb" maybeHeaders Nothing  = mempty@@ -93,4 +114,7 @@ bigintArrayEncoder :: V.Vector Int64 -> S.Snippet bigintArrayEncoder v =   "ARRAY[" <> M.mconcat (intersperse (S.sql ",") $ V.toList $ S.encoderAndParam (E.nonNullable E.int8) <$> v) <> "]::bigint[]"++columnsMessage :: T.Text+columnsMessage = "msg_id,read_ct,enqueued_at,last_read_at,vt,message,headers" 
src/Database/PostgreSQL/Stakhanov/Metrics.hs view
@@ -17,9 +17,9 @@  import           Data.Int import           Data.Time-import qualified Data.Vector                              as V+import qualified Data.Vector                                      as V import           Database.PostgreSQL.Stakhanov.Internal-import           Database.PostgreSQL.Stakhanov.Statements+import           Database.PostgreSQL.Stakhanov.Metrics.Statements import           Database.PostgreSQL.Stakhanov.Types import           Hasql.Connection import           Hasql.Errors
+ src/Database/PostgreSQL/Stakhanov/Metrics/Statements.hs view
@@ -0,0 +1,46 @@+module Database.PostgreSQL.Stakhanov.Metrics.Statements where++import           Data.Int+import qualified Data.Text       as T+import           Data.Time+import qualified Data.Vector     as V+import qualified Hasql.Decoders  as D+import qualified Hasql.Encoders  as E+import           Hasql.Statement+import           Prelude         hiding (pi)++getMetrics :: Statement T.Text (Int64, Maybe Int32, Maybe Int32, Int64, UTCTime, Int64)+getMetrics =+  preparable sql encoder decoder+  where+    sql = "select " <> columnsMetrics <> " from pgmq.metrics($1)"+    encoder = E.param (E.nonNullable E.text)+    decoder =+      D.singleRow $+        (,,,,,) <$>+          D.column (D.nonNullable D.int8) <*>+          D.column (D.nullable D.int4) <*>+          D.column (D.nullable D.int4) <*>+          D.column (D.nonNullable D.int8) <*>+          D.column (D.nonNullable D.timestamptz) <*>+          D.column (D.nonNullable D.int8)++getAllMetrics :: Statement () (V.Vector (T.Text, Int64, Maybe Int32, Maybe Int32, Int64, UTCTime, Int64))+getAllMetrics =+  preparable sql E.noParams decoder+  where+    sql = "select queue_name," <> columnsMetrics <> " from pgmq.metrics_all()"+    decoder =+      D.rowVector $+        (,,,,,,) <$>+          D.column (D.nonNullable D.text) <*>+          D.column (D.nonNullable D.int8) <*>+          D.column (D.nullable D.int4) <*>+          D.column (D.nullable D.int4) <*>+          D.column (D.nonNullable D.int8) <*>+          D.column (D.nonNullable D.timestamptz) <*>+          D.column (D.nonNullable D.int8)++columnsMetrics :: T.Text+columnsMetrics = "queue_length,newest_msg_age_sec,oldest_msg_age_sec,total_messages,scrape_time,queue_visible_length"+
src/Database/PostgreSQL/Stakhanov/Statements.hs view
@@ -1,6 +1,6 @@ module Database.PostgreSQL.Stakhanov.Statements where -import           Contravariant.Extras.Contrazip         (contrazip2, contrazip3)+import           Contravariant.Extras.Contrazip         (contrazip2) import           Data.Aeson import           Data.Int import qualified Data.List                              as L@@ -86,14 +86,9 @@  readMessages :: Statement (T.Text,Int32,Int32) (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value)) readMessages =-  preparable sql encoder tupleMessageDecoder+  preparable sql readTupleEncoder tupleMessageDecoder   where     sql = "select " <> columnsMessage <> " from pgmq.read($1,$2,$3)"-    encoder =-      contrazip3-        (E.param $ E.nonNullable E.text)-        (E.param $ E.nonNullable E.int4)-        (E.param $ E.nonNullable E.int4)  popMessages :: Statement (T.Text,Int32) (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value)) popMessages =@@ -160,21 +155,6 @@   let snippet = "select * from pgmq.set_vt(" <> S.param q <> ","                 <> bigintArrayEncoder v <> "," <> S.param s <> ")"   in S.toStatement snippet tupleMessageDecoder--tupleMessageDecoder :: D.Result (V.Vector (Int64, Int32, UTCTime, Maybe UTCTime, UTCTime, Value, Maybe Value))-tupleMessageDecoder =-  D.rowVector $-    (,,,,,,) <$>-      D.column (D.nonNullable D.int8) <*>-      D.column (D.nonNullable D.int4) <*>-      D.column (D.nonNullable D.timestamptz) <*>-      D.column (D.nullable D.timestamptz) <*>-      D.column (D.nonNullable D.timestamptz) <*>-      D.column (D.nonNullable D.jsonb) <*>-      D.column (D.nullable D.jsonb)--columnsMessage :: T.Text-columnsMessage = "msg_id,read_ct,enqueued_at,last_read_at,vt,message,headers"  columnsMetrics :: T.Text columnsMetrics = "queue_length,newest_msg_age_sec,oldest_msg_age_sec,total_messages,scrape_time,queue_visible_length"
stakhanov.cabal view
@@ -1,10 +1,8 @@ cabal-version:      3.0 name:               stakhanov-version:            0.1.0.1+version:            0.1.1.0 synopsis:           A Haskell PGMQ client-description:-  A fast Haskell PGMQ client for busy workers-+description:        A fast Haskell PGMQ client for busy workers homepage:           https://github.com/MichelBoucey/stakhanov license:            BSD-3-Clause license-file:       LICENSE@@ -15,8 +13,7 @@ build-type:         Simple extra-doc-files:    CHANGELOG.md extra-source-files: ReadMe.md-tested-with:-  GHC ==9.6.7 || ==9.8.4 || ==9.10.3 || ==9.12.4+tested-with:        GHC ==9.6.7 || ==9.8.4 || ==9.10.3 || ==9.12.4  source-repository head   type:     git@@ -30,11 +27,14 @@   exposed-modules:     Database.PostgreSQL.Stakhanov     Database.PostgreSQL.Stakhanov.Connection+    Database.PostgreSQL.Stakhanov.FIFO     Database.PostgreSQL.Stakhanov.Metrics     Database.PostgreSQL.Stakhanov.Types    other-modules:     Database.PostgreSQL.Stakhanov.Internal+    Database.PostgreSQL.Stakhanov.FIFO.Statements+    Database.PostgreSQL.Stakhanov.Metrics.Statements     Database.PostgreSQL.Stakhanov.Statements    default-extensions:@@ -44,15 +44,15 @@     RecordWildCards    build-depends:-    , aeson                     >=2.2.3    && <2.4-    , base                      >=4.8      && <5-    , contravariant-extras      >=0.3.5    && <0.3.6-    , hasql                     >=1.10      && <1.11-    , hasql-dynamic-statements  >=0.5      && <0.6-    , hasql-th                  >=0.5      && <0.6-    , text                      >=1.2.3    && <2.2-    , time                      >=1.9.3    && <1.16-    , vector                    >=0.12.2   && <0.14+    , aeson                     >=2.2.3  && <2.4+    , base                      >=4.8    && <5+    , contravariant-extras      >=0.3.5  && <0.3.6+    , hasql                     >=1.10   && <1.11+    , hasql-dynamic-statements  >=0.5    && <0.6+    , hasql-th                  >=0.5    && <0.6+    , text                      >=1.2.3  && <2.2+    , time                      >=1.9.3  && <1.17+    , vector                    >=0.12.2 && <0.14    hs-source-dirs:     src   default-language:   Haskell2010
tests/hspec.hs view
@@ -4,9 +4,11 @@ import           Data.Aeson import qualified Data.Aeson.KeyMap                        as K import           Data.Int-import qualified Data.Vector                              as V+import           Data.Traversable+import qualified Data.Vector                              as V hiding (forM) import qualified Database.PostgreSQL.Stakhanov            as S import           Database.PostgreSQL.Stakhanov.Connection+import           Database.PostgreSQL.Stakhanov.FIFO import           Database.PostgreSQL.Stakhanov.Metrics import           Database.PostgreSQL.Stakhanov.Types import           Test.Hspec@@ -163,5 +165,93 @@       it "Return True" $ do         Right c <- acquireLocalPGConn         let q = S.declare "HspecTestQueue" c+        S.drop q `shouldReturn` Right True++  describe "Create a Hspec FIFO test queue" $+      it "Return the record of the created queue" $ do+        Right c <- acquireLocalPGConn+        Right q <- S.create "HspecFIFOTestQueue" c+        q `shouldBe` (q :: Queue)++  describe "Send 10 messages with FIFO group ID 'A'" $+      it "Return the IDs of messages created" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        forM [1,2,3,4,5,6,7,8,9,10]+              (\n -> S.batchSend' q (V.fromList [Object (K.fromList [("order", Number n)])]) (Just $ V.fromList [Object (K.fromList [("x-pgmq-group", String "A")])]) Nothing)+        `shouldReturn`+        [Right (V.fromList [1]),Right (V.fromList [2]), Right (V.fromList [3]),Right (V.fromList [4]),Right (V.fromList [5]),+         Right (V.fromList [6]),Right (V.fromList [7]),Right (V.fromList [8]),Right (V.fromList [9]),Right (V.fromList [10])]++  describe "Send 10 messages with FIFO group ID 'B'" $+      it "Return the IDs of messages created" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        forM [11,12,13,14,15,16,17,18,19,20]+             (\n -> S.batchSend' q (V.fromList [Object (K.fromList [("order", Number n)])]) (Just $ V.fromList [Object (K.fromList [("x-pgmq-group", String "B")])]) Nothing)+        `shouldReturn`+        [Right (V.fromList [11]),Right (V.fromList [12]), Right (V.fromList [13]),Right (V.fromList [14]),Right (V.fromList [15]),+         Right (V.fromList [16]),Right (V.fromList [17]),Right (V.fromList [18]),Right (V.fromList [19]),Right (V.fromList [20])]++  describe "Send 10 messages with FIFO group ID 'C'" $+      it "Return the IDs of messages created" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        forM [21,22,23,24,25,26,27,28,29,30]+              (\n -> S.batchSend' q (V.fromList [Object (K.fromList [("order", Number n)])]) (Just $ V.fromList [Object (K.fromList [("x-pgmq-group", String "C")])]) Nothing)+        `shouldReturn`+        [Right (V.fromList [21]),Right (V.fromList [22]), Right (V.fromList [23]),Right (V.fromList [24]),Right (V.fromList [25]),+         Right (V.fromList [26]),Right (V.fromList [27]),Right (V.fromList [28]),Right (V.fromList [29]),Right (V.fromList [30])]++  describe "Read messages with readGroupedRR" $+      it "May return messages" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        Right vm <- readGroupedRR q 10 5+        vm `shouldBe` (vm::Maybe Messages)++  describe "Read messages with readGrouped" $+      it "May return messages" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        Right vm <- readGrouped q 10 5+        vm `shouldBe` (vm::Maybe Messages)++  describe "Read messages with readGroupedWithPoll" $+      it "May return messages" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        Right vm <- readGroupedWithPoll q 10 5 Nothing Nothing+        vm `shouldBe` (vm::Maybe Messages)++  describe "Read messages with readGroupedRRWithPoll" $+      it "May return messages" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        Right vm <- readGroupedRRWithPoll q 10 5 Nothing Nothing+        vm `shouldBe` (vm::Maybe Messages)++  describe "Read messages with readGroupedHead" $+      it "May return messages" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        Right vm <- readGroupedHead q 10 5+        vm `shouldBe` (vm::Maybe Messages)++  describe "Create a FIFO index" $+      it "Return Right ()" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c+        createFIFOIndex q `shouldReturn` Right ()++  describe "Create all FIFO indexes" $+      it "Return Right ()" $ do+        Right c <- acquireLocalPGConn+        createFIFOIndexesAll c `shouldReturn` Right ()++  describe "Drop HspecFIFOTestQueue" $+      it "Return True" $ do+        Right c <- acquireLocalPGConn+        let q = S.declare "HspecFIFOTestQueue" c         S.drop q `shouldReturn` Right True