packages feed

kiroku-metrics 0.2.0.0 → 0.3.0.0

raw patch · 24 files changed

+3724/−287 lines, 24 filesdep +case-insensitivedep +directorydep +effectful-coredep ~kiroku-clidep ~kiroku-storedep ~kiroku-test-supportnew-component:exe:kiroku-inspectPVP ok

version bump matches the API change (PVP)

Dependencies added: case-insensitive, directory, effectful-core, network, optparse-applicative, process, unix

Dependency ranges changed: kiroku-cli, kiroku-store, kiroku-test-support

API changes (from Hackage documentation)

+ Kiroku.Metrics.Browse: InvalidBrowseLimits :: !Int -> !Int -> BrowseLimitsError
+ Kiroku.Metrics.Browse: ReadBackward :: ReadDirection
+ Kiroku.Metrics.Browse: ReadForward :: ReadDirection
+ Kiroku.Metrics.Browse: StoreBrowser :: (forall a. () => Eff '[Store, Error StoreError, IOE] a -> IO (Either StoreError a)) -> !BrowseLimits -> StoreBrowser
+ Kiroku.Metrics.Browse: [limits] :: StoreBrowser -> !BrowseLimits
+ Kiroku.Metrics.Browse: [runStoreRead] :: StoreBrowser -> forall a. () => Eff '[Store, Error StoreError, IOE] a -> IO (Either StoreError a)
+ Kiroku.Metrics.Browse: browseApp :: StoreBrowser -> Application
+ Kiroku.Metrics.Browse: browseNotConfiguredApp :: Application
+ Kiroku.Metrics.Browse: data BrowseLimits
+ Kiroku.Metrics.Browse: data BrowseLimitsError
+ Kiroku.Metrics.Browse: data ReadDirection
+ Kiroku.Metrics.Browse: data StoreBrowser
+ Kiroku.Metrics.Browse: defaultBrowseLimits :: BrowseLimits
+ Kiroku.Metrics.Browse: defaultLimit :: BrowseLimits -> Int
+ Kiroku.Metrics.Browse: instance GHC.Classes.Eq Kiroku.Metrics.Browse.BrowseLimits
+ Kiroku.Metrics.Browse: instance GHC.Classes.Eq Kiroku.Metrics.Browse.BrowseLimitsError
+ Kiroku.Metrics.Browse: instance GHC.Classes.Eq Kiroku.Metrics.Browse.ReadDirection
+ Kiroku.Metrics.Browse: instance GHC.Internal.Show.Show Kiroku.Metrics.Browse.BrowseLimits
+ Kiroku.Metrics.Browse: instance GHC.Internal.Show.Show Kiroku.Metrics.Browse.BrowseLimitsError
+ Kiroku.Metrics.Browse: instance GHC.Internal.Show.Show Kiroku.Metrics.Browse.ReadDirection
+ Kiroku.Metrics.Browse: maxLimit :: BrowseLimits -> Int
+ Kiroku.Metrics.Browse: mkBrowseLimits :: Int -> Int -> Either BrowseLimitsError BrowseLimits
+ Kiroku.Metrics.Browse: storeBrowser :: KirokuStore -> StoreBrowser
+ Kiroku.Metrics.Browse: storeBrowserWith :: BrowseLimits -> KirokuStore -> StoreBrowser
+ Kiroku.Metrics.Browse: streamInfoToJSON :: StreamInfo -> Value
+ Kiroku.Metrics.Capabilities: Capabilities :: !Text -> !Text -> !RouteAvailability -> !Bool -> ![Text] -> Capabilities
+ Kiroku.Metrics.Capabilities: ProviderPresence :: !Bool -> !Bool -> !Bool -> !Bool -> !WebSocketChannels -> ProviderPresence
+ Kiroku.Metrics.Capabilities: RouteAvailability :: !Bool -> !Bool -> !Bool -> !Bool -> !Bool -> !Bool -> !Bool -> !Bool -> !Bool -> RouteAvailability
+ Kiroku.Metrics.Capabilities: WebSocketChannels :: !Bool -> !Bool -> WebSocketChannels
+ Kiroku.Metrics.Capabilities: [browse] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [corsIsEnabled] :: Capabilities -> !Bool
+ Kiroku.Metrics.Capabilities: [deadLetters] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [eventsChannel] :: WebSocketChannels -> !Bool
+ Kiroku.Metrics.Capabilities: [hasBrowser] :: ProviderPresence -> !Bool
+ Kiroku.Metrics.Capabilities: [hasCheckpointInventory] :: ProviderPresence -> !Bool
+ Kiroku.Metrics.Capabilities: [hasDeadLetters] :: ProviderPresence -> !Bool
+ Kiroku.Metrics.Capabilities: [hasSubscriptionStatus] :: ProviderPresence -> !Bool
+ Kiroku.Metrics.Capabilities: [health] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [metricsChannel] :: WebSocketChannels -> !Bool
+ Kiroku.Metrics.Capabilities: [metrics] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [package] :: Capabilities -> !Text
+ Kiroku.Metrics.Capabilities: [presentWebSocketChannels] :: ProviderPresence -> !WebSocketChannels
+ Kiroku.Metrics.Capabilities: [processLocal] :: Capabilities -> ![Text]
+ Kiroku.Metrics.Capabilities: [prometheus] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [routes] :: Capabilities -> !RouteAvailability
+ Kiroku.Metrics.Capabilities: [subscriptionsCheckpoints] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [subscriptionsLive] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [version] :: Capabilities -> !Text
+ Kiroku.Metrics.Capabilities: [websocketEvents] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: [websocketMetrics] :: RouteAvailability -> !Bool
+ Kiroku.Metrics.Capabilities: capabilitiesApp :: Capabilities -> Application
+ Kiroku.Metrics.Capabilities: capabilitiesFor :: MetricsServerConfig -> ProviderPresence -> Capabilities
+ Kiroku.Metrics.Capabilities: capabilitiesPath :: [Text]
+ Kiroku.Metrics.Capabilities: data Capabilities
+ Kiroku.Metrics.Capabilities: data ProviderPresence
+ Kiroku.Metrics.Capabilities: data RouteAvailability
+ Kiroku.Metrics.Capabilities: data WebSocketChannels
+ Kiroku.Metrics.Capabilities: instance Data.Aeson.Types.FromJSON.FromJSON Kiroku.Metrics.Capabilities.Capabilities
+ Kiroku.Metrics.Capabilities: instance Data.Aeson.Types.FromJSON.FromJSON Kiroku.Metrics.Capabilities.RouteAvailability
+ Kiroku.Metrics.Capabilities: instance Data.Aeson.Types.ToJSON.ToJSON Kiroku.Metrics.Capabilities.Capabilities
+ Kiroku.Metrics.Capabilities: instance Data.Aeson.Types.ToJSON.ToJSON Kiroku.Metrics.Capabilities.RouteAvailability
+ Kiroku.Metrics.Capabilities: instance GHC.Classes.Eq Kiroku.Metrics.Capabilities.Capabilities
+ Kiroku.Metrics.Capabilities: instance GHC.Classes.Eq Kiroku.Metrics.Capabilities.ProviderPresence
+ Kiroku.Metrics.Capabilities: instance GHC.Classes.Eq Kiroku.Metrics.Capabilities.RouteAvailability
+ Kiroku.Metrics.Capabilities: instance GHC.Classes.Eq Kiroku.Metrics.Capabilities.WebSocketChannels
+ Kiroku.Metrics.Capabilities: instance GHC.Internal.Show.Show Kiroku.Metrics.Capabilities.Capabilities
+ Kiroku.Metrics.Capabilities: instance GHC.Internal.Show.Show Kiroku.Metrics.Capabilities.ProviderPresence
+ Kiroku.Metrics.Capabilities: instance GHC.Internal.Show.Show Kiroku.Metrics.Capabilities.RouteAvailability
+ Kiroku.Metrics.Capabilities: instance GHC.Internal.Show.Show Kiroku.Metrics.Capabilities.WebSocketChannels
+ Kiroku.Metrics.Capabilities: kirokuMetricsVersion :: Text
+ Kiroku.Metrics.Capabilities: noWebSocketChannels :: WebSocketChannels
+ Kiroku.Metrics.Capabilities: processLocalRoutes :: [Text]
+ Kiroku.Metrics.Capabilities: storeWebSocketChannels :: WebSocketChannels
+ Kiroku.Metrics.Checkpoints: CheckpointInventoryResponse :: !Int64 -> ![CheckpointRow] -> CheckpointInventoryResponse
+ Kiroku.Metrics.Checkpoints: CheckpointRow :: !Text -> !Int32 -> !Int64 -> !UTCTime -> CheckpointRow
+ Kiroku.Metrics.Checkpoints: [checkpointPosition] :: CheckpointRow -> !Int64
+ Kiroku.Metrics.Checkpoints: [checkpoints] :: CheckpointInventoryResponse -> ![CheckpointRow]
+ Kiroku.Metrics.Checkpoints: [member] :: CheckpointRow -> !Int32
+ Kiroku.Metrics.Checkpoints: [storePosition] :: CheckpointInventoryResponse -> !Int64
+ Kiroku.Metrics.Checkpoints: [subscription] :: CheckpointRow -> !Text
+ Kiroku.Metrics.Checkpoints: [updatedAt] :: CheckpointRow -> !UTCTime
+ Kiroku.Metrics.Checkpoints: checkpointInventoryResponse :: SubscriptionCheckpointInventory -> CheckpointInventoryResponse
+ Kiroku.Metrics.Checkpoints: checkpointsApp :: CheckpointInventoryProvider -> Application
+ Kiroku.Metrics.Checkpoints: checkpointsNotConfiguredApp :: Application
+ Kiroku.Metrics.Checkpoints: checkpointsPath :: [Text]
+ Kiroku.Metrics.Checkpoints: data CheckpointInventoryResponse
+ Kiroku.Metrics.Checkpoints: data CheckpointRow
+ Kiroku.Metrics.Checkpoints: instance Data.Aeson.Types.FromJSON.FromJSON Kiroku.Metrics.Checkpoints.CheckpointInventoryResponse
+ Kiroku.Metrics.Checkpoints: instance Data.Aeson.Types.FromJSON.FromJSON Kiroku.Metrics.Checkpoints.CheckpointRow
+ Kiroku.Metrics.Checkpoints: instance Data.Aeson.Types.ToJSON.ToJSON Kiroku.Metrics.Checkpoints.CheckpointInventoryResponse
+ Kiroku.Metrics.Checkpoints: instance Data.Aeson.Types.ToJSON.ToJSON Kiroku.Metrics.Checkpoints.CheckpointRow
+ Kiroku.Metrics.Checkpoints: instance GHC.Classes.Eq Kiroku.Metrics.Checkpoints.CheckpointInventoryResponse
+ Kiroku.Metrics.Checkpoints: instance GHC.Classes.Eq Kiroku.Metrics.Checkpoints.CheckpointRow
+ Kiroku.Metrics.Checkpoints: instance GHC.Internal.Show.Show Kiroku.Metrics.Checkpoints.CheckpointInventoryResponse
+ Kiroku.Metrics.Checkpoints: instance GHC.Internal.Show.Show Kiroku.Metrics.Checkpoints.CheckpointRow
+ Kiroku.Metrics.Checkpoints: storeCheckpointInventory :: KirokuStore -> CheckpointInventoryProvider
+ Kiroku.Metrics.Checkpoints: type CheckpointInventoryProvider = IO Either StoreError SubscriptionCheckpointInventory
+ Kiroku.Metrics.Config: [cors] :: MetricsServerConfig -> !CorsPolicy
+ Kiroku.Metrics.Cors: CorsPolicy :: ![AllowedOrigin] -> !Bool -> !Maybe Int -> CorsPolicy
+ Kiroku.Metrics.Cors: EmptyHost :: Text -> OriginError
+ Kiroku.Metrics.Cors: HasPathQueryOrFragment :: Text -> OriginError
+ Kiroku.Metrics.Cors: MissingScheme :: Text -> OriginError
+ Kiroku.Metrics.Cors: NotAnOrigin :: Text -> OriginError
+ Kiroku.Metrics.Cors: OpaqueOrigin :: OriginError
+ Kiroku.Metrics.Cors: WildcardOrigin :: OriginError
+ Kiroku.Metrics.Cors: [allowCredentials] :: CorsPolicy -> !Bool
+ Kiroku.Metrics.Cors: [allowedOrigins] :: CorsPolicy -> ![AllowedOrigin]
+ Kiroku.Metrics.Cors: [maxAgeSeconds] :: CorsPolicy -> !Maybe Int
+ Kiroku.Metrics.Cors: allowedOrigin :: Text -> Either OriginError AllowedOrigin
+ Kiroku.Metrics.Cors: corsAllowOrigins :: [AllowedOrigin] -> CorsPolicy
+ Kiroku.Metrics.Cors: corsAllowedMethods :: ByteString
+ Kiroku.Metrics.Cors: corsDisabled :: CorsPolicy
+ Kiroku.Metrics.Cors: corsEnabled :: CorsPolicy -> Bool
+ Kiroku.Metrics.Cors: corsMiddleware :: CorsPolicy -> Middleware
+ Kiroku.Metrics.Cors: data AllowedOrigin
+ Kiroku.Metrics.Cors: data CorsPolicy
+ Kiroku.Metrics.Cors: data OriginError
+ Kiroku.Metrics.Cors: instance GHC.Classes.Eq Kiroku.Metrics.Cors.AllowedOrigin
+ Kiroku.Metrics.Cors: instance GHC.Classes.Eq Kiroku.Metrics.Cors.CorsPolicy
+ Kiroku.Metrics.Cors: instance GHC.Classes.Eq Kiroku.Metrics.Cors.OriginError
+ Kiroku.Metrics.Cors: instance GHC.Classes.Ord Kiroku.Metrics.Cors.AllowedOrigin
+ Kiroku.Metrics.Cors: instance GHC.Internal.Show.Show Kiroku.Metrics.Cors.AllowedOrigin
+ Kiroku.Metrics.Cors: instance GHC.Internal.Show.Show Kiroku.Metrics.Cors.CorsPolicy
+ Kiroku.Metrics.Cors: instance GHC.Internal.Show.Show Kiroku.Metrics.Cors.OriginError
+ Kiroku.Metrics.Cors: isPreflight :: Request -> Bool
+ Kiroku.Metrics.Cors: originAllowed :: CorsPolicy -> ByteString -> Bool
+ Kiroku.Metrics.Cors: renderAllowedOrigin :: AllowedOrigin -> Text
+ Kiroku.Metrics.DeadLetters: DeadLetterItem :: !Int64 -> !Text -> !Int32 -> !Int64 -> !UUID -> !Value -> !Text -> !Int32 -> !UTCTime -> DeadLetterItem
+ Kiroku.Metrics.DeadLetters: DeadLetterPageResponse :: ![DeadLetterItem] -> !Maybe Text -> DeadLetterPageResponse
+ Kiroku.Metrics.DeadLetters: DeadLetterRequest :: !Maybe Int32 -> !Maybe SubscriptionDeadLetterCursor -> !SubscriptionDeadLetterLimit -> DeadLetterRequest
+ Kiroku.Metrics.DeadLetters: [attemptCount] :: DeadLetterItem -> !Int32
+ Kiroku.Metrics.DeadLetters: [createdAt] :: DeadLetterItem -> !UTCTime
+ Kiroku.Metrics.DeadLetters: [deadLetterId] :: DeadLetterItem -> !Int64
+ Kiroku.Metrics.DeadLetters: [eventId] :: DeadLetterItem -> !UUID
+ Kiroku.Metrics.DeadLetters: [globalPosition] :: DeadLetterItem -> !Int64
+ Kiroku.Metrics.DeadLetters: [items] :: DeadLetterPageResponse -> ![DeadLetterItem]
+ Kiroku.Metrics.DeadLetters: [member] :: DeadLetterItem -> !Int32
+ Kiroku.Metrics.DeadLetters: [nextCursor] :: DeadLetterPageResponse -> !Maybe Text
+ Kiroku.Metrics.DeadLetters: [reasonSummary] :: DeadLetterItem -> !Text
+ Kiroku.Metrics.DeadLetters: [reason] :: DeadLetterItem -> !Value
+ Kiroku.Metrics.DeadLetters: [requestAfter] :: DeadLetterRequest -> !Maybe SubscriptionDeadLetterCursor
+ Kiroku.Metrics.DeadLetters: [requestLimit] :: DeadLetterRequest -> !SubscriptionDeadLetterLimit
+ Kiroku.Metrics.DeadLetters: [requestMember] :: DeadLetterRequest -> !Maybe Int32
+ Kiroku.Metrics.DeadLetters: [subscription] :: DeadLetterItem -> !Text
+ Kiroku.Metrics.DeadLetters: data DeadLetterItem
+ Kiroku.Metrics.DeadLetters: data DeadLetterPageResponse
+ Kiroku.Metrics.DeadLetters: data DeadLetterRequest
+ Kiroku.Metrics.DeadLetters: deadLetterPageResponse :: SubscriptionDeadLetterPage -> DeadLetterPageResponse
+ Kiroku.Metrics.DeadLetters: deadLettersApp :: DeadLetterProvider -> Application
+ Kiroku.Metrics.DeadLetters: deadLettersNotConfiguredApp :: Application
+ Kiroku.Metrics.DeadLetters: instance Data.Aeson.Types.FromJSON.FromJSON Kiroku.Metrics.DeadLetters.DeadLetterItem
+ Kiroku.Metrics.DeadLetters: instance Data.Aeson.Types.FromJSON.FromJSON Kiroku.Metrics.DeadLetters.DeadLetterPageResponse
+ Kiroku.Metrics.DeadLetters: instance Data.Aeson.Types.ToJSON.ToJSON Kiroku.Metrics.DeadLetters.DeadLetterItem
+ Kiroku.Metrics.DeadLetters: instance Data.Aeson.Types.ToJSON.ToJSON Kiroku.Metrics.DeadLetters.DeadLetterPageResponse
+ Kiroku.Metrics.DeadLetters: instance GHC.Classes.Eq Kiroku.Metrics.DeadLetters.DeadLetterItem
+ Kiroku.Metrics.DeadLetters: instance GHC.Classes.Eq Kiroku.Metrics.DeadLetters.DeadLetterPageResponse
+ Kiroku.Metrics.DeadLetters: instance GHC.Classes.Eq Kiroku.Metrics.DeadLetters.DeadLetterRequest
+ Kiroku.Metrics.DeadLetters: instance GHC.Internal.Show.Show Kiroku.Metrics.DeadLetters.DeadLetterItem
+ Kiroku.Metrics.DeadLetters: instance GHC.Internal.Show.Show Kiroku.Metrics.DeadLetters.DeadLetterPageResponse
+ Kiroku.Metrics.DeadLetters: instance GHC.Internal.Show.Show Kiroku.Metrics.DeadLetters.DeadLetterRequest
+ Kiroku.Metrics.DeadLetters: parseDeadLetterCursor :: Text -> Maybe SubscriptionDeadLetterCursor
+ Kiroku.Metrics.DeadLetters: parseDeadLetterRequest :: Query -> Either Failure DeadLetterRequest
+ Kiroku.Metrics.DeadLetters: renderDeadLetterCursor :: SubscriptionDeadLetterCursor -> Text
+ Kiroku.Metrics.DeadLetters: storeDeadLetters :: KirokuStore -> DeadLetterProvider
+ Kiroku.Metrics.DeadLetters: type DeadLetterProvider = SubscriptionDeadLetterQuery -> IO Either StoreError SubscriptionDeadLetterPage
+ Kiroku.Metrics.JSON: errorEnvelope :: Text -> Text -> Maybe Value -> Value
+ Kiroku.Metrics.JSON: errorResponse :: Status -> Text -> Text -> Maybe Value -> Response
+ Kiroku.Metrics.JSON: storeErrorResponse :: Text -> StoreError -> Response
+ Kiroku.Metrics.Server: ServerProviders :: !ServerApp -> !Maybe SubscriptionStatusProvider -> !Maybe CheckpointInventoryProvider -> !Maybe StoreBrowser -> !Maybe DeadLetterProvider -> !WebSocketChannels -> ServerProviders
+ Kiroku.Metrics.Server: [checkpointInventory] :: ServerProviders -> !Maybe CheckpointInventoryProvider
+ Kiroku.Metrics.Server: [deadLetters] :: ServerProviders -> !Maybe DeadLetterProvider
+ Kiroku.Metrics.Server: [storeBrowsing] :: ServerProviders -> !Maybe StoreBrowser
+ Kiroku.Metrics.Server: [subscriptionStatus] :: ServerProviders -> !Maybe SubscriptionStatusProvider
+ Kiroku.Metrics.Server: [webSocketChannels] :: ServerProviders -> !WebSocketChannels
+ Kiroku.Metrics.Server: [webSocketServer] :: ServerProviders -> !ServerApp
+ Kiroku.Metrics.Server: combinedAppWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> Application
+ Kiroku.Metrics.Server: data ServerProviders
+ Kiroku.Metrics.Server: defaultServerProviders :: ServerProviders
+ Kiroku.Metrics.Server: httpAppWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> Application
+ Kiroku.Metrics.Server: providerPresence :: ServerProviders -> ProviderPresence
+ Kiroku.Metrics.Server: startMetricsServerWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> IO MetricsServer
+ Kiroku.Metrics.Server: storeServerProviders :: MetricsServerConfig -> KirokuMetrics -> KirokuStore -> IO ServerProviders
+ Kiroku.Metrics.Server: withMetricsServerWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> (MetricsServer -> IO a) -> IO a
+ Kiroku.Metrics.Standalone: InspectHooks :: !Int -> Capabilities -> IO () -> !IO () -> InspectHooks
+ Kiroku.Metrics.Standalone: InspectOptions :: !Maybe Text -> !Maybe Text -> !Maybe Int -> !Maybe Int -> ![Text] -> !Maybe Bool -> !Maybe Int -> InspectOptions
+ Kiroku.Metrics.Standalone: InspectRuntime :: !Text -> !Text -> !Int -> !Int -> !CorsPolicy -> !Int -> InspectRuntime
+ Kiroku.Metrics.Standalone: [corsAllowCredentials] :: InspectOptions -> !Maybe Bool
+ Kiroku.Metrics.Standalone: [corsOrigins] :: InspectOptions -> ![Text]
+ Kiroku.Metrics.Standalone: [cors] :: InspectRuntime -> !CorsPolicy
+ Kiroku.Metrics.Standalone: [databaseUrl] :: InspectRuntime -> !Text
+ Kiroku.Metrics.Standalone: [onListening] :: InspectHooks -> !Int -> Capabilities -> IO ()
+ Kiroku.Metrics.Standalone: [poolSize] :: InspectRuntime -> !Int
+ Kiroku.Metrics.Standalone: [port] :: InspectRuntime -> !Int
+ Kiroku.Metrics.Standalone: [schema] :: InspectRuntime -> !Text
+ Kiroku.Metrics.Standalone: [waitForShutdown] :: InspectHooks -> !IO ()
+ Kiroku.Metrics.Standalone: [wsMaxConnections] :: InspectRuntime -> !Int
+ Kiroku.Metrics.Standalone: data InspectHooks
+ Kiroku.Metrics.Standalone: data InspectOptions
+ Kiroku.Metrics.Standalone: data InspectRuntime
+ Kiroku.Metrics.Standalone: inspectOptionsParser :: Parser InspectOptions
+ Kiroku.Metrics.Standalone: inspectParserInfo :: ParserInfo InspectOptions
+ Kiroku.Metrics.Standalone: instance GHC.Classes.Eq Kiroku.Metrics.Standalone.InspectOptions
+ Kiroku.Metrics.Standalone: instance GHC.Classes.Eq Kiroku.Metrics.Standalone.InspectRuntime
+ Kiroku.Metrics.Standalone: renderStartupBanner :: InspectRuntime -> Int -> Capabilities -> [Text]
+ Kiroku.Metrics.Standalone: resolveInspectOptions :: [(String, String)] -> InspectOptions -> Either Text InspectRuntime
+ Kiroku.Metrics.Standalone: runInspect :: InspectHooks -> InspectRuntime -> IO ()
+ Kiroku.Metrics.WebSocket: CodedError :: !Text -> !Text -> ServerMessage
+ Kiroku.Metrics.WebSocket: UnsubscribeMetrics :: ClientMessage
+ Kiroku.Metrics.WebSocket: broadcastEventsWith :: (ServerMessage -> IO ()) -> StreamNameCache -> ([StreamId] -> IO (Map StreamId StreamName)) -> PublisherSubscription -> (RecordedEvent -> Bool) -> IO ()
+ Kiroku.Metrics.WebSocket: data StreamNameCache
+ Kiroku.Metrics.WebSocket: errorCodeCategoryReadFailed :: Text
+ Kiroku.Metrics.WebSocket: errorCodeEventStreamOverflowed :: Text
+ Kiroku.Metrics.WebSocket: errorCodeLiveDecodeFailed :: Text
+ Kiroku.Metrics.WebSocket: errorCodeReplayFailed :: Text
+ Kiroku.Metrics.WebSocket: newStreamNameCache :: IO StreamNameCache
+ Kiroku.Metrics.WebSocket: overflowNotice :: Word64 -> Word64 -> Maybe ServerMessage
+ Kiroku.Metrics.WebSocket: recordedEventToJSONResolved :: Map StreamId StreamName -> RecordedEvent -> Value
+ Kiroku.Metrics.WebSocket: resolveEventNames :: StreamNameCache -> ([StreamId] -> IO (Map StreamId StreamName)) -> Vector RecordedEvent -> IO (Map StreamId StreamName)
+ Kiroku.Metrics.WebSocket: streamNameCacheSize :: StreamNameCache -> IO (Int, Int)
+ Kiroku.Metrics.WebSocket: withWorkerSlot :: ((IO () -> IO ()) -> IO () -> IO a) -> IO a
- Kiroku.Metrics.Config: MetricsServerConfig :: !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Natural -> !Int64 -> !Int -> MetricsServerConfig
+ Kiroku.Metrics.Config: MetricsServerConfig :: !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Natural -> !Int64 -> !Int -> !CorsPolicy -> MetricsServerConfig

Files

CHANGELOG.md view
@@ -1,5 +1,60 @@ # Revision history for kiroku-metrics +## 0.3.0.0 — 2026-10-11++### Breaking Changes++- New umbrella exports can make record labels ambiguous; qualify configuration record updates.++- `ClientMessage` gains `UnsubscribeMetrics` and `ServerMessage` gains `CodedError`; exhaustive Haskell matches must handle the additions. Published wire growth is additive.++* `MetricsServerConfig` gains `cors`, defaulting to `corsDisabled`. Use+  `defaultConfig` record updates; complete or positional construction must supply it.++### New Features++- `GET`/`HEAD /capabilities` reports actual wiring, declared WebSocket channels and process-local scope without database reads. `Kiroku.Metrics.Capabilities` exports codecs and the generated `kirokuMetricsVersion`.+- `kiroku-inspect` and `Kiroku.Metrics.Standalone` serve the store-backed inspection API from a database URL, with validated CLI/environment options, explicit CORS and joined SIGINT/SIGTERM shutdown. It runs no subscriptions.+- Add `optparse-applicative` as a library dependency for the reusable standalone parser.++- `unsubscribe_metrics` stops periodic snapshots; `subscribe_metrics` requests a fresh snapshot and resumes periodic delivery without duplicate workers.+- Tail errors carry stable `code` values: `replay_failed`, `category_read_failed`, `live_decode_failed`, `event_stream_overflowed`, with sanitized messages.+- Live, replay and category event frames add `original_stream_name` through one batched lookup and a bounded 4096-name FIFO cache per tail.+- Drop-oldest loss now sends an overflow notice before surviving events. Earlier versions documented that notice but never emitted it.++- Read-only `GET`/`HEAD /subscriptions/<name>/dead-letters`, structured reasons and errors, member filtering and opaque cursor pages. `Kiroku.Metrics.DeadLetters` exports the codec and provider; store-backed starters configure the new `deadLetters` provider automatically.++* Add bounded stream/category/event browsing with category-plus-literal-prefix filters, exclusive cursors, validated page limits and GET/HEAD support. Store-backed servers configure the provider automatically.+* Add `recordedEventToJSONResolved`, preserving existing event keys and adding `original_stream_name`.++* `GET /subscription-checkpoints` serves exact durable member checkpoints and+  the same-snapshot store position through `Kiroku.Metrics.Checkpoints`.+* `ServerProviders`, `defaultServerProviders`, `storeServerProviders` and four+  `...WithProviders` functions compose inspection sources without changing+  legacy starter signatures. Store-backed starters include durable inventory.+++* `Kiroku.Metrics.Cors` provides validated `AllowedOrigin`, `CorsPolicy` and+  `corsMiddleware`: explicit default-off browser access with cache-correct+  HTTP/preflight handling and WebSocket origin refusal before upgrade.+* Shared `errorEnvelope`, `errorResponse` and sanitized `storeErrorResponse`+  helpers for new inspection routes. CORS refusals use `origin_not_allowed`,+  `cors_method_not_allowed` and `invalid_cors_request` codes.++### Other Changes++- Require `kiroku-store ^>=0.11.0.0` and `kiroku-cli ^>=0.2.0.10`. Construct the new `ServerProviders` from `defaultServerProviders`; full construction supplies all six fields.++- Cache-miss detection inspects only current-batch stream IDs, avoiding a walk of every retained name on each small tail batch.++- Document discovery, standalone hosting and the complete client workflow in `docs/guides/building-an-inspection-ui.md`.++* Server acquisition waits for Warp readiness and propagates bind failures;+  bracketed lifetimes supervise server termination and release sockets.+* Prefix-mounted WebSocket dispatch uses escaped mount-relative paths and+  honors `enableWebSocket` before upgrades.+* Add a direct `network` dependency for explicit ephemeral-socket cleanup.+ ## 0.2.0.0 — 2026-10-10  ### Breaking Changes
+ app-inspect/Main.hs view
@@ -0,0 +1,38 @@+module Main (main) where++import Control.Concurrent.MVar (newEmptyMVar, takeMVar, tryPutMVar)+import Control.Exception (SomeAsyncException, SomeException, fromException, tryJust)+import Control.Monad (forM_, void)+import Data.Text.IO qualified as TIO+import Options.Applicative (execParser)+import System.Environment (getEnvironment)+import System.Exit (ExitCode (..), exitFailure, exitWith)+import System.IO (hFlush, stderr, stdout)+import System.Posix.Signals (Handler (..), installHandler, sigINT, sigTERM)++import Kiroku.Metrics.Standalone++main :: IO ()+main = do+    opts <- execParser inspectParserInfo+    env <- getEnvironment+    case resolveInspectOptions env opts of+        Left err -> TIO.hPutStrLn stderr err >> exitWith (ExitFailure 2)+        Right rt -> do+            done <- newEmptyMVar+            forM_ [sigINT, sigTERM] $ \signal ->+                void (installHandler signal (Catch (void (tryPutMVar done ()))) Nothing)+            let hooks =+                    InspectHooks+                        { onListening = \port caps -> mapM_ TIO.putStrLn (renderStartupBanner rt port caps) >> hFlush stdout+                        , waitForShutdown = takeMVar done >> TIO.putStrLn "kiroku-inspect: shutting down"+                        }+            result <- tryJust synchronousException (runInspect hooks rt)+            case result of+                Left _ -> TIO.hPutStrLn stderr "kiroku-inspect: startup or server failure (details redacted)" >> exitFailure+                Right () -> pure ()++synchronousException :: SomeException -> Maybe SomeException+synchronousException err = case fromException err :: Maybe SomeAsyncException of+    Just _ -> Nothing+    Nothing -> Just err
+ bench/WebSocketTail.hs view
@@ -0,0 +1,106 @@+{-# LANGUAGE ImportQualifiedPost #-}+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-}++module Main (main) where++import Control.Concurrent.Async (concurrently, withAsync)+import Control.Concurrent.STM+import Control.Exception (evaluate)+import Control.Monad (forM, replicateM_)+import Data.Aeson (encode, object)+import Data.ByteString.Lazy qualified as LBS+import Data.IORef+import Data.List (sort)+import Data.Maybe (fromJust)+import Data.Text qualified as T+import Data.Time.Clock+import Data.Vector qualified as V+import Kiroku.Metrics.WebSocket+import Kiroku.Store hiding (cancel, id)+import Kiroku.Store.Subscription.EventPublisher qualified as Pub+import Kiroku.Test.Postgres (withMigratedTestDatabase)+import System.Timeout (timeout)++main :: IO ()+main = withMigratedTestDatabase $ \conn -> withStore (defaultConnectionSettings conn) $ \store -> do+    -- Setup and publication complete before the fixed comparison.+    mapM_ (\n -> append store (StreamName ("names-" <> T.pack (show n))) 1) [1 .. 500 :: Int]+    append store (StreamName "appender") 1+    bounded (atomically (Pub.publisherPosition store.publisher >>= check . (>= GlobalPosition 501)))+    Right events <- runStoreIO store (readAllForward (GlobalPosition 0) 500)+    putStrLn "Focused local tail delivery: 500 distinct names, 50 batches/trial, 5 paired rounds."+    putStrLn "Each case concurrently appends 50 batches of 10 events to the same existing stream."+    putStrLn "Control uses frozen encoder, warm caches pre-resolve, cold starts one cache per trial. No remote acceptance inference."+    samples <- forM [0 .. 4 :: Int] $ \roundNumber -> do+        let modes = if even roundNumber then ["control", "warm", "cold"] else ["cold", "warm", "control"]+        measured <- traverse (runCase store events) modes+        pure [(mode, result) | (mode, result) <- zip modes measured]+    mapM_+        ( \name -> do+            let xs = map (fromJust . lookup name) samples+            putStrLn (name <> " raw (tail ms, append ms, lookups, retained map/FIFO): " <> show xs)+            putStrLn (name <> " medians tail/append ms: " <> show (median [t | (t, _, _, _) <- xs], median [a | (_, a, _, _) <- xs]))+        )+        ["control", "warm", "cold"]+    finalHead <- runStoreIO store visibleGlobalHeadPosition+    case finalHead of+        Right (GlobalPosition 8001) -> putStrLn "Verified 501 seeded + 7500 concurrently appended events; final visible head 8001."+        other -> fail ("unexpected final head: " <> show other)++runCase :: KirokuStore -> V.Vector RecordedEvent -> String -> IO (Double, Double, Int, (Int, Int))+runCase store events mode = do+    cache <- newStreamNameCache+    calls <- newIORef (0 :: Int)+    count <- newTVarIO (0 :: Int)+    sub <- Pub.PublisherSubscription <$> newTBQueueIO 1 <*> newTVarIO Pub.Active <*> newTVarIO 0 <*> pure (pure ())+    let lookupBatch ids = do+            modifyIORef' calls (+ 1)+            either (error . show) id <$> runStoreIO store (lookupStreamNames ids)+        send msg = do+            _ <- evaluate (LBS.length (encode msg))+            case msg of+                Event _ -> atomically (modifyTVar' count (+ 1))+                _ -> fail "unexpected loss/error"+        worker =+            if mode == "control"+                then replicateM_ 50 $ do+                    UnchangedBatch batch <- atomically (readTBQueue sub.subscriptionQueue)+                    V.mapM_ (send . Event . recordedEventToJSON) (V.filter (const True) batch)+                    -- Reproduce the original loop's post-delivery status sample.+                    _ <- atomically (readTVar sub.subscriptionStatus)+                    pure ()+                else broadcastEventsWith send cache lookupBatch sub (const True)+    if mode == "warm" then resolveEventNames cache lookupBatch events >> writeIORef calls 0 else pure ()+    withAsync worker $ \_ -> do+        (tailMs, appendMs) <-+            concurrently+                ( timed $ do+                    replicateM_ 50 (atomically (writeTBQueue sub.subscriptionQueue (UnchangedBatch events)))+                    bounded (atomically (readTVar count >>= check . (== 25000)))+                )+                (timed $ replicateM_ 50 (append store (StreamName "appender") 10))+        lookups <- readIORef calls+        retained <- streamNameCacheSize cache+        if mode == "control"+            then pure ()+            else do+                if lookups == (if mode == "warm" then 0 else 1) && retained == (500, 500) then pure () else fail "cache invariant failed"+        pure (tailMs, appendMs, lookups, retained)++append :: KirokuStore -> StreamName -> Int -> IO ()+append store name n = do+    result <- runStoreIO store (appendToStream name AnyVersion (replicate n (EventData Nothing (EventType "Bench") (object []) Nothing Nothing Nothing)))+    either (fail . show) (const (pure ())) result++timed :: IO a -> IO Double+timed action = do+    t0 <- getCurrentTime+    _ <- action+    t1 <- getCurrentTime+    pure (realToFrac (diffUTCTime t1 t0) * 1000 :: Double)++bounded :: IO a -> IO a+bounded action = timeout 30_000_000 action >>= maybe (fail "focused tail timeout") pure+median :: [Double] -> Double+median xs = sort xs !! (length xs `div` 2)
example/Main.hs view
@@ -34,9 +34,12 @@     Manager,     defaultManagerSettings,     httpLbs,+    method,     newManager,     parseRequest,+    requestHeaders,     responseBody,+    responseHeaders,     responseStatus,  ) import Network.HTTP.Types (statusCode)@@ -46,9 +49,14 @@ import System.Timeout (timeout)  import Kiroku.Metrics (+    Capabilities (..),+    CheckpointInventoryResponse (..),     MetricsServer (..),     MetricsSnapshot (..),+    RouteAvailability (..),     StoreGauges (..),+    allowedOrigin,+    corsAllowOrigins,     defaultConfig,     metricsEventHandler,     metricsObservationHandler,@@ -75,7 +83,7 @@  main :: IO () main = withMigratedTestDatabase $ \connStr -> do-    step "[1/6] ephemeral postgres ready"+    step "[1/11] ephemeral postgres ready"      -- The collector must observe events from the first append, so its callbacks     -- go on ConnectionSettings BEFORE withStore. But snapshots read store-level@@ -90,14 +98,15 @@      withStore settings $ \store -> do         atomically (writeTVar storeVar (Just store))-        withMetricsServerWithStore (defaultConfig{port = 0}) metrics store [postgresPing store] $ \srv -> do+        origin <- either (fail . show) pure (allowedOrigin "https://ops.example.com")+        withMetricsServerWithStore (defaultConfig{port = 0, cors = corsAllowOrigins [origin]}) metrics store [postgresPing store] $ \srv -> do             threadDelay 300_000             let port = srv.serverPort                 base = "http://127.0.0.1:" <> show port-            step ("[2/6] store + collector + metrics server on port " <> show port)+            step ("[2/11] store + collector + metrics server on port " <> show port)              appendEvents store (StreamName "orders-1") ["OrderCreated", "OrderPaid", "OrderShipped"]-            step "[3/6] appended 3 events to orders-1"+            step "[3/11] appended 3 events to orders-1"              mgr <- newManager defaultManagerSettings @@ -116,8 +125,62 @@             check "GET /health/live is 200" (sLive == 200)             (sReady, _) <- httpGet mgr (base <> "/health/ready")             check "GET /health/ready is 200" (sReady == 200)-            step "[4/6] HTTP /metrics, /prometheus, /health/live, /health/ready all OK"+            step "[4/11] HTTP /metrics, /prometheus, /health/live, /health/ready all OK" +            -- Browser preflight and ordinary reads share the host's allowlist.+            request <- parseRequest (base <> "/metrics")+            let ops = "https://ops.example.com"+                evil = "https://evil.example.com"+            preflight <- httpLbs (request{method = "OPTIONS", requestHeaders = [("Origin", ops), ("Access-Control-Request-Method", "GET")]}) mgr+            check "CORS preflight is 204" (statusCode (responseStatus preflight) == 204)+            check "CORS preflight echoes origin" (lookup "Access-Control-Allow-Origin" (responseHeaders preflight) == Just ops)+            allowed <- httpLbs (request{requestHeaders = [("Origin", ops)]}) mgr+            check "CORS GET echoes origin" (lookup "Access-Control-Allow-Origin" (responseHeaders allowed) == Just ops)+            denied <- httpLbs (request{requestHeaders = [("Origin", evil)]}) mgr+            check "CORS denied GET is unchanged" (statusCode (responseStatus denied) == 200)+            check "CORS denied GET grants nothing" (all (\(name, _) -> name `notElem` ["Access-Control-Allow-Origin", "Access-Control-Allow-Credentials", "Access-Control-Allow-Methods", "Access-Control-Allow-Headers", "Access-Control-Max-Age"]) (responseHeaders denied))+            step "[5/11] CORS: preflight and GET from https://ops.example.com allowed; https://evil.example.com undecorated"++            (sCheckpoints, bCheckpoints) <- httpGet mgr (base <> "/subscription-checkpoints")+            check "GET /subscription-checkpoints is 200" (sCheckpoints == 200)+            case decode bCheckpoints of+                Just (CheckpointInventoryResponse position rows) -> do+                    check "durable store_position >= 3" (position >= 3)+                    check "no durable checkpoints without subscriptions" (null rows)+                    step ("[6/11] GET /subscription-checkpoints store_position=" <> show position <> " with no durable checkpoints (this example runs no subscription)")+                Nothing -> check "durable inventory decodes" False++            (sStreams, bStreams) <- httpGet mgr (base <> "/streams?category=orders")+            check "GET /streams is 200" (sStreams == 200)+            check "stream page has items" (maybe False (\v -> case lookKey ["items"] v of Just (Array rows) -> not (null rows); _ -> False) (decode bStreams))+            (sEvents, bEvents) <- httpGet mgr (base <> "/events?from=0&limit=10")+            check "GET /events is 200" (sEvents == 200)+            case decode bEvents of+                Just value | Just (Array rows) <- lookKey ["items"] value -> do+                    check "event page includes original names" (all (\row -> lookKey ["original_stream_name"] row == Just (String "orders-1")) rows)+                    case foldr (:) [] rows of+                        first : _ | Just (String uuid) <- lookKey ["eventId"] first -> do+                            (status, _) <- httpGet mgr (base <> "/events/" <> T.unpack uuid)+                            check "GET event by id is 200" (status == 200)+                        _ -> check "event page is nonempty" False+                _ -> check "event page decodes" False+            step "[7/11] Stream browsing and historical events resolve original stream names"++            (sDeadLetters, bDeadLetters) <- httpGet mgr (base <> "/subscriptions/example/dead-letters?limit=10")+            check "dead-letter page status" (sDeadLetters == 200)+            check "dead-letter page is empty without a cursor" (decode bDeadLetters == Just (object ["items" .= ([] :: [Value])]))+            step "[8/11] GET /subscriptions/example/dead-letters returned an empty page (this example runs no subscription)"++            (sCaps, bCaps) <- httpGet mgr (base <> "/capabilities")+            check "discovery status" (sCaps == 200)+            case decode bCaps of+                Just (caps :: Capabilities) ->+                    check+                        "discovery reflects legacy store wiring"+                        (caps.routes.browse && caps.routes.subscriptionsCheckpoints && caps.routes.deadLetters && caps.routes.websocketEvents && not caps.routes.subscriptionsLive)+                Nothing -> check "discovery decodes" False+            step "[9/11] GET /capabilities reports browse, checkpoints, dead letters, and the event tail"+             -- WebSocket: subscribe to the live event tail, then append one more             -- event and assert it arrives over the socket as a JSON event message.             evType <-@@ -131,10 +194,10 @@                             wait appendThread                             pure (eventTypeOf ev)             check ("WebSocket event eventType == OrderRefunded (got " <> T.unpack evType <> ")") (evType == "OrderRefunded")-            step ("[5/6] WebSocket /ws/events received event eventType=" <> T.unpack evType)+            step ("[10/11] WebSocket /ws/events received event eventType=" <> T.unpack evType)              snap <- snapshotMetrics metrics-            step ("[6/6] kiroku-metrics-example: all checks passed (snapshot global position = " <> show snap.store.globalPosition <> ")")+            step ("[11/11] kiroku-metrics-example: all checks passed (snapshot global position = " <> show snap.store.globalPosition <> ")")  -------------------------------------------------------------------------------- -- Helpers
kiroku-metrics.cabal view
@@ -1,6 +1,6 @@ cabal-version:   3.0 name:            kiroku-metrics-version:         0.2.0.0+version:         0.3.0.0 synopsis:   Metrics, health, and event-streaming HTTP endpoints for Kiroku @@ -41,39 +41,62 @@   import:          common   exposed-modules:     Kiroku.Metrics+    Kiroku.Metrics.Browse+    Kiroku.Metrics.Capabilities+    Kiroku.Metrics.Checkpoints     Kiroku.Metrics.Collector     Kiroku.Metrics.Config+    Kiroku.Metrics.Cors+    Kiroku.Metrics.DeadLetters     Kiroku.Metrics.Health     Kiroku.Metrics.JSON     Kiroku.Metrics.Prometheus     Kiroku.Metrics.Server+    Kiroku.Metrics.Standalone     Kiroku.Metrics.Subscriptions     Kiroku.Metrics.Types     Kiroku.Metrics.WebSocket +  other-modules:   Paths_kiroku_metrics+  autogen-modules: Paths_kiroku_metrics   build-depends:-    , aeson           >=2.1       && <2.3-    , async           >=2.2       && <2.3-    , base            >=4.18      && <5-    , bytestring      >=0.11      && <0.13-    , containers      >=0.6       && <0.8-    , hasql           >=1.10      && <1.11-    , hasql-pool      >=1.2       && <1.5-    , http-types      >=0.12      && <0.13-    , kiroku-cli      ^>=0.2.0.9-    , kiroku-store    ^>=0.10.0.0-    , stm             >=2.5       && <2.6-    , text            >=2.0       && <2.2-    , time            >=1.12      && <1.15-    , uuid            >=1.3       && <1.4-    , vector          >=0.13      && <0.14-    , wai             >=3.2       && <3.3-    , wai-websockets  >=3.0       && <3.1-    , warp            >=3.4       && <3.5-    , websockets      >=0.13      && <0.14+    , aeson                 >=2.1       && <2.3+    , async                 >=2.2       && <2.3+    , base                  >=4.18      && <5+    , bytestring            >=0.11      && <0.13+    , containers            >=0.6       && <0.8+    , effectful-core        >=2.6.1     && <2.7  || >=2.7.1.1 && <2.8+    , hasql                 >=1.10      && <1.11+    , hasql-pool            >=1.2       && <1.5+    , http-types            >=0.12      && <0.13+    , kiroku-cli            ^>=0.2.0.10+    , kiroku-store          ^>=0.11.0.0+    , network               >=3.1       && <3.3+    , optparse-applicative  >=0.19      && <0.20+    , stm                   >=2.5       && <2.6+    , text                  >=2.0       && <2.2+    , time                  >=1.12      && <1.15+    , uuid                  >=1.3       && <1.4+    , vector                >=0.13      && <0.14+    , wai                   >=3.2       && <3.3+    , wai-websockets        >=3.0       && <3.1+    , warp                  >=3.4       && <3.5+    , websockets            >=0.13      && <0.14    hs-source-dirs:  src +executable kiroku-inspect+  import:         common+  main-is:        Main.hs+  hs-source-dirs: app-inspect+  ghc-options:    -threaded -rtsopts -with-rtsopts=-N+  build-depends:+    , base                  >=4.18 && <5+    , kiroku-metrics+    , optparse-applicative  >=0.19 && <0.20+    , text                  >=2.0  && <2.2+    , unix                  >=2.8  && <2.9+ flag example   description:     Build the self-verifying example executable. Off by default because it@@ -103,7 +126,7 @@     , http-client          >=0.7       && <0.8     , http-types           >=0.12      && <0.13     , kiroku-metrics-    , kiroku-store         ^>=0.10.0.0+    , kiroku-store         ^>=0.11.0.0     , kiroku-test-support  ^>=0.1     , lens                 >=5.2       && <5.4     , scientific           >=0.3       && <0.4@@ -112,38 +135,75 @@     , websockets           >=0.13      && <0.14  test-suite kiroku-metrics-test-  import:         common-  type:           exitcode-stdio-1.0-  main-is:        Main.hs+  import:             common+  type:               exitcode-stdio-1.0+  build-tool-depends: kiroku-metrics:kiroku-inspect+  main-is:            Main.hs   other-modules:+    Test.BrowseSpec+    Test.CapabilitiesSpec+    Test.CheckpointsSpec     Test.CollectorSpec+    Test.CorsSpec+    Test.DeadLettersSpec     Test.IntegrationSpec     Test.ServerSpec+    Test.StandaloneSpec     Test.SubscriptionsSpec+    Test.WebSocketConvergenceSpec     Test.WebSocketSpec -  hs-source-dirs: test-  ghc-options:    -threaded -rtsopts -with-rtsopts=-N+  hs-source-dirs:     test+  ghc-options:        -threaded -rtsopts -with-rtsopts=-N   build-depends:-    , aeson                >=2.1       && <2.3-    , async                >=2.2       && <2.3-    , base                 >=4.18      && <5-    , bytestring           >=0.11      && <0.13-    , containers           >=0.6       && <0.8-    , generic-lens         >=2.2       && <2.4-    , hasql                >=1.10      && <1.11-    , hasql-pool           >=1.2       && <1.5-    , hspec                >=2.10      && <2.12-    , http-client          >=0.7       && <0.8-    , http-types           >=0.12      && <0.13-    , kiroku-cli           ^>=0.2.0.9+    , aeson                 >=2.1       && <2.3+    , async                 >=2.2       && <2.3+    , base                  >=4.18      && <5+    , bytestring            >=0.11      && <0.13+    , case-insensitive      >=1.2       && <1.3+    , containers            >=0.6       && <0.8+    , directory             >=1.3       && <1.4+    , effectful-core        >=2.6.1     && <2.7  || >=2.7.1.1 && <2.8+    , generic-lens          >=2.2       && <2.4+    , hasql                 >=1.10      && <1.11+    , hasql-pool            >=1.2       && <1.5+    , hspec                 >=2.10      && <2.12+    , http-client           >=0.7       && <0.8+    , http-types            >=0.12      && <0.13+    , kiroku-cli            ^>=0.2.0.10     , kiroku-metrics-    , kiroku-store         ^>=0.10.0.0+    , kiroku-store          ^>=0.11.0.0     , kiroku-test-support-    , lens                 >=5.2       && <5.4-    , scientific           >=0.3       && <0.4-    , stm                  >=2.5       && <2.6-    , text                 >=2.0       && <2.2-    , uuid                 >=1.3       && <1.4-    , warp                 >=3.4       && <3.5-    , websockets           >=0.13      && <0.14+    , lens                  >=5.2       && <5.4+    , network               >=3.1       && <3.3+    , optparse-applicative  >=0.19      && <0.20+    , process               >=1.6       && <1.7+    , scientific            >=0.3       && <0.4+    , stm                   >=2.5       && <2.6+    , text                  >=2.0       && <2.2+    , time                  >=1.12      && <1.15+    , unix                  >=2.8       && <2.9+    , uuid                  >=1.3       && <1.4+    , vector                >=0.13      && <0.14+    , wai                   >=3.2       && <3.3+    , warp                  >=3.4       && <3.5+    , websockets            >=0.13      && <0.14++benchmark kiroku-websocket-tail+  import:         common+  type:           exitcode-stdio-1.0+  main-is:        WebSocketTail.hs+  hs-source-dirs: bench+  ghc-options:    -threaded -rtsopts "-with-rtsopts=-N -A32m"+  build-depends:+    , aeson                >=2.1  && <2.3+    , async                >=2.2  && <2.3+    , base                 >=4.18 && <5+    , bytestring           >=0.11 && <0.13+    , kiroku-metrics+    , kiroku-store+    , kiroku-test-support+    , stm                  >=2.5  && <2.6+    , text                 >=2.0  && <2.2+    , time                 >=1.12 && <1.15+    , vector               >=0.13 && <0.14
src/Kiroku/Metrics.hs view
@@ -4,19 +4,31 @@ types into scope. -} module Kiroku.Metrics (+    module Kiroku.Metrics.Browse,+    module Kiroku.Metrics.Capabilities,     module Kiroku.Metrics.Types,     module Kiroku.Metrics.Collector,+    module Kiroku.Metrics.Checkpoints,     module Kiroku.Metrics.Config,+    module Kiroku.Metrics.Cors,+    module Kiroku.Metrics.DeadLetters,     module Kiroku.Metrics.Health,     module Kiroku.Metrics.Server,+    module Kiroku.Metrics.Standalone,     module Kiroku.Metrics.Subscriptions,     module Kiroku.Metrics.WebSocket, ) where +import Kiroku.Metrics.Browse+import Kiroku.Metrics.Capabilities+import Kiroku.Metrics.Checkpoints import Kiroku.Metrics.Collector import Kiroku.Metrics.Config+import Kiroku.Metrics.Cors+import Kiroku.Metrics.DeadLetters import Kiroku.Metrics.Health import Kiroku.Metrics.Server+import Kiroku.Metrics.Standalone import Kiroku.Metrics.Subscriptions import Kiroku.Metrics.Types import Kiroku.Metrics.WebSocket
+ src/Kiroku/Metrics/Browse.hs view
@@ -0,0 +1,246 @@+{- | Bounded inspection reads. Stream names use stable UTF-8 byte order;+category enumeration retains the database's category collation. Cursors are+exclusive and describe live pages rather than a cross-request snapshot.+-}+module Kiroku.Metrics.Browse (+    StoreBrowser (..),+    storeBrowser,+    storeBrowserWith,+    BrowseLimits,+    BrowseLimitsError (..),+    mkBrowseLimits,+    defaultLimit,+    maxLimit,+    defaultBrowseLimits,+    ReadDirection (..),+    browseApp,+    browseNotConfiguredApp,+    streamInfoToJSON,+) where++import Data.Aeson (Value, encode, object, toJSON, (.=))+import Data.ByteString qualified as BS+import Data.Int (Int32, Int64)+import Data.Map.Strict (Map)+import Data.Map.Strict qualified as Map+import Data.Set qualified as Set+import Data.Text (Text)+import Data.Text qualified as T+import Data.Text.Encoding qualified as TE+import Data.UUID qualified as UUID+import Data.Vector (Vector)+import Data.Vector qualified as V+import Effectful (Eff, IOE)+import Effectful.Error.Static (Error)+import Kiroku.Metrics.JSON (errorResponse, jsonResponse, storeErrorResponse)+import Kiroku.Metrics.WebSocket (recordedEventToJSONResolved)+import Kiroku.Store+import Network.HTTP.Types (Status, status200, status400, status404, status405)+import Network.Wai (Application, Request, Response, mapResponseHeaders, pathInfo, queryString, requestMethod, responseHeaders, responseLBS, responseStatus)+import Text.Read (readMaybe)++data StoreBrowser = StoreBrowser+    { runStoreRead :: forall a. Eff '[Store, Error StoreError, IOE] a -> IO (Either StoreError a)+    , limits :: !BrowseLimits+    }++data BrowseLimits = BrowseLimits {defaultLimit :: !Int, maxLimit :: !Int}+    deriving stock (Eq, Show)+data BrowseLimitsError = InvalidBrowseLimits !Int !Int deriving stock (Eq, Show)+mkBrowseLimits :: Int -> Int -> Either BrowseLimitsError BrowseLimits+mkBrowseLimits def cap+    | 1 <= def && def <= cap && cap <= 1000 = Right (BrowseLimits def cap)+    | otherwise = Left (InvalidBrowseLimits def cap)+defaultBrowseLimits :: BrowseLimits+defaultBrowseLimits = BrowseLimits 100 1000+storeBrowser :: KirokuStore -> StoreBrowser+storeBrowser = storeBrowserWith defaultBrowseLimits+storeBrowserWith :: BrowseLimits -> KirokuStore -> StoreBrowser+storeBrowserWith lims store = StoreBrowser (runStoreIO store) lims++data ReadDirection = ReadForward | ReadBackward deriving stock (Eq, Show)+type ReadProgram a = Eff '[Store, Error StoreError, IOE] a+type Failure = (Status, Text, Text, Maybe Value)++-- | Preserve GET status/headers for HEAD. Other methods never run a read.+readMethods :: (Request -> IO Response) -> Application+readMethods action req respond+    | requestMethod req == "GET" = action req >>= respond+    | requestMethod req == "HEAD" = do+        response <- action req+        respond (responseLBS (responseStatus response) (responseHeaders response) "")+    | otherwise =+        respond $+            mapResponseHeaders (("Allow", "GET, HEAD") :) $+                errorResponse status405 "method_not_allowed" "Only GET and HEAD are supported." Nothing++browseNotConfiguredApp :: Application+browseNotConfiguredApp = readMethods $ \_ -> pure (errorResponse status404 "store_browsing_not_configured" "Store browsing is not configured." Nothing)++browseApp :: StoreBrowser -> Application+browseApp browser = readMethods $ \req -> case route browser req of+    Left (status, code, message, details) -> pure (errorResponse status code message details)+    Right program ->+        runStoreRead browser program >>= \case+            Left err -> pure (storeErrorResponse "store_unavailable" err)+            Right (Left (status, code, message, details)) -> pure (errorResponse status code message details)+            Right (Right value) -> pure (jsonResponse status200 (encode value))++route :: StoreBrowser -> Request -> Either Failure (ReadProgram (Either Failure Value))+route browser req = case pathInfo req of+    ["streams"] -> do+        limit <- pageLimit+        category <- fmap CategoryName <$> textParam "category"+        prefix <- textParam "prefix"+        cursor <- fmap StreamName <$> textParam "from"+        size <- browseSize limit+        pure $ do+            rows <- listStreams category prefix cursor size+            pure $ Right $ pageJSON limit streamInfoToJSON (\row -> let StreamName name = row.name in AesonText name) rows+    ["streams", name] -> do+        stream <- validStream name+        pure $ getStream stream >>= pure . maybe (Left (missing "stream_not_found" "Stream not found.")) (Right . streamInfoToJSON)+    ["streams", name, "events"] -> do+        stream <- validStream name+        limit <- pageLimit+        cursor <- position+        direction <- readDirection+        pure $+            getStream stream >>= \case+                Nothing -> pure (Left (missing "stream_not_found" "Stream not found."))+                Just _ ->+                    Right+                        <$> eventPage+                            limit+                            (\event -> let StreamVersion v = event.streamVersion in AesonNumber v)+                            (case direction of ReadForward -> readStreamForward stream (StreamVersion cursor) (overfetch limit); ReadBackward -> readStreamBackward stream (StreamVersion cursor) (overfetch limit))+    ["categories"] -> do+        limit <- pageLimit+        cursor <- fmap CategoryName <$> textParam "from"+        size <- browseSize limit+        pure $ do+            rows <- listCategories cursor size+            pure $ Right $ pageJSON limit (\(CategoryName name) -> object ["name" .= name]) (\(CategoryName name) -> AesonText name) rows+    ["categories", name, "events"] -> do+        validText "category" name+        limit <- pageLimit+        cursor <- position+        -- Category reads have one supported direction; do not silently ignore it.+        direction <- readDirection+        if direction == ReadBackward+            then Left (invalid "direction" "backward" "Category events support forward reads.")+            else+                pure $ Right <$> eventPage limit globalCursor (readCategory (CategoryName name) (GlobalPosition cursor) (overfetch limit))+    ["events"] -> do+        limit <- pageLimit+        cursor <- position+        direction <- readDirection+        pure $+            Right+                <$> eventPage+                    limit+                    globalCursor+                    (case direction of ReadForward -> readAllForward (GlobalPosition cursor) (overfetch limit); ReadBackward -> readAllBackward (GlobalPosition cursor) (overfetch limit))+    ["events", rawId] -> case UUID.fromText rawId of+        Nothing -> Left (status400, "invalid_event_id", "The event id must be a UUID.", Nothing)+        Just uuid ->+            pure $+                getEvent (EventId uuid) >>= \case+                    Nothing -> pure (Left (missing "event_not_found" "Event not found."))+                    Just event -> do+                        mapping <- lookupStreamNames [event.originalStreamId]+                        pure (Right (recordedEventToJSONResolved mapping event))+    _ -> Left (missing "not_found" "Not found.")+  where+    pageLimit = do+        value <- rawParam req "limit"+        case value of+            Nothing -> Right browser.limits.defaultLimit+            Just raw -> case decimal raw of+                Just n | 1 <= n && n <= toInteger browser.limits.maxLimit -> Right (fromInteger n)+                _ -> Left (invalid "limit" raw "Expected a decimal integer within the configured page limit.")+    textParam key = do+        value <- rawParam req key+        mapM_ (validText key) value+        pure value+    position = do+        value <- rawParam req "from"+        case value of+            Nothing -> Right 0+            Just raw -> case decimal raw of+                Just n | n <= toInteger (maxBound :: Int64) -> Right (fromInteger n)+                _ -> Left (invalid "from" raw "Expected a non-negative Int64 decimal integer.")+    readDirection =+        rawParam req "direction" >>= \case+            Nothing -> Right ReadForward+            Just "forward" -> Right ReadForward+            Just "backward" -> Right ReadBackward+            Just raw -> Left (invalid "direction" raw "Expected forward or backward.")++-- Closed bounds make over-fetch safe before conversion to the store's Int32.+overfetch :: Int -> Int32+overfetch limit = fromIntegral (limit + 1)+browseSize :: Int -> Either Failure BrowsePageSize+browseSize limit = case mkBrowsePageSize (limit + 1) of+    Right size -> Right size+    Left _ -> Left (invalid "limit" (T.pack (show limit)) "Invalid page limit.")++data Cursor = AesonText Text | AesonNumber Int64+cursorJSON :: Cursor -> Value+cursorJSON (AesonText text) = toJSON text+cursorJSON (AesonNumber n) = toJSON n++pageJSON :: Int -> (a -> Value) -> (a -> Cursor) -> Vector a -> Value+pageJSON limit encodeItem cursor rows =+    object $ ["items" .= V.map encodeItem items] <> ["next_cursor" .= cursorJSON (cursor (V.last items)) | V.length rows > limit && not (V.null items)]+  where+    items = V.take limit rows++eventPage :: Int -> (RecordedEvent -> Cursor) -> ReadProgram (Vector RecordedEvent) -> ReadProgram Value+eventPage limit cursor readPage = do+    rows <- readPage+    mapping <- resolveNames (V.take limit rows)+    pure (pageJSON limit (recordedEventToJSONResolved mapping) cursor rows)+resolveNames :: Vector RecordedEvent -> ReadProgram (Map StreamId StreamName)+resolveNames rows+    | V.null rows = pure Map.empty+    | otherwise = lookupStreamNames (Set.toList (Set.fromList (map (\event -> event.originalStreamId) (V.toList rows))))+globalCursor :: RecordedEvent -> Cursor+globalCursor event = let GlobalPosition n = event.globalPosition in AesonNumber n++streamInfoToJSON :: StreamInfo -> Value+streamInfoToJSON stream =+    object+        [ "stream_id" .= (let StreamId n = stream.id in n)+        , "name" .= (let StreamName name = stream.name in name)+        , "category" .= (let CategoryName category = categoryName stream.name in category)+        , "version" .= (let StreamVersion n = stream.version in n)+        , "created_at" .= stream.createdAt+        , "deleted_at" .= stream.deletedAt+        , "truncate_before" .= (let StreamVersion n = stream.truncateBefore in n)+        ]++missing :: Text -> Text -> Failure+missing code message = (status404, code, message, Nothing)+invalid :: Text -> Text -> Text -> Failure+invalid key raw reason = (status400, "invalid_query_parameter", "Invalid query parameter.", Just (object ["parameter" .= key, "value" .= raw, "reason" .= reason]))+rawParam :: Request -> BS.ByteString -> Either Failure (Maybe Text)+rawParam req key = case lookup key (queryString req) of+    Nothing -> Right Nothing+    Just Nothing -> Right (Just "")+    Just (Just bytes) -> case TE.decodeUtf8' bytes of+        Right value -> Right (Just value)+        Left _ -> Left (invalid (TE.decodeUtf8 key) "<invalid UTF-8>" "Expected valid UTF-8.")+validText :: BS.ByteString -> Text -> Either Failure ()+validText key value+    | T.any (== '\0') value || BS.length (TE.encodeUtf8 value) > 512 = Left (invalid (TE.decodeUtf8 key) value "Expected at most 512 UTF-8 bytes without NUL.")+    | otherwise = Right ()+validStream :: Text -> Either Failure StreamName+validStream value = case validateStreamName (StreamName value) of+    Right () | not (T.any (== '\0') value) -> Right (StreamName value)+    _ -> Left (status400, "invalid_stream_name", "Invalid stream name.", Just (object ["stream_name" .= value]))+decimal :: Text -> Maybe Integer+decimal value+    | T.null value || not (T.all (\c -> '0' <= c && c <= '9') value) = Nothing+    | T.length value > 20 = Nothing+    | otherwise = readMaybe (T.unpack value)
+ src/Kiroku/Metrics/Capabilities.hs view
@@ -0,0 +1,152 @@+{-# LANGUAGE NoFieldSelectors #-}++-- | Mount-relative discovery, computed from configuration and actual provider wiring.+module Kiroku.Metrics.Capabilities (+    WebSocketChannels (..),+    noWebSocketChannels,+    storeWebSocketChannels,+    ProviderPresence (..),+    RouteAvailability (..),+    Capabilities (..),+    capabilitiesFor,+    kirokuMetricsVersion,+    processLocalRoutes,+    capabilitiesApp,+    capabilitiesPath,+) where++import Data.Aeson (FromJSON (..), ToJSON (..), encode, object, withObject, (.:), (.=))+import Data.Text (Text)+import Data.Text qualified as T+import Data.Version (showVersion)+import Network.HTTP.Types (status200, status404, status405)+import Network.Wai (Application, mapResponseHeaders, pathInfo, requestMethod, responseHeaders, responseLBS, responseStatus)++import Kiroku.Metrics.Config (MetricsServerConfig (..))+import Kiroku.Metrics.Cors qualified as Cors+import Kiroku.Metrics.JSON (errorResponse, jsonResponse)+import Paths_kiroku_metrics qualified as Paths++-- | Declaration by the host choosing the opaque WebSocket application.+data WebSocketChannels = WebSocketChannels+    {metricsChannel :: !Bool, eventsChannel :: !Bool}+    deriving stock (Eq, Show)++noWebSocketChannels, storeWebSocketChannels :: WebSocketChannels+noWebSocketChannels = WebSocketChannels False False+storeWebSocketChannels = WebSocketChannels True True++data ProviderPresence = ProviderPresence+    { hasSubscriptionStatus :: !Bool+    , hasCheckpointInventory :: !Bool+    , hasBrowser :: !Bool+    , hasDeadLetters :: !Bool+    , presentWebSocketChannels :: !WebSocketChannels+    }+    deriving stock (Eq, Show)++data RouteAvailability = RouteAvailability+    { metrics :: !Bool+    , prometheus :: !Bool+    , health :: !Bool+    , subscriptionsLive :: !Bool+    , subscriptionsCheckpoints :: !Bool+    , deadLetters :: !Bool+    , browse :: !Bool+    , websocketMetrics :: !Bool+    , websocketEvents :: !Bool+    }+    deriving stock (Eq, Show)++data Capabilities = Capabilities+    { package :: !Text+    , version :: !Text+    , routes :: !RouteAvailability+    , corsIsEnabled :: !Bool+    , processLocal :: ![Text]+    }+    deriving stock (Eq, Show)++kirokuMetricsVersion :: Text+kirokuMetricsVersion = T.pack (showVersion Paths.version)++-- | These answers describe only the process answering the request.+processLocalRoutes :: [Text]+processLocalRoutes = ["metrics", "prometheus", "health", "subscriptions_live", "websocket_metrics"]++capabilitiesFor :: MetricsServerConfig -> ProviderPresence -> Capabilities+capabilitiesFor cfg presence = Capabilities "kiroku-metrics" kirokuMetricsVersion availability (Cors.corsEnabled cfg.cors) processLocalRoutes+  where+    availability =+        RouteAvailability+            cfg.enableJSON+            cfg.enablePrometheus+            cfg.enableJSON+            presence.hasSubscriptionStatus+            presence.hasCheckpointInventory+            presence.hasDeadLetters+            presence.hasBrowser+            (cfg.enableWebSocket && presence.presentWebSocketChannels.metricsChannel)+            (cfg.enableWebSocket && presence.presentWebSocketChannels.eventsChannel)++instance ToJSON RouteAvailability where+    toJSON r =+        object+            [ "metrics" .= r.metrics+            , "prometheus" .= r.prometheus+            , "health" .= r.health+            , "subscriptions_live" .= r.subscriptionsLive+            , "subscriptions_checkpoints" .= r.subscriptionsCheckpoints+            , "dead_letters" .= r.deadLetters+            , "browse" .= r.browse+            , "websocket_metrics" .= r.websocketMetrics+            , "websocket_events" .= r.websocketEvents+            ]++instance FromJSON RouteAvailability where+    parseJSON = withObject "RouteAvailability" $ \o ->+        RouteAvailability+            <$> o .: "metrics"+            <*> o .: "prometheus"+            <*> o .: "health"+            <*> o .: "subscriptions_live"+            <*> o .: "subscriptions_checkpoints"+            <*> o .: "dead_letters"+            <*> o .: "browse"+            <*> o .: "websocket_metrics"+            <*> o .: "websocket_events"++instance ToJSON Capabilities where+    toJSON c =+        object+            [ "package" .= c.package+            , "version" .= c.version+            , "routes" .= c.routes+            , "cors" .= object ["enabled" .= c.corsIsEnabled]+            , "process_local" .= c.processLocal+            ]++instance FromJSON Capabilities where+    parseJSON = withObject "Capabilities" $ \o -> do+        corsObject <- o .: "cors"+        enabled <- withObject "cors" (.: "enabled") corsObject+        Capabilities <$> o .: "package" <*> o .: "version" <*> o .: "routes" <*> pure enabled <*> o .: "process_local"++capabilitiesPath :: [Text]+capabilitiesPath = ["capabilities"]++-- | No store access; the encoded body is shared by requests to this application.+capabilitiesApp :: Capabilities -> Application+capabilitiesApp caps = \req respond -> do+    let response+            | pathInfo req /= capabilitiesPath = errorResponse status404 "not_found" "Not found" Nothing+            | requestMethod req == "GET" || requestMethod req == "HEAD" = jsonResponse status200 body+            | otherwise =+                mapResponseHeaders (("Allow", "GET, HEAD") :) $+                    errorResponse status405 "method_not_allowed" "Use GET or HEAD." Nothing+    respond $+        if requestMethod req == "HEAD"+            then responseLBS (responseStatus response) (responseHeaders response) ""+            else response+  where+    body = encode caps
+ src/Kiroku/Metrics/Checkpoints.hs view
@@ -0,0 +1,89 @@+-- | Read-only durable checkpoint inventory, independent of the live registry.+module Kiroku.Metrics.Checkpoints (+    CheckpointInventoryProvider,+    storeCheckpointInventory,+    CheckpointInventoryResponse (..),+    CheckpointRow (..),+    checkpointInventoryResponse,+    checkpointsPath,+    checkpointsApp,+    checkpointsNotConfiguredApp,+) where++import Data.Aeson (FromJSON (..), ToJSON (..), encode, object, withObject, (.:), (.=))+import Data.Int (Int32, Int64)+import Data.Text (Text)+import Data.Time (UTCTime)+import Data.Vector qualified as V+import Network.HTTP.Types (status200, status404, status405)+import Network.Wai (Application, Response, mapResponseHeaders, pathInfo, requestMethod, responseHeaders, responseLBS, responseStatus)++import Kiroku.Metrics.JSON (errorResponse, jsonResponse, storeErrorResponse)+import Kiroku.Store (GlobalPosition (..), KirokuStore, StoreError, SubscriptionCheckpoint (..), SubscriptionCheckpointInventory (..), SubscriptionName (..), runStoreIO, subscriptionCheckpointInventory)++{- | One call reads the frontier and all checkpoint rows in one SQL snapshot.+Inventory work is proportional to row count; clients should not overlap polls.+-}+type CheckpointInventoryProvider = IO (Either StoreError SubscriptionCheckpointInventory)++storeCheckpointInventory :: KirokuStore -> CheckpointInventoryProvider+storeCheckpointInventory store = runStoreIO store subscriptionCheckpointInventory++data CheckpointRow = CheckpointRow+    { subscription :: !Text+    , member :: !Int32+    , checkpointPosition :: !Int64+    , updatedAt :: !UTCTime+    }+    deriving stock (Eq, Show)++data CheckpointInventoryResponse = CheckpointInventoryResponse+    { storePosition :: !Int64+    , checkpoints :: ![CheckpointRow]+    }+    deriving stock (Eq, Show)++instance ToJSON CheckpointRow where+    toJSON row = object ["subscription" .= row.subscription, "member" .= row.member, "checkpoint_position" .= row.checkpointPosition, "updated_at" .= row.updatedAt]++instance FromJSON CheckpointRow where+    parseJSON = withObject "CheckpointRow" $ \o -> CheckpointRow <$> o .: "subscription" <*> o .: "member" <*> o .: "checkpoint_position" <*> o .: "updated_at"++instance ToJSON CheckpointInventoryResponse where+    toJSON inventory = object ["store_position" .= inventory.storePosition, "checkpoints" .= inventory.checkpoints]++instance FromJSON CheckpointInventoryResponse where+    parseJSON = withObject "CheckpointInventoryResponse" $ \o -> CheckpointInventoryResponse <$> o .: "store_position" <*> o .: "checkpoints"++checkpointInventoryResponse :: SubscriptionCheckpointInventory -> CheckpointInventoryResponse+checkpointInventoryResponse (SubscriptionCheckpointInventory (GlobalPosition position) rows) =+    CheckpointInventoryResponse position (map toRow (V.toList rows))+  where+    toRow (SubscriptionCheckpoint (SubscriptionName name) index (GlobalPosition cp) updated) = CheckpointRow name index cp updated++checkpointsPath :: [Text]+checkpointsPath = ["subscription-checkpoints"]++{- | GET and HEAD only. Unknown query parameters are ignored. Expected store+failures are sanitized; thrown exceptions (including cancellation) propagate.+-}+checkpointsApp :: CheckpointInventoryProvider -> Application+checkpointsApp provider = checkpointsResponseApp $ either (storeErrorResponse "checkpoint_inventory_unavailable") (jsonResponse status200 . encode . checkpointInventoryResponse) <$> provider++-- | Same method and HEAD behavior when a server has no inventory provider.+checkpointsNotConfiguredApp :: Application+checkpointsNotConfiguredApp = checkpointsResponseApp $ pure $ errorResponse status404 "checkpoint_inventory_not_configured" "Durable checkpoint inventory is not configured." Nothing++checkpointsResponseApp :: IO Response -> Application+checkpointsResponseApp readResponse req respond = do+    response <-+        if pathInfo req /= checkpointsPath+            then pure $ errorResponse status404 "not_found" "Not found" Nothing+            else+                if requestMethod req `notElem` ["GET", "HEAD"]+                    then pure $ mapResponseHeaders (("Allow", "GET, HEAD") :) $ errorResponse status405 "method_not_allowed" "Use GET or HEAD." Nothing+                    else readResponse+    respond $+        if requestMethod req == "HEAD"+            then responseLBS (responseStatus response) (responseHeaders response) ""+            else response
src/Kiroku/Metrics/Config.hs view
@@ -7,6 +7,8 @@ import Data.Int (Int64) import Numeric.Natural (Natural) +import Kiroku.Metrics.Cors (CorsPolicy, corsDisabled)+ {- | Configuration for the metrics web server. The @ws*@ fields are consumed by the WebSocket endpoint (EP-3); they exist here so that plan needs no config change. @readinessMaxLag@ is the Kiroku analogue of Marten's @maxEventLag@; it@@ -36,6 +38,8 @@     -- ^ A subscription lagging beyond this fails readiness (default: 10_000).     , livenessTimeoutUs :: !Int     -- ^ Timeout for the liveness snapshot in microseconds (default: 1_000_000 = 1s).+    , cors :: !CorsPolicy+    -- ^ Explicit browser origins; disabled by default for HTTP and WebSocket alike.     }     deriving stock (Eq, Show) @@ -52,4 +56,5 @@         , wsEventQueueCap = 256         , readinessMaxLag = 10_000         , livenessTimeoutUs = 1_000_000+        , cors = corsDisabled         }
+ src/Kiroku/Metrics/Cors.hs view
@@ -0,0 +1,295 @@+{-# LANGUAGE BangPatterns #-}++-- | Default-off browser access for HTTP and WebSocket inspection surfaces.+module Kiroku.Metrics.Cors (+    AllowedOrigin,+    OriginError (..),+    allowedOrigin,+    renderAllowedOrigin,+    CorsPolicy (..),+    corsDisabled,+    corsAllowOrigins,+    corsEnabled,+    corsAllowedMethods,+    originAllowed,+    isPreflight,+    corsMiddleware,+) where++import Control.Monad (guard)+import Data.ByteString (ByteString)+import Data.ByteString.Char8 qualified as BS+import Data.Char (digitToInt, isAscii, isAsciiLower, isAsciiUpper, isHexDigit, ord, toLower)+import Data.List (nubBy)+import Data.Maybe (isJust)+import Data.Set qualified as Set+import Data.Text (Text)+import Data.Text qualified as T+import Data.Text.Encoding qualified as TE+import Network.HTTP.Types (HeaderName, RequestHeaders, ResponseHeaders, methodOptions, status204, status400, status403)+import Network.Wai (Middleware, Request, mapResponseHeaders, requestHeaders, requestMethod, responseLBS)+import Network.Wai.Handler.WebSockets (isWebSocketsReq)+import Numeric (showHex)++import Kiroku.Metrics.JSON (errorResponse)++-- | An explicit normalized HTTP(S) origin. Construct with 'allowedOrigin'.+newtype AllowedOrigin = AllowedOrigin Text+    deriving stock (Eq, Ord, Show)++-- | Configuration errors; input is retained only for the host's diagnostic use.+data OriginError+    = WildcardOrigin+    | OpaqueOrigin+    | MissingScheme Text+    | EmptyHost Text+    | HasPathQueryOrFragment Text+    | NotAnOrigin Text+    deriving stock (Eq, Show)++{- | Validate configuration, trimming whitespace and tolerating one final slash.+DNS hosts must be ASCII; use the ASCII form of internationalized names.+-}+allowedOrigin :: Text -> Either OriginError AllowedOrigin+allowedOrigin input = parseOrigin (maybe trimmed id (T.stripSuffix "/" trimmed))+  where+    trimmed = T.strip input++renderAllowedOrigin :: AllowedOrigin -> Text+renderAllowedOrigin (AllowedOrigin value) = value++parseOrigin :: Text -> Either OriginError AllowedOrigin+parseOrigin input+    | input == "*" = Left WildcardOrigin+    | input == "null" = Left OpaqueOrigin+    | T.any (\c -> not (isAscii c) || ord c <= 32 || ord c == 127) input = invalid+    | T.null rest = Left (MissingScheme input)+    | scheme /= "http" && scheme /= "https" = invalid+    | T.null authority = Left (EmptyHost input)+    | T.any (`elem` ("/?#" :: String)) authority = Left (HasPathQueryOrFragment input)+    | otherwise = case normalizeAuthority authority of+        Nothing -> invalid+        Just (host, mPort) ->+            let defaultPort = if scheme == "http" then 80 else 443+                suffix = maybe "" (\p -> if p == defaultPort then "" else ":" <> T.pack (show p)) mPort+             in Right (AllowedOrigin (scheme <> "://" <> host <> suffix))+  where+    (rawScheme, rest) = T.breakOn "://" input+    scheme = T.map toLower rawScheme+    authority = T.drop 3 rest+    invalid = Left (NotAnOrigin input)++normalizeAuthority :: Text -> Maybe (Text, Maybe Int)+normalizeAuthority authority = do+    guard (not (T.any (`elem` ("@%*\\" :: String)) authority))+    if T.isPrefixOf "[" authority+        then do+            let (literal, close) = T.breakOn "]" (T.drop 1 authority)+            guard (not (T.null close))+            host <- ipv6 literal+            p <- portSuffix (T.drop 1 close)+            pure ("[" <> host <> "]", p)+        else do+            let (host, suffix) = T.breakOn ":" authority+            guard (validHost host)+            p <- portSuffix suffix+            pure (T.map toLower host, p)++portSuffix :: Text -> Maybe (Maybe Int)+portSuffix "" = Just Nothing+portSuffix suffix = do+    digits <- T.stripPrefix ":" suffix+    n <- decimal digits+    guard (n <= 65535)+    pure (Just (fromInteger n))++decimal :: Text -> Maybe Integer+decimal digits = do+    guard (not (T.null digits) && T.all (\c -> c >= '0' && c <= '9') digits)+    -- Both uses are bounded (ports and IPv4 octets). Stop growing the number+    -- once it cannot be a port, even for a very long untrusted header.+    T.foldl' step (Just 0) digits+  where+    step number c = do+        n <- number+        let next = n * 10 + fromIntegral (ord c - ord '0')+        guard (next <= 65535)+        pure next++validHost :: Text -> Bool+validHost host+    | T.null host || T.length host > 253 = False+    | T.all (\c -> c == '.' || (c >= '0' && c <= '9')) host = isJust (ipv4 host)+    | otherwise = all validLabel (T.splitOn "." host)+  where+    validLabel label =+        not (T.null label)+            && T.length label <= 63+            && T.head label /= '-'+            && T.last label /= '-'+            && T.all (\c -> isAsciiLower c || isAsciiUpper c || (c >= '0' && c <= '9') || c == '-') label++ipv4 :: Text -> Maybe [Int]+ipv4 host = do+    let parts = T.splitOn "." host+    guard (length parts == 4)+    traverse octet parts+  where+    octet part = do+        guard (T.length part <= 3 && (T.length part == 1 || not (T.isPrefixOf "0" part)))+        n <- decimal part+        guard (n <= 255)+        pure (fromInteger n)++-- Normalize IPv6 by expanding the groups; compressed and embedded-IPv4 forms+-- compare by address, without DNS lookup or platform-dependent parsing.+ipv6 :: Text -> Maybe Text+ipv6 literal = do+    expanded <- case T.splitOn "::" literal of+        [whole] -> do+            groups <- side whole+            guard (length groups == 8)+            pure groups+        [left, right] -> do+            -- An embedded IPv4 address must occupy the final two groups.+            guard (not (T.any (== '.') left))+            l <- side left+            r <- side right+            guard (length l + length r < 8)+            pure (l <> replicate (8 - length l - length r) 0 <> r)+        _ -> Nothing+    pure (T.intercalate ":" (map (\n -> T.pack (showHex n "")) expanded))+  where+    side "" = Just []+    side value = do+        let parts = T.splitOn ":" value+        case reverse parts of+            final : remaining | T.any (== '.') final -> do+                octets <- ipv4 final+                case octets of+                    [a, b, c, d] -> do+                        preceding <- traverse hexGroup (reverse remaining)+                        pure (preceding <> [a * 256 + b, c * 256 + d])+                    _ -> Nothing+            _ -> traverse hexGroup parts+    hexGroup value = do+        guard (not (T.null value) && T.length value <= 4 && T.all isHexDigit value)+        pure (T.foldl' (\n c -> n * 16 + digitToInt c) 0 value)++-- | Empty origins disable the middleware. Negative max age is treated as absent.+data CorsPolicy = CorsPolicy+    { allowedOrigins :: ![AllowedOrigin]+    , allowCredentials :: !Bool+    , maxAgeSeconds :: !(Maybe Int)+    }+    deriving stock (Eq, Show)++corsDisabled :: CorsPolicy+corsDisabled = CorsPolicy [] False Nothing++corsAllowOrigins :: [AllowedOrigin] -> CorsPolicy+corsAllowOrigins origins = CorsPolicy origins False Nothing++corsEnabled :: CorsPolicy -> Bool+corsEnabled = not . null . (.allowedOrigins)++corsAllowedMethods :: ByteString+corsAllowedMethods = "GET, HEAD, OPTIONS"++originSet :: CorsPolicy -> Set.Set AllowedOrigin+originSet = Set.fromList . (.allowedOrigins)++-- Request parsing deliberately does not trim or strip a configuration slash.+requestOrigin :: ByteString -> Maybe AllowedOrigin+requestOrigin raw = either (const Nothing) (either (const Nothing) Just . parseOrigin) (TE.decodeUtf8' raw)++originAllowed :: CorsPolicy -> ByteString -> Bool+originAllowed policy raw = maybe False (`Set.member` originSet policy) (requestOrigin raw)++-- Header literals work across the supported http-types range, including older+-- umbrella modules which do not re-export the named constants.+originHeader, varyHeader :: HeaderName+originHeader = "Origin"+varyHeader = "Vary"++headerValues :: HeaderName -> RequestHeaders -> [ByteString]+headerValues name = map snd . filter ((== name) . fst)++isPreflight :: Request -> Bool+isPreflight req =+    requestMethod req == methodOptions+        && not (null (headerValues originHeader (requestHeaders req)))+        && not (null (headerValues "Access-Control-Request-Method" (requestHeaders req)))++{- | Disabled policy returns the application itself. Enabled policy captures one+normalized set, protects upgrades before dispatch, and varies all HTTP responses.+Raw WebSocket responses are preserved by WAI's 'mapResponseHeaders'.+-}+corsMiddleware :: CorsPolicy -> Middleware+corsMiddleware policy+    | not (corsEnabled policy) = id+    | otherwise =+        let !origins = originSet policy+         in \app req respond ->+                let values = headerValues originHeader (requestHeaders req)+                    grant = case values of+                        [raw] | maybe False (`Set.member` origins) (requestOrigin raw) -> Just raw+                        _ -> Nothing+                    varyKeys = if isPreflight req then ["Origin", "Access-Control-Request-Method", "Access-Control-Request-Headers"] else ["Origin"]+                    decorate = mapResponseHeaders (decorateHeaders policy grant varyKeys)+                    reply = respond . decorate+                 in if isWebSocketsReq req+                        then case (values, grant) of+                            ([], _) -> app req reply+                            (_, Just _) -> app req reply+                            _ -> reply (errorResponse status403 "origin_not_allowed" "The request origin is not allowed." Nothing)+                        else case (isPreflight req, grant) of+                            (True, Just _) -> case headerValues "Access-Control-Request-Method" (requestHeaders req) of+                                [method] | method == "GET" || method == "HEAD" ->+                                    case requestedHeaders (requestHeaders req) of+                                        Nothing -> reply (errorResponse status400 "invalid_cors_request" "Requested header names must be HTTP tokens." Nothing)+                                        Just headers ->+                                            reply (responseLBS status204 (preflightHeaders policy headers) "")+                                [_] -> reply (errorResponse status403 "cors_method_not_allowed" "The requested method is not allowed." Nothing)+                                _ -> reply (errorResponse status400 "invalid_cors_request" "A single requested method is required." Nothing)+                            _ -> app req reply++requestedHeaders :: RequestHeaders -> Maybe (Maybe ByteString)+requestedHeaders headers = case headerValues "Access-Control-Request-Headers" headers of+    [] -> Just Nothing+    values -> do+        guard (all (not . BS.null . trimOWS) values)+        let tokens = map trimOWS (concatMap (BS.split ',') values)+        guard (not (null tokens) && all (\t -> not (BS.null t) && BS.all tokenChar t) tokens)+        pure (Just (BS.intercalate ", " tokens))+  where+    tokenChar c = isAsciiLower c || isAsciiUpper c || (c >= '0' && c <= '9') || c `elem` ("!#$%&'*+-.^_`|~" :: String)++trimOWS :: ByteString -> ByteString+trimOWS = BS.dropWhileEnd space . BS.dropWhile space+  where+    space c = c == ' ' || c == '\t'++preflightHeaders :: CorsPolicy -> Maybe ByteString -> ResponseHeaders+preflightHeaders policy headers =+    [("Access-Control-Allow-Methods", corsAllowedMethods)]+        <> maybe [] (\v -> [("Access-Control-Allow-Headers", v)]) headers+        <> case policy.maxAgeSeconds of+            Just n | n >= 0 -> [("Access-Control-Max-Age", BS.pack (show n))]+            _ -> []++decorateHeaders :: CorsPolicy -> Maybe ByteString -> [ByteString] -> ResponseHeaders -> ResponseHeaders+decorateHeaders policy grant keys headers =+    mergeVary keys (filter (not . isGrant . fst) headers)+        <> maybe [] (\raw -> [("Access-Control-Allow-Origin", raw)] <> [("Access-Control-Allow-Credentials", "true") | policy.allowCredentials]) grant+  where+    isGrant name = name == "Access-Control-Allow-Origin" || name == "Access-Control-Allow-Credentials"++mergeVary :: [ByteString] -> ResponseHeaders -> ResponseHeaders+mergeVary keys headers =+    (varyHeader, BS.intercalate ", " tokens) : filter ((/= varyHeader) . fst) headers+  where+    existing = filter (not . BS.null) (map trimOWS (concatMap (BS.split ',' . snd) (filter ((== varyHeader) . fst) headers)))+    tokens+        | "*" `elem` existing = ["*"]+        | otherwise = nubBy (\a b -> BS.map toLower a == BS.map toLower b) (existing <> keys)
+ src/Kiroku/Metrics/DeadLetters.hs view
@@ -0,0 +1,127 @@+-- | Read-only subscription dead-letter inspection with exclusive keyset pages.+module Kiroku.Metrics.DeadLetters (+    DeadLetterProvider,+    storeDeadLetters,+    DeadLetterItem (..),+    DeadLetterPageResponse (..),+    deadLetterPageResponse,+    renderDeadLetterCursor,+    parseDeadLetterCursor,+    DeadLetterRequest (..),+    parseDeadLetterRequest,+    deadLettersApp,+    deadLettersNotConfiguredApp,+) where++import Data.Aeson (FromJSON (..), ToJSON (..), Value, encode, object, withObject, (.:), (.:?), (.=))+import Data.ByteString qualified as BS+import Data.Int (Int32, Int64)+import Data.Text (Text)+import Data.Text qualified as T+import Data.Text.Encoding qualified as TE+import Data.Time.Clock (UTCTime)+import Data.UUID (UUID)+import Data.Vector qualified as V+import Kiroku.Metrics.JSON (errorResponse, jsonResponse, storeErrorResponse)+import Kiroku.Store qualified as Store+import Network.HTTP.Types (Query, Status, status200, status400, status404, status405)+import Network.Wai (Application, Response, mapResponseHeaders, pathInfo, queryString, requestMethod, responseHeaders, responseLBS, responseStatus)+import Text.Read (readMaybe)++type DeadLetterProvider = Store.SubscriptionDeadLetterQuery -> IO (Either Store.StoreError Store.SubscriptionDeadLetterPage)+storeDeadLetters :: Store.KirokuStore -> DeadLetterProvider+storeDeadLetters store = Store.runStoreIO store . Store.subscriptionDeadLetters++data DeadLetterItem = DeadLetterItem+    { deadLetterId :: !Int64+    , subscription :: !Text+    , member :: !Int32+    , globalPosition :: !Int64+    , eventId :: !UUID+    , reason :: !Value+    , reasonSummary :: !Text+    , attemptCount :: !Int32+    , createdAt :: !UTCTime+    }+    deriving stock (Eq, Show)+data DeadLetterPageResponse = DeadLetterPageResponse+    { items :: ![DeadLetterItem]+    , nextCursor :: !(Maybe Text)+    }+    deriving stock (Eq, Show)+instance ToJSON DeadLetterItem where+    toJSON row = object ["dead_letter_id" .= row.deadLetterId, "subscription" .= row.subscription, "member" .= row.member, "global_position" .= row.globalPosition, "event_id" .= row.eventId, "reason" .= row.reason, "reason_summary" .= row.reasonSummary, "attempt_count" .= row.attemptCount, "created_at" .= row.createdAt]+instance FromJSON DeadLetterItem where+    parseJSON = withObject "DeadLetterItem" $ \o -> DeadLetterItem <$> o .: "dead_letter_id" <*> o .: "subscription" <*> o .: "member" <*> o .: "global_position" <*> o .: "event_id" <*> o .: "reason" <*> o .: "reason_summary" <*> o .: "attempt_count" <*> o .: "created_at"+instance ToJSON DeadLetterPageResponse where+    toJSON page = object (["items" .= page.items] <> maybe [] (\cursor -> ["next_cursor" .= cursor]) page.nextCursor)+instance FromJSON DeadLetterPageResponse where+    parseJSON = withObject "DeadLetterPageResponse" $ \o -> DeadLetterPageResponse <$> o .: "items" <*> o .:? "next_cursor"+deadLetterPageResponse :: Store.SubscriptionDeadLetterPage -> DeadLetterPageResponse+deadLetterPageResponse (Store.SubscriptionDeadLetterPage rows cursor) = DeadLetterPageResponse (map item (V.toList rows)) (renderDeadLetterCursor <$> cursor)+  where+    item (Store.SubscriptionDeadLetter ident (Store.SubscriptionName name) member (Store.GlobalPosition position) (Store.EventId eid) reason summary attempts date) = DeadLetterItem ident name member position eid reason summary attempts date++renderDeadLetterCursor :: Store.SubscriptionDeadLetterCursor -> Text+renderDeadLetterCursor (Store.SubscriptionDeadLetterCursor (Store.GlobalPosition position) ident) = T.pack (show position) <> ":" <> T.pack (show ident)+parseDeadLetterCursor :: Text -> Maybe Store.SubscriptionDeadLetterCursor+parseDeadLetterCursor raw = case T.splitOn ":" raw of+    [p, i] -> Store.SubscriptionDeadLetterCursor . Store.GlobalPosition <$> unsigned p <*> unsigned i+    _ -> Nothing+unsigned :: forall a. (Integral a, Bounded a) => Text -> Maybe a+unsigned raw+    | T.null raw || T.length raw > 20 || not (T.all (\c -> c >= '0' && c <= '9') raw) = Nothing+    | otherwise = do+        value <- readMaybe (T.unpack raw) :: Maybe Integer+        if value <= toInteger (maxBound @a) then Just (fromInteger value) else Nothing++data DeadLetterRequest = DeadLetterRequest+    { requestMember :: !(Maybe Int32)+    , requestAfter :: !(Maybe Store.SubscriptionDeadLetterCursor)+    , requestLimit :: !Store.SubscriptionDeadLetterLimit+    }+    deriving stock (Eq, Show)+type Failure = (Status, Text, Text, Maybe Value)+invalid :: Text -> Text -> Text -> Failure+invalid parameter value reason = (status400, "invalid_query_parameter", "Invalid query parameter.", Just (object ["parameter" .= parameter, "value" .= value, "reason" .= reason]))+parseDeadLetterRequest :: Query -> Either Failure DeadLetterRequest+parseDeadLetterRequest query = do+    member <- parameter "member" unsigned "Expected a non-negative Int32 decimal integer."+    cursor <- parameter "from" parseDeadLetterCursor "Expected an opaque position:id cursor with non-negative Int64 components."+    size <- parameter "limit" validLimit "Expected a decimal page size from 1 through 1000."+    pure (DeadLetterRequest member cursor (maybe Store.defaultSubscriptionDeadLetterLimit id size))+  where+    validLimit raw = unsigned raw >>= either (const Nothing) Just . Store.mkSubscriptionDeadLetterLimit+    parameter :: Text -> (Text -> Maybe a) -> Text -> Either Failure (Maybe a)+    parameter key parser reason = case [v | (k, v) <- query, k == TE.encodeUtf8 key] of+        [] -> Right Nothing+        [Just bytes] -> case TE.decodeUtf8' bytes of+            Left _ -> Left (invalid key "<invalid UTF-8>" "Expected valid UTF-8.")+            Right raw -> maybe (Left (invalid key raw reason)) (Right . Just) (parser raw)+        [Nothing] -> Left (invalid key "" reason)+        _ -> Left (invalid key "<duplicate>" "Parameter must occur at most once.")++-- | Applies GET/HEAD/405 at WAI level, including when mounted without Warp.+readMethods :: IO Response -> Application+readMethods action req respond+    | requestMethod req == "GET" = action >>= respond+    | requestMethod req == "HEAD" = action >>= \response -> respond (responseLBS (responseStatus response) (responseHeaders response) "")+    | otherwise = respond $ mapResponseHeaders (("Allow", "GET, HEAD") :) $ errorResponse status405 "method_not_allowed" "Only GET and HEAD are supported." Nothing++deadLettersNotConfiguredApp :: Application+deadLettersNotConfiguredApp req = readMethods (pure (errorResponse status404 "dead_letters_not_configured" "Dead letters are not configured." Nothing)) req++deadLettersApp :: DeadLetterProvider -> Application+deadLettersApp provider req = readMethods action req+  where+    action = case pathInfo req of+        ["subscriptions", name, "dead-letters"]+            | T.any (== '\0') name || BS.length (TE.encodeUtf8 name) > 512 -> pure $ failure (invalid "subscription" name "Expected at most 512 UTF-8 bytes without NUL.")+            | otherwise -> case parseDeadLetterRequest (queryString req) of+                Left err -> pure (failure err)+                Right (DeadLetterRequest member cursor size) ->+                    provider (Store.SubscriptionDeadLetterQuery (Store.SubscriptionName name) member cursor size) >>= \case+                        Left err -> pure (storeErrorResponse "dead_letters_unavailable" err)+                        Right page -> pure (jsonResponse status200 (encode (deadLetterPageResponse page)))+        _ -> pure (errorResponse status404 "not_found" "Not found" Nothing)+    failure (status, code, message, details) = errorResponse status code message details
src/Kiroku/Metrics/JSON.hs view
@@ -4,17 +4,21 @@ module Kiroku.Metrics.JSON (     jsonApp,     jsonResponse,+    errorEnvelope,+    errorResponse,+    storeErrorResponse, ) where -import Data.Aeson (encode, object, (.=))+import Data.Aeson (Value, encode, object, (.=)) import Data.ByteString.Lazy qualified as LBS import Data.Map.Strict qualified as Map import Data.Text (Text)-import Network.HTTP.Types (Status, hContentType, status200, status404)+import Network.HTTP.Types (Status, hContentType, status200, status404, status500, status503) import Network.Wai (Application, Response, pathInfo, responseLBS)  import Kiroku.Metrics.Collector (KirokuMetrics, snapshotMetrics) import Kiroku.Metrics.Types (MetricsSnapshot (..))+import Kiroku.Store.Error (StoreError (..))  {- | WAI application for the JSON metrics endpoints. Routes @/metrics@ and @/metrics/\<name\>@; any other path returns a 404 JSON body.@@ -42,3 +46,21 @@ -- | Build an @application/json@ response with the given status and body. jsonResponse :: Status -> LBS.ByteString -> Response jsonResponse status = responseLBS status [(hContentType, "application/json")]++-- | Shared structured error for new inspection routes; legacy errors stay unchanged.+errorEnvelope :: Text -> Text -> Maybe Value -> Value+errorEnvelope code message details =+    object ["error" .= object (["code" .= code, "message" .= message] <> maybe [] (\v -> ["details" .= v]) details)]++-- | Build a JSON error, omitting @details@ when absent.+errorResponse :: Status -> Text -> Text -> Maybe Value -> Response+errorResponse status code message = jsonResponse status . encode . errorEnvelope code message++{- | Sanitize store failures without exposing connection strings or event payloads.+This function handles values only and never catches asynchronous exceptions.+-}+storeErrorResponse :: Text -> StoreError -> Response+storeErrorResponse unavailableCode = \case+    ConnectionError _ -> errorResponse status503 unavailableCode "The event store is unavailable." Nothing+    EventDecodeFailed _ -> errorResponse status500 "event_decode_failed" "An event could not be decoded." Nothing+    _ -> errorResponse status500 "store_error" "The event store operation failed." Nothing
src/Kiroku/Metrics/Server.hs view
@@ -1,17 +1,18 @@-{- | The combined metrics web server.--Builds a single WAI 'Application' with 'WaiWS.websocketsOr': WebSocket upgrades-go to a 'WS.ServerApp' seam, everything else to the HTTP router. EP-2 supplies a-rejecting stub for the seam ('stubWebSocketApp'); EP-3 replaces it with the real-event-streaming app via 'startMetricsServerWith' without changing this module.--The server takes 'KirokuMetrics' plus a list of 'DependencyCheck's and a-'WS.ServerApp' — it does /not/ take the 'KirokuStore' directly. Everything-store-specific is captured in caller-built closures ('postgresPing' over the-pool, and EP-3's WebSocket app over the store), keeping the server store-agnostic.+{- | The shared inspection composition. Store-backed behavior enters through+provider closures; HTTP and WebSocket dispatch share a mount-relative path and+one outer CORS policy. /capabilities reports this wiring without store access.+Legacy starters retain their signatures. -} module Kiroku.Metrics.Server (     MetricsServer (..),+    ServerProviders (..),+    defaultServerProviders,+    providerPresence,+    storeServerProviders,+    startMetricsServerWithProviders,+    withMetricsServerWithProviders,+    combinedAppWithProviders,+    httpAppWithProviders,     startMetricsServer,     startMetricsServerWith,     startMetricsServerWith',@@ -25,18 +26,29 @@     stubWebSocketApp, ) where -import Control.Concurrent.Async (Async, async, cancel)-import Control.Exception (bracket)+import Control.Concurrent.Async (Async, asyncWithUnmask, cancel, race, wait, waitCatchSTM)+import Control.Concurrent.STM (atomically, newEmptyTMVarIO, orElse, putTMVar, readTMVar)+import Control.Exception (bracket, finally, mask, onException, throwIO) import Data.Aeson (encode, object, (.=))+import Data.ByteString.Builder (toLazyByteString)+import Data.ByteString.Lazy qualified as LBS+import Data.Maybe (isJust) import Data.Text (Text) import Network.HTTP.Types (status200, status404, status503)-import Network.Wai (Application, pathInfo)+import Network.HTTP.Types.URI (encodePathSegments)+import Network.Socket qualified as Socket+import Network.Wai (Application, pathInfo, rawPathInfo) import Network.Wai.Handler.Warp qualified as Warp import Network.Wai.Handler.WebSockets qualified as WaiWS import Network.WebSockets qualified as WS +import Kiroku.Metrics.Browse (StoreBrowser, browseApp, browseNotConfiguredApp, storeBrowser)+import Kiroku.Metrics.Capabilities (ProviderPresence (..), WebSocketChannels, capabilitiesApp, capabilitiesFor, noWebSocketChannels, storeWebSocketChannels)+import Kiroku.Metrics.Checkpoints (CheckpointInventoryProvider, checkpointsApp, checkpointsNotConfiguredApp, storeCheckpointInventory) import Kiroku.Metrics.Collector (KirokuMetrics) import Kiroku.Metrics.Config (MetricsServerConfig (..))+import Kiroku.Metrics.Cors (corsMiddleware)+import Kiroku.Metrics.DeadLetters (DeadLetterProvider, deadLettersApp, deadLettersNotConfiguredApp, storeDeadLetters) import Kiroku.Metrics.Health (     DependencyCheck,     LivenessStatus (..),@@ -47,7 +59,7 @@  ) import Kiroku.Metrics.JSON (jsonApp, jsonResponse) import Kiroku.Metrics.Prometheus (prometheusApp)-import Kiroku.Metrics.Subscriptions (SubscriptionStatusProvider, subscriptionsApp)+import Kiroku.Metrics.Subscriptions (SubscriptionStatusProvider, storeSubscriptionStatus, subscriptionsApp) import Kiroku.Metrics.WebSocket (newWebSocketState, websocketApp) import Kiroku.Store (KirokuStore) @@ -57,166 +69,182 @@     , serverPort :: !Int     } -{- | Start the server with the rejecting WebSocket stub. Use this until EP-3's-real WebSocket app is wired via 'startMetricsServerWith'.--}-startMetricsServer :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> IO MetricsServer-startMetricsServer cfg m deps = startMetricsServerWith' cfg m deps Nothing stubWebSocketApp+-- | Optional data sources for the shared inspection application.+data ServerProviders = ServerProviders+    { webSocketServer :: !WS.ServerApp+    , subscriptionStatus :: !(Maybe SubscriptionStatusProvider)+    , checkpointInventory :: !(Maybe CheckpointInventoryProvider)+    , storeBrowsing :: !(Maybe StoreBrowser)+    -- ^ Backs the stream, category and event inspection routes.+    , deadLetters :: !(Maybe DeadLetterProvider)+    -- ^ Backs GET /subscriptions/<name>/dead-letters.+    , webSocketChannels :: !WebSocketChannels+    -- ^ Channels served by the opaque WebSocket app, declared for /capabilities.+    } -{- | Start the server with an explicit WebSocket app (the IP-3 seam) and no-subscription-status provider. When @cfg.port == 0@ an OS-assigned free port is-used and reported in 'serverPort'.+-- | Reject upgrades and leave all optional providers unconfigured.+defaultServerProviders :: ServerProviders+defaultServerProviders = ServerProviders stubWebSocketApp Nothing Nothing Nothing Nothing noWebSocketChannels++{- | Build every store-backed provider, including the process-local live registry.+Bind this action first, then use 'withMetricsServerWithProviders'. -}-startMetricsServerWith ::-    MetricsServerConfig ->-    KirokuMetrics ->-    [DependencyCheck] ->-    WS.ServerApp ->-    IO MetricsServer-startMetricsServerWith cfg m deps = startMetricsServerWith' cfg m deps Nothing+storeServerProviders :: MetricsServerConfig -> KirokuMetrics -> KirokuStore -> IO ServerProviders+storeServerProviders cfg m store = do+    wsState <- newWebSocketState cfg.wsMaxConnections+    pure $ ServerProviders (websocketApp cfg m store wsState) (Just (storeSubscriptionStatus store)) (Just (storeCheckpointInventory store)) (Just (storeBrowser store)) (Just (storeDeadLetters store)) storeWebSocketChannels -{- | Start the server with an explicit WebSocket app (the IP-3 seam) /and/ an-optional subscription-status provider (the IP-5 seam, EP-5). The provider, when-@Just@, serves @GET /subscriptions@; when @Nothing@, that route returns a-configured-404. All EP-2/EP-3 starters delegate here with @Nothing@.+{- | Return only after Warp is ready; bind/setup failures are rethrown.+Ephemeral sockets are explicitly closed on every exit, including cancellation. -}-startMetricsServerWith' ::-    MetricsServerConfig ->-    KirokuMetrics ->-    [DependencyCheck] ->-    Maybe SubscriptionStatusProvider ->-    WS.ServerApp ->-    IO MetricsServer-startMetricsServerWith' cfg m deps mProvider wsApp = do-    let app = combinedApp cfg m deps mProvider wsApp+startMetricsServerWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> IO MetricsServer+startMetricsServerWithProviders cfg m deps providers = mask $ \restore -> do+    ready <- newEmptyTMVarIO+    let app = combinedAppWithProviders cfg m deps providers+        settings port = Warp.setBeforeMainLoop (atomically $ putTMVar ready ()) $ Warp.setHost "*" $ Warp.setPort port Warp.defaultSettings+        await thread port = do+            result <- restore (atomically $ (Left <$> waitCatchSTM thread) `orElse` (Right <$> readTMVar ready)) `onException` cancel thread+            case result of+                Left (Left err) -> throwIO err+                Left (Right ()) -> fail "Metrics server terminated before readiness."+                Right () -> pure (MetricsServer thread port)     if cfg.port == 0         then do-            (actualPort, sock) <- Warp.openFreePort-            let settings = Warp.setPort actualPort Warp.defaultSettings-            thread <- async (Warp.runSettingsSocket settings sock app)-            pure (MetricsServer thread actualPort)+            (port, sock) <- Warp.openFreePort+            thread <- asyncWithUnmask (\unmask -> unmask (Warp.runSettingsSocket (settings port) sock app) `finally` Socket.close sock) `onException` Socket.close sock+            await thread port         else do-            let settings = Warp.setHost "*" (Warp.setPort cfg.port Warp.defaultSettings)-            thread <- async (Warp.runSettings settings app)-            pure (MetricsServer thread cfg.port)+            thread <- asyncWithUnmask (\unmask -> unmask $ Warp.runSettings (settings cfg.port) app)+            await thread cfg.port -{- | Start the server with the real WebSocket app (EP-3), which streams live-metrics and events out of the given 'KirokuStore'. This is the recommended entry-point once event streaming is wanted: it allocates one shared connection-limiting-state (bounded by @cfg.wsMaxConnections@) and wires-'Kiroku.Metrics.WebSocket.websocketApp'. EP-2's 'startMetricsServer' (stub) is-unchanged for callers who do not want the WebSocket.+{- | Supervise the callback and the server together, then release both.+Unexpected server termination cancels the callback and is rethrown to the owner. -}-startMetricsServerWithStore ::-    MetricsServerConfig ->-    KirokuMetrics ->-    KirokuStore ->-    [DependencyCheck] ->-    IO MetricsServer+withMetricsServerWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> (MetricsServer -> IO a) -> IO a+withMetricsServerWithProviders cfg m deps providers action =+    withRunningServer (startMetricsServerWithProviders cfg m deps providers) action++withRunningServer :: IO MetricsServer -> (MetricsServer -> IO a) -> IO a+withRunningServer acquire action = bracket acquire stopMetricsServer $ \server -> do+    result <- race (wait server.serverThread) (action server)+    either (\() -> fail "Metrics server terminated unexpectedly.") pure result++-- | Start with the rejecting WebSocket stub and no optional providers.+startMetricsServer :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> IO MetricsServer+startMetricsServer cfg m deps = startMetricsServerWithProviders cfg m deps defaultServerProviders++-- | Start with a caller-supplied WebSocket app and no optional providers.+startMetricsServerWith :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> WS.ServerApp -> IO MetricsServer+startMetricsServerWith cfg m deps = startMetricsServerWith' cfg m deps Nothing++-- | Legacy binding of a WebSocket app and optional live status provider.+startMetricsServerWith' :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> Maybe SubscriptionStatusProvider -> WS.ServerApp -> IO MetricsServer+startMetricsServerWith' cfg m deps mProvider wsApp =+    startMetricsServerWithProviders cfg m deps defaultServerProviders{webSocketServer = wsApp, subscriptionStatus = mProvider}++{- | Serve the real WebSocket and durable inventory; the legacy live route stays+unconfigured. New hosts wanting every provider use 'storeServerProviders'.+-}+startMetricsServerWithStore :: MetricsServerConfig -> KirokuMetrics -> KirokuStore -> [DependencyCheck] -> IO MetricsServer startMetricsServerWithStore cfg m store deps = do-    wsState <- newWebSocketState cfg.wsMaxConnections-    startMetricsServerWith cfg m deps (websocketApp cfg m store wsState)+    providers <- storeServerProviders cfg m store+    startMetricsServerWithProviders cfg m deps providers{subscriptionStatus = Nothing} --- | Stop the server by cancelling its Warp thread. stopMetricsServer :: MetricsServer -> IO () stopMetricsServer server = cancel server.serverThread --- | Run an action with a running server, tearing it down afterwards.-withMetricsServer ::-    MetricsServerConfig ->-    KirokuMetrics ->-    [DependencyCheck] ->-    (MetricsServer -> IO a) ->-    IO a-withMetricsServer cfg m deps =-    bracket (startMetricsServer cfg m deps) stopMetricsServer+withMetricsServer :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> (MetricsServer -> IO a) -> IO a+withMetricsServer cfg m deps = withMetricsServerWithProviders cfg m deps defaultServerProviders -{- | Run an action with a running store-aware server (EP-3 WebSocket), tearing-it down afterwards.--}-withMetricsServerWithStore ::-    MetricsServerConfig ->-    KirokuMetrics ->-    KirokuStore ->-    [DependencyCheck] ->-    (MetricsServer -> IO a) ->-    IO a-withMetricsServerWithStore cfg m store deps =-    bracket (startMetricsServerWithStore cfg m store deps) stopMetricsServer+withMetricsServerWithStore :: MetricsServerConfig -> KirokuMetrics -> KirokuStore -> [DependencyCheck] -> (MetricsServer -> IO a) -> IO a+withMetricsServerWithStore cfg m store deps = withRunningServer (startMetricsServerWithStore cfg m store deps) -{- | Run an action with a server that serves @GET /subscriptions@ from the given-provider (EP-5), using the rejecting WebSocket stub. The common case for a worker-that wants remote subscription introspection but not the event-streaming socket.--}-withMetricsServerSubscriptions ::-    MetricsServerConfig ->-    KirokuMetrics ->-    [DependencyCheck] ->-    SubscriptionStatusProvider ->-    (MetricsServer -> IO a) ->-    IO a+withMetricsServerSubscriptions :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> SubscriptionStatusProvider -> (MetricsServer -> IO a) -> IO a withMetricsServerSubscriptions cfg m deps provider =-    bracket-        (startMetricsServerWith' cfg m deps (Just provider) stubWebSocketApp)-        stopMetricsServer+    withMetricsServerWithProviders cfg m deps defaultServerProviders{subscriptionStatus = Just provider} --- | The combined WAI app: WebSocket upgrades to @wsApp@, everything else to the HTTP router.-combinedApp ::-    MetricsServerConfig ->-    KirokuMetrics ->-    [DependencyCheck] ->-    Maybe SubscriptionStatusProvider ->-    WS.ServerApp ->-    Application+-- | Legacy composition binding. CORS is applied once by the general composition.+combinedApp :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> Maybe SubscriptionStatusProvider -> WS.ServerApp -> Application combinedApp cfg m deps mProvider wsApp =-    WaiWS.websocketsOr WS.defaultConnectionOptions wsApp (httpApp cfg m deps mProvider)+    combinedAppWithProviders cfg m deps defaultServerProviders{webSocketServer = wsApp, subscriptionStatus = mProvider} -{- | The EP-2 WebSocket stub: reject the upgrade with a clear message. EP-3-replaces this with the real event-streaming app.+{- | Mountable composition. Respect the WebSocket switch before upgrade dispatch.+Normalize only the dispatch copy's raw path; keep its original query string. -}+combinedAppWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> Application+combinedAppWithProviders cfg m deps providers = corsMiddleware cfg.cors dispatch+  where+    http = httpAppWithProviders cfg m deps providers+    dispatch req respond+        | cfg.enableWebSocket =+            let relativePath = LBS.toStrict $ toLazyByteString $ encodePathSegments (pathInfo req)+             in WaiWS.websocketsOr WS.defaultConnectionOptions providers.webSocketServer http (req{rawPathInfo = relativePath}) respond+        | otherwise = http req respond+ stubWebSocketApp :: WS.ServerApp stubWebSocketApp pending = WS.rejectRequest pending "WebSocket endpoint not yet implemented" --- | HTTP router. Matches @/metrics/prometheus@ before @/metrics/\<name\>@.-httpApp ::-    MetricsServerConfig ->-    KirokuMetrics ->-    [DependencyCheck] ->-    Maybe SubscriptionStatusProvider ->-    Application-httpApp cfg m deps mProvider req respond =-    case pathInfo req of-        ["metrics", "prometheus"] | cfg.enablePrometheus -> prometheusApp m req respond-        ["metrics"] | cfg.enableJSON -> jsonApp m req respond-        ["metrics", _] | cfg.enableJSON -> jsonApp m req respond-        ["subscriptions"] -> subscriptionsRoute-        ["subscriptions", _] -> subscriptionsRoute-        ["health"] | cfg.enableJSON -> do-            (readiness, snap) <- checkDetailedHealth cfg m deps-            respond $-                jsonResponse-                    (statusFor readiness.ready)-                    (encode (object ["status" .= readiness, "metrics" .= snap]))-        ["health", "live"] | cfg.enableJSON -> do-            liveness <- checkLiveness cfg m-            respond (jsonResponse (statusFor liveness.alive) (encode liveness))-        ["health", "ready"] | cfg.enableJSON -> do-            readiness <- checkReadiness cfg m deps-            respond (jsonResponse (statusFor readiness.ready) (encode readiness))-        ["ws"]-            | cfg.enableWebSocket ->+-- | Legacy unwrapped HTTP binding.+httpApp :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> Maybe SubscriptionStatusProvider -> Application+httpApp cfg m deps mProvider = httpAppWithProviders cfg m deps defaultServerProviders{subscriptionStatus = mProvider}++-- | Unwrapped HTTP router, matching paths relative to the host's mount.+httpAppWithProviders :: MetricsServerConfig -> KirokuMetrics -> [DependencyCheck] -> ServerProviders -> Application+httpAppWithProviders cfg m deps providers = dispatch+  where+    discovery = capabilitiesApp (capabilitiesFor cfg (providerPresence providers))+    dispatch req respond =+        case pathInfo req of+            ["capabilities"] -> discovery req respond+            ["metrics", "prometheus"] | cfg.enablePrometheus -> prometheusApp m req respond+            ["metrics"] | cfg.enableJSON -> jsonApp m req respond+            ["metrics", _] | cfg.enableJSON -> jsonApp m req respond+            prefix : _ | prefix `elem` ["streams", "categories", "events"] -> browseRoute+            ["subscription-checkpoints"] -> checkpointsRoute+            ["subscriptions", _, "dead-letters"] -> deadLettersRoute+            ["subscriptions"] -> subscriptionsRoute+            ["subscriptions", _] -> subscriptionsRoute+            ["health"] | cfg.enableJSON -> do+                (readiness, snap) <- checkDetailedHealth cfg m deps                 respond $                     jsonResponse+                        (statusFor readiness.ready)+                        (encode (object ["status" .= readiness, "metrics" .= snap]))+            ["health", "live"] | cfg.enableJSON -> do+                liveness <- checkLiveness cfg m+                respond (jsonResponse (statusFor liveness.alive) (encode liveness))+            ["health", "ready"] | cfg.enableJSON -> do+                readiness <- checkReadiness cfg m deps+                respond (jsonResponse (statusFor readiness.ready) (encode readiness))+            ["ws"]+                | cfg.enableWebSocket ->+                    respond $+                        jsonResponse+                            status404+                            (encode (object ["error" .= ("WebSocket endpoint - use ws:// protocol" :: Text)]))+            _ ->+                respond (jsonResponse status404 (encode (object ["error" .= ("Not found" :: Text)])))+      where+        statusFor ok = if ok then status200 else status503+        deadLettersRoute = maybe deadLettersNotConfiguredApp deadLettersApp providers.deadLetters req respond+        browseRoute = maybe browseNotConfiguredApp browseApp providers.storeBrowsing req respond+        checkpointsRoute = case providers.checkpointInventory of+            Just provider -> checkpointsApp provider req respond+            Nothing -> checkpointsNotConfiguredApp req respond+        subscriptionsRoute = case providers.subscriptionStatus of+            Just provider -> subscriptionsApp provider req respond+            Nothing ->+                respond $+                    jsonResponse                         status404-                        (encode (object ["error" .= ("WebSocket endpoint - use ws:// protocol" :: Text)]))-        _ ->-            respond (jsonResponse status404 (encode (object ["error" .= ("Not found" :: Text)])))-  where-    statusFor ok = if ok then status200 else status503-    subscriptionsRoute = case mProvider of-        Just provider -> subscriptionsApp provider req respond-        Nothing ->-            respond $-                jsonResponse-                    status404-                    (encode (object ["error" .= ("subscription status not configured" :: Text)]))+                        (encode (object ["error" .= ("subscription status not configured" :: Text)]))++-- | Pure wiring summary; evaluating it never invokes a provider.+providerPresence :: ServerProviders -> ProviderPresence+providerPresence providers =+    ProviderPresence+        (isJust providers.subscriptionStatus)+        (isJust providers.checkpointInventory)+        (isJust providers.storeBrowsing)+        (isJust providers.deadLetters)+        providers.webSocketChannels
+ src/Kiroku/Metrics/Standalone.hs view
@@ -0,0 +1,192 @@+{-# LANGUAGE NoFieldSelectors #-}++-- | Open a store and serve its inspection surface without a Haskell host program.+module Kiroku.Metrics.Standalone (+    InspectOptions (..),+    inspectOptionsParser,+    inspectParserInfo,+    InspectRuntime (..),+    resolveInspectOptions,+    InspectHooks (..),+    runInspect,+    renderStartupBanner,+) where++import Control.Applicative (many, optional, (<|>))+import Control.Concurrent.STM (STM, TVar, atomically, newTVarIO, readTVar, writeTVar)+import Data.IntMap.Strict qualified as IntMap+import Data.Text (Text)+import Data.Text qualified as T+import Options.Applicative qualified as O++import Kiroku.Metrics.Capabilities+import Kiroku.Metrics.Collector+import Kiroku.Metrics.Config (defaultConfig)+import Kiroku.Metrics.Config qualified as Config+import Kiroku.Metrics.Cors+import Kiroku.Metrics.Health (postgresPing)+import Kiroku.Metrics.Server+import Kiroku.Store qualified as Store+import Kiroku.Store.Subscription.EventPublisher (EventPublisher (..), publisherPosition)++-- No Show instances: a database URL can contain credentials.+data InspectOptions = InspectOptions+    { databaseUrl :: !(Maybe Text)+    , schema :: !(Maybe Text)+    , poolSize :: !(Maybe Int)+    , port :: !(Maybe Int)+    , corsOrigins :: ![Text]+    , corsAllowCredentials :: !(Maybe Bool)+    , wsMaxConnections :: !(Maybe Int)+    }+    deriving stock (Eq)++data InspectRuntime = InspectRuntime+    { databaseUrl :: !Text+    , schema :: !Text+    , poolSize :: !Int+    , port :: !Int+    , cors :: !CorsPolicy+    , wsMaxConnections :: !Int+    }+    deriving stock (Eq)++data InspectHooks = InspectHooks+    { onListening :: !(Int -> Capabilities -> IO ())+    -- ^ Called once after successful binding, with the actual port and wiring.+    , waitForShutdown :: !(IO ())+    }++inspectOptionsParser :: O.Parser InspectOptions+inspectOptionsParser =+    InspectOptions+        <$> optional (textOption "database-url" "URL" "Database connection string (or DATABASE_URL)")+        <*> optional (textOption "schema" "NAME" "Migrated schema (default kiroku)")+        <*> optional (numberOption "pool-size" 1 intMaximum "Connection pool size (default 10)")+        <*> optional (numberOption "port" 0 65535 "Listen port (default 9091; 0 selects a free port)")+        <*> many (textOption "cors-origin" "ORIGIN" "Allowed HTTP(S) browser origin; repeatable")+        <*> optional+            ( O.flag' True (O.long "cors-allow-credentials" <> O.help "Allow credentials for explicit origins")+                <|> O.flag' False (O.long "no-cors-allow-credentials" <> O.help "Disable credentials, overriding the environment")+            )+        <*> optional (numberOption "ws-max-connections" 1 intMaximum "Maximum WebSocket connections (default 100)")+  where+    textOption name metavar help = O.strOption (O.long name <> O.metavar metavar <> O.help help)+    numberOption name lower upper help =+        O.option+            (O.eitherReader (either (Left . T.unpack) Right . boundedDecimal lower upper . T.pack))+            (O.long name <> O.metavar "N" <> O.help help)++inspectParserInfo :: O.ParserInfo InspectOptions+inspectParserInfo =+    O.info+        (inspectOptionsParser O.<**> O.helper)+        ( O.failureCode 2+            <> O.fullDesc+            <> O.header "kiroku-inspect - standalone inspection server"+            <> O.progDesc "Serve the read-only inspection API from a migrated Kiroku database. Runs no subscriptions."+        )++intMaximum :: Integer+intMaximum = toInteger (maxBound :: Int)++-- Parse into Integer and stop growing at the upper bound, before narrowing to Int.+boundedDecimal :: Integer -> Integer -> Text -> Either Text Int+boundedDecimal lower upper value+    | T.null value || not (T.all (\c -> c >= '0' && c <= '9') value) = Left "expected ASCII decimal digits"+    | otherwise = do+        n <- T.foldl' step (Right 0) value+        if n < lower then Left "value is below the allowed range" else Right (fromInteger n)+  where+    step result c = do+        n <- result+        let next = n * 10 + toInteger (fromEnum c - fromEnum '0')+        if next > upper then Left "value exceeds the allowed range" else Right next++resolveInspectOptions :: [(String, String)] -> InspectOptions -> Either Text InspectRuntime+resolveInspectOptions env opts = do+    url <- maybe (Left "kiroku-inspect: no database; pass --database-url or set DATABASE_URL (a libpq URI such as postgresql://user@host/db)") Right (opts.databaseUrl <|> variable "DATABASE_URL")+    nonempty "--database-url" url+    let schemaName = maybe "kiroku" id (opts.schema <|> variable "KIROKU_INSPECT_SCHEMA")+    nonempty "--schema" schemaName+    pool <- number opts.poolSize "KIROKU_INSPECT_POOL_SIZE" 1 intMaximum (Store.defaultConnectionSettings url).poolSize+    listenPort <- number opts.port "KIROKU_INSPECT_PORT" 0 65535 defaultConfig.port+    maxConnections <- number opts.wsMaxConnections "KIROKU_INSPECT_WS_MAX_CONNECTIONS" 1 intMaximum defaultConfig.wsMaxConnections+    credentials <- case opts.corsAllowCredentials of+        Just b -> Right b+        Nothing -> case variable "KIROKU_INSPECT_CORS_ALLOW_CREDENTIALS" of+            Nothing -> Right False+            Just "true" -> Right True+            Just "1" -> Right True+            Just "false" -> Right False+            Just "0" -> Right False+            Just _ -> Left "kiroku-inspect: KIROKU_INSPECT_CORS_ALLOW_CREDENTIALS must be true, false, 1 or 0"+    let originTexts = if null opts.corsOrigins then maybe [] (map T.strip . T.splitOn ",") (variable "KIROKU_INSPECT_CORS_ORIGINS") else opts.corsOrigins+    origins <-+        traverse+            ( \origin -> case allowedOrigin origin of+                Right o -> Right o+                Left WildcardOrigin -> Left "kiroku-inspect: wildcard '*' is not an allowed CORS origin"+                Left _ -> Left "kiroku-inspect: invalid CORS origin; use an explicit HTTP(S) origin"+            )+            originTexts+    let policy = if null origins then corsDisabled else (corsAllowOrigins origins){allowCredentials = credentials}+    pure (InspectRuntime url schemaName pool listenPort policy maxConnections)+  where+    variable name = case T.pack <$> lookup name env of+        Just value | not (T.null value) -> Just value+        _ -> Nothing+    nonempty name value = if T.null (T.strip value) then Left ("kiroku-inspect: " <> name <> " must not be empty") else Right ()+    number explicit name lower upper fallback = case explicit of+        Just value | toInteger value >= lower && toInteger value <= upper -> Right value+        Just _ -> Left ("kiroku-inspect: " <> T.pack name <> " flag is outside the allowed range")+        Nothing -> case variable name of+            Nothing -> Right fallback+            Just value -> either (\message -> Left ("kiroku-inspect: " <> T.pack name <> ": " <> message)) Right (boundedDecimal lower upper value)++runInspect :: InspectHooks -> InspectRuntime -> IO ()+runInspect hooks rt = do+    storeVar <- newTVarIO Nothing+    metrics <- newKirokuMetricsWith (readPosition storeVar) (readSubscribers storeVar)+    let settings =+            (Store.defaultConnectionSettings rt.databaseUrl)+                { Store.schema = rt.schema+                , Store.poolSize = rt.poolSize+                , Store.eventHandler = Just (metricsEventHandler metrics Nothing)+                , Store.observationHandler = Just (metricsObservationHandler metrics Nothing)+                }+        cfg = defaultConfig{Config.port = rt.port, Config.cors = rt.cors, Config.wsMaxConnections = rt.wsMaxConnections}+    Store.withStore settings $ \store -> do+        atomically (writeTVar storeVar (Just store))+        providers <- storeServerProviders cfg metrics store+        withMetricsServerWithProviders cfg metrics [postgresPing store] providers $ \server -> do+            hooks.onListening server.serverPort (capabilitiesFor cfg (providerPresence providers))+            hooks.waitForShutdown++readPosition :: TVar (Maybe Store.KirokuStore) -> STM Store.GlobalPosition+readPosition storeVar = readTVar storeVar >>= maybe (pure (Store.GlobalPosition 0)) (publisherPosition . (.publisher))++readSubscribers :: TVar (Maybe Store.KirokuStore) -> STM Int+readSubscribers storeVar = readTVar storeVar >>= maybe (pure 0) (\store -> IntMap.size <$> readTVar (subscribers store.publisher))++-- | Never includes the database URL. The port is the actual bound port.+renderStartupBanner :: InspectRuntime -> Int -> Capabilities -> [Text]+renderStartupBanner rt boundPort caps =+    [ "kiroku-inspect: connected to schema " <> T.pack (show rt.schema) <> "; listening on port " <> T.pack (show boundPort)+    , "kiroku-inspect: routes browse="+        <> enabled caps.routes.browse+        <> " subscriptions_checkpoints="+        <> enabled caps.routes.subscriptionsCheckpoints+        <> " dead_letters="+        <> enabled caps.routes.deadLetters+        <> " subscriptions_live="+        <> enabled caps.routes.subscriptionsLive+        <> " websocket_events="+        <> enabled caps.routes.websocketEvents+        <> " cors="+        <> enabled caps.corsIsEnabled+    , "kiroku-inspect: this process runs no subscriptions; /subscriptions, /metrics, and /health reflect only this process"+    ]+  where+    enabled True = "on"+    enabled False = "off"
src/Kiroku/Metrics/WebSocket.hs view
@@ -29,7 +29,21 @@     ClientMessage (..),     ServerMessage (..),     recordedEventToJSON,+    recordedEventToJSONResolved,+    errorCodeReplayFailed,+    errorCodeCategoryReadFailed,+    errorCodeEventStreamOverflowed,+    errorCodeLiveDecodeFailed,+    overflowNotice, +    -- * Tail delivery building blocks+    StreamNameCache,+    newStreamNameCache,+    streamNameCacheSize,+    resolveEventNames,+    broadcastEventsWith,+    withWorkerSlot,+     -- * Connection limiting     WebSocketState (..),     newWebSocketState,@@ -39,10 +53,9 @@ ) where  import Control.Concurrent (threadDelay)-import Control.Concurrent.Async (async, cancel, link)+import Control.Concurrent.Async (asyncWithUnmask, link, uninterruptibleCancel) import Control.Concurrent.STM (     STM,-    TBQueue,     TVar,     atomically,     check,@@ -52,12 +65,12 @@     readTVar,     writeTVar,  )-import Control.Exception (catch, finally)+import Control.Exception (bracket, catch, finally, mask, mask_) import Control.Monad (forever) import Data.Aeson (     FromJSON (..),     ToJSON (..),-    Value,+    Value (..),     eitherDecode',     encode,     object,@@ -66,14 +79,22 @@     (.:?),     (.=),  )+import Data.Aeson.KeyMap qualified as KeyMap import Data.ByteString qualified as BS import Data.ByteString.Lazy qualified as LBS-import Data.Foldable (for_)+import Data.Foldable (foldl', for_)+import Data.IORef (IORef, newIORef, readIORef, writeIORef) import Data.Int (Int64)+import Data.Map.Strict (Map)+import Data.Map.Strict qualified as Map+import Data.Sequence (Seq, ViewL (..), (|>))+import Data.Sequence qualified as Seq+import Data.Set qualified as Set import Data.Text (Text) import Data.Text qualified as T import Data.Vector (Vector) import Data.Vector qualified as V+import Data.Word (Word64) import Network.WebSockets qualified as WS  import Kiroku.Metrics.Collector (KirokuMetrics, snapshotMetrics)@@ -84,21 +105,24 @@     GlobalPosition (..),     KirokuStore (..),     RecordedEvent (..),+    lookupStreamNames,     readAllForward,     readCategory,     runStoreIO,  ) import Kiroku.Store.Settings (DecodedBatch (..), DecodedEvent (..), decodedEventRecorded) import Kiroku.Store.Subscription.EventPublisher (+    PublisherSubscription (..),     SubscriberStatus (..),     publisherPosition,-    subscribePublisher,+    subscribePublisherWith,  ) import Kiroku.Store.Subscription.Types (OverflowPolicy (..)) import Kiroku.Store.Types (     EventId (..),     EventType (..),     StreamId (..),+    StreamName (..),     StreamVersion (..),  ) @@ -124,6 +148,8 @@       SubscribeEvents !(Maybe Int64) !(Maybe Text)     | -- | (Event channel) stop the current tail.       UnsubscribeEvents+    | -- | Stop periodic metrics snapshots; subscribe_metrics resumes them.+      UnsubscribeMetrics     deriving stock (Eq, Show)  -- | Messages from the server to a WebSocket client (tagged on a @"type"@ field).@@ -142,6 +168,8 @@       Goodbye     | -- | A non-fatal error message for the client.       ErrorMsg !Text+    | -- | An error with a stable machine-readable code.+      CodedError !Text !Text     deriving stock (Eq, Show)  instance FromJSON ClientMessage where@@ -150,6 +178,7 @@         case msgType :: Text of             "ping" -> pure Ping             "subscribe_metrics" -> pure SubscribeMetrics+            "unsubscribe_metrics" -> pure UnsubscribeMetrics             "subscribe_events" ->                 SubscribeEvents <$> v .:? "from_position" <*> v .:? "category"             "unsubscribe_events" -> pure UnsubscribeEvents@@ -163,7 +192,29 @@         object ["type" .= ("event_stream_started" :: Text), "from_position" .= p]     toJSON Goodbye = object ["type" .= ("goodbye" :: Text)]     toJSON (ErrorMsg msg) = object ["type" .= ("error" :: Text), "message" .= msg]+    toJSON (CodedError code msg) = object ["type" .= ("error" :: Text), "code" .= code, "message" .= msg] +-- | Published error codes (human messages may change).+errorCodeReplayFailed, errorCodeCategoryReadFailed, errorCodeEventStreamOverflowed, errorCodeLiveDecodeFailed :: Text+errorCodeReplayFailed = "replay_failed"+errorCodeCategoryReadFailed = "category_read_failed"+errorCodeEventStreamOverflowed = "event_stream_overflowed"+errorCodeLiveDecodeFailed = "live_decode_failed"++-- | A notice on any counter change, including a Word64 wrap (modular subtraction).+overflowNotice :: Word64 -> Word64 -> Maybe ServerMessage+overflowNotice previous current+    | previous /= current =+        Just+            ( CodedError+                errorCodeEventStreamOverflowed+                ( "event stream overflowed; "+                    <> T.pack (show (current - previous))+                    <> " undelivered batch(es) dropped since the last notice; re-read from your last position"+                )+            )+    | otherwise = Nothing+ {- | Encode a 'RecordedEvent' to a JSON 'Value' (IP-4). An explicit function rather than a @ToJSON@ instance: 'RecordedEvent' has no instance today and a library-level orphan is undesirable. EP-4's user guide documents this shape.@@ -278,14 +329,34 @@ handleMetrics :: MetricsServerConfig -> KirokuMetrics -> WS.PendingConnection -> IO () handleMetrics cfg m pending = do     conn <- WS.acceptRequest pending-    WS.withPingThread conn 30 (pure ()) $ do-        sendMsg conn . Snapshot =<< snapshotMetrics m-        pushThread <- async (metricsPushLoop cfg m conn)-        link pushThread-        finally-            (metricsReceiveLoop m conn)-            (cancel pushThread >> sendMsg conn Goodbye)+    WS.withPingThread conn 30 (pure ()) $+        withWorkerSlot $ \start stop -> do+            sendMsg conn . Snapshot =<< snapshotMetrics m+            start (metricsPushLoop cfg m conn)+            metricsReceiveLoop m conn (start (metricsPushLoop cfg m conn)) stop+                `finally` (stop >> sendMsg conn Goodbye) +{- | Scope one linked worker. Creation and registration are masked, the worker+body is unmasked, and cleanup cancels and joins even during acquisition. Starting+an already occupied slot is a no-op. The receive loop is the sole slot owner.+-}+withWorkerSlot :: ((IO () -> IO ()) -> IO () -> IO a) -> IO a+withWorkerSlot body = mask $ \restore -> do+    workerVar <- newTVarIO Nothing+    let stop = mask_ $ do+            worker <- atomically (readTVar workerVar)+            for_ worker uninterruptibleCancel+            atomically (writeTVar workerVar Nothing)+        start action = mask_ $ do+            worker <- atomically (readTVar workerVar)+            case worker of+                Just _ -> pure ()+                Nothing -> do+                    child <- asyncWithUnmask (\unmask -> unmask action)+                    atomically (writeTVar workerVar (Just child))+                    link child+    restore (body start stop) `finally` stop+ -- | Periodically push a fresh snapshot every @wsPushIntervalUs@. metricsPushLoop :: MetricsServerConfig -> KirokuMetrics -> WS.Connection -> IO () metricsPushLoop cfg m conn = forever $ do@@ -293,12 +364,15 @@     sendMsg conn . Snapshot =<< snapshotMetrics m  -- | Answer @ping@ with @pong@ and @subscribe_metrics@ with a fresh snapshot.-metricsReceiveLoop :: KirokuMetrics -> WS.Connection -> IO ()-metricsReceiveLoop m conn = forever $ do+metricsReceiveLoop :: KirokuMetrics -> WS.Connection -> IO () -> IO () -> IO ()+metricsReceiveLoop m conn start stop = forever $ do     cmd <- recvMsg conn     case cmd of         Just Ping -> sendMsg conn Pong-        Just SubscribeMetrics -> sendMsg conn . Snapshot =<< snapshotMetrics m+        Just SubscribeMetrics -> do+            sendMsg conn . Snapshot =<< snapshotMetrics m+            start+        Just UnsubscribeMetrics -> stop         _ -> pure ()  --------------------------------------------------------------------------------@@ -313,26 +387,18 @@ handleEvents :: MetricsServerConfig -> KirokuStore -> WS.PendingConnection -> IO () handleEvents cfg store pending = do     conn <- WS.acceptRequest pending-    WS.withPingThread conn 30 (pure ()) $ do-        tailVar <- newTVarIO Nothing-        let stopTail = do-                mt <- atomically (readTVar tailVar)-                for_ mt cancel-                atomically (writeTVar tailVar Nothing)-            startTail from cat = do-                stopTail-                t <- async (eventTail cfg store conn from cat)-                atomically (writeTVar tailVar (Just t))-        finally-            ( forever $ do-                cmd <- recvMsg conn-                case cmd of-                    Just Ping -> sendMsg conn Pong-                    Just (SubscribeEvents from cat) -> startTail from cat-                    Just UnsubscribeEvents -> stopTail-                    _ -> pure ()-            )-            (stopTail >> sendMsg conn Goodbye)+    WS.withPingThread conn 30 (pure ()) $+        withWorkerSlot $ \start stop ->+            finally+                ( forever $ do+                    cmd <- recvMsg conn+                    case cmd of+                        Just Ping -> sendMsg conn Pong+                        Just (SubscribeEvents from cat) -> stop >> start (eventTail cfg store conn from cat)+                        Just UnsubscribeEvents -> stop+                        _ -> pure ()+                )+                (stop >> sendMsg conn Goodbye)  -- | A reasonable replay/category page size. eventReadLimit :: Int@@ -356,28 +422,33 @@     Maybe Int64 ->     Maybe Text ->     IO ()-eventTail cfg store conn mFrom mCategory =+eventTail cfg store conn mFrom mCategory = do+    cache <- newStreamNameCache     case mCategory of         Just cat -> do             start <- case mFrom of                 Just p -> pure p                 Nothing -> unGP <$> atomically (publisherPosition store.publisher)             sendMsg conn (EventStreamStarted start)-            categoryLoop store conn (CategoryName cat) start-        Nothing -> do-            (queue, statusVar, unsubscribe) <--                atomically (subscribePublisher store.publisher cfg.wsEventQueueCap DropOldest)-            attachPos <- unGP <$> atomically (publisherPosition store.publisher)-            flip finally unsubscribe $-                case mFrom of-                    Nothing -> do-                        sendMsg conn (EventStreamStarted attachPos)-                        broadcastLoop conn queue statusVar (const True)-                    Just p -> do-                        sendMsg conn (EventStreamStarted p)-                        mCovered <- replayHistory store conn p attachPos-                        for_ mCovered $ \covered ->-                            broadcastLoop conn queue statusVar (\e -> unGP e.globalPosition > covered)+            categoryLoop store cache conn (CategoryName cat) start+        Nothing ->+            bracket+                ( atomically $ do+                    sub <- subscribePublisherWith store.publisher cfg.wsEventQueueCap DropOldest+                    attachPos <- unGP <$> publisherPosition store.publisher+                    pure (sub, attachPos)+                )+                (unsubscribe . fst)+                $ \(sub, attachPos) ->+                    case mFrom of+                        Nothing -> do+                            sendMsg conn (EventStreamStarted attachPos)+                            broadcastEventsWith (sendMsg conn) cache (lookupNames store) sub (const True)+                        Just p -> do+                            sendMsg conn (EventStreamStarted p)+                            mCovered <- replayHistory store cache conn p attachPos+                            for_ mCovered $ \covered ->+                                broadcastEventsWith (sendMsg conn) cache (lookupNames store) sub (\e -> unGP e.globalPosition > covered)  {- | Page history from the requested position up to @attachPos@ with 'readAllForward'. Returns @Just covered@, the highest global position the@@ -386,60 +457,68 @@ been surfaced to the client as an 'ErrorMsg'. The caller must terminate the tail on @Nothing@ rather than continue live with a gap. -}-replayHistory :: KirokuStore -> WS.Connection -> Int64 -> Int64 -> IO (Maybe Int64)-replayHistory store conn from attachPos = go from attachPos+replayHistory :: KirokuStore -> StreamNameCache -> WS.Connection -> Int64 -> Int64 -> IO (Maybe Int64)+replayHistory store cache conn from attachPos = go from attachPos   where     go cursor covered         | cursor >= attachPos = pure (Just covered)         | otherwise = do             res <- runStoreIO store (readAllForward (GlobalPosition cursor) (fromIntegral eventReadLimit))             case res of-                Left err -> do-                    sendMsg conn (ErrorMsg (T.pack ("replay error: " <> show err)))+                Left _ -> do+                    sendMsg conn (CodedError errorCodeReplayFailed "replay error: history unavailable")                     pure Nothing                 Right evs                     | V.null evs -> pure (Just covered)                     | otherwise -> do-                        sendEvents conn evs+                        sendEvents store cache conn evs                         let lastPos = unGP (V.last evs).globalPosition                         go lastPos (max covered lastPos) -{- | Drain the broadcast queue forever, sending each kept event. Defensively-surfaces an @Overflowed@ status (not set under 'DropOldest', but handled).+{- | Production live delivery with an injectable frame writer and name lookup.+Queue, status and counter are sampled atomically. Loss is signalled before any+surviving event, preserving the client's last contiguous recovery cursor. -}-broadcastLoop ::-    WS.Connection ->-    -- | broadcast queue-    TBQueue DecodedBatch ->-    TVar SubscriberStatus ->+broadcastEventsWith ::+    (ServerMessage -> IO ()) ->+    StreamNameCache ->+    ([StreamId] -> IO (Map StreamId StreamName)) ->+    PublisherSubscription ->     (RecordedEvent -> Bool) ->     IO ()-broadcastLoop conn queue statusVar keep = go+broadcastEventsWith send cache lookupBatch sub keep = go 0 False   where-    go = do-        batch <- atomically (readTBQueue queue)+    go previous warned = do+        (batch, status, dropped) <-+            atomically $+                (,,)+                    <$> readTBQueue sub.subscriptionQueue+                    <*> readTVar sub.subscriptionStatus+                    <*> readTVar sub.subscriptionDropped+        let notice = case overflowNotice previous dropped of+                Just msg -> Just msg+                Nothing | status == Overflowed && not warned -> Just (CodedError errorCodeEventStreamOverflowed "event stream overflowed; some events dropped")+                Nothing -> Nothing+        for_ notice send         let decoded = case batch of                 UnchangedBatch events -> Right (V.filter keep events)                 TransformedBatch events -> V.mapM unwrap (V.filter (keep . decodedEventRecorded) events)             unwrap (Decoded event) = Right event             unwrap (Undecodable _ failure) = Left failure         case decoded of-            Left failure -> sendMsg conn (ErrorMsg (T.pack ("live decode error: " <> show failure)))+            Left _ -> send (CodedError errorCodeLiveDecodeFailed "live event decoding failed")             Right events -> do-                sendEvents conn events-                status <- atomically (readTVar statusVar)-                case status of-                    Overflowed -> sendMsg conn (ErrorMsg "event stream overflowed; some events dropped")-                    _ -> pure ()-                go+                names <- resolveEventNames cache lookupBatch events+                V.mapM_ (send . Event . recordedEventToJSONResolved names) events+                go dropped (status == Overflowed)  {- | DB-driven category live loop. Mirrors the subscription worker's @liveLoopDbDriven@: gate on the publisher advancing past the /last observed/ position (not the cursor) so an unmatched category does not busy-spin, then drain the category to empty before waiting again. -}-categoryLoop :: KirokuStore -> WS.Connection -> CategoryName -> Int64 -> IO ()-categoryLoop store conn cat startPos = go startPos 0+categoryLoop :: KirokuStore -> StreamNameCache -> WS.Connection -> CategoryName -> Int64 -> IO ()+categoryLoop store cache conn cat startPos = go startPos 0   where     go cursor waitFrom = do         pubPos <- atomically $ do@@ -453,13 +532,13 @@     drainTo cursor = do         res <- runStoreIO store (readCategory cat (GlobalPosition cursor) (fromIntegral eventReadLimit))         case res of-            Left err -> do-                sendMsg conn (ErrorMsg (T.pack ("category read error: " <> show err)))+            Left _ -> do+                sendMsg conn (CodedError errorCodeCategoryReadFailed "category read error: events unavailable")                 pure Nothing             Right evs                 | V.null evs -> pure (Just cursor)                 | otherwise -> do-                    sendEvents conn evs+                    sendEvents store cache conn evs                     drainTo (unGP (V.last evs).globalPosition)  --------------------------------------------------------------------------------@@ -467,9 +546,54 @@ --------------------------------------------------------------------------------  -- | Send each event in a batch as an 'Event' message.-sendEvents :: WS.Connection -> Vector RecordedEvent -> IO ()-sendEvents conn = V.mapM_ (sendMsg conn . Event . recordedEventToJSON)+sendEvents :: KirokuStore -> StreamNameCache -> WS.Connection -> Vector RecordedEvent -> IO ()+sendEvents store cache conn events = do+    names <- resolveEventNames cache (lookupNames store) events+    V.mapM_ (sendMsg conn . Event . recordedEventToJSONResolved names) events +-- Typed lookup failures fall back to null names; thrown exceptions propagate.+lookupNames :: KirokuStore -> [StreamId] -> IO (Map StreamId StreamName)+lookupNames store ids = either (const Map.empty) id <$> runStoreIO store (lookupStreamNames ids)++data NameCache = NameCache !(Map StreamId StreamName) !(Seq StreamId)++-- | Per-tail FIFO name cache. Missing names are not retained. Capacity: 4096.+newtype StreamNameCache = StreamNameCache (IORef NameCache)++-- | Allocate once for a tail; dispose on unsubscribe/disconnect.+newStreamNameCache :: IO StreamNameCache+newStreamNameCache = StreamNameCache <$> newIORef (NameCache Map.empty Seq.empty)++-- | Retained map and FIFO sizes (both bounded by 4096).+streamNameCacheSize :: StreamNameCache -> IO (Int, Int)+streamNameCacheSize (StreamNameCache ref) = do+    NameCache names order <- readIORef ref+    pure (Map.size names, Seq.length order)++{- | Resolve distinct misses in at most one lookup. The temporary encoding map+contains all current-batch names even when that batch exceeds cache capacity.+The tail owns the cache; calls on one cache must be serialized.+-}+resolveEventNames :: StreamNameCache -> ([StreamId] -> IO (Map StreamId StreamName)) -> Vector RecordedEvent -> IO (Map StreamId StreamName)+resolveEventNames (StreamNameCache ref) lookupBatch events = do+    NameCache cached order <- readIORef ref+    let wanted = Set.fromList (V.toList (V.map (.originalStreamId) events))+        -- Inspect the requested ids rather than materializing all retained keys+        -- for each small batch after the cache has filled.+        missing = Set.filter (`Map.notMember` cached) wanted+    found <- if Set.null missing then pure Map.empty else Map.restrictKeys <$> lookupBatch (Set.toList missing) <*> pure missing+    let current = Map.union cached found+        inserted = foldl' (|>) order (Map.keys found)+        trimmed = evict current inserted+    trimmed `seq` writeIORef ref trimmed+    pure (Map.restrictKeys current wanted)+  where+    evict names order+        | Map.size names <= 4096 = NameCache names order+        | otherwise = case Seq.viewl order of+            oldest :< rest -> evict (Map.delete oldest names) rest+            EmptyL -> NameCache Map.empty Seq.empty+ {- | Send a 'ServerMessage', swallowing a closed-connection exception so cleanup in a @finally@ never re-throws on an already-dead socket. -}@@ -489,3 +613,9 @@ -- | Unwrap a 'GlobalPosition' to its underlying 'Int64'. unGP :: GlobalPosition -> Int64 unGP (GlobalPosition n) = n++-- | The frozen event shape plus the resolved original name (or null).+recordedEventToJSONResolved :: Map StreamId StreamName -> RecordedEvent -> Value+recordedEventToJSONResolved names event = case recordedEventToJSON event of+    Object fields -> Object (KeyMap.insert "original_stream_name" (toJSON (fmap (\(StreamName name) -> name) (Map.lookup event.originalStreamId names))) fields)+    _ -> error "recordedEventToJSON must return an object"
test/Main.hs view
@@ -3,16 +3,30 @@ import Test.Hspec (hspec)  import Kiroku.Test.Postgres (withSharedMigratedPostgres)+import Test.BrowseSpec qualified as BrowseSpec+import Test.CapabilitiesSpec qualified as CapabilitiesSpec+import Test.CheckpointsSpec qualified as CheckpointsSpec import Test.CollectorSpec qualified as CollectorSpec+import Test.CorsSpec qualified as CorsSpec+import Test.DeadLettersSpec qualified as DeadLettersSpec import Test.IntegrationSpec qualified as IntegrationSpec import Test.ServerSpec qualified as ServerSpec+import Test.StandaloneSpec qualified as StandaloneSpec import Test.SubscriptionsSpec qualified as SubscriptionsSpec+import Test.WebSocketConvergenceSpec qualified as WebSocketConvergenceSpec import Test.WebSocketSpec qualified as WebSocketSpec  main :: IO () main = withSharedMigratedPostgres $ hspec $ do+    StandaloneSpec.spec+    CapabilitiesSpec.spec+    BrowseSpec.spec+    CheckpointsSpec.spec+    DeadLettersSpec.spec+    CorsSpec.spec     CollectorSpec.spec     IntegrationSpec.spec     ServerSpec.spec+    WebSocketConvergenceSpec.spec     WebSocketSpec.spec     SubscriptionsSpec.spec
+ test/Test/BrowseSpec.hs view
@@ -0,0 +1,173 @@+{-# 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
+ test/Test/CapabilitiesSpec.hs view
@@ -0,0 +1,167 @@+module Test.CapabilitiesSpec (spec) where++import Control.Exception (SomeException, try)+import Control.Monad (forM_)+import Data.Aeson qualified as A+import Data.Aeson.KeyMap qualified as KM+import Data.ByteString.Builder (toLazyByteString)+import Data.ByteString.Lazy qualified as LBS+import Data.Either (isLeft)+import Data.IORef (modifyIORef', newIORef, readIORef, writeIORef)+import Data.Text (Text)+import Data.Text qualified as T+import Network.HTTP.Client qualified as HTTP+import Network.HTTP.Types (Method, ResponseHeaders, Status, status200, status404, status405)+import Network.Wai qualified as Wai+import Network.Wai.Handler.Warp qualified as Warp+import Network.Wai.Internal (ResponseReceived (..))+import Network.WebSockets qualified as WS+import System.Timeout (timeout)+import Test.Hspec++-- The umbrella import also checks that record labels have no export collisions.+import Kiroku.Metrics hiding (InspectOptions (..), InspectRuntime (..))+import Kiroku.Store qualified as Store+import Kiroku.Test.Postgres (withMigratedTestDatabase)++none, allPresent :: ProviderPresence+none = ProviderPresence False False False False noWebSocketChannels+allPresent = ProviderPresence True True True True storeWebSocketChannels++spec :: Spec+spec = describe "Kiroku.Metrics.Capabilities" $ do+    it "pins every wire key, includes Prometheus as process-local, and round-trips" $ do+        let caps = capabilitiesFor defaultConfig allPresent+            keys = ["metrics", "prometheus", "health", "subscriptions_live", "subscriptions_checkpoints", "dead_letters", "browse", "websocket_metrics", "websocket_events"]+        A.toJSON caps+            `shouldBe` A.object+                [ "package" A..= ("kiroku-metrics" :: Text)+                , "version" A..= kirokuMetricsVersion+                , "routes" A..= A.object [key A..= True | key <- keys]+                , "cors" A..= A.object ["enabled" A..= False]+                , "process_local" A..= (["metrics", "prometheus", "health", "subscriptions_live", "websocket_metrics"] :: [Text])+                ]+        A.eitherDecode (A.encode caps) `shouldBe` Right caps+        kirokuMetricsVersion `shouldSatisfy` (T.isInfixOf ".")+    it "derives availability from wiring and switches without invoking any provider" $ do+        let plain = capabilitiesFor defaultConfig none+            off = capabilitiesFor defaultConfig{enableJSON = False, enablePrometheus = False, enableWebSocket = False} allPresent+        plain.routes `shouldBe` RouteAvailability True True True False False False False False False+        off.routes `shouldBe` RouteAvailability False False False True True True True False False+        origin <- either (fail . show) pure (allowedOrigin "http://a")+        (capabilitiesFor defaultConfig{cors = corsAllowOrigins [origin]} none).corsIsEnabled `shouldBe` True+        let providers =+                defaultServerProviders+                    { subscriptionStatus = Just (fail "Discovery invoked live provider")+                    , checkpointInventory = Just (fail "Discovery invoked checkpoint provider")+                    , storeBrowsing = Just (StoreBrowser (\_ -> fail "Discovery invoked browse provider") defaultBrowseLimits)+                    , deadLetters = Just (\_ -> fail "Discovery invoked dead-letter provider")+                    }+        m <- emptyMetrics+        Warp.testWithApplication (pure (httpAppWithProviders defaultConfig m [] providers)) $ \port -> do+            caps <- getCaps port+            caps.routes.subscriptionsLive `shouldBe` True+            caps.routes.subscriptionsCheckpoints `shouldBe` True+            caps.routes.browse `shouldBe` True+            caps.routes.deadLetters `shouldBe` True+    it "implements GET, bodyless HEAD, exact path and 405 in direct WAI" $ do+        let app = capabilitiesApp (capabilitiesFor defaultConfig none)+        (s, h, b) <- capture app "GET" ["capabilities"]+        s `shouldBe` status200+        capture app "HEAD" ["capabilities"] `shouldReturn` (s, h, "")+        b `shouldBe` A.encode (capabilitiesFor defaultConfig none)+        forM_ ["POST", "PUT", "DELETE", "OPTIONS"] $ \method -> do+            (status, headers, errorBody) <- capture app method ["capabilities"]+            status `shouldBe` status405+            errorCode errorBody `shouldBe` Just "method_not_allowed"+            lookup "Allow" headers `shouldBe` Just "GET, HEAD"+        (unknown, unknownHeaders, unknownBody) <- capture app "GET" ["capabilities", "x"]+        unknown `shouldBe` status404+        errorCode unknownBody `shouldBe` Just "not_found"+        capture app "HEAD" ["capabilities", "x"] `shouldReturn` (unknown, unknownHeaders, "")+    it "reports plain/stub servers honestly and remains reachable with every switch off" $ do+        m <- emptyMetrics+        forM_ [defaultConfig{port = 0}, defaultConfig{port = 0, enableJSON = False, enablePrometheus = False, enableWebSocket = False}] $ \cfg ->+            withMetricsServer cfg m [] $ \server -> do+                caps <- getCaps server.serverPort+                caps `shouldBe` capabilitiesFor cfg none+    it "reports all store providers and the legacy store starter's absent live registry" $+        withMigratedTestDatabase $ \url -> Store.withStore (Store.defaultConnectionSettings url) $ \store -> do+            m <- emptyMetrics+            let cfg = defaultConfig{port = 0}+            providers <- storeServerProviders cfg m store+            withMetricsServerWithProviders cfg m [] providers $ \server -> do+                getCaps server.serverPort `shouldReturn` capabilitiesFor cfg allPresent+                forM_ ["/metrics", "/metrics/prometheus", "/health", "/subscriptions", "/subscription-checkpoints", "/streams", "/subscriptions/missing/dead-letters"] $ \path ->+                    HTTP.responseStatus <$> get server.serverPort path `shouldReturn` status200+                forM_ ["/ws/metrics", "/ws/events"] $ \path ->+                    bounded (WS.runClient "127.0.0.1" server.serverPort path (\conn -> WS.sendClose conn ("done" :: Text)))+            withMetricsServerWithStore cfg m store [] $ \server -> do+                caps <- getCaps server.serverPort+                caps.routes.subscriptionsLive `shouldBe` False+                caps.routes.websocketMetrics `shouldBe` True+                caps.routes.websocketEvents `shouldBe` True+    it "honours custom declarations, refuses disabled upgrades and keeps legacy opaque apps conservative" $ do+        m <- emptyMetrics+        let ws connectionRequest = WS.acceptRequest connectionRequest >>= \conn -> WS.sendTextData conn ("ok" :: Text)+            providers = defaultServerProviders{webSocketServer = ws, webSocketChannels = WebSocketChannels False True}+        forM_ [True, False] $ \enabled -> do+            let cfg = defaultConfig{port = 0, enableWebSocket = enabled}+            withMetricsServerWithProviders cfg m [] providers $ \server -> do+                caps <- getCaps server.serverPort+                caps.routes.websocketEvents `shouldBe` enabled+                result <- bounded (try (WS.runClient "127.0.0.1" server.serverPort "/ws/events" (\conn -> WS.receiveData conn :: IO Text)) :: IO (Either SomeException Text))+                if enabled then either (fail . show) (`shouldBe` "ok") result else result `shouldSatisfy` isLeft+        let cfg = defaultConfig{port = 0}+        Warp.testWithApplication (pure (combinedApp cfg m [] Nothing ws)) $ \port -> do+            caps <- getCaps port+            caps.routes.websocketEvents `shouldBe` False+    it "serves discovery behind a mount prefix and preserves the configured CORS policy" $ do+        m <- emptyMetrics+        origin <- either (fail . show) pure (allowedOrigin "http://a")+        let cfg = defaultConfig{cors = corsAllowOrigins [origin]}+            mounted req = combinedAppWithProviders cfg m [] defaultServerProviders req{Wai.pathInfo = drop 1 (Wai.pathInfo req)}+        Warp.testWithApplication (pure mounted) $ \port -> do+            manager <- HTTP.newManager HTTP.defaultManagerSettings+            req <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> "/kiroku/capabilities")+            response <- HTTP.httpLbs req{HTTP.requestHeaders = [("Origin", "http://a")]} manager+            HTTP.responseStatus response `shouldBe` status200+            lookup "Access-Control-Allow-Origin" (HTTP.responseHeaders response) `shouldBe` Just "http://a"+            A.eitherDecode (HTTP.responseBody response) `shouldBe` Right (capabilitiesFor cfg none)++emptyMetrics :: IO KirokuMetrics+emptyMetrics = newKirokuMetricsWith (pure (Store.GlobalPosition 0)) (pure 0)++get :: Int -> String -> IO (HTTP.Response LBS.ByteString)+get port path = do+    manager <- HTTP.newManager HTTP.defaultManagerSettings+    req <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> path)+    HTTP.httpLbs req manager++getCaps :: Int -> IO Capabilities+getCaps port = do+    response <- get port "/capabilities"+    HTTP.responseStatus response `shouldBe` status200+    either fail pure (A.eitherDecode (HTTP.responseBody response))++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+        chunks <- newIORef mempty+        stream $ \body -> body (\builder -> modifyIORef' chunks (<> builder)) (pure ())+        bytes <- toLazyByteString <$> readIORef chunks+        writeIORef result (Just (status, headers, bytes))+        pure ResponseReceived+    readIORef result >>= maybe (fail "No response") pure++bounded :: IO a -> IO a+bounded action = timeout 15_000_000 action >>= maybe (fail "Timed out") pure++errorCode :: LBS.ByteString -> Maybe Text+errorCode body = do+    A.Object root <- A.decode body+    A.Object err <- KM.lookup "error" root+    A.String code <- KM.lookup "code" err+    pure code
+ test/Test/CheckpointsSpec.hs view
@@ -0,0 +1,402 @@+{-# LANGUAGE ScopedTypeVariables #-}++module Test.CheckpointsSpec (spec) where++import Control.Concurrent (threadDelay)+import Control.Concurrent.Async qualified as Async+import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)+import Control.Exception (SomeException, bracket, finally, throwIO, try)+import Control.Monad (forM_)+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.ByteString.Lazy.Char8 qualified as LBSC+import Data.Either (isLeft)+import Data.IORef (modifyIORef', newIORef, readIORef, writeIORef)+import Data.Int (Int32, Int64)+import Data.Map.Strict qualified as Map+import Data.Maybe (isJust)+import Data.Text (Text)+import Data.Text qualified as T+import Data.Time (UTCTime (..), fromGregorian)+import Data.Vector qualified as V+import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Network.HTTP.Client qualified as HTTP+import Network.HTTP.Types (Method, ResponseHeaders, Status, status200, status404, status405, status500, status503)+import Network.Socket qualified as Socket+import Network.Wai qualified as Wai+import Network.Wai.Handler.Warp qualified as Warp+import Network.Wai.Internal (ResponseReceived (..))+import Network.WebSockets qualified as WS+import System.Timeout (timeout)+import Test.Hspec hiding (pending)++import Data.UUID qualified as UUID+import Kiroku.Cli.Subscription.Status (SubscriptionStatusRow (..))+import Kiroku.Metrics.Checkpoints+import Kiroku.Metrics.Collector (KirokuMetrics, newKirokuMetrics, newKirokuMetricsWith)+import Kiroku.Metrics.Config (MetricsServerConfig (..), defaultConfig)+import Kiroku.Metrics.Cors (allowedOrigin, corsAllowOrigins)+import Kiroku.Metrics.Server+import Kiroku.Store (GlobalPosition (..), SubscriptionCheckpoint (..), SubscriptionCheckpointInventory (..), SubscriptionName (..))+import Kiroku.Store qualified as Store+import Kiroku.Store.SQL qualified as SQL+import Kiroku.Store.Settings (DecodeFailure (..))+import Kiroku.Test.Postgres (withMigratedTestDatabase)++fixture :: SubscriptionCheckpointInventory+fixture = SubscriptionCheckpointInventory (GlobalPosition 17) $ V.fromList [SubscriptionCheckpoint (SubscriptionName "beta") 1 (GlobalPosition 11) date, SubscriptionCheckpoint (SubscriptionName "alpha") 0 (GlobalPosition 7) date]+  where+    date = UTCTime (fromGregorian 2026 9 10) 0++spec :: Spec+spec = describe "Kiroku.Metrics.Checkpoints (/subscription-checkpoints)" $ do+    it "pins exact snake_case keys, preserves row order and round-trips lossless Int64 positions" $ do+        let response = checkpointInventoryResponse fixture+            row name index cp = Aeson.object ["subscription" Aeson..= (name :: Text), "member" Aeson..= (index :: Int32), "checkpoint_position" Aeson..= (cp :: Int64), "updated_at" Aeson..= ("2026-09-10T00:00:00Z" :: Text)]+        Aeson.toJSON response `shouldBe` Aeson.object ["store_position" Aeson..= (17 :: Int64), "checkpoints" Aeson..= [row "beta" 1 11, row "alpha" 0 7]]+        Aeson.eitherDecode (Aeson.encode response) `shouldBe` Right response+        let large = CheckpointInventoryResponse maxBound [CheckpointRow "large" 0 9007199254740993 (UTCTime (fromGregorian 2026 9 10) 0)]+        Aeson.eitherDecode (Aeson.encode large) `shouldBe` Right large+    it "serves an inventory and a structured standalone 404" $+        Warp.testWithApplication (pure (checkpointsApp (pure (Right fixture)))) $ \port -> do+            manager <- HTTP.newManager HTTP.defaultManagerSettings+            request <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> "/subscription-checkpoints")+            response <- HTTP.httpLbs request manager+            HTTP.responseStatus response `shouldBe` status200+            Aeson.eitherDecode (HTTP.responseBody response) `shouldBe` Right (checkpointInventoryResponse fixture)+            unknown <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> "/unknown") >>= flip HTTP.httpLbs manager+            HTTP.responseStatus unknown `shouldBe` status404+            Aeson.decode (HTTP.responseBody unknown) `shouldBe` Just (Aeson.object ["error" Aeson..= Aeson.object ["code" Aeson..= ("not_found" :: Text), "message" Aeson..= ("Not found" :: Text)]])++    it "implements HEAD and 405 directly in WAI and calls the provider once per read" $ do+        calls <- newIORef (0 :: Int)+        let app = checkpointsApp (modifyIORef' calls (+ 1) >> pure (Right fixture))+        getResponse <- capture app "GET"+        (headStatus, headHeaders, headBody) <- capture app "HEAD"+        let (getStatus, getHeaders, _) = getResponse+        headStatus `shouldBe` getStatus+        headHeaders `shouldBe` getHeaders+        headBody `shouldBe` ""+        forM_ ["POST", "PUT", "DELETE", "OPTIONS"] $ \method -> do+            (status, headers, body) <- capture app method+            status `shouldBe` status405+            lookup "Allow" headers `shouldBe` Just "GET, HEAD"+            code body `shouldBe` Just "method_not_allowed"+        readIORef calls `shouldReturn` 2++    it "sanitizes typed failures and preserves HEAD error status and headers" $ do+        let failures = [(Store.ConnectionError "postgres://secret", status503, "checkpoint_inventory_unavailable"), (Store.StreamNotFound (Store.StreamName "secret"), status500, "store_error"), (Store.EventDecodeFailed (DecodeFailure (Store.EventId UUID.nil) "secret"), status500, "event_decode_failed")]+        forM_ failures $ \(err, expected, expectedCode) -> do+            let app = checkpointsApp (pure (Left err))+            (status, headers, body) <- capture app "GET"+            status `shouldBe` expected+            code body `shouldBe` Just expectedCode+            body `shouldSatisfy` (not . T.isInfixOf "secret" . T.pack . LBSC.unpack)+            capture app "HEAD" `shouldReturn` (status, headers, "")++    it "does not disguise thrown provider exceptions or cancellation as store errors" $ do+        capture (checkpointsApp (throwIO (userError "provider failed"))) "GET" `shouldThrow` anyIOException+        capture (checkpointsApp (throwIO Async.AsyncCancelled)) "GET" `shouldThrow` (\(_ :: Async.AsyncCancelled) -> True)++    it "keeps a stopped worker durable after its live registry entry disappears" $+        withTestStore $ \store -> do+            append store 3+            let name = Store.SubscriptionName "durable-vs-live"+            handle <- Store.subscribe store (Store.defaultSubscriptionConfig name Store.AllStreams (\_ -> pure Store.Continue))+            awaitLive store name+            Store.cancel handle+            stopped <- timeout 10_000_000 (Store.wait handle)+            stopped `shouldSatisfy` isJust+            states <- Store.subscriptionStates store+            Map.member (name, 0) states `shouldBe` False+            metrics <- newKirokuMetrics store+            let cfg = defaultConfig{port = 0}+            providers <- storeServerProviders cfg metrics store+            withMetricsServerWithProviders cfg metrics [] providers $ \server -> do+                live <- get server.serverPort "/subscriptions"+                HTTP.responseStatus live `shouldBe` status200+                Aeson.decode (HTTP.responseBody live) `shouldBe` Just ([] :: [SubscriptionStatusRow])+                inventory <- getInventory server.serverPort+                inventory.storePosition `shouldBe` 3+                triples inventory `shouldBe` [("durable-vs-live", 0, 3)]++    it "preserves SQL name and numeric member ordering with the captured frontier" $+        withTestStore $ \store -> do+            append store 20+            seed store "zeta" 2 7+            seed store "alpha" 10 3+            seed store "alpha" 2 5+            metrics <- newKirokuMetrics store+            withMetricsServerWithStore defaultConfig{port = 0} metrics store [] $ \server -> do+                inventory <- getInventory server.serverPort+                inventory.storePosition `shouldBe` 20+                triples inventory `shouldBe` [("alpha", 2, 5), ("alpha", 10, 3), ("zeta", 2, 7)]++    it "returns equal quiescent inventories across handles with different live providers" $+        withMigratedTestDatabase $ \connStr ->+            Store.withStore (Store.defaultConnectionSettings connStr) $ \storeA ->+                Store.withStore (Store.defaultConnectionSettings connStr) $ \storeB -> do+                    append storeA 3+                    let name = Store.SubscriptionName "worker-a"+                    handle <- Store.subscribe storeA (Store.defaultSubscriptionConfig name Store.AllStreams (\_ -> pure Store.Continue))+                    awaitLive storeA name+                    seed storeB "worker-b" 0 2+                    metricsA <- newKirokuMetrics storeA+                    metricsB <- newKirokuMetrics storeB+                    let cfg = defaultConfig{port = 0}+                    providersA <- storeServerProviders cfg metricsA storeA+                    providersB <- storeServerProviders cfg metricsB storeB+                    withMetricsServerWithProviders cfg metricsA [] providersA $ \serverA ->+                        withMetricsServerWithProviders cfg metricsB [] providersB $ \serverB -> do+                            liveBodyA <- HTTP.responseBody <$> get serverA.serverPort "/subscriptions"+                            liveBodyB <- HTTP.responseBody <$> get serverB.serverPort "/subscriptions"+                            liveBodyA `shouldNotBe` liveBodyB+                            -- Stop the worker before comparing separate SQL snapshots.+                            Store.cancel handle+                            _ <- Store.wait handle+                            a <- getInventory serverA.serverPort+                            b <- getInventory serverB.serverPort+                            a `shouldBe` b++    it "serves an empty store without changing the legacy unconfigured live response" $+        withTestStore $ \store -> do+            metrics <- newKirokuMetrics store+            withMetricsServerWithStore defaultConfig{port = 0} metrics store [] $ \server -> do+                getInventory server.serverPort `shouldReturn` CheckpointInventoryResponse 0 []+                live <- get server.serverPort "/subscriptions"+                HTTP.responseStatus live `shouldBe` status404+                Aeson.decode (HTTP.responseBody live) `shouldBe` Just (Aeson.object ["error" Aeson..= ("subscription status not configured" :: Text)])++    it "returns structured not-configured errors without shadowing a live name checkpoints" $ do+        metrics <- emptyMetrics+        let live = pure [SubscriptionStatusRow "checkpoints" 0 "live" 7]+        withMetricsServerSubscriptions defaultConfig{port = 0} metrics [] live $ \server -> do+            inventory <- get server.serverPort "/subscription-checkpoints"+            HTTP.responseStatus inventory `shouldBe` status404+            code (HTTP.responseBody inventory) `shouldBe` Just "checkpoint_inventory_not_configured"+            named <- get server.serverPort "/subscriptions/checkpoints"+            HTTP.responseStatus named `shouldBe` status200+            Aeson.decode (HTTP.responseBody named) `shouldBe` Just [SubscriptionStatusRow "checkpoints" 0 "live" 7]+        let app = httpAppWithProviders defaultConfig metrics [] defaultServerProviders+        (status, headers, _) <- capture app "GET"+        capture app "HEAD" `shouldReturn` (status, headers, "")+        (methodStatus, methodHeaders, _) <- capture app "POST"+        methodStatus `shouldBe` status405+        lookup "Allow" methodHeaders `shouldBe` Just "GET, HEAD"++    it "mounts HTTP and the real event WebSocket when the host strips only pathInfo" $+        withTestStore $ \store -> do+            metrics <- newKirokuMetrics store+            providers <- storeServerProviders defaultConfig metrics store+            let composed = combinedAppWithProviders defaultConfig metrics [] providers+                mounted req respond = composed req{Wai.pathInfo = drop 1 (Wai.pathInfo req)} respond+            Warp.testWithApplication (pure mounted) $ \port -> do+                response <- get port "/kiroku/subscription-checkpoints"+                HTTP.responseStatus response `shouldBe` status200+                bounded $ WS.runClient "127.0.0.1" port "/kiroku/ws/events" $ \conn -> do+                    WS.sendTextData conn (Aeson.encode $ Aeson.object ["type" Aeson..= ("subscribe_events" :: Text)])+                    raw <- WS.receiveData conn :: IO LBS.ByteString+                    frameType raw `shouldBe` Just "event_stream_started"++    it "escapes mounted WebSocket segments and retains the raw query string" $ do+        metrics <- emptyMetrics+        observed <- newEmptyMVar+        let ws pending = do+                putMVar observed (WS.requestPath (WS.pendingRequest pending))+                conn <- WS.acceptRequest pending+                WS.sendTextData conn ("ok" :: Text)+            app = combinedAppWithProviders defaultConfig metrics [] defaultServerProviders{webSocketServer = ws}+            mounted req respond = app req{Wai.pathInfo = drop 1 (Wai.pathInfo req)} respond+        Warp.testWithApplication (pure mounted) $ \port -> do+            bounded $ WS.runClient "127.0.0.1" port "/kiroku/ws/a%20b%2Fc?token=a%2Fb" $ \conn -> do+                WS.receiveData conn `shouldReturn` ("ok" :: Text)+            takeMVar observed `shouldReturn` "/ws/a%20b%2Fc?token=a%2Fb"++    it "enforces disabled WebSockets through both legacy and general composition" $ do+        metrics <- emptyMetrics+        calls <- newIORef (0 :: Int)+        let ws pending = modifyIORef' calls (+ 1) >> WS.rejectRequest pending "called"+            cfg = defaultConfig{enableWebSocket = False}+        forM_ [combinedApp cfg metrics [] Nothing ws, combinedAppWithProviders cfg metrics [] defaultServerProviders{webSocketServer = ws}] $ \app ->+            Warp.testWithApplication (pure app) $ \port -> do+                result <- bounded $ try (WS.runClient "127.0.0.1" port "/ws/metrics" (\_ -> pure ()))+                (result :: Either WS.HandshakeException ()) `shouldSatisfy` isLeft+        readIORef calls `shouldReturn` 0++    it "does no checkpoint work on legacy metrics reads and invokes inventory only on demand" $ do+        metrics <- emptyMetrics+        calls <- newIORef (0 :: Int)+        let provider = modifyIORef' calls (+ 1) >> pure (Right fixture)+            providers = defaultServerProviders{checkpointInventory = Just provider}+            app = httpAppWithProviders defaultConfig metrics [] providers+            metricsApp req respond = app req{Wai.pathInfo = ["metrics"]} respond+        forM_ [1 .. 20 :: Int] $ \_ -> do+            (status, _, _) <- capture metricsApp "GET"+            status `shouldBe` status200+        readIORef calls `shouldReturn` 0+        _ <- capture app "GET"+        readIORef calls `shouldReturn` 1++    it "inherits exactly one CORS wrap on the provider composition" $ do+        metrics <- emptyMetrics+        let origin = either (error . show) id (allowedOrigin "https://ops.example.com")+            cfg = defaultConfig{cors = corsAllowOrigins [origin]}+            app = combinedAppWithProviders cfg metrics [] defaultServerProviders{checkpointInventory = Just (pure (Right fixture))}+        Warp.testWithApplication (pure app) $ \port -> do+            manager <- HTTP.newManager HTTP.defaultManagerSettings+            req <- HTTP.parseRequest (url port "/subscription-checkpoints")+            response <- HTTP.httpLbs req{HTTP.requestHeaders = [("Origin", "https://ops.example.com")]} manager+            HTTP.responseStatus response `shouldBe` status200+            filter ((== "Access-Control-Allow-Origin") . fst) (HTTP.responseHeaders response) `shouldBe` [("Access-Control-Allow-Origin", "https://ops.example.com")]++    it "is immediately ready on ephemeral and fixed ports and releases each socket" $ do+        metrics <- emptyMetrics+        withMetricsServer defaultConfig{port = 0} metrics [] $ \server -> do+            HTTP.responseStatus <$> get server.serverPort "/metrics" `shouldReturn` status200+        (port, socket) <- Warp.openFreePort+        Socket.close socket+        withMetricsServer defaultConfig{port = port} metrics [] $ \server ->+            HTTP.responseStatus <$> get server.serverPort "/metrics" `shouldReturn` status200+        assertPortReleased port++    it "reports an occupied-port failure to acquisition without invoking the callback" $ do+        metrics <- emptyMetrics+        (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+                -- On macOS a loopback listener can coexist with a wildcard listener.+                -- Occupy both wildcard addresses, as the server binds host "*".+                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+                called <- newIORef False+                result <- bounded $ try $ withMetricsServer defaultConfig{port = port} metrics [] (\_ -> writeIORef called True)+                (result :: Either SomeException ()) `shouldSatisfy` isLeft+                readIORef called `shouldReturn` False++    it "releases the server after callback failure and propagates unexpected termination" $ do+        metrics <- emptyMetrics+        portRef <- newIORef 0+        result <- bounded $ try $ withMetricsServer defaultConfig{port = 0} metrics [] $ \server -> do+            writeIORef portRef server.serverPort+            throwIO (userError "callback failed")+        (result :: Either SomeException ()) `shouldSatisfy` isLeft+        readIORef portRef >>= assertPortReleased+        callbackStopped <- newEmptyMVar+        result2 <- bounded $ try $ withMetricsServer defaultConfig{port = 0} metrics [] $ \server ->+            ( do+                writeIORef portRef server.serverPort+                Async.cancel server.serverThread+                threadDelay 10_000_000+            )+                `finally` putMVar callbackStopped ()+        (result2 :: Either SomeException ()) `shouldSatisfy` isLeft+        bounded (takeMVar callbackStopped)+        readIORef portRef >>= assertPortReleased++    it "cancels acquisition or its returned lifetime without leaking a fixed-port listener" $ do+        metrics <- emptyMetrics+        (port, socket) <- Warp.openFreePort+        Socket.close socket+        -- Exercise the acquisition boundary and the readiness/callback boundary.+        forM_ [0, 100, 1000] $ \delay -> do+            entered <- newEmptyMVar+            thread <- Async.async $ do+                putMVar entered ()+                withMetricsServer defaultConfig{port = port} metrics [] (\_ -> threadDelay 10_000_000)+            takeMVar entered+            threadDelay delay+            bounded (Async.cancel thread)+            assertPortReleased port++emptyMetrics :: IO KirokuMetrics+emptyMetrics = newKirokuMetricsWith (pure (GlobalPosition 0)) (pure 0)++withTestStore :: (Store.KirokuStore -> IO a) -> IO a+withTestStore action = withMigratedTestDatabase $ \connection -> Store.withStore (Store.defaultConnectionSettings connection) action++append :: Store.KirokuStore -> Int -> IO ()+append store count = do+    let ev = Store.EventData Nothing (Store.EventType "E") Aeson.Null Nothing Nothing Nothing+    result <- Store.runStoreIO store $ Store.appendToStream (Store.StreamName "orders-1") Store.NoStream (replicate count ev)+    result `shouldSatisfy` either (const False) (const True)++seed :: Store.KirokuStore -> Text -> Int32 -> Int64 -> IO ()+seed store name index position = do+    result <- Pool.use store.pool $ Session.statement (name, index, position, max 1 (index + 1), "unbound", Nothing) SQL.saveCheckpointMemberStmt+    result `shouldBe` Right ()++awaitLive :: Store.KirokuStore -> Store.SubscriptionName -> IO ()+awaitLive store name = bounded loop+  where+    loop = do+        states <- Store.subscriptionStates store+        case Map.lookup (name, 0) states of+            Just view | view.statePhase == "live" -> pure ()+            _ -> threadDelay 20_000 >> loop++triples :: CheckpointInventoryResponse -> [(Text, Int32, Int64)]+triples response = [(row.subscription, row.member, row.checkpointPosition) | row <- response.checkpoints]++url :: Int -> String -> String+url port path = "http://127.0.0.1:" <> show port <> path++get :: Int -> String -> IO (HTTP.Response LBS.ByteString)+get port path = do+    manager <- HTTP.newManager HTTP.defaultManagerSettings+    request <- HTTP.parseRequest (url port path)+    HTTP.httpLbs request manager++getInventory :: Int -> IO CheckpointInventoryResponse+getInventory port = do+    response <- get port "/subscription-checkpoints"+    HTTP.responseStatus response `shouldBe` status200+    either fail pure (Aeson.eitherDecode (HTTP.responseBody response))++code :: LBS.ByteString -> Maybe Text+code bytes = do+    Aeson.Object root <- Aeson.decode bytes+    Aeson.Object err <- KM.lookup "error" root+    Aeson.String value <- KM.lookup "code" err+    pure value++frameType :: LBS.ByteString -> Maybe Text+frameType bytes = do+    Aeson.Object root <- Aeson.decode bytes+    Aeson.String value <- KM.lookup "type" root+    pure value++capture :: Wai.Application -> Method -> IO (Status, ResponseHeaders, LBS.ByteString)+capture app method = do+    result <- newIORef Nothing+    _ <- app Wai.defaultRequest{Wai.requestMethod = method, Wai.pathInfo = checkpointsPath} $ \response -> do+        let (status, headers, stream) = Wai.responseToStream response+        chunks <- newIORef mempty+        stream $ \body -> body (\builder -> modifyIORef' chunks (<> builder)) (pure ())+        bytes <- toLazyByteString <$> readIORef chunks+        writeIORef result (Just (status, headers, bytes))+        pure ResponseReceived+    readIORef result >>= maybe (fail "No WAI response") pure++bounded :: IO a -> IO a+bounded action = timeout 15_000_000 action >>= maybe (fail "Timed out") pure++assertPortReleased :: Int -> IO ()+assertPortReleased port = forM_ addresses $ \(family, address) ->+    bracket (Socket.socket family Socket.Stream Socket.defaultProtocol) Socket.close $ \sock -> do+        Socket.setSocketOption sock Socket.ReuseAddr 1+        if family == Socket.AF_INET6 then Socket.setSocketOption sock Socket.IPv6Only 1 else pure ()+        Socket.bind sock address+        Socket.listen sock 1+  where+    addresses =+        [ (Socket.AF_INET, Socket.SockAddrInet (fromIntegral port) (Socket.tupleToHostAddress (127, 0, 0, 1)))+        , (Socket.AF_INET, Socket.SockAddrInet (fromIntegral port) (Socket.tupleToHostAddress (0, 0, 0, 0)))+        , (Socket.AF_INET6, Socket.SockAddrInet6 (fromIntegral port) 0 (0, 0, 0, 0) 0)+        ]
+ test/Test/CorsSpec.hs view
@@ -0,0 +1,265 @@+{-# LANGUAGE ScopedTypeVariables #-}++module Test.CorsSpec (spec) where++import Control.Concurrent (threadDelay)+import Control.Exception (try)+import Control.Monad (forM_)+import Data.Aeson (Value (..), decode, object, (.=))+import Data.Aeson.KeyMap qualified as KM+import Data.ByteString (ByteString)+import Data.ByteString.Builder (toLazyByteString)+import Data.ByteString.Char8 qualified as BS+import Data.ByteString.Lazy qualified as LBS+import Data.CaseInsensitive qualified as CI+import Data.Either (isLeft)+import Data.IORef (modifyIORef', newIORef, readIORef, writeIORef)+import Data.Text (Text)+import Data.UUID qualified as UUID+import Network.HTTP.Client qualified as HTTP+import Network.HTTP.Types hiding (hOrigin, hVary)+import Network.Wai qualified as Wai+import Network.Wai.Handler.Warp qualified as Warp+import Network.Wai.Internal (Response (ResponseRaw), ResponseReceived (..))+import Network.WebSockets qualified as WS+import System.Timeout (timeout)+import Test.Hspec++import Kiroku.Metrics+import Kiroku.Metrics.Config qualified as Config+import Kiroku.Metrics.JSON (errorEnvelope, storeErrorResponse)+import Kiroku.Store (StreamName (..), defaultConnectionSettings, withStore)+import Kiroku.Store.Error (StoreError (..))+import Kiroku.Store.Settings (DecodeFailure (..))+import Kiroku.Store.Types (EventId (..), GlobalPosition (..))+import Kiroku.Test.Postgres (withMigratedTestDatabase)++hOrigin, hVary :: HeaderName+hOrigin = "Origin"+hVary = "Vary"++ops, evil :: ByteString+ops = "https://ops.example.com"+evil = "https://evil.example.com"++policy :: CorsPolicy+policy = corsAllowOrigins [either (error . show) id (allowedOrigin "https://ops.example.com")]++spec :: Spec+spec = do+    describe "Kiroku.Metrics.Cors (configuration)" $ do+        it "defaults to disabled and validates normalized HTTP(S) origins" $ do+            cors defaultConfig `shouldBe` corsDisabled+            corsEnabled corsDisabled `shouldBe` False+            allowedOrigin " HTTPS://OPS.Example.Com:443/ " `shouldBe` allowedOrigin "https://ops.example.com"+            allowedOrigin "http://localhost:80" `shouldBe` allowedOrigin "http://localhost"+            fmap renderAllowedOrigin (allowedOrigin "http://127.0.0.1:5173") `shouldBe` Right "http://127.0.0.1:5173"+            originAllowed policy "HTTPS://OPS.EXAMPLE.COM:443" `shouldBe` True+            originAllowed policy "https://ops.example.com:8443" `shouldBe` False+        it "normalizes IPv6 literals including embedded IPv4" $ do+            allowedOrigin "http://[::1]" `shouldBe` allowedOrigin "http://[0:0:0:0:0:0:0:1]:80"+            allowedOrigin "https://[2001:DB8::1]:8443" `shouldBe` allowedOrigin "https://[2001:db8:0:0:0:0:0:1]:8443"+            allowedOrigin "http://[::ffff:192.0.2.1]" `shouldBe` allowedOrigin "http://[0:0:0:0:0:ffff:c000:201]"+        it "rejects wildcards, opaque origins, malformed authorities, paths and Unicode" $+            forM_ ["*", "null", "", "example.com", "://example.com", "ftp://example.com", "http://", "http:///", "https://a/path", "https://a//", "https://a?x", "https://a#x", "https://a b", "https://a\nb", "https://user@a", "https://*.a", "https://a%2eb", "https://a\\b", "https://é.example", "http://:80", "http://a:", "http://a:-1", "http://a:+80", "http://a:65536", "http://a:999999999999999999999999", "http://a:80:90", "http://-a", "http://a-", "http://a..b", "http://256.1.2.3", "http://01.2.3.4", "http://[:::1]", "http://[::1", "http://[1:2:3:4:5:6:7]", "http://[1:2:3:4:5:6:7:8:9]", "http://[::1]:", "http://[::1]x", "http://[fe80::1%eth0]", "http://::1", "http://[192.0.2.1::]"] $ \input ->+                allowedOrigin input `shouldSatisfy` isLeft+        it "tolerates configuration slash/whitespace but refuses them in request origins" $ do+            allowedOrigin " https://ops.example.com/ " `shouldBe` allowedOrigin "https://ops.example.com"+            forM_ ["https://ops.example.com/", " https://ops.example.com", "https://ops.example.com ", "https://ops.example.com https://evil.example.com", "https://ops.example.com,https://evil.example.com", "\xff"] $ \raw ->+                originAllowed policy raw `shouldBe` False++    describe "Kiroku.Metrics.Cors (middleware, standalone)" $ do+        it "preserves every response byte and header when disabled, including upgrades and preflights" $ do+            let original = Wai.responseLBS status201 [("Vary", "Accept"), ("X-Custom", "yes")] "unchanged"+                app _ respond = respond original+            forM_ [[], [(hOrigin, ops)], [(hOrigin, evil), ("Upgrade", "websocket")], preflight ops "DELETE"] $ \headers -> do+                unwrapped <- capture app "OPTIONS" headers+                wrapped <- capture (corsMiddleware corsDisabled app) "OPTIONS" headers+                wrapped `shouldBe` unwrapped+        it "decorates GET and HEAD and varies allowed, disallowed, absent, malformed and duplicate origins" $ do+            forM_ ["GET", "HEAD"] $ \method -> do+                (_, headers, body) <- capture (corsMiddleware policy baseApp) method [(hOrigin, ops)]+                lookup "Access-Control-Allow-Origin" headers `shouldBe` Just ops+                lookup hVary headers `shouldBe` Just "Origin"+                body `shouldBe` "legacy"+            forM_ [[], [(hOrigin, evil)], [(hOrigin, "null")], [(hOrigin, ops <> "/")], [(hOrigin, ops), (hOrigin, ops)]] $ \headers -> do+                (status, hs, body) <- capture (corsMiddleware policy baseApp) "GET" headers+                status `shouldBe` status200+                body `shouldBe` "legacy"+                grants hs `shouldBe` []+                lookup hVary hs `shouldBe` Just "Origin"+        it "answers GET/HEAD preflights, reflects validated tokens, and varies on all inputs" $ do+            forM_ ["GET", "HEAD"] $ \method -> do+                (status, headers, body) <- capture (corsMiddleware policy baseApp) "OPTIONS" (preflight ops method <> [("Access-Control-Request-Headers", "Authorization, X-Trace")])+                status `shouldBe` status204+                body `shouldBe` ""+                lookup "Access-Control-Allow-Methods" headers `shouldBe` Just "GET, HEAD, OPTIONS"+                lookup "Access-Control-Allow-Headers" headers `shouldBe` Just "Authorization, X-Trace"+                lookup hVary headers `shouldBe` Just "Origin, Access-Control-Request-Method, Access-Control-Request-Headers"+        it "rejects unsupported/duplicate methods and malformed requested header tokens" $ do+            (status, _, body) <- capture (corsMiddleware policy baseApp) "OPTIONS" (preflight ops "POST")+            status `shouldBe` status403+            errorCode body `shouldBe` Just (String "cors_method_not_allowed")+            forM_ ["", "X Header", "X:Header", "X-Good,", "X-Good, \xff", "X\r\nInjected"] $ \header -> do+                (s, _, b) <- capture (corsMiddleware policy baseApp) "OPTIONS" (preflight ops "GET" <> [("Access-Control-Request-Headers", header)])+                s `shouldBe` status400+                errorCode b `shouldBe` Just (String "invalid_cors_request")+            (emptyStatus, _, _) <- capture (corsMiddleware policy baseApp) "OPTIONS" (preflight ops "GET" <> [("Access-Control-Request-Headers", ""), ("Access-Control-Request-Headers", "Authorization")])+            emptyStatus `shouldBe` status400+            (s, _, _) <- capture (corsMiddleware policy baseApp) "OPTIONS" (preflight ops "GET" <> [("Access-Control-Request-Method", "HEAD")])+            s `shouldBe` status400+        it "passes plain OPTIONS and disallowed preflights through without grants" $ do+            forM_ [[], [(hOrigin, ops)], preflight evil "GET", [(hOrigin, ops), (hOrigin, ops), ("Access-Control-Request-Method", "GET")]] $ \headers -> do+                (s, hs, body) <- capture (corsMiddleware policy baseApp) "OPTIONS" headers+                s `shouldBe` status200+                body `shouldBe` "legacy"+                if headers == [(hOrigin, ops)] then lookup "Access-Control-Allow-Origin" hs `shouldBe` Just ops else grants hs `shouldBe` []+        it "merges existing Vary case-insensitively and preserves wildcard variation" $ do+            let app headers _ respond = respond (Wai.responseLBS status200 headers "")+            (_, hs, _) <- capture (corsMiddleware policy (app [(hVary, "Accept, origin"), (hVary, "ACCEPT, X-Foo")])) "GET" []+            lookup hVary hs `shouldBe` Just "Accept, origin, X-Foo"+            (_, wildcard, _) <- capture (corsMiddleware policy (app [(hVary, "Accept, *")])) "GET" [(hOrigin, ops)]+            lookup hVary wildcard `shouldBe` Just "*"+        it "owns grants once, adds optional credentials and ignores negative max age" $ do+            let app _ respond = respond (Wai.responseLBS status200 [("Access-Control-Allow-Origin", "*"), ("Access-Control-Allow-Origin", evil), ("Access-Control-Allow-Credentials", "true")] "")+            (_, hs, _) <- capture (corsMiddleware policy app) "GET" [(hOrigin, ops)]+            grants hs `shouldBe` [("Access-Control-Allow-Origin", ops)]+            (_, denied, _) <- capture (corsMiddleware policy app) "GET" [(hOrigin, evil)]+            grants denied `shouldBe` []+            forM_ [Nothing, Just (-1), Just 0, Just 3600] $ \age -> do+                (_, headers, _) <- capture (corsMiddleware (policy{allowCredentials = True, maxAgeSeconds = age}) baseApp) "OPTIONS" (preflight ops "GET")+                lookup "Access-Control-Allow-Credentials" headers `shouldBe` Just "true"+                lookup "Access-Control-Max-Age" headers `shouldBe` case age of+                    Just n | n >= 0 -> Just (fromStringInt n)+                    _ -> Nothing+        it "refuses malformed/duplicate/disallowed upgrade origins before invoking the inner app" $ do+            calls <- newIORef (0 :: Int)+            let app req respond = modifyIORef' calls (+ 1) >> baseApp req respond+            forM_ [[(hOrigin, evil)], [(hOrigin, "null")], [(hOrigin, ops <> "/")], [(hOrigin, ops), (hOrigin, ops)]] $ \headers -> do+                (s, hs, body) <- capture (corsMiddleware policy app) "GET" (("Upgrade", "websocket") : headers)+                s `shouldBe` status403+                lookup hContentType hs `shouldBe` Just "application/json"+                errorCode body `shouldBe` Just (String "origin_not_allowed")+            readIORef calls `shouldReturn` 0+            forM_ [[], [(hOrigin, ops)]] $ \headers -> do+                _ <- capture (corsMiddleware policy app) "GET" (("Upgrade", "websocket") : headers)+                pure ()+            readIORef calls `shouldReturn` 2+        it "never reuses grants between sequential requests with different origins" $ do+            forM_ [[], [(hOrigin, ops)], [(hOrigin, evil)], [(hOrigin, ops)], []] $ \headers -> do+                (_, hs, _) <- capture (corsMiddleware policy baseApp) "GET" headers+                lookup hVary hs `shouldBe` Just "Origin"+                lookup "Access-Control-Allow-Origin" hs `shouldBe` if headers == [(hOrigin, ops)] then Just ops else Nothing+        it "leaves raw upgrade responses untouched" $ do+            let original = Wai.responseRaw (\_ _ -> pure ()) (Wai.responseLBS status500 [] "fallback")+                app _ respond = respond original+            ref <- newIORef Nothing+            _ <- corsMiddleware policy app (Wai.defaultRequest{Wai.requestHeaders = [("Upgrade", "websocket"), (hOrigin, ops)]}) (\r -> writeIORef ref (Just r) >> pure ResponseReceived)+            response <- readIORef ref+            case response of+                Just (ResponseRaw _ _) -> pure ()+                _ -> expectationFailure "raw response was rewritten"+        it "strips HEAD bodies on a real Warp server while preserving grants and status" $+            Warp.testWithApplication (pure (corsMiddleware policy baseApp)) $ \port -> do+                manager <- HTTP.newManager HTTP.defaultManagerSettings+                get <- networkRequest manager port "GET" [(hOrigin, ops)]+                headResponse <- networkRequest manager port "HEAD" [(hOrigin, ops)]+                HTTP.responseStatus headResponse `shouldBe` HTTP.responseStatus get+                lookup "Access-Control-Allow-Origin" (HTTP.responseHeaders headResponse) `shouldBe` Just ops+                HTTP.responseBody headResponse `shouldBe` ""++    describe "Kiroku.Metrics.JSON (shared inspection errors)" $ do+        it "pins envelope keys and omits optional details" $ do+            errorEnvelope "origin_not_allowed" "Denied." Nothing `shouldBe` object ["error" .= object ["code" .= ("origin_not_allowed" :: Text), "message" .= ("Denied." :: Text)]]+            errorEnvelope "invalid_query_parameter" "Invalid." (Just (object ["parameter" .= ("limit" :: Text)])) `shouldBe` object ["error" .= object ["code" .= ("invalid_query_parameter" :: Text), "message" .= ("Invalid." :: Text), "details" .= object ["parameter" .= ("limit" :: Text)]]]+        it "sanitizes unavailable and other store errors" $ do+            forM_ [(ConnectionError "postgres://secret", status503, "store_unavailable"), (StreamNotFound (StreamName "private"), status500, "store_error"), (EventDecodeFailed (DecodeFailure (EventId UUID.nil) "secret payload"), status500, "event_decode_failed")] $ \(err, expected, code) -> do+                (status, _, body) <- capture (\_ respond -> respond (storeErrorResponse "store_unavailable" err)) "GET" []+                status `shouldBe` expected+                errorCode body `shouldBe` Just (String code)+                body `shouldSatisfy` (not . BS.isInfixOf "secret" . LBS.toStrict)++    describe "Kiroku.Metrics.Cors (real server)" $ do+        it "decorates store-backed metrics, answers preflights and keeps denied bodies unchanged" $+            withInspection policy $ \port -> do+                manager <- HTTP.newManager HTTP.defaultManagerSettings+                allowed <- networkRequest manager port "GET" [(hOrigin, ops)]+                HTTP.responseStatus allowed `shouldBe` status200+                lookup "Access-Control-Allow-Origin" (HTTP.responseHeaders allowed) `shouldBe` Just ops+                denied <- networkRequest manager port "GET" [(hOrigin, evil)]+                HTTP.responseStatus denied `shouldBe` status200+                grants (HTTP.responseHeaders denied) `shouldBe` []+                pre <- networkRequest manager port "OPTIONS" (preflight ops "GET")+                HTTP.responseStatus pre `shouldBe` status204+        it "upgrades allowed and absent origins but refuses an unlisted browser origin" $+            withInspection policy $ \port -> do+                assertSnapshot port [(hOrigin, ops)]+                refused <- wsSnapshot port [(hOrigin, evil)]+                case refused of+                    Left (WS.MalformedResponse _ _) -> pure ()+                    other -> expectationFailure ("expected 403 MalformedResponse, got " <> show other)+                assertSnapshot port []+        it "keeps upgrades open to any origin under the default disabled policy" $+            withInspection corsDisabled $ \port ->+                assertSnapshot port [(hOrigin, evil)]++baseApp :: Wai.Application+baseApp _ respond = respond (Wai.responseLBS status200 [("X-Legacy", "kept")] "legacy")++preflight :: ByteString -> ByteString -> RequestHeaders+preflight origin method = [(hOrigin, origin), ("Access-Control-Request-Method", method)]++grants :: ResponseHeaders -> ResponseHeaders+grants = filter (\(name, _) -> "access-control-" `BS.isPrefixOf` CI.foldedCase name)++capture :: Wai.Application -> Method -> RequestHeaders -> IO (Status, ResponseHeaders, LBS.ByteString)+capture app method headers = do+    ref <- newIORef Nothing+    _ <- app (Wai.defaultRequest{Wai.requestMethod = method, Wai.requestHeaders = headers}) $ \response -> do+        let (status, hs, stream) = Wai.responseToStream response+        chunks <- newIORef mempty+        stream $ \body -> body (\builder -> modifyIORef' chunks (<> builder)) (pure ())+        bytes <- toLazyByteString <$> readIORef chunks+        writeIORef ref (Just (status, hs, bytes))+        pure ResponseReceived+    readIORef ref >>= maybe (fail "application did not respond") pure++errorCode :: LBS.ByteString -> Maybe Value+errorCode bytes = do+    Object root <- decode bytes+    Object err <- KM.lookup "error" root+    KM.lookup "code" err++networkRequest :: HTTP.Manager -> Int -> Method -> RequestHeaders -> IO (HTTP.Response LBS.ByteString)+networkRequest manager port method headers = do+    request <- HTTP.parseRequest ("http://127.0.0.1:" <> show port <> "/metrics")+    HTTP.httpLbs (request{HTTP.method = method, HTTP.requestHeaders = headers}) manager++withInspection :: CorsPolicy -> (Int -> IO a) -> IO a+withInspection cors action = withMigratedTestDatabase $ \connStr -> do+    metrics <- newKirokuMetricsWith (pure (GlobalPosition 0)) (pure 0)+    withStore (defaultConnectionSettings connStr) $ \store ->+        withMetricsServerWithStore (defaultConfig{Config.port = 0, Config.cors = cors}) metrics store [] $ \server -> do+            threadDelay 300_000+            action server.serverPort++wsSnapshot :: Int -> RequestHeaders -> IO (Either WS.HandshakeException (Maybe Value))+wsSnapshot port headers = do+    result <- timeout 15_000_000 $+        try $+            WS.runClientWith "127.0.0.1" port "/ws/metrics" WS.defaultConnectionOptions headers $ \conn -> do+                raw <- WS.receiveData conn :: IO LBS.ByteString+                pure $ do+                    Object fields <- decode raw+                    KM.lookup "type" fields+    maybe (fail "WebSocket snapshot timed out") pure result++fromStringInt :: Int -> ByteString+fromStringInt = BS.pack . show++assertSnapshot :: Int -> RequestHeaders -> Expectation+assertSnapshot port headers = do+    result <- wsSnapshot port headers+    case result of+        Right value -> value `shouldBe` Just (String "snapshot")+        Left err -> expectationFailure (show err)
+ test/Test/DeadLettersSpec.hs view
@@ -0,0 +1,190 @@+{-# LANGUAGE OverloadedLabels #-}++module Test.DeadLettersSpec (spec) where++import Control.Exception (SomeException, throwIO, try)+import Control.Lens ((^.))+import Control.Monad (forM_, void)+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.ByteString.Lazy.Char8 qualified as LBSC+import Data.Generics.Labels ()+import Data.IORef+import Data.Int (Int32)+import Data.Text (Text)+import Data.Text qualified as T+import Data.Time (UTCTime (..), fromGregorian)+import Data.UUID qualified as UUID+import Data.Vector qualified as V+import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Kiroku.Metrics.Collector (newKirokuMetrics, newKirokuMetricsWith)+import Kiroku.Metrics.Config (MetricsServerConfig (..), defaultConfig)+import Kiroku.Metrics.Cors (allowedOrigin, corsAllowOrigins)+import Kiroku.Metrics.DeadLetters+import Kiroku.Metrics.Server+import Kiroku.Store qualified as Store+import Kiroku.Store.SQL qualified as SQL+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 System.Timeout (timeout)+import Test.Hspec++fixture :: Store.SubscriptionDeadLetterPage+fixture = Store.SubscriptionDeadLetterPage (V.singleton (Store.SubscriptionDeadLetter 7 (Store.SubscriptionName "orders") 9 (Store.GlobalPosition 9007199254740993) (Store.EventId UUID.nil) (object ["kind" .= ("poison" :: Text), "detail" .= ("unknown SKU" :: Text)]) "poison: unknown SKU" 1 (UTCTime (fromGregorian 2026 9 10) 0))) (Just (Store.SubscriptionDeadLetterCursor (Store.GlobalPosition 9007199254740993) 7))++spec :: Spec+spec = describe "DeadLetters" $ do+    it "pins exact snake_case keys, preserves reason JSON and round-trips lossless positions" $ do+        let response = deadLetterPageResponse fixture+            item = object ["dead_letter_id" .= (7 :: Int), "subscription" .= ("orders" :: Text), "member" .= (9 :: Int), "global_position" .= (9007199254740993 :: Integer), "event_id" .= UUID.nil, "reason" .= object ["kind" .= ("poison" :: Text), "detail" .= ("unknown SKU" :: Text)], "reason_summary" .= ("poison: unknown SKU" :: Text), "attempt_count" .= (1 :: Int), "created_at" .= ("2026-09-10T00:00:00Z" :: Text)]+        Aeson.toJSON response `shouldBe` object ["items" .= [item], "next_cursor" .= ("9007199254740993:7" :: Text)]+        Aeson.eitherDecode (Aeson.encode response) `shouldBe` Right response+        Aeson.toJSON (deadLetterPageResponse (Store.SubscriptionDeadLetterPage V.empty Nothing)) `shouldBe` object ["items" .= ([] :: [Value])]+    it "validates decimal cursor components without signed or overflow narrowing" $ do+        forM_ ["", "1", "1:", ":1", "1:2:3", "-1:2", "1:-2", "1.5:2", "18446744073709551616:1", "9223372036854775808:1", "1:18446744073709551616", "+1:2", "1:2"] $ \raw -> parseDeadLetterCursor raw `shouldBe` Nothing+        forM_ [Store.SubscriptionDeadLetterCursor (Store.GlobalPosition 0) 0, Store.SubscriptionDeadLetterCursor (Store.GlobalPosition maxBound) maxBound] $ \cursor -> parseDeadLetterCursor (renderDeadLetterCursor cursor) `shouldBe` Just cursor+    it "passes the complete query once and uses the documented defaults" $ do+        calls <- newIORef []+        let provider query = modifyIORef' calls (<> [query]) >> pure (Right fixture)+        void $ capture (deadLettersApp provider) "GET" ["subscriptions", "orders", "dead-letters"] [("member", Just "9"), ("from", Just "4211:7"), ("limit", Just "5")]+        void $ capture (deadLettersApp provider) "GET" ["subscriptions", "orders", "dead-letters"] []+        recorded <- readIORef calls+        recorded `shouldBe` [Store.SubscriptionDeadLetterQuery (Store.SubscriptionName "orders") (Just 9) (Just (Store.SubscriptionDeadLetterCursor (Store.GlobalPosition 4211) 7)) (either (error . show) (\v -> v) (Store.mkSubscriptionDeadLetterLimit 5)), Store.defaultSubscriptionDeadLetterQuery (Store.SubscriptionName "orders")]+    it "rejects invalid, duplicate, missing and malformed UTF-8 values before calling a provider" $ do+        calls <- newIORef (0 :: Int)+        let app = deadLettersApp (\_ -> modifyIORef' calls (+ 1) >> pure (Right fixture))+        forM_ [[("limit", Just raw)] | raw <- ["0", "1001", "abc", "", "-1", "18446744073709551616"]] $ \query -> invalidRequest calls app query+        forM_ [[("member", Just "2147483648")], [("member", Just "-1")], [("from", Just "1:9223372036854775808")], [("limit", Nothing)], [("limit", Just "2"), ("limit", Just "3")], [("member", Just "\255")]] $ invalidRequest calls app+        readIORef calls `shouldReturn` 0+        (status, _, _) <- capture app "GET" ["subscriptions", "orders", "dead-letters"] [("unknown", Nothing)]+        status `shouldBe` status200+    it "implements HEAD and 405 in WAI, including unavailable and unconfigured responses" $ do+        calls <- newIORef (0 :: Int)+        let app = deadLettersApp (\_ -> modifyIORef' calls (+ 1) >> pure (Right fixture))+        (status, headers, _) <- capture app "GET" route []+        capture app "HEAD" route [] `shouldReturn` (status, headers, "")+        forM_ ["POST", "PUT", "DELETE", "OPTIONS"] $ \method -> do+            (actual, hs, body) <- capture app method route []+            actual `shouldBe` status405+            lookup "Allow" hs `shouldBe` Just "GET, HEAD"+            code body `shouldBe` Just "method_not_allowed"+        readIORef calls `shouldReturn` 2+        (missing, _, body) <- capture deadLettersNotConfiguredApp "GET" route []+        missing `shouldBe` status404+        code body `shouldBe` Just "dead_letters_not_configured"+    it "sanitizes typed store failures and preserves HEAD error headers" $ do+        forM_ [(Store.ConnectionError "postgres://secret", status503, "dead_letters_unavailable"), (Store.StreamNotFound (Store.StreamName "secret"), status500, "store_error"), (Store.EventDecodeFailed (Store.DecodeFailure (Store.EventId UUID.nil) "secret"), status500, "event_decode_failed")] $ \(err, expected, expectedCode) -> do+            let app = deadLettersApp (\_ -> pure (Left err))+            (status, headers, body) <- capture app "GET" route []+            status `shouldBe` expected+            code body `shouldBe` Just expectedCode+            body `shouldSatisfy` (not . T.isInfixOf "secret" . T.pack . LBSC.unpack)+            capture app "HEAD" route [] `shouldReturn` (status, headers, "")+    it "propagates thrown provider exceptions and serves structured unknown-path errors" $ do+        result <- try @SomeException $ capture (deadLettersApp (\_ -> throwIO (userError "failure"))) "GET" route []+        result `shouldSatisfy` either (const True) (const False)+        (status, _, body) <- capture (deadLettersApp (\_ -> error "must not run")) "GET" ["unknown"] []+        status `shouldBe` status404+        code body `shouldBe` Just "not_found"+    it "serves real pages, member filters and legacy behavior from the store-aware server" $ withTestStore $ \store -> do+        let event = Store.EventData Nothing (Store.EventType "Created") (object []) Nothing Nothing Nothing+        Right _ <- Store.runStoreIO store $ Store.appendToStream (Store.StreamName "orders-1") Store.NoStream (replicate 5 event)+        Right events <- Store.runStoreIO store $ Store.readAllForward (Store.GlobalPosition 0) 10+        forM_ (V.toList events) $ \recorded -> seed store "paged/name" 0 recorded+        seed store "other" 0 (V.head events)+        handle <- Store.subscribe store $ Store.defaultSubscriptionConfig (Store.SubscriptionName "worker") Store.AllStreams $ \recorded -> pure $ case recorded ^. #globalPosition of+            Store.GlobalPosition 2 -> Store.DeadLetter (Store.DeadLetterPoison "unknown SKU")+            Store.GlobalPosition 5 -> Store.Stop+            _ -> Store.Continue+        stopped <- timeout 10_000_000 (Store.wait handle)+        case stopped of+            Just (Right ()) -> pure ()+            other -> Store.cancel handle >> expectationFailure (show other)+        metrics <- newKirokuMetrics store+        withMetricsServerWithStore defaultConfig{port = 0} metrics store [] $ \server -> do+            workerPage <- get server.serverPort "/subscriptions/worker/dead-letters" >>= decodePage+            map (.globalPosition) workerPage.items `shouldBe` [2]+            map (.reason) workerPage.items `shouldBe` [object ["kind" .= ("poison" :: Text), "detail" .= ("unknown SKU" :: Text)]]+            putStrLn ("Dead-letter worker response: " <> LBSC.unpack (Aeson.encode workerPage))+            first <- get server.serverPort "/subscriptions/paged%2Fname/dead-letters?limit=2"+            p1 <- decodePage first+            map (.globalPosition) p1.items `shouldBe` [5, 4]+            p2 <- get server.serverPort ("/subscriptions/paged%2Fname/dead-letters?limit=2&from=" <> maybe "" T.unpack p1.nextCursor) >>= decodePage+            p3 <- get server.serverPort ("/subscriptions/paged%2Fname/dead-letters?limit=2&from=" <> maybe "" T.unpack p2.nextCursor) >>= decodePage+            map (.globalPosition) p2.items `shouldBe` [3, 2]+            map (.globalPosition) p3.items `shouldBe` [1]+            p3.nextCursor `shouldBe` Nothing+            empty <- get server.serverPort "/subscriptions/never/dead-letters" >>= decodePage+            empty `shouldBe` DeadLetterPageResponse [] Nothing+            filtered <- get server.serverPort "/subscriptions/paged%2Fname/dead-letters?member=9" >>= decodePage+            filtered `shouldBe` DeadLetterPageResponse [] Nothing+            old <- get server.serverPort "/subscriptions"+            Aeson.decode (HTTP.responseBody old) `shouldBe` Just (object ["error" .= ("subscription status not configured" :: Text)])+            unknown <- get server.serverPort "/nope"+            Aeson.decode (HTTP.responseBody unknown) `shouldBe` Just (object ["error" .= ("Not found" :: Text)])+        providers <- storeServerProviders defaultConfig metrics store+        withMetricsServerWithProviders defaultConfig{port = 0} metrics [] providers $ \server -> do+            live <- get server.serverPort "/subscriptions"+            HTTP.responseStatus live `shouldBe` status200+            void $ get server.serverPort "/subscriptions/paged%2Fname/dead-letters" >>= decodePage+    it "inherits mount-relative routing and CORS with legacy switches disabled" $ do+        metrics <- newKirokuMetricsWith (pure (Store.GlobalPosition 0)) (pure 0)+        let origin = either (error . show) (\v -> v) (allowedOrigin "https://ops.example.com")+            config = defaultConfig{cors = corsAllowOrigins [origin], enableJSON = False, enableWebSocket = False}+            providers = defaultServerProviders{deadLetters = Just (\_ -> pure (Right fixture))}+            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/subscriptions/orders/dead-letters")+            response <- HTTP.httpLbs request{HTTP.requestHeaders = [("Origin", "https://ops.example.com")]} manager+            HTTP.responseStatus response `shouldBe` status200+            lookup "Access-Control-Allow-Origin" (HTTP.responseHeaders response) `shouldBe` Just "https://ops.example.com"++route :: [Text]+route = ["subscriptions", "orders", "dead-letters"]+invalidRequest :: IORef Int -> Wai.Application -> Query -> IO ()+invalidRequest calls app query = do+    priorCalls <- readIORef calls+    (status, _, body) <- capture app "GET" route query+    status `shouldBe` status400+    code body `shouldBe` Just "invalid_query_parameter"+    readIORef calls `shouldReturn` priorCalls+code :: LBS.ByteString -> Maybe Text+code body = case Aeson.decode body of Just (Object fields) | Just (Object err) <- KM.lookup "error" fields, Just (String value) <- KM.lookup "code" err -> Just value; _ -> Nothing+capture :: Wai.Application -> Method -> [Text] -> Query -> IO (Status, ResponseHeaders, LBS.ByteString)+capture app method path query = do+    result <- newIORef Nothing+    _ <- app Wai.defaultRequest{Wai.requestMethod = method, Wai.pathInfo = path, Wai.queryString = query} $ \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+            writeIORef result (Just (status, headers, bytes))+        pure ResponseReceived+    maybe (fail "No response") pure =<< readIORef result+withTestStore :: (Store.KirokuStore -> IO ()) -> IO ()+withTestStore action = withMigratedTestDatabase $ \connection -> Store.withStore (Store.defaultConnectionSettings connection) action+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+decodePage :: HTTP.Response LBS.ByteString -> IO DeadLetterPageResponse+decodePage response = do+    HTTP.responseStatus response `shouldBe` status200+    either fail pure (Aeson.eitherDecode (HTTP.responseBody response))+seed :: Store.KirokuStore -> Text -> Int32 -> Store.RecordedEvent -> IO ()+seed store name member recorded = do+    let Store.EventId eid = recorded ^. #eventId+        Store.GlobalPosition position = recorded ^. #globalPosition+        params = SQL.DeadLetterParams name member "unbound" Nothing (member + 1) position eid (object ["kind" .= ("poison" :: Text), "detail" .= ("unknown SKU" :: Text)]) "poison: unknown SKU" 1+    Pool.use (store ^. #pool) (Session.statement params SQL.insertDeadLetterAndCheckpointStmt) >>= either (fail . show) pure
+ test/Test/StandaloneSpec.hs view
@@ -0,0 +1,250 @@+{-# 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
+ test/Test/WebSocketConvergenceSpec.hs view
@@ -0,0 +1,356 @@+{-# LANGUAGE OverloadedLabels #-}+{-# LANGUAGE ScopedTypeVariables #-}++module Test.WebSocketConvergenceSpec (spec) where++import Control.Concurrent.Async (cancel, withAsync)+import Control.Concurrent.STM+import Control.Exception (MaskingState (..), SomeException, bracket, bracket_, getMaskingState, throwIO, try)+import Control.Lens ((&), (.~))+import Control.Monad (replicateM, replicateM_)+import Data.Aeson (Value (..), eitherDecode, encode, object, toJSON, (.=))+import Data.Aeson.Key qualified as Key+import Data.Aeson.KeyMap qualified as KM+import Data.ByteString.Lazy qualified as LBS+import Data.IORef (atomicModifyIORef', newIORef, readIORef)+import Data.Int (Int64)+import Data.IntMap.Strict qualified as IntMap+import Data.List (nub, sort)+import Data.Map.Strict qualified as Map+import Data.Text (Text)+import Data.Text qualified as T+import Data.Time (UTCTime (..), fromGregorian)+import Data.UUID qualified as UUID+import Data.Vector qualified as V+import Data.Word (Word64)+import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Network.WebSockets qualified as WS+import System.Timeout (timeout)+import Test.Hspec++import Kiroku.Metrics (+    MetricsServer (..),+    MetricsServerConfig (..),+    defaultConfig,+    newKirokuMetricsWith,+    snapshotMetrics,+    startMetricsServerWithStore,+    stopMetricsServer,+ )+import Kiroku.Metrics.WebSocket+import Kiroku.Store hiding (cancel, id)+import Kiroku.Store.Subscription.EventPublisher qualified as Pub+import Kiroku.Test.Postgres (withMigratedTestDatabase)++spec :: Spec+spec = do+    describe "Kiroku.Metrics.WebSocket (frames)" $ do+        it "preserves every existing server frame" $ do+            toJSON Pong `shouldBe` object ["type" .= ("pong" :: Text)]+            toJSON Goodbye `shouldBe` object ["type" .= ("goodbye" :: Text)]+            toJSON (EventStreamStarted 7) `shouldBe` object ["type" .= ("event_stream_started" :: Text), "from_position" .= (7 :: Int)]+            toJSON (ErrorMsg "x") `shouldBe` object ["type" .= ("error" :: Text), "message" .= ("x" :: Text)]+            toJSON (Event (object ["k" .= (1 :: Int)])) `shouldBe` object ["type" .= ("event" :: Text), "event" .= object ["k" .= (1 :: Int)]]+            km <- newKirokuMetricsWith (pure (GlobalPosition 0)) (pure 0)+            snap <- snapshotMetrics km+            look ["type"] (toJSON (Snapshot snap)) `shouldBe` Just (String "snapshot")+            look ["metrics"] (toJSON (Snapshot snap)) `shouldBe` Just (toJSON snap)+        it "pins coded errors and all four code spellings" $ do+            toJSON (CodedError "replay_failed" "boom") `shouldBe` object ["type" .= ("error" :: Text), "code" .= ("replay_failed" :: Text), "message" .= ("boom" :: Text)]+            [errorCodeReplayFailed, errorCodeCategoryReadFailed, errorCodeEventStreamOverflowed, errorCodeLiveDecodeFailed] `shouldBe` ["replay_failed", "category_read_failed", "event_stream_overflowed", "live_decode_failed"]+        it "accepts every old client frame and unsubscribe_metrics" $ do+            let parse raw = eitherDecode raw :: Either String ClientMessage+            map parse ["{\"type\":\"ping\"}", "{\"type\":\"subscribe_metrics\"}", "{\"type\":\"unsubscribe_events\"}", "{\"type\":\"unsubscribe_metrics\"}", "{\"type\":\"subscribe_events\"}", "{\"type\":\"subscribe_events\",\"from_position\":7,\"category\":\"orders\"}"] `shouldBe` map Right [Ping, SubscribeMetrics, UnsubscribeEvents, UnsubscribeMetrics, SubscribeEvents Nothing Nothing, SubscribeEvents (Just 7) (Just "orders")]+        it "adds exactly original_stream_name and preserves all old event fields" $ do+            let event = fixture 9+                original = objectFields (recordedEventToJSON event)+                resolved = objectFields (recordedEventToJSONResolved (Map.singleton (StreamId 9) (StreamName "orders-7")) event)+            KM.delete "original_stream_name" resolved `shouldBe` original+            length (KM.keys resolved) `shouldBe` 12+            KM.lookup "original_stream_name" resolved `shouldBe` Just (String "orders-7")+            look ["original_stream_name"] (recordedEventToJSONResolved Map.empty event) `shouldBe` Just Null+        it "detects only counter changes, including modular wraparound" $ do+            overflowNotice 0 0 `shouldBe` Nothing+            overflowNotice 3 3 `shouldBe` Nothing+            assertNotice 0 2 "2 undelivered"+            assertNotice 2 5 "3 undelivered"+            assertNotice maxBound 1 "2 undelivered"++    describe "Kiroku.Metrics.WebSocket (bounded delivery)" $ do+        it "resolves one cold batch, no warm or empty batches, and bounds FIFO retention" $ do+            cache <- newStreamNameCache+            calls <- newIORef ([] :: [[StreamId]])+            let lookupBatch ids = do+                    atomicModifyIORef' calls (\xs -> (xs <> [ids], ()))+                    pure (Map.fromList [(sid, streamName sid) | sid <- ids])+                events = V.fromList (map fixture [1 .. 5000])+            names <- resolveEventNames cache lookupBatch events+            Map.size names `shouldBe` 5000+            Map.lookup (StreamId 1) names `shouldBe` Just (StreamName "stream-1")+            Map.lookup (StreamId 5000) names `shouldBe` Just (StreamName "stream-5000")+            streamNameCacheSize cache `shouldReturn` (4096, 4096)+            _ <- resolveEventNames cache lookupBatch (V.fromList [fixture 5000, fixture 5000])+            _ <- resolveEventNames cache lookupBatch V.empty+            length <$> readIORef calls `shouldReturn` 1+            _ <- resolveEventNames cache lookupBatch (V.singleton (fixture 1))+            length <$> readIORef calls `shouldReturn` 2+            streamNameCacheSize cache `shouldReturn` (4096, 4096)+        it "does not retain misses and propagates cancellation from a name lookup" $ do+            cache <- newStreamNameCache+            resolveEventNames cache (const (pure Map.empty)) (V.singleton (fixture 1)) `shouldReturn` Map.empty+            streamNameCacheSize cache `shouldReturn` (0, 0)+            entered <- newEmptyTMVarIO+            withAsync (resolveEventNames cache (\_ -> atomically (putTMVar entered ()) >> atomically retry) (V.singleton (fixture 1))) $ \worker -> do+                bounded (atomically (takeTMVar entered))+                cancel worker+            streamNameCacheSize cache `shouldReturn` (0, 0)+        it "delivers unchanged and transformed batches without false overflow notices" $ do+            sub <- fakeSubscription+            cache <- newStreamNameCache+            frames <- newTVarIO []+            atomically $ do+                writeTBQueue sub.subscriptionQueue (UnchangedBatch (V.singleton (fixture 1)))+                writeTBQueue sub.subscriptionQueue (TransformedBatch (V.singleton (Decoded (fixture 2))))+            withAsync (broadcastEventsWith (capture frames) cache (pure . Map.fromList . map (\sid -> (sid, streamName sid))) sub (const True)) $ \_ -> do+                bounded (atomically (readTVar frames >>= \xs -> check (length xs == 2)))+            xs <- readTVarIO frames+            map framePosition xs `shouldBe` [Just 1, Just 2]+        it "never sends partial data and terminates on an applicable typed decode failure" $ do+            sub <- fakeSubscription+            cache <- newStreamNameCache+            frames <- newTVarIO []+            let failed = fixture 2+            atomically (writeTBQueue sub.subscriptionQueue (TransformedBatch (V.fromList [Decoded (fixture 1), Undecodable failed (DecodeFailure failed.eventId "secret payload")])))+            bounded (broadcastEventsWith (capture frames) cache (\_ -> fail "must not look up partial batch") sub (const True))+            readTVarIO frames `shouldReturn` [CodedError errorCodeLiveDecodeFailed "live event decoding failed"]+        it "filters typed failures below the covered replay boundary" $ do+            sub <- fakeSubscription+            cache <- newStreamNameCache+            frames <- newTVarIO []+            let old = fixture 1+            atomically (writeTBQueue sub.subscriptionQueue (TransformedBatch (V.fromList [Undecodable old (DecodeFailure old.eventId "old failure"), Decoded (fixture 2)])))+            withAsync (broadcastEventsWith (capture frames) cache (const (pure Map.empty)) sub (\e -> e.globalPosition > GlobalPosition 1)) $ \_ ->+                bounded (atomically (readTVar frames >>= \xs -> check (length xs == 1)))+            map framePosition <$> readTVarIO frames `shouldReturn` [Just 2]+        it "signals real publisher loss before survivors and recovers from the pre-notice cursor" $ withBareStore $ \store ->+            bracket (atomically (Pub.subscribePublisherWith store.publisher 1 DropOldest)) Pub.unsubscribe $ \sub -> do+                cache <- newStreamNameCache+                frames <- newTVarIO []+                entered <- newEmptyTMVarIO+                release <- newEmptyTMVarIO+                let writer msg = do+                        case msg of+                            Event ev | look ["globalPosition"] ev == Just (toJSON (1 :: Int)) -> atomically (putTMVar entered ()) >> atomically (takeTMVar release)+                            _ -> pure ()+                        capture frames msg+                    lookupBatch ids = either (error . show) id <$> runStoreIO store (lookupStreamNames ids)+                withAsync (broadcastEventsWith writer cache lookupBatch sub (const True)) $ \_ -> do+                    appendEvents store "loss-1" 1+                    bounded (atomically (takeTMVar entered))+                    appendEvents store "loss-2" 1+                    waitPosition store 2+                    appendEvents store "loss-3" 1+                    waitPosition store 3+                    readTVarIO sub.subscriptionDropped `shouldReturn` 1+                    atomically (putTMVar release ())+                    bounded (atomically (readTVar frames >>= \xs -> check (length xs == 3)))+                xs <- readTVarIO frames+                map framePosition xs `shouldBe` [Just 1, Nothing, Just 3]+                case xs of+                    [_, CodedError code _, _] -> code `shouldBe` errorCodeEventStreamOverflowed+                    _ -> expectationFailure (show xs)+                Right recovered <- runStoreIO store (readAllForward (GlobalPosition 1) 10)+                map (.globalPosition) (V.toList recovered) `shouldBe` [GlobalPosition 2, GlobalPosition 3]+                sort (nub ([p | Just p <- map framePosition xs] <> map (\e -> let GlobalPosition p = e.globalPosition in p) (V.toList recovered))) `shouldBe` [1, 2, 3]+        it "coalesces defensive status overflow and a drop notice in the same delivery" $ do+            sub <- fakeSubscription+            cache <- newStreamNameCache+            frames <- newTVarIO []+            atomically $ do+                writeTVar sub.subscriptionStatus Pub.Overflowed+                writeTVar sub.subscriptionDropped 1+                writeTBQueue sub.subscriptionQueue (UnchangedBatch (V.singleton (fixture 1)))+                writeTBQueue sub.subscriptionQueue (UnchangedBatch (V.singleton (fixture 2)))+            withAsync (broadcastEventsWith (capture frames) cache (const (pure Map.empty)) sub (const True)) $ \_ ->+                bounded (atomically (readTVar frames >>= \xs -> check (length xs == 3)))+            map framePosition <$> readTVarIO frames `shouldReturn` [Nothing, Just 1, Just 2]++    describe "Kiroku.Metrics.WebSocket (worker lifecycle)" $ do+        it "keeps one unmasked worker through repeated start/stop and joins on exit" $ do+            active <- newTVarIO (0 :: Int)+            masking <- newTVarIO Nothing+            let worker = bracket_ (atomically (modifyTVar' active (+ 1))) (atomically (modifyTVar' active (subtract 1))) $ do+                    state <- getMaskingState+                    atomically (writeTVar masking (Just state))+                    atomically retry+            withWorkerSlot $ \start stop -> do+                replicateM_ 5 $ do+                    replicateM_ 3 (start worker)+                    bounded (atomically (readTVar active >>= check . (== 1)))+                    bounded (atomically (readTVar masking >>= check . (== Just Unmasked)))+                    stop >> stop+                    readTVarIO active `shouldReturn` 0+                start worker+                bounded (atomically (readTVar active >>= check . (== 1)))+            readTVarIO active `shouldReturn` 0+        it "joins a worker when the connection is cancelled during worker startup" $ do+            active <- newTVarIO (0 :: Int)+            let worker = bracket_ (atomically (modifyTVar' active (+ 1))) (atomically (modifyTVar' active (subtract 1))) (atomically retry)+            withAsync (withWorkerSlot $ \start _ -> start worker >> atomically retry) $ \owner -> do+                bounded (atomically (readTVar active >>= check . (== 1)))+                cancel owner+            readTVarIO active `shouldReturn` 0+        it "propagates unexpected worker failure to the owner" $ do+            result <- try (bounded (withWorkerSlot $ \start _ -> start (throwIO (userError "worker failed")) >> atomically retry)) :: IO (Either SomeException ())+            result `shouldSatisfy` either (T.isInfixOf "worker failed" . T.pack . show) (const False)++    describe "Kiroku.Metrics.WebSocket (convergence, real server)" $ do+        it "stops metrics pushes, resumes them and preserves ping" $ withServer id $ \_ srv -> client srv "/ws/metrics" $ \conn -> do+            _ <- waitForType conn "snapshot"+            _ <- waitForType conn "snapshot"+            replicateM_ 3 (command conn "unsubscribe_metrics")+            -- Ping is an ordering barrier: the receive loop has completed cancellation.+            command conn "ping"+            _ <- waitForType conn "pong"+            timeout 700_000 (WS.receiveData conn :: IO LBS.ByteString) `shouldReturn` Nothing+            replicateM_ 3 (command conn "subscribe_metrics")+            replicateM_ 4 (waitForType conn "snapshot")+            command conn "ping"+            _ <- waitForType conn "pong"+            pure ()+        it "labels live events on both new and warm streams" $ withServer id $ \store srv -> do+            names <- client srv "/ws/events" $ \conn -> do+                command conn "subscribe_events"+                _ <- waitForType conn "event_stream_started"+                appendEvents store "conv-live-1" 2+                appendEvents store "conv-live-2" 1+                replicateM 3 (readEventField conn "original_stream_name")+            names `shouldBe` map (Just . String) ["conv-live-1", "conv-live-1", "conv-live-2"]+            waitSubscriberCount store 0+        it "labels replay and category events" $ withServer id $ \store srv -> do+            appendEvents store "convcat-1" 2+            waitPosition store 2+            replayed <- client srv "/ws/events" $ \conn -> do+                sendJSON conn (object ["type" .= ("subscribe_events" :: Text), "from_position" .= (0 :: Int)])+                _ <- waitForType conn "event_stream_started"+                replicateM 2 (readEventField conn "original_stream_name")+            replayed `shouldBe` replicate 2 (Just (String "convcat-1"))+            names <- client srv "/ws/events" $ \conn -> do+                sendJSON conn (object ["type" .= ("subscribe_events" :: Text), "category" .= ("convcat" :: Text)])+                _ <- waitForType conn "event_stream_started"+                appendEvents store "convcat-2" 1+                readEventField conn "original_stream_name"+            names `shouldBe` Just (String "convcat-2")+            waitSubscriberCount store 0+        it "codes and sanitizes a replay failure and ends its tail" $ withServer id $ \store srv -> do+            appendEvents store "conv-fail-1" 1+            waitPosition store 1+            Pool.use store.pool (Session.script "ALTER TABLE events RENAME TO events_hidden") `shouldReturn` Right ()+            client srv "/ws/events" $ \conn -> do+                sendJSON conn (object ["type" .= ("subscribe_events" :: Text), "from_position" .= (0 :: Int)])+                err <- waitForType conn "error"+                look ["code"] err `shouldBe` Just (String "replay_failed")+                look ["message"] err `shouldBe` Just (String "replay error: history unavailable")+                timeout 300_000 (WS.receiveData conn :: IO LBS.ByteString) `shouldReturn` Nothing+            waitSubscriberCount store 0+        it "codes a category failure without appending to the renamed table" $ withServer id $ \store srv -> do+            appendEvents store "convcat-1" 1+            waitPosition store 1+            Pool.use store.pool (Session.script "ALTER TABLE events RENAME TO events_hidden") `shouldReturn` Right ()+            client srv "/ws/events" $ \conn -> do+                sendJSON conn (object ["type" .= ("subscribe_events" :: Text), "from_position" .= (0 :: Int), "category" .= ("convcat" :: Text)])+                err <- waitForType conn "error"+                look ["code"] err `shouldBe` Just (String "category_read_failed")+                look ["message"] err `shouldBe` Just (String "category read error: events unavailable")+                timeout 300_000 (WS.receiveData conn :: IO LBS.ByteString) `shouldReturn` Nothing+            waitSubscriberCount store 0+        it "codes typed live decode failures without leaking hook details" $ withServer (\settings -> settings & #storeSettings . #decodeHook .~ Just (\e -> pure (Left (DecodeFailure e.eventId "secret")))) $ \store srv -> do+            client srv "/ws/events" $ \conn -> do+                command conn "subscribe_events"+                _ <- waitForType conn "event_stream_started"+                appendEvents store "conv-decode" 1+                err <- waitForType conn "error"+                look ["code"] err `shouldBe` Just (String "live_decode_failed")+                look ["message"] err `shouldBe` Just (String "live event decoding failed")+            waitSubscriberCount store 0++fixture :: Int64 -> RecordedEvent+fixture n = RecordedEvent (EventId UUID.nil) (EventType "E") (StreamVersion 1) (GlobalPosition n) (StreamId n) (StreamVersion 1) (object []) Nothing Nothing Nothing (UTCTime (fromGregorian 2026 10 10) 0)++streamName :: StreamId -> StreamName+streamName (StreamId n) = StreamName ("stream-" <> T.pack (show n))++assertNotice :: Word64 -> Word64 -> Text -> Expectation+assertNotice previous current text = case overflowNotice previous current of+    Just (CodedError code msg) -> do+        code `shouldBe` errorCodeEventStreamOverflowed+        msg `shouldSatisfy` T.isInfixOf text+    other -> expectationFailure (show other)++fakeSubscription :: IO Pub.PublisherSubscription+fakeSubscription = Pub.PublisherSubscription <$> newTBQueueIO 16 <*> newTVarIO Pub.Active <*> newTVarIO 0 <*> pure (pure ())++capture :: TVar [ServerMessage] -> ServerMessage -> IO ()+capture frames msg = atomically (modifyTVar' frames (<> [msg]))++framePosition :: ServerMessage -> Maybe Int64+framePosition (Event ev) = case look ["globalPosition"] ev of+    Just (Number n) -> Just (round n)+    _ -> Nothing+framePosition _ = Nothing++bounded :: IO a -> IO a+bounded action = timeout 15_000_000 action >>= maybe (fail "convergence timeout") pure++withBareStore :: (KirokuStore -> IO a) -> IO a+withBareStore action = withMigratedTestDatabase $ \conn -> withStore (defaultConnectionSettings conn) action++withServer :: (ConnectionSettings -> ConnectionSettings) -> (KirokuStore -> MetricsServer -> IO a) -> IO a+withServer tweak action = withMigratedTestDatabase $ \conn -> do+    var <- newTVarIO Nothing+    km <-+        newKirokuMetricsWith+            (readTVar var >>= maybe (pure (GlobalPosition 0)) (Pub.publisherPosition . (.publisher)))+            (readTVar var >>= maybe (pure 0) (fmap IntMap.size . readTVar . Pub.subscribers . (.publisher)))+    withStore (tweak (defaultConnectionSettings conn)) $ \store -> do+        atomically (writeTVar var (Just store))+        bracket (startMetricsServerWithStore (defaultConfig{port = 0, wsPushIntervalUs = 200_000}) km store []) stopMetricsServer (action store)++client :: MetricsServer -> String -> (WS.Connection -> IO a) -> IO a+client srv path = bounded . WS.runClient "127.0.0.1" srv.serverPort path++sendJSON :: WS.Connection -> Value -> IO ()+sendJSON conn = WS.sendTextData conn . encode++command :: WS.Connection -> Text -> IO ()+command conn name = sendJSON conn (object ["type" .= name])++waitForType :: WS.Connection -> Text -> IO Value+waitForType conn name = do+    raw <- WS.receiveData conn :: IO LBS.ByteString+    value <- either fail pure (eitherDecode raw)+    if look ["type"] value == Just (String name) then pure value else waitForType conn name++readEventField :: WS.Connection -> Text -> IO (Maybe Value)+readEventField conn key = look ["event", key] <$> waitForType conn "event"++look :: [Text] -> Value -> Maybe Value+look [] value = Just value+look (key : keys) (Object fields) = KM.lookup (Key.fromText key) fields >>= look keys+look _ _ = Nothing++appendEvents :: KirokuStore -> Text -> Int -> IO ()+appendEvents store name count = do+    result <- runStoreIO store (appendToStream (StreamName name) NoStream [EventData Nothing (EventType ("E" <> T.pack (show n))) (object []) Nothing Nothing Nothing | n <- [1 .. count]])+    either (fail . show) (const (pure ())) result++waitPosition :: KirokuStore -> Int64 -> IO ()+waitPosition store target = bounded (atomically (Pub.publisherPosition store.publisher >>= check . (>= GlobalPosition target)))++waitSubscriberCount :: KirokuStore -> Int -> IO ()+waitSubscriberCount store count = bounded (atomically (readTVar (Pub.subscribers store.publisher) >>= check . (== count) . IntMap.size))++objectFields :: Value -> KM.KeyMap Value+objectFields (Object fields) = fields+objectFields other = error ("expected object: " <> show other)