packages feed

kiroku-metrics-0.3.0.0: test/Test/StandaloneSpec.hs

{-# LANGUAGE ScopedTypeVariables #-}

module Test.StandaloneSpec (spec) where

import Control.Concurrent.Async qualified as Async
import Control.Concurrent.MVar
import Control.Exception (SomeException, bracket, throwIO, try)
import Control.Monad (forM_, void)
import Data.Aeson qualified as A
import Data.Aeson.Key qualified as K
import Data.Aeson.KeyMap qualified as KM
import Data.ByteString.Lazy qualified as LBS
import Data.ByteString.Lazy.Char8 qualified as LBSC
import Data.Either (isLeft, isRight)
import Data.IORef (newIORef, readIORef, writeIORef)
import Data.Text (Text)
import Data.Text qualified as T
import Network.HTTP.Client qualified as HTTP
import Network.HTTP.Types (RequestHeaders, status200)
import Network.Socket qualified as Socket
import Network.Wai.Handler.Warp qualified as Warp
import Network.WebSockets qualified as WS
import Options.Applicative qualified as O
import System.Directory (findExecutable)
import System.Environment (getEnvironment)
import System.Exit (ExitCode (..))
import System.IO (hGetLine)
import System.Posix.Signals (sigINT, sigTERM, signalProcess)
import System.Process qualified as Process
import System.Timeout (timeout)
import Test.Hspec

import Kiroku.Metrics
import Kiroku.Store qualified as Store
import Kiroku.Test.Postgres (withMigratedTestDatabase)

spec :: Spec
spec = do
    describe "Kiroku.Metrics.Standalone (options)" $ do
        it "parses defaults, repeated origins, and all options" $ do
            isRight (parsed []) `shouldBe` True
            opts <- requireParsed ["--database-url", "postgresql://x", "--schema", "tenant", "--pool-size", "2", "--port", "0", "--cors-origin", "http://a", "--cors-origin", "http://b", "--no-cors-allow-credentials", "--ws-max-connections", "3"]
            opts.databaseUrl `shouldBe` Just "postgresql://x"
            opts.schema `shouldBe` Just "tenant"
            opts.poolSize `shouldBe` Just 2
            opts.port `shouldBe` Just 0
            opts.corsOrigins `shouldBe` ["http://a", "http://b"]
            opts.corsAllowCredentials `shouldBe` Just False
            opts.wsMaxConnections `shouldBe` Just 3
        it "rejects signed, non-ASCII, overflowing and out-of-range CLI numbers and conflicting flags" $ do
            forM_ ["-1", "+1", "", "1", "18446744073709551616", "65536"] $ \value -> isLeft (parsed ["--port", value]) `shouldBe` True
            forM_ ["--pool-size", "--ws-max-connections"] $ \flag -> isLeft (parsed [flag, "0"]) `shouldBe` True
            isLeft (parsed ["--cors-allow-credentials", "--no-cors-allow-credentials"]) `shouldBe` True
        it "resolves environment defaults and reports a missing or empty database/schema" $ do
            opts <- requireParsed []
            isLeft (resolveInspectOptions [] opts) `shouldBe` True
            rt <- resolved [("DATABASE_URL", "postgresql://env")] opts
            rt.databaseUrl `shouldBe` "postgresql://env"
            rt.schema `shouldBe` "kiroku"
            rt.poolSize `shouldBe` (Store.defaultConnectionSettings "").poolSize
            rt.port `shouldBe` 9091
            rt.cors `shouldBe` corsDisabled
            rt.wsMaxConnections `shouldBe` 100
            forM_ [["--database-url", ""], ["--database-url", "x", "--schema", " "]] $ \args -> requireParsed args >>= \o -> isLeft (resolveInspectOptions [] o) `shouldBe` True
        it "valid flags override malformed environment values including explicit credentials False" $ do
            opts <- requireParsed ["--database-url", "flag", "--port", "0", "--pool-size", "1", "--schema", "tenant", "--ws-max-connections", "1", "--cors-origin", "http://a", "--no-cors-allow-credentials"]
            rt <- resolved [("DATABASE_URL", "env"), ("KIROKU_INSPECT_PORT", "bad"), ("KIROKU_INSPECT_POOL_SIZE", "bad"), ("KIROKU_INSPECT_SCHEMA", "bad"), ("KIROKU_INSPECT_WS_MAX_CONNECTIONS", "bad"), ("KIROKU_INSPECT_CORS_ORIGINS", "*"), ("KIROKU_INSPECT_CORS_ALLOW_CREDENTIALS", "bad")] opts
            rt.databaseUrl `shouldBe` "flag"
            rt.port `shouldBe` 0
            rt.schema `shouldBe` "tenant"
            rt.cors.allowCredentials `shouldBe` False
            rt2 <- resolved [("KIROKU_INSPECT_CORS_ALLOW_CREDENTIALS", "true")] opts
            rt2.cors.allowCredentials `shouldBe` False
        it "validates environment numbers and booleans, normalizes origins and refuses wildcard" $ do
            opts <- requireParsed ["--database-url", "x"]
            forM_ ["KIROKU_INSPECT_POOL_SIZE", "KIROKU_INSPECT_WS_MAX_CONNECTIONS", "KIROKU_INSPECT_PORT"] $ \name ->
                forM_ ["abc", "-1", "+1", "9", "18446744073709551616"] $ \value -> isLeft (resolveInspectOptions [(name, value)] opts) `shouldBe` True
            isLeft (resolveInspectOptions [("KIROKU_INSPECT_CORS_ALLOW_CREDENTIALS", "yes")] opts) `shouldBe` True
            rt <- resolved [("KIROKU_INSPECT_CORS_ORIGINS", "https://Ops.example.com/, http://localhost:5173"), ("KIROKU_INSPECT_CORS_ALLOW_CREDENTIALS", "true")] opts
            map renderAllowedOrigin rt.cors.allowedOrigins `shouldBe` ["https://ops.example.com", "http://localhost:5173"]
            rt.cors.allowCredentials `shouldBe` True
            wild <- requireParsed ["--database-url", "x", "--cors-origin", "*"]
            either (T.isInfixOf "wildcard") (const False) (resolveInspectOptions [] wild) `shouldBe` True
            empty <- resolved [("KIROKU_INSPECT_PORT", ""), ("KIROKU_INSPECT_SCHEMA", "")] opts
            empty.port `shouldBe` 9091
            empty.schema `shouldBe` "kiroku"
    describe "Kiroku.Metrics.Standalone (end to end)" $ do
        it "serves durable reads, an empty live registry, CORS, and real event tail from a database URL" $
            withMigratedTestDatabase $ \url -> do
                rt <- runtime url
                ready <- newEmptyMVar
                done <- newEmptyMVar
                let hooks = InspectHooks (\port caps -> putMVar ready (port, caps)) (takeMVar done)
                Async.withAsync (runInspect hooks rt) $ \server -> do
                    (port, caps) <- bounded (takeMVar ready)
                    caps.routes `shouldBe` RouteAvailability True True True True True True True True True
                    caps.corsIsEnabled `shouldBe` True
                    Store.withStore (Store.defaultConnectionSettings url) $ \writer -> do
                        append writer "orders-1" 3
                        response <- get port "/capabilities" []
                        A.eitherDecode (HTTP.responseBody response) `shouldBe` Right caps
                        LBSC.putStrLn ("standalone discovery capture: " <> HTTP.responseBody response)
                        streams <- get port "/streams?category=orders" []
                        HTTP.responseStatus streams `shouldBe` status200
                        (A.decode (HTTP.responseBody streams) >>= key "items") `shouldSatisfy` maybe False (\case A.Array rows -> any (\row -> key "name" row == Just (A.String "orders-1")) rows; _ -> False)
                        live <- get port "/subscriptions" []
                        A.decode (HTTP.responseBody live) `shouldBe` Just (A.toJSON ([] :: [A.Value]))
                        inventory <- get port "/subscription-checkpoints" []
                        case A.eitherDecode (HTTP.responseBody inventory) of
                            Right (CheckpointInventoryResponse pos rows) -> do
                                pos `shouldSatisfy` (>= 3)
                                rows `shouldBe` []
                            Left err -> expectationFailure err
                        forM_ ["/health/ready", "/subscriptions/missing/dead-letters", "/events", "/categories"] $ \path -> HTTP.responseStatus <$> get port path [] `shouldReturn` status200
                        forM_ [("http://localhost:5173", Just "http://localhost:5173"), ("http://evil.example", Nothing)] $ \(origin, expected) -> do
                            response2 <- get port "/metrics" [("Origin", origin)]
                            lookup "Access-Control-Allow-Origin" (HTTP.responseHeaders response2) `shouldBe` expected
                            (A.decode (HTTP.responseBody response2) >>= key "subscriptions") `shouldBe` Just (A.object [])
                        bounded $ WS.runClient "127.0.0.1" port "/ws/events" $ \conn -> do
                            WS.sendTextData conn (A.encode (A.object ["type" A..= ("subscribe_events" :: Text)]))
                            void (frame conn "event_stream_started")
                            append writer "orders-2" 1
                            ev <- frame conn "event"
                            (key "event" ev >>= key "original_stream_name") `shouldBe` Just (A.String "orders-2")
                            (key "event" ev >>= key "eventType") `shouldBe` Just (A.String "OrderCreated")
                    putMVar done ()
                    bounded (Async.wait server)
                    assertReleased port
        it "releases its listener on immediate shutdown, hook exception and cancellation" $
            withMigratedTestDatabase $ \url -> do
                rt <- runtime url
                forM_ [False, True] $ \shouldFail -> do
                    portRef <- newIORef 0
                    let hooks = InspectHooks (\port _ -> writeIORef portRef port >> if shouldFail then throwIO (userError "hook failed") else pure ()) (pure ())
                    result <- bounded (try (runInspect hooks rt) :: IO (Either SomeException ()))
                    result `shouldSatisfy` (if shouldFail then isLeft else isRight)
                    readIORef portRef >>= assertReleased
                ready <- newEmptyMVar
                never <- newEmptyMVar
                Async.withAsync (runInspect (InspectHooks (\port _ -> putMVar ready port) (takeMVar never)) rt) $ \server -> do
                    port <- bounded (takeMVar ready)
                    bounded (Async.cancel server)
                    assertReleased port
        it "does not call onListening on an occupied port" $
            withMigratedTestDatabase $ \url -> do
                rt <- runtime url
                withOccupiedPort $ \port -> do
                    called <- newIORef False
                    result <- bounded (try (runInspect (InspectHooks (\_ _ -> writeIORef called True) (pure ())) (InspectRuntime rt.databaseUrl rt.schema rt.poolSize port rt.cors rt.wsMaxConnections)) :: IO (Either SomeException ()))
                    result `shouldSatisfy` isLeft
                    readIORef called `shouldReturn` False
    describe "kiroku-inspect (executable)" $ do
        it "exits 2 for usage/resolution errors and redacts runtime connection failures" $ do
            exe <- executable
            env <- cleanEnvironment
            forM_ [[], ["--unknown"], ["--database-url", "x", "--cors-origin", "*"]] $ \args -> do
                (exit, _, _) <- bounded (Process.readCreateProcessWithExitCode (Process.proc exe args){Process.env = Just env} "")
                exit `shouldBe` ExitFailure 2
            (exit, out, err) <- bounded (Process.readCreateProcessWithExitCode (Process.proc exe ["--database-url", "postgresql://user:secret@127.0.0.1:1/db?connect_timeout=1"]){Process.env = Just env} "")
            exit `shouldBe` ExitFailure 1
            out `shouldBe` ""
            T.pack err `shouldSatisfy` (not . T.isInfixOf "secret")
        it "prints no success banner and exits 1 on bind failure" $
            withMigratedTestDatabase $ \url -> withOccupiedPort $ \port -> do
                exe <- executable
                env <- cleanEnvironment
                (exit, out, err) <- bounded (Process.readCreateProcessWithExitCode (Process.proc exe ["--database-url", T.unpack url, "--port", show port]){Process.env = Just env} "")
                exit `shouldBe` ExitFailure 1
                out `shouldBe` ""
                err `shouldSatisfy` (not . null)
        it "exits 0 after SIGINT or SIGTERM, including repeated signals" $
            withMigratedTestDatabase $ \url -> do
                exe <- executable
                env <- cleanEnvironment
                forM_ [sigINT, sigTERM] $ \signal ->
                    bracket (Process.createProcess (Process.proc exe ["--database-url", T.unpack url, "--port", "0"]){Process.env = Just env, Process.std_out = Process.CreatePipe, Process.std_err = Process.CreatePipe}) Process.cleanupProcess $ \(_, out, _, process) -> do
                        handle <- maybe (fail "No stdout") pure out
                        first <- bounded (hGetLine handle)
                        first `shouldSatisfy` (T.isInfixOf "listening on port" . T.pack)
                        void (bounded (hGetLine handle))
                        void (bounded (hGetLine handle))
                        pid <- Process.getPid process >>= maybe (fail "No PID") pure
                        signalProcess signal pid
                        signalProcess signal pid
                        shutting <- bounded (hGetLine handle)
                        shutting `shouldBe` "kiroku-inspect: shutting down"
                        bounded (Process.waitForProcess process) `shouldReturn` ExitSuccess

parsed :: [String] -> Either String InspectOptions
parsed args = maybe (Left "Parse failed") Right (O.getParseResult (O.execParserPure O.defaultPrefs inspectParserInfo args))

requireParsed :: [String] -> IO InspectOptions
requireParsed = either fail pure . parsed

resolved :: [(String, String)] -> InspectOptions -> IO InspectRuntime
resolved env opts = either (fail . T.unpack) pure (resolveInspectOptions env opts)

runtime :: Text -> IO InspectRuntime
runtime url = requireParsed ["--database-url", T.unpack url, "--port", "0", "--cors-origin", "http://localhost:5173"] >>= resolved []

key :: Text -> A.Value -> Maybe A.Value
key name (A.Object o) = KM.lookup (K.fromText name) o
key _ _ = Nothing

frame :: WS.Connection -> Text -> IO A.Value
frame conn wanted = do
    raw <- WS.receiveData conn :: IO LBS.ByteString
    case A.decode raw of
        Just value | key "type" value == Just (A.String wanted) -> pure value
        _ -> frame conn wanted

get :: Int -> String -> RequestHeaders -> IO (HTTP.Response LBS.ByteString)
get port path headers = do
    manager <- HTTP.newManager HTTP.defaultManagerSettings
    req <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> path)
    HTTP.httpLbs req{HTTP.requestHeaders = headers} manager

