-- | Reproduces the @hasql@ pipeline parity benchmark scenarios at the @pqi@
-- level: the same sequence of queries is run both sequentially and inside a
-- pipeline, and the two ways of executing them must produce identical
-- observations.
--
-- This is a scenario test for 'Pqi.pipelineSync': it marks the sync point that
-- lets the pipelined batch complete.
module Pqi.Conformance.Operation.PipelineSync.Parity
  ( spec,
  )
where

import qualified Pqi
import qualified Pqi as Lq
import Pqi.Conformance.Harness
import Pqi.Conformance.Observation
import Pqi.Conformance.Prelude
import Pqi.Conformance.Scenario (observed, takeCommandResults, takeResult)
import Test.Hspec

spec :: Pqi.Adapter -> SpecWith ByteString
spec :: Adapter -> SpecWith ByteString
spec Adapter
adapter =
  String -> SpecWith ByteString -> SpecWith ByteString
forall a. HasCallStack => String -> SpecWith a -> SpecWith a
describe String
"parity" do
    String
-> (ByteString -> IO ()) -> SpecWith (Arg (ByteString -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"manySmallResults matches sequential execution" \ByteString
conninfo ->
      Adapter
-> ByteString
-> (Connection
    -> IO
         (Bool, [Bool], Bool,
          [(Maybe ResultObservation, Maybe ResultObservation)],
          Maybe ResultObservation, Maybe ResultObservation, Bool,
          [Maybe ResultObservation]))
-> IO ()
forall a.
(Eq a, Show a, HasCallStack) =>
Adapter -> ByteString -> (Connection -> IO a) -> IO ()
differential Adapter
adapter ByteString
conninfo \Connection
connection -> do
        let query :: a
query = a
"SELECT 1, 2"
        sequential <- Int -> IO (Maybe ResultObservation) -> IO [Maybe ResultObservation]
forall (m :: * -> *) a. Applicative m => Int -> m a -> m [a]
replicateM Int
100 (ByteString
-> [Maybe (Word32, ByteString, Format)]
-> Format
-> Connection
-> IO (Maybe ResultObservation)
observed ByteString
forall {a}. IsString a => a
query [] Format
Lq.Text Connection
connection)
        entered <- Lq.enterPipelineMode connection
        sent <- replicateM 100 (Lq.sendQueryParams connection query [] Lq.Text)
        synced <- Lq.pipelineSync connection
        pipeline <- replicateM 100 (takeCommandResults connection)
        syncResult <- takeResult connection
        trailing <- takeResult connection
        exited <- Lq.exitPipelineMode connection
        let pipelineResults = ((Maybe ResultObservation, Maybe ResultObservation)
 -> Maybe ResultObservation)
-> [(Maybe ResultObservation, Maybe ResultObservation)]
-> [Maybe ResultObservation]
forall a b. (a -> b) -> [a] -> [b]
map (Maybe ResultObservation, Maybe ResultObservation)
-> Maybe ResultObservation
forall a b. (a, b) -> a
fst [(Maybe ResultObservation, Maybe ResultObservation)]
pipeline
        sequential `shouldBe` pipelineResults
        pure (entered, sent, synced, pipeline, syncResult, trailing, exited, sequential)

    String
-> (ByteString -> IO ()) -> SpecWith (Arg (ByteString -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"manyLargeResults matches sequential execution" \ByteString
conninfo ->
      Adapter
-> ByteString
-> (Connection
    -> IO
         (Bool, [Bool], Bool,
          [(Maybe ResultObservation, Maybe ResultObservation)],
          Maybe ResultObservation, Maybe ResultObservation, Bool,
          [Maybe ResultObservation]))
-> IO ()
forall a.
(Eq a, Show a, HasCallStack) =>
Adapter -> ByteString -> (Connection -> IO a) -> IO ()
differential Adapter
adapter ByteString
conninfo \Connection
connection -> do
        let query :: a
query = a
"SELECT generate_series(0,1000) as a, generate_series(1000,2000) as b"
        sequential <- Int -> IO (Maybe ResultObservation) -> IO [Maybe ResultObservation]
forall (m :: * -> *) a. Applicative m => Int -> m a -> m [a]
replicateM Int
100 (ByteString
-> [Maybe (Word32, ByteString, Format)]
-> Format
-> Connection
-> IO (Maybe ResultObservation)
observed ByteString
forall {a}. IsString a => a
query [] Format
Lq.Text Connection
connection)
        entered <- Lq.enterPipelineMode connection
        sent <- replicateM 100 (Lq.sendQueryParams connection query [] Lq.Text)
        synced <- Lq.pipelineSync connection
        pipeline <- replicateM 100 (takeCommandResults connection)
        syncResult <- takeResult connection
        trailing <- takeResult connection
        exited <- Lq.exitPipelineMode connection
        let pipelineResults = ((Maybe ResultObservation, Maybe ResultObservation)
 -> Maybe ResultObservation)
-> [(Maybe ResultObservation, Maybe ResultObservation)]
-> [Maybe ResultObservation]
forall a b. (a -> b) -> [a] -> [b]
map (Maybe ResultObservation, Maybe ResultObservation)
-> Maybe ResultObservation
forall a b. (a, b) -> a
fst [(Maybe ResultObservation, Maybe ResultObservation)]
pipeline
        sequential `shouldBe` pipelineResults
        pure (entered, sent, synced, pipeline, syncResult, trailing, exited, sequential)