packages feed

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

{-# LANGUAGE OverloadedLabels #-}

module Test.BrowseSpec (spec) where

import Control.Lens ((^.))
import Control.Monad (forM_)
import Control.Monad.IO.Class (liftIO)
import Data.Aeson (Value (..), object, (.=))
import Data.Aeson qualified as Aeson
import Data.Aeson.KeyMap qualified as KM
import Data.ByteString.Builder (toLazyByteString)
import Data.ByteString.Lazy qualified as LBS
import Data.Generics.Labels ()
import Data.IORef (modifyIORef', newIORef, readIORef)
import Data.Map.Strict qualified as Map
import Data.Text (Text)
import Data.Vector qualified as V
import Effectful (Eff, IOE, runEff)
import Effectful.Dispatch.Dynamic (interpret_)
import Effectful.Error.Static (Error, runErrorNoCallStack)
import Kiroku.Metrics hiding (items)
import Kiroku.Metrics.Config qualified as Config
import Kiroku.Store qualified as Store
import Kiroku.Store.Effect (Store (..))
import Kiroku.Test.Postgres (withMigratedTestDatabase)
import Network.HTTP.Client qualified as HTTP
import Network.HTTP.Types
import Network.Wai qualified as Wai
import Network.Wai.Handler.Warp qualified as Warp
import Network.Wai.Internal (ResponseReceived (..))
import Test.Hspec

spec :: Spec
spec = describe "Kiroku.Metrics.Browse" $ do
    it "validates page limits before construction" $ do
        forM_ [(0, 10), (10, 1), (1, 1001)] $ \(def, cap) -> mkBrowseLimits def cap `shouldBe` Left (InvalidBrowseLimits def cap)
        fmap maxLimit (mkBrowseLimits 1 1000) `shouldBe` Right 1000
    it "rejects malformed queries before invoking the store" $ do
        calls <- newIORef (0 :: Int)
        let browser = StoreBrowser (\_ -> modifyIORef' calls (+ 1) >> pure (Left (Store.ConnectionError "secret"))) defaultBrowseLimits
        Warp.testWithApplication (pure (browseApp browser)) $ \port -> do
            forM_ ["/streams?limit=0", "/streams?limit=1001", "/streams?limit=abc", "/streams?prefix=%00", "/streams?from=%FF", "/events?from=-1", "/events?from=9223372036854775808", "/events?direction=sideways", "/streams/$all/events", "/events/not-a-uuid"] $ \path -> do
                response <- get port path
                HTTP.responseStatus response `shouldBe` status400
            response <- get port "/streams?limit=1"
            HTTP.responseStatus response `shouldBe` status503
            code (HTTP.responseBody response) `shouldBe` Just "store_unavailable"
        readIORef calls `shouldReturn` 1
    it "preserves HEAD status/headers, refuses mutations and sanitizes errors" $ do
        let app = browseApp (StoreBrowser (\_ -> pure (Left (Store.ConnectionError "postgres://secret"))) defaultBrowseLimits)
        (status, headers, responseBody) <- capture app "GET" ["events"]
        status `shouldBe` status503
        code responseBody `shouldBe` Just "store_unavailable"
        capture app "HEAD" ["events"] `shouldReturn` (status, headers, "")
        forM_ ["POST", "PUT", "DELETE", "OPTIONS"] $ \method -> do
            (actual, responseHeaders, _) <- capture app method ["events"]
            actual `shouldBe` status405
            lookup "Allow" responseHeaders `shouldBe` Just "GET, HEAD"
        (disabled, _, disabledBody) <- capture browseNotConfiguredApp "GET" ["streams"]
        disabled `shouldBe` status404
        code disabledBody `shouldBe` Just "store_browsing_not_configured"
    it "serves all browse routes, names and exclusive pages through the shared server" $ withTestStore $ \store -> do
        append store "orders-1" 3
        append store "orders-2" 1
        append store "shipments-1" 1
        metrics <- newKirokuMetrics store
        providers <- storeServerProviders defaultConfig metrics store
        withMetricsServerWithProviders defaultConfig{Config.port = 0} metrics [] providers $ \server -> do
            let fetch = get server.serverPort
            first <- fetch "/streams?category=orders&limit=1"
            field "next_cursor" (body first) `shouldBe` Just (String "orders-1")
            second <- fetch "/streams?category=orders&prefix=orders-&from=orders-1&limit=1"
            map (field "name") (items (body second)) `shouldBe` [Just (String "orders-2")]
            field "next_cursor" (body second) `shouldBe` Nothing
            summary <- fetch "/streams/orders-1"
            field "category" (body summary) `shouldBe` Just (String "orders")
            streamPage <- fetch "/streams/orders-1/events?limit=2"
            map (field "streamVersion") (items (body streamPage)) `shouldBe` [Just (Number 1), Just (Number 2)]
            field "next_cursor" (body streamPage) `shouldBe` Just (Number 2)
            continuation <- fetch "/streams/orders-1/events?from=2&limit=2"
            map (field "streamVersion") (items (body continuation)) `shouldBe` [Just (Number 3)]
            backward <- fetch "/streams/orders-1/events?direction=backward"
            map (field "streamVersion") (items (body backward)) `shouldBe` map (Just . Number) [3, 2, 1]
            cats <- fetch "/categories"
            map (field "name") (items (body cats)) `shouldBe` [Just (String "orders"), Just (String "shipments")]
            categoryEvents <- fetch "/categories/orders/events"
            length (items (body categoryEvents)) `shouldBe` 4
            globals <- fetch "/events?from=3"
            map (field "globalPosition") (items (body globals)) `shouldBe` [Just (Number 4), Just (Number 5)]
            map (field "original_stream_name") (items (body globals)) `shouldBe` [Just (String "orders-2"), Just (String "shipments-1")]
            reverseGlobal <- fetch "/events?direction=backward&limit=2"
            map (field "globalPosition") (items (body reverseGlobal)) `shouldBe` [Just (Number 5), Just (Number 4)]
            Right recorded <- Store.runStoreIO store $ Store.readAllForward (Store.GlobalPosition 0) 10
            let Store.EventId uuid = V.head recorded ^. #eventId
            event <- fetch ("/events/" <> show uuid)
            HTTP.responseStatus event `shouldBe` status200
            field "original_stream_name" (body event) `shouldBe` Just (String "orders-1")
            forM_ ["/streams/missing", "/streams/missing/events", "/events/00000000-0000-0000-0000-000000000000", "/streams/unknown/path"] $ \path -> do
                response <- fetch path
                HTTP.responseStatus response `shouldBe` status404
            -- Summary JSON and legacy endpoints remain independently mounted.
            legacy <- fetch "/nope"
            body legacy `shouldBe` object ["error" .= ("Not found" :: Text)]
    it "keeps prefix-mounted browse routes and CORS active when legacy JSON and WebSockets are disabled" $ do
        metrics <- newKirokuMetricsWith (pure (Store.GlobalPosition 0)) (pure 0)
        let origin = either (error . show) (\value -> value) (allowedOrigin "https://ops.example.com")
            cors = corsAllowOrigins [origin]
            browser = StoreBrowser (\_ -> pure (Left (Store.ConnectionError "secret"))) defaultBrowseLimits
            config = defaultConfig{cors, enableJSON = False, enableWebSocket = False}
            providers = defaultServerProviders{storeBrowsing = Just browser}
            mounted req respond = combinedAppWithProviders config metrics [] providers (req{Wai.pathInfo = drop 1 (Wai.pathInfo req)}) respond
        Warp.testWithApplication (pure mounted) $ \port -> do
            manager <- HTTP.newManager HTTP.defaultManagerSettings
            request <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> "/inspect/events?from=0")
            response <- HTTP.httpLbs request{HTTP.requestHeaders = [("Origin", "https://ops.example.com")]} manager
            HTTP.responseStatus response `shouldBe` status503
            lookup "Access-Control-Allow-Origin" (HTTP.responseHeaders response) `shouldBe` Just "https://ops.example.com"
            code (HTTP.responseBody response) `shouldBe` Just "store_unavailable"
    it "trims over-fetch before one distinct-name lookup and skips it on empty pages" $ withTestStore $ \store -> do
        append store "orders-1" 2
        append store "orders-2" 1
        Right events <- Store.runStoreIO store $ Store.readAllForward (Store.GlobalPosition 0) 10
        calls <- newIORef ([] :: [[Store.StreamId]])
        let interpreter :: forall a. Eff '[Store, Error Store.StoreError, IOE] a -> Eff '[Error Store.StoreError, IOE] a
            interpreter = interpret_ $ \case
                ReadAllForward (Store.GlobalPosition cursor) limit -> pure $ V.take (fromIntegral limit) $ V.filter (\event -> let Store.GlobalPosition n = event ^. #globalPosition in n > cursor) events
                LookupStreamNames ids -> liftIO (modifyIORef' calls (<> [ids])) >> pure Map.empty
                _ -> error "unexpected browse mock operation"
            browser = StoreBrowser (runEff . runErrorNoCallStack . interpreter) defaultBrowseLimits
        Warp.testWithApplication (pure (browseApp browser)) $ \port -> do
            response <- get port "/events?limit=2"
            map (field "original_stream_name") (items (body response)) `shouldBe` [Just Null, Just Null]
            _ <- get port "/events?from=3"
            pure ()
        readIORef calls `shouldReturn` [[V.head events ^. #originalStreamId]]

withTestStore :: (Store.KirokuStore -> IO ()) -> IO ()
withTestStore action = withMigratedTestDatabase $ \connection -> Store.withStore (Store.defaultConnectionSettings connection) action
append :: Store.KirokuStore -> Text -> Int -> IO ()
append store name count = do
    let event = Store.EventData Nothing (Store.EventType "Created") (object []) Nothing Nothing Nothing
    result <- Store.runStoreIO store $ Store.appendToStream (Store.StreamName name) Store.NoStream (replicate count event)
    result `shouldSatisfy` either (const False) (const True)
get :: Int -> String -> IO (HTTP.Response LBS.ByteString)
get port path = do
    manager <- HTTP.newManager HTTP.defaultManagerSettings
    request <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> path)
    HTTP.httpLbs request manager
body :: HTTP.Response LBS.ByteString -> Value
body response = maybe (error "Invalid JSON response") (\value -> value) (Aeson.decode (HTTP.responseBody response))
field :: Aeson.Key -> Value -> Maybe Value
field key (Object objectValue) = KM.lookup key objectValue
field _ _ = Nothing
items :: Value -> [Value]
items value = case field "items" value of Just (Array rows) -> V.toList rows; _ -> error "Expected items"
code :: LBS.ByteString -> Maybe Text
code bytes = do
    value <- Aeson.decode bytes
    err <- field "error" value
    String actual <- field "code" err
    pure actual
capture :: Wai.Application -> Method -> [Text] -> IO (Status, ResponseHeaders, LBS.ByteString)
capture app method path = do
    result <- newIORef Nothing
    _ <- app Wai.defaultRequest{Wai.requestMethod = method, Wai.pathInfo = path} $ \response -> do
        let (status, headers, stream) = Wai.responseToStream response
        stream $ \send -> do
            chunks <- newIORef []
            send (\chunk -> modifyIORef' chunks (<> [toLazyByteString chunk])) (pure ())
            bytes <- mconcat <$> readIORef chunks
            modifyIORef' result (const (Just (status, headers, bytes)))
        pure ResponseReceived
    maybe (error "No response") (\value -> pure value) =<< readIORef result