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 +8/−4
- src/Database/PostgreSQL/Stakhanov.hs +2/−1
- src/Database/PostgreSQL/Stakhanov/FIFO.hs +87/−0
- src/Database/PostgreSQL/Stakhanov/FIFO/Statements.hs +66/−0
- src/Database/PostgreSQL/Stakhanov/Internal.hs +24/−0
- src/Database/PostgreSQL/Stakhanov/Metrics.hs +2/−2
- src/Database/PostgreSQL/Stakhanov/Metrics/Statements.hs +46/−0
- src/Database/PostgreSQL/Stakhanov/Statements.hs +2/−22
- stakhanov.cabal +15/−15
- tests/hspec.hs +91/−1
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