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 +55/−0
- app-inspect/Main.hs +38/−0
- bench/WebSocketTail.hs +106/−0
- example/Main.hs +70/−7
- kiroku-metrics.cabal +106/−46
- src/Kiroku/Metrics.hs +12/−0
- src/Kiroku/Metrics/Browse.hs +246/−0
- src/Kiroku/Metrics/Capabilities.hs +152/−0
- src/Kiroku/Metrics/Checkpoints.hs +89/−0
- src/Kiroku/Metrics/Config.hs +5/−0
- src/Kiroku/Metrics/Cors.hs +295/−0
- src/Kiroku/Metrics/DeadLetters.hs +127/−0
- src/Kiroku/Metrics/JSON.hs +24/−2
- src/Kiroku/Metrics/Server.hs +179/−151
- src/Kiroku/Metrics/Standalone.hs +192/−0
- src/Kiroku/Metrics/WebSocket.hs +211/−81
- test/Main.hs +14/−0
- test/Test/BrowseSpec.hs +173/−0
- test/Test/CapabilitiesSpec.hs +167/−0
- test/Test/CheckpointsSpec.hs +402/−0
- test/Test/CorsSpec.hs +265/−0
- test/Test/DeadLettersSpec.hs +190/−0
- test/Test/StandaloneSpec.hs +250/−0
- test/Test/WebSocketConvergenceSpec.hs +356/−0
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)