append :: Store.KirokuStore -> Text -> Int -> IO ()
append store name count = do
    result <- Store.runStoreIO store (Store.appendToStream (Store.StreamName name) Store.NoStream (replicate count (Store.EventData Nothing (Store.EventType "OrderCreated") A.Null Nothing Nothing Nothing)))
    result `shouldSatisfy` isRight

bounded :: IO a -> IO a
bounded action = timeout 15_000_000 action >>= maybe (fail "Timed out") pure

assertReleased :: Int -> IO ()
assertReleased port =
    bracket (Socket.socket Socket.AF_INET Socket.Stream Socket.defaultProtocol) Socket.close $ \socket -> do
        Socket.setSocketOption socket Socket.ReuseAddr 1
        Socket.bind socket (Socket.SockAddrInet (fromIntegral port) (Socket.tupleToHostAddress (0, 0, 0, 0)))
        Socket.listen socket 1

withOccupiedPort :: (Int -> IO a) -> IO a
withOccupiedPort action = do
    (port, reserved) <- Warp.openFreePort
    Socket.close reserved
    bracket (Socket.socket Socket.AF_INET Socket.Stream Socket.defaultProtocol) Socket.close $ \ipv4 ->
        bracket (Socket.socket Socket.AF_INET6 Socket.Stream Socket.defaultProtocol) Socket.close $ \ipv6 -> do
            Socket.bind ipv4 (Socket.SockAddrInet (fromIntegral port) (Socket.tupleToHostAddress (0, 0, 0, 0)))
            Socket.listen ipv4 1
            Socket.setSocketOption ipv6 Socket.IPv6Only 1
            Socket.bind ipv6 (Socket.SockAddrInet6 (fromIntegral port) 0 (0, 0, 0, 0) 0)
            Socket.listen ipv6 1
            action port

executable :: IO FilePath
executable = findExecutable "kiroku-inspect" >>= maybe (fail "Cabal build-tool-depends did not supply kiroku-inspect") pure

cleanEnvironment :: IO [(String, String)]
cleanEnvironment = filter (\(name, _) -> name /= "DATABASE_URL" && not ("KIROKU_INSPECT_" `T.isPrefixOf` T.pack name)) <$> getEnvironment