module Pqi.Conformance.Operation.ExitPipelineMode
( spec,
)
where
import qualified Pqi as Lq
import Pqi.Conformance.Harness
import Pqi.Conformance.Prelude
import Pqi.Conformance.Scenario (drainResults, execScenario, float8Oid, takeCommandResults, takeResult)
import System.Timeout (timeout)
import Test.Hspec
spec :: Lq.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
"exitPipelineMode" do
String
-> (ByteString -> IO ()) -> SpecWith (Arg (ByteString -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"returns the connection to its non-pipeline status" \ByteString
conninfo ->
Adapter
-> ByteString -> (Connection -> IO (Bool, PipelineStatus)) -> IO ()
forall a.
(Eq a, Show a, HasCallStack) =>
Adapter -> ByteString -> (Connection -> IO a) -> IO ()
differential Adapter
adapter ByteString
conninfo \Connection
connection -> do
_ <- Connection -> IO Bool
Lq.enterPipelineMode Connection
connection
exited <- Lq.exitPipelineMode connection
after <- Lq.pipelineStatus connection
pure (exited, after)
String
-> (ByteString -> IO ()) -> SpecWith (Arg (ByteString -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"fails with work pending and succeeds once drained" \ByteString
conninfo ->
Adapter
-> ByteString
-> (Connection
-> IO
(Bool, Bool, Bool, Bool,
(Maybe ResultObservation, Maybe ResultObservation),
Maybe ResultObservation, Bool))
-> IO ()
forall a.
(Eq a, Show a, HasCallStack) =>
Adapter -> ByteString -> (Connection -> IO a) -> IO ()
differential Adapter
adapter ByteString
conninfo \Connection
connection -> do
entered <- Connection -> IO Bool
Lq.enterPipelineMode Connection
connection
sent <- Lq.sendQueryParams connection "select 1" [] Lq.Text
prematureExit <- Lq.exitPipelineMode connection
synced <- Lq.pipelineSync connection
results <- takeCommandResults connection
syncResult <- takeResult connection
exited <- Lq.exitPipelineMode connection
pure (entered, sent, prematureExit, synced, results, syncResult, exited)
String
-> (ByteString -> IO ()) -> SpecWith (Arg (ByteString -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"recovers after mid-pipeline cancel (mirrors cleanUpAfterInterruption)" \ByteString
conninfo ->
Adapter -> ByteString -> (Connection -> IO Bool) -> IO ()
forall a.
(Eq a, Show a, HasCallStack) =>
Adapter -> ByteString -> (Connection -> IO a) -> IO ()
differential Adapter
adapter ByteString
conninfo \Connection
connection -> do
_ <- Connection -> IO Bool
Lq.enterPipelineMode Connection
connection
_ <- Lq.sendPrepare connection "s1" "select 1" Nothing
_ <- Lq.sendQueryPrepared connection "s1" [] Lq.Text
_ <- Lq.sendPrepare connection "s2" "select pg_sleep($1)" (Just [float8Oid])
_ <- Lq.sendQueryPrepared connection "s2" [Just ("0.5", Lq.Text)] Lq.Text
_ <- Lq.pipelineSync connection
_ <- Lq.sendFlushRequest connection
_ <- takeCommandResults connection
_ <- takeCommandResults connection
_ <- takeCommandResults connection
mHandle <- Lq.getCancel connection
_ <- for mHandle Lq.cancel
_ <- drainResults connection
_ <- drainResults connection
_ <- Lq.pipelineSync connection
_ <- drainResults connection
_ <- Lq.sendFlushRequest connection
_ <- drainResults connection
exited <- Lq.exitPipelineMode connection
pure exited
String
-> (ByteString -> IO ()) -> SpecWith (Arg (ByteString -> IO ()))
forall a.
(HasCallStack, Example a) =>
String -> a -> SpecWith (Arg a)
it String
"recovers after timeout interrupts mid-pipeline read" \ByteString
conninfo ->
Adapter
-> ByteString
-> (Connection
-> IO (Bool, PipelineStatus, Maybe ResultObservation))
-> IO ()
forall a.
(Eq a, Show a, HasCallStack) =>
Adapter -> ByteString -> (Connection -> IO a) -> IO ()
differential Adapter
adapter ByteString
conninfo \Connection
connection -> do
_ <- Connection -> IO Bool
Lq.enterPipelineMode Connection
connection
_ <- Lq.sendPrepare connection "s1" "select $1::int" Nothing
_ <- Lq.sendQueryPrepared connection "s1" [Just ("42", Lq.Text)] Lq.Text
_ <- Lq.sendPrepare connection "s2" "select pg_sleep($1)" (Just [float8Oid])
_ <- Lq.sendQueryPrepared connection "s2" [Just ("0.1", Lq.Text)] Lq.Text
_ <- Lq.pipelineSync connection
_ <- timeout 50000 (drainResults connection)
_ <- drainResults connection
mHandle <- Lq.getCancel connection
_ <- for mHandle Lq.cancel
_ <- drainResults connection
pipelineStatusBefore <- Lq.pipelineStatus connection
exited <-
if pipelineStatusBefore == Lq.PipelineOn
then do
_ <- Lq.pipelineSync connection
_ <- drainResults connection
_ <- Lq.sendFlushRequest connection
_ <- drainResults connection
ok <- Lq.exitPipelineMode connection
if ok
then pure True
else do
_ <- drainResults connection
Lq.exitPipelineMode connection
else pure True
afterStatus <- Lq.pipelineStatus connection
usable <- execScenario "select 99" connection
pure (exited, afterStatus, usable)