shibuya-metrics 0.9.0.3 → 0.10.0.0
raw patch · 20 files changed
+1558/−116 lines, 20 filesdep +atomic-primopsdep +effectfuldep +hspecdep ~aesondep ~asyncdep ~bytestringnew-component:exe:metrics-wire-loadPVP ok
version bump matches the API change (PVP)
Dependencies added: atomic-primops, effectful, hspec, http-client, nqe, shibuya-metrics, wai-extra
Dependency ranges changed: aeson, async, bytestring, containers, http-types, shibuya-core, stm, text, time, wai, warp, websockets
API changes (from Hackage documentation)
+ Shibuya.Metrics: ApplicationFailed :: ApplicationStatus
+ Shibuya.Metrics: ApplicationStopped :: ApplicationStatus
+ Shibuya.Metrics: ConfiguredEmpty :: ApplicationStatus
+ Shibuya.Metrics: Draining :: ApplicationStatus
+ Shibuya.Metrics: ProcessorTerminal :: !ProcessorId -> !ProcessorTerminalStatus -> ServerMessage
+ Shibuya.Metrics: Running :: ApplicationStatus
+ Shibuya.Metrics: Starting :: ApplicationStatus
+ Shibuya.Metrics: TerminalFailed :: !Text -> !Maybe Text -> ProcessorTerminalStatus
+ Shibuya.Metrics: TerminalStopped :: ProcessorTerminalStatus
+ Shibuya.Metrics: [application] :: ReadinessStatus -> !ApplicationStatus
+ Shibuya.Metrics: [dependencyTimeoutMicros] :: HealthConfig -> !Int
+ Shibuya.Metrics: [host] :: MetricsServerConfig -> !String
+ Shibuya.Metrics: [wsMaxSubscriptions] :: MetricsServerConfig -> !Int
+ Shibuya.Metrics: data ApplicationStatus
+ Shibuya.Metrics: data ProcessorTerminalStatus
+ Shibuya.Metrics.Config: [dependencyTimeoutMicros] :: MetricsServerConfig -> !Int
+ Shibuya.Metrics.Config: [host] :: MetricsServerConfig -> !String
+ Shibuya.Metrics.Config: [wsMaxSubscriptions] :: MetricsServerConfig -> !Int
+ Shibuya.Metrics.Health: ApplicationFailed :: ApplicationStatus
+ Shibuya.Metrics.Health: ApplicationStopped :: ApplicationStatus
+ Shibuya.Metrics.Health: ConfiguredEmpty :: ApplicationStatus
+ Shibuya.Metrics.Health: Draining :: ApplicationStatus
+ Shibuya.Metrics.Health: Running :: ApplicationStatus
+ Shibuya.Metrics.Health: Starting :: ApplicationStatus
+ Shibuya.Metrics.Health: [application] :: ReadinessStatus -> !ApplicationStatus
+ Shibuya.Metrics.Health: [dependencyTimeoutMicros] :: HealthConfig -> !Int
+ Shibuya.Metrics.Health: data ApplicationStatus
+ Shibuya.Metrics.Health: instance Data.Aeson.Types.ToJSON.ToJSON Shibuya.Metrics.Health.ApplicationStatus
+ Shibuya.Metrics.Health: instance GHC.Classes.Eq Shibuya.Metrics.Health.ApplicationStatus
+ Shibuya.Metrics.Health: instance GHC.Internal.Show.Show Shibuya.Metrics.Health.ApplicationStatus
+ Shibuya.Metrics.Server: [dependencyTimeoutMicros] :: MetricsServerConfig -> !Int
+ Shibuya.Metrics.Server: [host] :: MetricsServerConfig -> !String
+ Shibuya.Metrics.Server: [wsMaxSubscriptions] :: MetricsServerConfig -> !Int
+ Shibuya.Metrics.Server: combinedApp :: MetricsServerConfig -> Master -> WebSocketState -> [DependencyCheck] -> Application
+ Shibuya.Metrics.Types: ProcessorTerminal :: !ProcessorId -> !ProcessorTerminalStatus -> ServerMessage
+ Shibuya.Metrics.Types: TerminalFailed :: !Text -> !Maybe Text -> ProcessorTerminalStatus
+ Shibuya.Metrics.Types: TerminalStopped :: ProcessorTerminalStatus
+ Shibuya.Metrics.Types: data ProcessorTerminalStatus
+ Shibuya.Metrics.Types: instance GHC.Classes.Eq Shibuya.Metrics.Types.ProcessorTerminalStatus
+ Shibuya.Metrics.Types: instance GHC.Internal.Generics.Generic Shibuya.Metrics.Types.ProcessorTerminalStatus
+ Shibuya.Metrics.Types: instance GHC.Internal.Show.Show Shibuya.Metrics.Types.ProcessorTerminalStatus
+ Shibuya.Metrics.WebSocket: [shutdownRequested] :: WebSocketState -> !TVar Bool
+ Shibuya.Metrics.WebSocket: shutdownWebSockets :: WebSocketState -> IO ()
- Shibuya.Metrics: HealthConfig :: !Int -> !NominalDiffTime -> HealthConfig
+ Shibuya.Metrics: HealthConfig :: !Int -> !Int -> !NominalDiffTime -> HealthConfig
- Shibuya.Metrics: MetricsServerConfig :: !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Int -> !NominalDiffTime -> MetricsServerConfig
+ Shibuya.Metrics: MetricsServerConfig :: !String -> !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Int -> !Int -> !Int -> !NominalDiffTime -> MetricsServerConfig
- Shibuya.Metrics: ReadinessStatus :: !Bool -> !ProcessorHealth -> ![DependencyStatus] -> ReadinessStatus
+ Shibuya.Metrics: ReadinessStatus :: !Bool -> !ApplicationStatus -> !ProcessorHealth -> ![DependencyStatus] -> ReadinessStatus
- Shibuya.Metrics.Config: MetricsServerConfig :: !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Int -> !NominalDiffTime -> MetricsServerConfig
+ Shibuya.Metrics.Config: MetricsServerConfig :: !String -> !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Int -> !Int -> !Int -> !NominalDiffTime -> MetricsServerConfig
- Shibuya.Metrics.Health: HealthConfig :: !Int -> !NominalDiffTime -> HealthConfig
+ Shibuya.Metrics.Health: HealthConfig :: !Int -> !Int -> !NominalDiffTime -> HealthConfig
- Shibuya.Metrics.Health: ReadinessStatus :: !Bool -> !ProcessorHealth -> ![DependencyStatus] -> ReadinessStatus
+ Shibuya.Metrics.Health: ReadinessStatus :: !Bool -> !ApplicationStatus -> !ProcessorHealth -> ![DependencyStatus] -> ReadinessStatus
- Shibuya.Metrics.Server: MetricsServerConfig :: !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Int -> !NominalDiffTime -> MetricsServerConfig
+ Shibuya.Metrics.Server: MetricsServerConfig :: !String -> !Int -> !Bool -> !Bool -> !Bool -> !Int -> !Int -> !Int -> !Int -> !Int -> !NominalDiffTime -> MetricsServerConfig
- Shibuya.Metrics.WebSocket: WebSocketState :: !TVar Int -> !Int -> WebSocketState
+ Shibuya.Metrics.WebSocket: WebSocketState :: !TVar Int -> !Int -> !TVar Bool -> WebSocketState
Files
- CHANGELOG.md +42/−0
- bench/WireLoad.hs +177/−0
- shibuya-metrics.cabal +93/−2
- src/Shibuya/Metrics.hs +6/−2
- src/Shibuya/Metrics/Config.hs +15/−2
- src/Shibuya/Metrics/Health.hs +133/−23
- src/Shibuya/Metrics/Prometheus.hs +2/−2
- src/Shibuya/Metrics/Server.hs +29/−9
- src/Shibuya/Metrics/Types.hs +32/−0
- src/Shibuya/Metrics/WebSocket.hs +167/−76
- test/Main.hs +18/−0
- test/Shibuya/Metrics/HealthSpec.hs +236/−0
- test/Shibuya/Metrics/JSONSpec.hs +11/−0
- test/Shibuya/Metrics/PrometheusSpec.hs +24/−0
- test/Shibuya/Metrics/ServerSpec.hs +116/−0
- test/Shibuya/Metrics/TestSupport.hs +122/−0
- test/Shibuya/Metrics/TypesSpec.hs +67/−0
- test/Shibuya/Metrics/WebSocketSpec.hs +237/−0
- test/golden/processor-metrics.json.golden +1/−0
- test/golden/prometheus.golden +30/−0
CHANGELOG.md view
@@ -1,5 +1,47 @@ # Changelog +## 0.10.0.0 — 2026-09-21++### Breaking Changes++- `MetricsServerConfig` and `HealthConfig` gain `dependencyTimeoutMicros`; direct record+ construction must choose a per-check deadline. `ReadinessStatus` gains `application`, and+ its JSON object gains the corresponding lifecycle status field.+- `MetricsServerConfig` gains `host` and `wsMaxSubscriptions`. The built-in server now+ defaults to loopback instead of all interfaces, validates resource limits, and bounds+ retained WebSocket processor selections and exclusions.+- The processing-state JSON object gains `lastProgress`, paired with the new third field of+ `shibuya-core`'s `ProcessorState.Processing` constructor.+- `ServerMessage` gains `ProcessorTerminal`, with the new public+ `ProcessorTerminalStatus` type, and the exposed `WebSocketState` record gains a shutdown+ cell. Exhaustive matches and direct record construction must handle these additions.++### New Features++- Export `combinedApp` so callers can mount the unified metrics WAI application on+ an externally managed server.+- Send an additive `terminal` WebSocket frame when a formerly visible processor leaves the+ live registry in a retained stopped or failed state.++### Other Changes++- Add `shibuya-metrics-test`, a release-gated Hspec suite covering every published+ HTTP route and WebSocket frame plus exact JSON and Prometheus golden contracts.+- Declare the generated `Paths_shibuya_metrics` test module and bound the wire-load+ executable's Effectful and NQE dependencies so `cabal check` accepts the source+ distribution.++### Bug Fixes++- Base stuck detection on sampled progress instead of burst age, retain failed configured+ processors after live metrics unregister, report stopped masters not live, and bound each+ dependency readiness check while normalizing its synchronous exceptions as unhealthy.+- Release WebSocket connection slots on every setup and connection exit, reject upgrades+ when WebSockets are disabled, support exclusions from subscribe-all, and deliver `goodbye`+ when server shutdown begins.+- Close WebSocket clients with policy code 1008 before their retained processor subscription+ or subscribe-all exclusion set can exceed `wsMaxSubscriptions`.+ ## 0.9.0.3 — 2026-09-20 Version bumped to track `shibuya-core` 0.9.0.3. The core dependency bound is
+ bench/WireLoad.hs view
@@ -0,0 +1,177 @@+-- | EP-45 real-wire health polling and WebSocket churn fixture.+module Main (main) where++import Control.Concurrent.NQE.Supervisor (Strategy (IgnoreAll))+import Control.Concurrent.STM (atomically, check, readTVar)+import Control.Exception (bracket, evaluate)+import Control.Monad (replicateM)+import Data.Aeson (ToJSON, encode)+import Data.ByteString.Lazy qualified as LBS+import Data.List (sort)+import Data.Time.Clock (diffUTCTime, getCurrentTime)+import Data.Word (Word64)+import Effectful (runEff)+import GHC.Clock (getMonotonicTimeNSec)+import GHC.Generics (Generic)+import GHC.Stats (GCDetails (..), RTSStats (..), getRTSStats, getRTSStatsEnabled)+import Network.HTTP.Client qualified as HTTP+import Network.HTTP.Types.Status (status200)+import Network.Wai.Handler.Warp qualified as Warp+import Network.WebSockets qualified as WS+import Shibuya.App (Master, ProcessorId (..))+import Shibuya.Core.Metrics (newMetricsHandle)+import Shibuya.Internal.Runner.Master (markMasterRunning, registerProcessor, startMaster, stopMaster)+import Shibuya.Metrics.Config (MetricsServerConfig (..), defaultConfig)+import Shibuya.Metrics.Server (combinedApp)+import Shibuya.Metrics.WebSocket (WebSocketState (..), newWebSocketState)+import System.Environment (lookupEnv)+import System.Exit (exitFailure)+import System.Mem (performMajorGC)+import System.Timeout (timeout)+import Text.Read (readMaybe)++data Scenario = HealthPolling | WebSocketChurn+ deriving stock (Eq, Show)++data Report = Report+ { schemaVersion :: !Int,+ scenario :: !String,+ iterations :: !Int,+ completed :: !Int,+ errors :: !Int,+ elapsedSeconds :: !Double,+ operationsPerSecond :: !Double,+ latencyP50Micros :: !Double,+ latencyP95Micros :: !Double,+ latencyP99Micros :: !Double,+ retainedBytes :: !Word64,+ maxLiveBytes :: !Word64,+ finalWebSocketConnections :: !Int+ }+ deriving stock (Generic)+ deriving anyclass (ToJSON)++main :: IO ()+main = do+ selected <- loadScenario+ defaultIterations <- pure $ if selected == HealthPolling then 5_000 else 500+ count <- envInt "ITERATIONS" defaultIterations+ output <- maybe "metrics-wire-load.json" id <$> lookupEnv "OUTPUT_JSON"+ report <- withMaster $ \master -> do+ registerIdleProcessor master+ wsState <- newWebSocketState defaultConfig.wsMaxConnections+ let config = defaultConfig {wsPushIntervalUs = 10_000}+ app = combinedApp config master wsState []+ Warp.testWithApplication (pure app) $ \port -> runScenario selected count port wsState+ LBS.writeFile output (encode report <> "\n")+ LBS.putStr (encode report <> "\n")+ if report.completed == report.iterations && report.errors == 0 && report.finalWebSocketConnections == 0+ then pure ()+ else exitFailure++loadScenario :: IO Scenario+loadScenario = do+ value <- maybe "health" id <$> lookupEnv "SCENARIO"+ case value of+ "health" -> pure HealthPolling+ "websocket" -> pure WebSocketChurn+ other -> error $ "SCENARIO must be health or websocket, got: " <> other++envInt :: String -> Int -> IO Int+envInt key fallback = maybe fallback id . (>>= readMaybe) <$> lookupEnv key++withMaster :: (Master -> IO a) -> IO a+withMaster = bracket acquire release+ where+ acquire = runEff $ do+ master <- startMaster IgnoreAll+ markMasterRunning master+ pure master+ release master = runEff $ stopMaster master++registerIdleProcessor :: Master -> IO ()+registerIdleProcessor master = do+ now <- getCurrentTime+ metrics <- newMetricsHandle now+ runEff $ registerProcessor master (ProcessorId "wire-load") metrics++runScenario :: Scenario -> Int -> Int -> WebSocketState -> IO Report+runScenario selected count port wsState = do+ start <- getCurrentTime+ latencies <- case selected of+ HealthPolling -> runHealthPolling count port+ WebSocketChurn -> runWebSocketChurn count port+ finish <- getCurrentTime+ remaining <- waitForNoConnections wsState+ (retained, highWater) <- getMemoryBytes+ let elapsed = realToFrac (diffUTCTime finish start)+ completedCount = length latencies+ pure+ Report+ { schemaVersion = 1,+ scenario = case selected of HealthPolling -> "health-polling"; WebSocketChurn -> "websocket-churn",+ iterations = count,+ completed = completedCount,+ errors = count - completedCount,+ elapsedSeconds = elapsed,+ operationsPerSecond = fromIntegral completedCount / max 0.000_001 elapsed,+ latencyP50Micros = percentile 0.50 latencies,+ latencyP95Micros = percentile 0.95 latencies,+ latencyP99Micros = percentile 0.99 latencies,+ retainedBytes = retained,+ maxLiveBytes = highWater,+ finalWebSocketConnections = remaining+ }++runHealthPolling :: Int -> Int -> IO [Word64]+runHealthPolling count port = do+ manager <- HTTP.newManager HTTP.defaultManagerSettings+ let request = HTTP.parseRequest_ $ "http://127.0.0.1:" <> show port <> "/health/ready"+ replicateM count $ timedMicros $ do+ response <- HTTP.httpLbs request manager+ if HTTP.responseStatus response /= status200+ then error $ "Unexpected health status: " <> show (HTTP.responseStatus response)+ else evaluate (LBS.length (HTTP.responseBody response)) >> pure ()++runWebSocketChurn :: Int -> Int -> IO [Word64]+runWebSocketChurn count port =+ replicateM count $+ timedMicros $+ WS.runClient "127.0.0.1" port "/ws" $ \connection -> do+ payload <- WS.receiveData connection :: IO LBS.ByteString+ evaluate (LBS.length payload) >> pure ()++timedMicros :: IO a -> IO Word64+timedMicros action = do+ start <- getMonotonicTimeNSec+ _ <- action+ finish <- getMonotonicTimeNSec+ pure $ (finish - start) `div` 1_000++percentile :: Double -> [Word64] -> Double+percentile _ [] = 0+percentile quantile values =+ let ordered = sort values+ index = min (length ordered - 1) (ceiling (quantile * fromIntegral (length ordered)) - 1)+ in fromIntegral (ordered !! max 0 index)++waitForNoConnections :: WebSocketState -> IO Int+waitForNoConnections wsState = do+ released <- timeout 5_000_000 $ atomically $ do+ count <- readTVar wsState.connectionCount+ check $ count == 0+ case released of+ Nothing -> readTVarIO wsState.connectionCount+ Just () -> pure 0+ where+ readTVarIO variable = atomically $ readTVar variable++getMemoryBytes :: IO (Word64, Word64)+getMemoryBytes = do+ enabled <- getRTSStatsEnabled+ if enabled+ then do+ performMajorGC+ stats <- getRTSStats+ pure (gcdetails_live_bytes stats.gc, max_live_bytes stats)+ else pure (0, 0)
shibuya-metrics.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.12 name: shibuya-metrics-version: 0.9.0.3+version: 0.10.0.0 synopsis: Metrics web server for Shibuya queue processing framework description: Provides HTTP/JSON, Prometheus, and WebSocket endpoints for@@ -12,6 +12,7 @@ build-type: Simple category: Concurrency extra-doc-files: CHANGELOG.md+data-files: test/golden/*.golden common warnings ghc-options: -Wall@@ -48,7 +49,7 @@ containers ^>=0.7, http-types ^>=0.12, prometheus-client ^>=1.1,- shibuya-core ^>=0.9.0.3,+ shibuya-core ^>=0.10.0.0, stm ^>=2.5, text ^>=2.1, time ^>=1.14,@@ -59,3 +60,93 @@ hs-source-dirs: src default-language: GHC2024++test-suite shibuya-metrics-test+ import: warnings+ default-language: GHC2024+ type: exitcode-stdio-1.0+ hs-source-dirs: test+ main-is: Main.hs+ default-extensions:+ DerivingStrategies+ DuplicateRecordFields+ LambdaCase+ NoFieldSelectors+ OverloadedLabels+ OverloadedRecordDot+ OverloadedStrings++ ghc-options:+ -threaded+ -rtsopts+ -with-rtsopts=-N2++ other-modules:+ Paths_shibuya_metrics+ Shibuya.Metrics.HealthSpec+ Shibuya.Metrics.JSONSpec+ Shibuya.Metrics.PrometheusSpec+ Shibuya.Metrics.ServerSpec+ Shibuya.Metrics.TestSupport+ Shibuya.Metrics.TypesSpec+ Shibuya.Metrics.WebSocketSpec++ autogen-modules:+ Paths_shibuya_metrics++ build-depends:+ aeson,+ async,+ atomic-primops,+ base ^>=4.21.0.0,+ bytestring,+ containers,+ effectful >=2.6.1 && <2.8,+ hspec ^>=2.11.17,+ http-types,+ nqe ^>=0.6,+ shibuya-core,+ shibuya-metrics,+ stm,+ text,+ time,+ wai,+ wai-extra ^>=3.1.18,+ warp,+ websockets,++-- EP-45 real-wire HTTP health polling and WebSocket churn fixture.+executable metrics-wire-load+ import: warnings+ default-language: GHC2024+ hs-source-dirs: bench+ main-is: WireLoad.hs+ default-extensions:+ DeriveAnyClass+ DerivingStrategies+ DuplicateRecordFields+ LambdaCase+ NoFieldSelectors+ OverloadedRecordDot+ OverloadedStrings++ ghc-options:+ -threaded+ -rtsopts+ "-with-rtsopts=-N4 -T -A32m"+ -O2++ build-depends:+ aeson,+ base ^>=4.21.0.0,+ bytestring,+ effectful >=2.6.1 && <2.8,+ http-client ^>=0.7.19,+ http-types,+ nqe ^>=0.6,+ shibuya-core,+ shibuya-metrics,+ stm,+ time,+ warp,+ websockets,
src/Shibuya/Metrics.hs view
@@ -43,6 +43,7 @@ -- -- * @{"type": "snapshot", "metrics": {...}}@ - Full metrics snapshot -- * @{"type": "update", "processor": "id", "metrics": {...}}@ - Single processor update+-- * @{"type": "terminal", "processor": "id", "status": "stopped"}@ - Processor left the live registry -- * @{"type": "pong"}@ - Response to ping -- * @{"type": "goodbye"}@ - Server shutting down module Shibuya.Metrics@@ -64,19 +65,22 @@ DependencyStatus (..), LivenessStatus (..), ReadinessStatus (..),+ ApplicationStatus (..), ProcessorHealth (..), HealthConfig (..), defaultHealthConfig, -- * WebSocket Protocol Types ClientMessage (..),+ ProcessorTerminalStatus (..), ServerMessage (..), ) where import Shibuya.Metrics.Config (MetricsServerConfig (..), defaultConfig) import Shibuya.Metrics.Health- ( DependencyCheck,+ ( ApplicationStatus (..),+ DependencyCheck, DependencyStatus (..), HealthConfig (..), LivenessStatus (..),@@ -85,4 +89,4 @@ defaultHealthConfig, ) import Shibuya.Metrics.Server (startMetricsServer, startMetricsServerWithDeps, stopMetricsServer, withMetricsServer)-import Shibuya.Metrics.Types (ClientMessage (..), MetricsServer (..), ServerMessage (..))+import Shibuya.Metrics.Types (ClientMessage (..), MetricsServer (..), ProcessorTerminalStatus (..), ServerMessage (..))
src/Shibuya/Metrics/Config.hs view
@@ -10,7 +10,12 @@ -- | Configuration for the metrics web server. data MetricsServerConfig = MetricsServerConfig- { -- | Port to listen on (default: 9090)+ { -- | Host to listen on (default: loopback only, @127.0.0.1@).+ --+ -- Set this to @*@ only behind an authentication and authorization boundary;+ -- metrics and terminal failure details are operationally sensitive.+ host :: !String,+ -- | Port to listen on (default: 9090) port :: !Int, -- | Enable JSON endpoints (default: True) enableJSON :: !Bool,@@ -22,8 +27,13 @@ wsPushIntervalUs :: !Int, -- | Maximum WebSocket connections (default: 100) wsMaxConnections :: !Int,+ -- | Maximum retained processor identifiers per WebSocket connection+ -- (default: 1,000). A client exceeding this policy is disconnected.+ wsMaxSubscriptions :: !Int, -- | Timeout for liveness check in microseconds (default: 1_000_000 = 1s) livenessTimeoutMicros :: !Int,+ -- | Timeout for each dependency readiness check in microseconds (default: 1s)+ dependencyTimeoutMicros :: !Int, -- | How long a processor can be in Processing state before considered stuck (default: 60s) stuckThreshold :: !NominalDiffTime }@@ -33,12 +43,15 @@ defaultConfig :: MetricsServerConfig defaultConfig = MetricsServerConfig- { port = 9090,+ { host = "127.0.0.1",+ port = 9090, enableJSON = True, enablePrometheus = True, enableWebSocket = True, wsPushIntervalUs = 100_000, -- 100ms wsMaxConnections = 100,+ wsMaxSubscriptions = 1_000, livenessTimeoutMicros = 1_000_000, -- 1 second+ dependencyTimeoutMicros = 1_000_000, -- 1 second per dependency stuckThreshold = 60 -- 60 seconds }
src/Shibuya/Metrics/Health.hs view
@@ -8,6 +8,7 @@ ( -- * Health Status Types LivenessStatus (..), ReadinessStatus (..),+ ApplicationStatus (..), ProcessorHealth (..), DependencyStatus (..), @@ -25,17 +26,32 @@ ) where +import Control.Exception+ ( SomeAsyncException,+ SomeException,+ displayException,+ fromException,+ tryJust,+ ) import Data.Aeson (ToJSON (..), object, (.=)) import Data.Map.Strict qualified as Map-import Data.Maybe (isJust) import Data.Text (Text)+import Data.Text qualified as Text import Data.Time.Clock (NominalDiffTime, UTCTime, diffUTCTime, getCurrentTime) import Shibuya.App (Master, getAllMetricsIO) import Shibuya.Core.Metrics ( MetricsMap,+ ProcessorId, ProcessorMetrics (..), ProcessorState (..), )+import Shibuya.Internal.Runner.Master+ ( LifecycleSnapshot,+ MasterPhase (..),+ ProcessorLifecycle (..),+ getLifecycleSnapshotIO,+ getMasterPhaseIO,+ ) import System.Timeout (timeout) --------------------------------------------------------------------------------@@ -46,6 +62,8 @@ data HealthConfig = HealthConfig { -- | Timeout for liveness check (microseconds) livenessTimeoutMicros :: !Int,+ -- | Timeout for each dependency check (microseconds)+ dependencyTimeoutMicros :: !Int, -- | How long a processor can be in Processing state before considered stuck stuckThreshold :: !NominalDiffTime }@@ -58,6 +76,7 @@ defaultHealthConfig = HealthConfig { livenessTimeoutMicros = 1_000_000,+ dependencyTimeoutMicros = 1_000_000, stuckThreshold = 60 } @@ -82,6 +101,7 @@ -- Indicates whether the system is ready to handle traffic. data ReadinessStatus = ReadinessStatus { ready :: !Bool,+ application :: !ApplicationStatus, processors :: !ProcessorHealth, dependencies :: ![DependencyStatus] }@@ -91,10 +111,31 @@ toJSON status = object [ "ready" .= status.ready,+ "application" .= status.application, "processors" .= status.processors, "dependencies" .= status.dependencies ] +-- | Health-level application lifecycle derived from the master phase and the+-- retained configured-processor lifecycle snapshot.+data ApplicationStatus+ = ConfiguredEmpty+ | Starting+ | Running+ | Draining+ | ApplicationStopped+ | ApplicationFailed+ deriving stock (Eq, Show)++instance ToJSON ApplicationStatus where+ toJSON = \case+ ConfiguredEmpty -> toJSON ("configured_empty" :: Text)+ Starting -> toJSON ("starting" :: Text)+ Running -> toJSON ("running" :: Text)+ Draining -> toJSON ("draining" :: Text)+ ApplicationStopped -> toJSON ("stopped" :: Text)+ ApplicationFailed -> toJSON ("failed" :: Text)+ -- | Summary of processor health across all processors. data ProcessorHealth = ProcessorHealth { total :: !Int,@@ -132,6 +173,8 @@ ] -- | A dependency check is an IO action that returns the dependency's status.+-- Synchronous exceptions become an unhealthy status; asynchronous exceptions+-- remain cancellation signals and are rethrown. type DependencyCheck = IO DependencyStatus --------------------------------------------------------------------------------@@ -142,9 +185,14 @@ -- This is a fast check suitable for Kubernetes liveness probes. checkLiveness :: HealthConfig -> Master -> IO LivenessStatus checkLiveness config master = do- -- Try to query metrics with timeout- result <- timeout config.livenessTimeoutMicros $ getAllMetricsIO master- pure $ LivenessStatus {alive = isJust result}+ result <- timeout config.livenessTimeoutMicros $ getMasterPhaseIO master+ pure $+ LivenessStatus+ { alive = case result of+ Just MasterStopped -> False+ Just _ -> True+ Nothing -> False+ } -- | Check readiness - are all processors healthy and dependencies available? -- This is suitable for Kubernetes readiness probes.@@ -156,17 +204,26 @@ checkReadiness config master depChecks = do now <- getCurrentTime metrics <- getAllMetricsIO master- let procHealth = analyzeProcessorHealth config now metrics- depStatus <- sequence depChecks+ lifecycles <- getLifecycleSnapshotIO master+ masterPhase <- getMasterPhaseIO master+ let procHealth = analyzeProcessorHealth config now metrics lifecycles+ application = classifyApplication masterPhase lifecycles+ allRunningVisible =+ all+ (\(pid, lifecycle) -> lifecycle /= LifecycleRunning || Map.member pid metrics)+ (Map.toList lifecycles)+ depStatus <- traverse (runDependencyCheck config) depChecks let allDepsHealthy = all (.healthy) depStatus noFailedProcessors = procHealth.failed == 0 noStuckProcessors = procHealth.stuck == 0- isReady = allDepsHealthy && noFailedProcessors && noStuckProcessors+ acceptsWork = application == Running || application == ConfiguredEmpty+ isReady = acceptsWork && allRunningVisible && allDepsHealthy && noFailedProcessors && noStuckProcessors pure ReadinessStatus { ready = isReady,+ application, processors = procHealth, dependencies = depStatus }@@ -188,11 +245,15 @@ -------------------------------------------------------------------------------- -- | Analyze processor health from metrics.-analyzeProcessorHealth :: HealthConfig -> UTCTime -> MetricsMap -> ProcessorHealth-analyzeProcessorHealth config now metrics =- let processors = Map.elems metrics- total = length processors- (healthy, failed, stuck) = foldr (categorize config now) (0, 0, 0) processors+analyzeProcessorHealth :: HealthConfig -> UTCTime -> MetricsMap -> LifecycleSnapshot -> ProcessorHealth+analyzeProcessorHealth config now metrics lifecycles =+ let processorIds = Map.keysSet metrics <> Map.keysSet lifecycles+ total = length processorIds+ (healthy, failed, stuck) =+ foldr+ (categorize config now metrics lifecycles)+ (0, 0, 0)+ processorIds in ProcessorHealth { total = total, healthy = healthy,@@ -204,16 +265,65 @@ categorize :: HealthConfig -> UTCTime ->- ProcessorMetrics ->+ MetricsMap ->+ LifecycleSnapshot ->+ ProcessorId -> (Int, Int, Int) -> (Int, Int, Int)-categorize config now pm (h, f, s) =- case pm.state of- Idle -> (h + 1, f, s)- Stopped -> (h, f, s) -- Stopped is neither healthy nor failed- Failed _ _ -> (h, f + 1, s)- Processing _ lastActivity ->- let timeSinceActivity = diffUTCTime now lastActivity- in if timeSinceActivity > config.stuckThreshold- then (h, f, s + 1) -- Stuck- else (h + 1, f, s) -- Healthy (actively processing)+categorize config now metrics lifecycles pid counts@(h, f, s) =+ case Map.lookup pid lifecycles of+ Just LifecycleFailed {} -> (h, f + 1, s)+ _ -> case Map.lookup pid metrics of+ Nothing -> counts+ Just pm -> case pm.state of+ Idle -> (h + 1, f, s)+ Stopped -> counts+ Failed _ _ -> (h, f + 1, s)+ Processing _ _ lastProgress ->+ let timeSinceProgress = diffUTCTime now lastProgress+ in if timeSinceProgress > config.stuckThreshold+ then (h, f, s + 1)+ else (h + 1, f, s)++classifyApplication :: MasterPhase -> LifecycleSnapshot -> ApplicationStatus+classifyApplication masterPhase lifecycles+ | any isFailed (Map.elems lifecycles) = ApplicationFailed+ | masterPhase == MasterStopped = ApplicationStopped+ | masterPhase == MasterDraining || any (== LifecycleDraining) (Map.elems lifecycles) = Draining+ | masterPhase == MasterStarting = Starting+ | Map.null lifecycles = ConfiguredEmpty+ | all (== LifecycleStopped) (Map.elems lifecycles) = ApplicationStopped+ | otherwise = Running+ where+ isFailed LifecycleFailed {} = True+ isFailed _ = False++runDependencyCheck :: HealthConfig -> DependencyCheck -> IO DependencyStatus+runDependencyCheck config check = do+ result <- timeout config.dependencyTimeoutMicros $ tryJust synchronousException check+ pure $ case result of+ Just (Right status) -> status+ Just (Left err) ->+ DependencyStatus+ { name = "unknown",+ healthy = False,+ latencyMs = Nothing,+ errorMsg = Just $ Text.pack $ displayException err+ }+ Nothing ->+ DependencyStatus+ { name = "unknown",+ healthy = False,+ latencyMs = Nothing,+ errorMsg =+ Just $+ "Dependency check timed out after "+ <> Text.pack (show config.dependencyTimeoutMicros)+ <> " microseconds"+ }++synchronousException :: SomeException -> Maybe SomeException+synchronousException exception =+ case fromException exception :: Maybe SomeAsyncException of+ Just _ -> Nothing+ Nothing -> Just exception
src/Shibuya/Metrics/Prometheus.hs view
@@ -94,11 +94,11 @@ -- | Convert processor state to integer for Prometheus. stateToInt :: ProcessorState -> Int stateToInt Idle = 1-stateToInt (Processing _ _) = 2+stateToInt (Processing _ _ _) = 2 stateToInt (Failed _ _) = 3 stateToInt Stopped = 4 -- | Get in-flight count from processor state. inFlightCount :: ProcessorState -> Int-inFlightCount (Processing info _) = info.inFlight+inFlightCount (Processing info _ _) = info.inFlight inFlightCount _ = 0
src/Shibuya/Metrics/Server.hs view
@@ -5,6 +5,7 @@ startMetricsServerWithDeps, stopMetricsServer, withMetricsServer,+ combinedApp, -- * Re-exports MetricsServer (..),@@ -15,8 +16,9 @@ where import Control.Concurrent.Async (async, cancel)-import Control.Exception (bracket)+import Control.Exception (bracket, finally) import Data.Aeson (encode, object, (.=))+import Data.String (fromString) import Data.Text (Text) import Network.HTTP.Types (hContentType, status404) import Network.Wai (Application, Response, pathInfo, responseLBS)@@ -29,7 +31,7 @@ import Shibuya.Metrics.JSON (jsonAppWithHealth) import Shibuya.Metrics.Prometheus (prometheusApp) import Shibuya.Metrics.Types (MetricsServer (..))-import Shibuya.Metrics.WebSocket (WebSocketState, newWebSocketState, websocketApp)+import Shibuya.Metrics.WebSocket (WebSocketState, newWebSocketState, shutdownWebSockets, websocketApp) -- | Start the metrics server without dependency checks. -- Returns a handle that can be used to stop the server.@@ -44,20 +46,33 @@ [DependencyCheck] -> IO MetricsServer startMetricsServerWithDeps config master depChecks = do+ validateConfig config wsState <- newWebSocketState config.wsMaxConnections let app = combinedApp config master wsState depChecks settings = Warp.setPort config.port $ Warp.setHost- "*"+ (fromString config.host) Warp.defaultSettings- serverAsync <- async $ Warp.runSettings settings app+ serverAsync <- async $ Warp.runSettings settings app `finally` shutdownWebSockets wsState pure MetricsServer { serverThread = serverAsync, serverPort = config.port } +validateConfig :: MetricsServerConfig -> IO ()+validateConfig config+ | null config.host = fail "MetricsServerConfig.host must not be empty"+ | config.port < 0 = fail "MetricsServerConfig.port must be non-negative"+ | config.wsPushIntervalUs <= 0 = fail "MetricsServerConfig.wsPushIntervalUs must be positive"+ | config.wsMaxConnections <= 0 = fail "MetricsServerConfig.wsMaxConnections must be positive"+ | config.wsMaxSubscriptions <= 0 = fail "MetricsServerConfig.wsMaxSubscriptions must be positive"+ | config.livenessTimeoutMicros <= 0 = fail "MetricsServerConfig.livenessTimeoutMicros must be positive"+ | config.dependencyTimeoutMicros <= 0 = fail "MetricsServerConfig.dependencyTimeoutMicros must be positive"+ | config.stuckThreshold <= 0 = fail "MetricsServerConfig.stuckThreshold must be positive"+ | otherwise = pure ()+ -- | Stop the metrics server. stopMetricsServer :: MetricsServer -> IO () stopMetricsServer server = cancel server.serverThread@@ -81,11 +96,15 @@ [DependencyCheck] -> Application combinedApp config master wsState depChecks =- -- Handle WebSocket upgrade first- WaiWS.websocketsOr- WS.defaultConnectionOptions- (websocketApp config master wsState)- (httpApp config master depChecks)+ if config.enableWebSocket+ then+ WaiWS.websocketsOr+ WS.defaultConnectionOptions+ (websocketApp config master wsState)+ fallback+ else fallback+ where+ fallback = httpApp config master depChecks -- | HTTP application routing based on path. httpApp :: MetricsServerConfig -> Master -> [DependencyCheck] -> Application@@ -94,6 +113,7 @@ healthConfig = HealthConfig { livenessTimeoutMicros = config.livenessTimeoutMicros,+ dependencyTimeoutMicros = config.dependencyTimeoutMicros, stuckThreshold = config.stuckThreshold } jsonHandler = jsonAppWithHealth healthConfig master depChecks
src/Shibuya/Metrics/Types.hs view
@@ -2,6 +2,7 @@ module Shibuya.Metrics.Types ( -- * WebSocket Protocol ClientMessage (..),+ ProcessorTerminalStatus (..), ServerMessage (..), -- * Server Handle@@ -16,6 +17,7 @@ object, withObject, (.:),+ (.:?), (.=), ) import Data.Text (Text)@@ -64,16 +66,38 @@ ProcessorUpdate !ProcessorId !ProcessorMetrics | -- | Pong response to ping Pong+ | -- | A formerly visible processor reached a retained terminal state+ ProcessorTerminal !ProcessorId !ProcessorTerminalStatus | -- | Server is shutting down Goodbye deriving stock (Eq, Show, Generic) +-- | Retained terminal state sent when a processor leaves the live registry.+data ProcessorTerminalStatus+ = TerminalStopped+ | TerminalFailed !Text !(Maybe Text)+ deriving stock (Eq, Show, Generic)+ instance ToJSON ServerMessage where toJSON (MetricsSnapshot metrics) = object ["type" .= ("snapshot" :: Text), "metrics" .= metrics] toJSON (ProcessorUpdate pid pm) = object ["type" .= ("update" :: Text), "processor" .= pid, "metrics" .= pm] toJSON Pong = object ["type" .= ("pong" :: Text)]+ toJSON (ProcessorTerminal pid TerminalStopped) =+ object+ [ "type" .= ("terminal" :: Text),+ "processor" .= pid,+ "status" .= ("stopped" :: Text)+ ]+ toJSON (ProcessorTerminal pid (TerminalFailed failure messageId)) =+ object+ [ "type" .= ("terminal" :: Text),+ "processor" .= pid,+ "status" .= ("failed" :: Text),+ "error" .= failure,+ "messageId" .= messageId+ ] toJSON Goodbye = object ["type" .= ("goodbye" :: Text)] instance FromJSON ServerMessage where@@ -83,6 +107,14 @@ "snapshot" -> MetricsSnapshot <$> v .: "metrics" "update" -> ProcessorUpdate <$> v .: "processor" <*> v .: "metrics" "pong" -> pure Pong+ "terminal" -> do+ pid <- v .: "processor"+ status <- v .: "status"+ terminal <- case status :: Text of+ "stopped" -> pure TerminalStopped+ "failed" -> TerminalFailed <$> v .: "error" <*> v .:? "messageId"+ other -> fail $ "Unknown terminal status: " <> Text.unpack other+ pure $ ProcessorTerminal pid terminal "goodbye" -> pure Goodbye other -> fail $ "Unknown message type: " <> Text.unpack other
src/Shibuya/Metrics/WebSocket.hs view
@@ -3,31 +3,46 @@ ( websocketApp, WebSocketState (..), newWebSocketState,+ shutdownWebSockets, ) where -import Control.Concurrent (threadDelay)-import Control.Concurrent.Async (async, cancel, link)+import Control.Concurrent.Async (race_) import Control.Concurrent.STM ( STM, TVar, atomically,+ check, modifyTVar', newTVarIO,+ orElse, readTVar,+ readTVarIO,+ registerDelay, writeTVar, )-import Control.Exception (finally)+import Control.Exception (catch, finally, mask, throwIO) import Control.Monad (forever, when) import Data.Aeson (decode, encode) import Data.Map.Strict qualified as Map import Data.Set (Set) import Data.Set qualified as Set+import Data.Text (Text)+import Data.Text qualified as Text import Network.WebSockets qualified as WS import Shibuya.App (Master, getAllMetricsIO) import Shibuya.Core.Metrics (MetricsMap, ProcessorId (..), ProcessorMetrics)+import Shibuya.Core.Types (MessageId (..))+import Shibuya.Internal.Runner.Master+ ( ProcessorLifecycle (..),+ getLifecycleSnapshotIO,+ ) import Shibuya.Metrics.Config (MetricsServerConfig (..))-import Shibuya.Metrics.Types (ClientMessage (..), ServerMessage (..))+import Shibuya.Metrics.Types+ ( ClientMessage (..),+ ProcessorTerminalStatus (..),+ ServerMessage (..),+ ) -------------------------------------------------------------------------------- -- WebSocket State@@ -38,28 +53,46 @@ { -- | Current number of connections connectionCount :: !(TVar Int), -- | Maximum allowed connections- maxConnections :: !Int+ maxConnections :: !Int,+ -- | Whether server shutdown has begun+ shutdownRequested :: !(TVar Bool) } -- | Create new WebSocket state. newWebSocketState :: Int -> IO WebSocketState newWebSocketState maxConns = do countVar <- newTVarIO 0+ shutdownVar <- newTVarIO False pure WebSocketState { connectionCount = countVar,- maxConnections = maxConns+ maxConnections = maxConns,+ shutdownRequested = shutdownVar } --- | Try to acquire a connection slot. Returns True if successful.-acquireConnection :: WebSocketState -> STM Bool+-- | Ask every active connection to send 'Goodbye' and finish.+shutdownWebSockets :: WebSocketState -> IO ()+shutdownWebSockets wsState =+ atomically $ writeTVar wsState.shutdownRequested True++data AcquireResult+ = Acquired+ | AtCapacity+ | ServerShuttingDown++-- | Try to acquire a connection slot.+acquireConnection :: WebSocketState -> STM AcquireResult acquireConnection wsState = do+ shuttingDown <- readTVar wsState.shutdownRequested count <- readTVar wsState.connectionCount- if count < wsState.maxConnections- then do- writeTVar wsState.connectionCount (count + 1)- pure True- else pure False+ if shuttingDown+ then pure ServerShuttingDown+ else+ if count >= wsState.maxConnections+ then pure AtCapacity+ else do+ writeTVar wsState.connectionCount (count + 1)+ pure Acquired -- | Release a connection slot. releaseConnection :: WebSocketState -> STM ()@@ -72,16 +105,20 @@ -- | State for a single WebSocket connection. data ConnectionState = ConnectionState- { -- | Subscribed processors (Nothing = all)- subscriptions :: !(TVar (Maybe (Set ProcessorId))),+ { -- | Processor selection, including exclusions from subscribe-all+ subscriptions :: !(TVar Subscription), -- | Last sent metrics for delta detection lastMetrics :: !(TVar MetricsMap) } +data Subscription+ = AllProcessors !(Set ProcessorId)+ | SelectedProcessors !(Set ProcessorId)+ -- | Create new connection state. newConnectionState :: IO ConnectionState newConnectionState = do- subsVar <- newTVarIO Nothing -- Start subscribed to all+ subsVar <- newTVarIO $ AllProcessors Set.empty lastVar <- newTVarIO Map.empty pure ConnectionState@@ -99,83 +136,103 @@ Master -> WebSocketState -> WS.ServerApp-websocketApp config master wsState pending = do- -- Try to acquire a connection slot- acquired <- atomically $ acquireConnection wsState- if not acquired- then WS.rejectRequest pending "Too many connections"- else do- conn <- WS.acceptRequest pending- -- Set up connection with ping/pong for keepalive- WS.withPingThread conn 30 (pure ()) $ do- connState <- newConnectionState- -- Send initial snapshot- metrics <- getAllMetricsIO master- WS.sendTextData conn $ encode $ MetricsSnapshot metrics- atomically $ writeTVar connState.lastMetrics metrics- -- Run receive and push loops concurrently- pushThread <- async $ pushLoop config master connState conn- link pushThread- finally- (receiveLoop master connState conn)- ( do- cancel pushThread- WS.sendTextData conn $ encode Goodbye- atomically $ releaseConnection wsState- )+websocketApp config master wsState pending =+ mask $ \restore -> do+ outcome <- atomically $ acquireConnection wsState+ case outcome of+ AtCapacity -> restore $ WS.rejectRequest pending "Too many connections"+ ServerShuttingDown -> restore $ WS.rejectRequest pending "Server shutting down"+ Acquired ->+ restore (serveConnection config master wsState pending `catch` normalPeerClosure)+ `finally` atomically (releaseConnection wsState) +normalPeerClosure :: WS.ConnectionException -> IO ()+normalPeerClosure = \case+ WS.ConnectionClosed -> pure ()+ WS.CloseRequest _ _ -> pure ()+ unexpected -> throwIO unexpected++serveConnection :: MetricsServerConfig -> Master -> WebSocketState -> WS.PendingConnection -> IO ()+serveConnection config master wsState pending = do+ conn <- WS.acceptRequest pending+ WS.withPingThread conn 30 (pure ()) $ do+ connState <- newConnectionState+ metrics <- getAllMetricsIO master+ WS.sendTextData conn $ encode $ MetricsSnapshot metrics+ atomically $ writeTVar connState.lastMetrics metrics+ race_+ (receiveLoop config master connState conn)+ (pushLoop config master wsState connState conn)+ -------------------------------------------------------------------------------- -- Receive Loop -------------------------------------------------------------------------------- -- | Handle incoming messages from client.-receiveLoop :: Master -> ConnectionState -> WS.Connection -> IO ()-receiveLoop master connState conn = forever $ do+receiveLoop :: MetricsServerConfig -> Master -> ConnectionState -> WS.Connection -> IO ()+receiveLoop config master connState conn = forever $ do msg <- WS.receiveData conn case decode msg of Nothing -> pure () -- Ignore invalid messages- Just clientMsg -> handleClientMessage master connState conn clientMsg+ Just clientMsg -> handleClientMessage config.wsMaxSubscriptions master connState conn clientMsg -- | Handle a client message. handleClientMessage ::+ Int -> Master -> ConnectionState -> WS.Connection -> ClientMessage -> IO ()-handleClientMessage master connState conn = \case+handleClientMessage maxSubscriptions master connState conn = \case SubscribeAll -> do- atomically $ writeTVar connState.subscriptions Nothing+ atomically $ writeTVar connState.subscriptions $ AllProcessors Set.empty -- Send snapshot of all metrics metrics <- getAllMetricsIO master WS.sendTextData conn $ encode $ MetricsSnapshot metrics atomically $ writeTVar connState.lastMetrics metrics Subscribe pids -> do- atomically $ do+ subscription <- atomically $ do current <- readTVar connState.subscriptions let newSubs = case current of- Nothing -> Just $ Set.fromList pids- Just existing -> Just $ existing <> Set.fromList pids- writeTVar connState.subscriptions newSubs- -- Send snapshot of subscribed processors- allMetrics <- getAllMetricsIO master- let filtered = Map.filterWithKey (\pid _ -> pid `elem` pids) allMetrics- WS.sendTextData conn $ encode $ MetricsSnapshot filtered+ AllProcessors _ -> SelectedProcessors $ Set.fromList pids+ SelectedProcessors existing -> SelectedProcessors $ existing <> Set.fromList pids+ if subscriptionSize newSubs > maxSubscriptions+ then pure Nothing+ else writeTVar connState.subscriptions newSubs >> pure (Just newSubs)+ case subscription of+ Nothing -> rejectOversizedSubscription conn maxSubscriptions+ Just accepted -> do+ allMetrics <- getAllMetricsIO master+ let filtered = filterMetrics accepted allMetrics+ WS.sendTextData conn $ encode $ MetricsSnapshot filtered+ atomically $ writeTVar connState.lastMetrics filtered Unsubscribe pids -> do- atomically $ do+ accepted <- atomically $ do current <- readTVar connState.subscriptions- case current of- Nothing -> do- -- Was subscribed to all, now remove these- -- We need all processor IDs to calculate the new set- pure () -- Keep as Nothing, will filter in push- Just existing ->- writeTVar connState.subscriptions $- Just $- Set.difference existing (Set.fromList pids)+ let removed = Set.fromList pids+ newSubs = case current of+ AllProcessors excluded -> AllProcessors $ excluded <> removed+ SelectedProcessors existing -> SelectedProcessors $ Set.difference existing removed+ if subscriptionSize newSubs > maxSubscriptions+ then pure False+ else writeTVar connState.subscriptions newSubs >> pure True+ when (not accepted) $ rejectOversizedSubscription conn maxSubscriptions Ping -> WS.sendTextData conn $ encode Pong +subscriptionSize :: Subscription -> Int+subscriptionSize = \case+ AllProcessors excluded -> Set.size excluded+ SelectedProcessors selected -> Set.size selected++rejectOversizedSubscription :: WS.Connection -> Int -> IO ()+rejectOversizedSubscription conn limit =+ WS.sendCloseCode conn 1008 $+ "WebSocket processor subscription limit exceeded (maximum "+ <> Text.pack (show limit)+ <> ")"+ -------------------------------------------------------------------------------- -- Push Loop --------------------------------------------------------------------------------@@ -184,25 +241,59 @@ pushLoop :: MetricsServerConfig -> Master ->+ WebSocketState -> ConnectionState -> WS.Connection -> IO ()-pushLoop config master connState conn = forever $ do- threadDelay config.wsPushIntervalUs- -- Get current metrics+pushLoop config master wsState connState conn = loop+ where+ loop = do+ shuttingDown <- waitForPushOrShutdown config.wsPushIntervalUs wsState+ if shuttingDown+ then WS.sendTextData conn $ encode Goodbye+ else pushUpdates master connState conn >> loop++waitForPushOrShutdown :: Int -> WebSocketState -> IO Bool+waitForPushOrShutdown intervalUs wsState = do+ intervalElapsed <- registerDelay intervalUs+ atomically $+ (readTVar wsState.shutdownRequested >>= \requested -> check requested >> pure True)+ `orElse` (readTVar intervalElapsed >>= \elapsed -> check elapsed >> pure False)++pushUpdates :: Master -> ConnectionState -> WS.Connection -> IO ()+pushUpdates master connState conn = do currentMetrics <- getAllMetricsIO master- -- Get subscription filter- mSubs <- atomically $ readTVar connState.subscriptions- -- Get last sent metrics- lastSent <- atomically $ readTVar connState.lastMetrics- -- Filter metrics based on subscriptions- let filteredMetrics = case mSubs of- Nothing -> currentMetrics- Just subs -> Map.filterWithKey (\pid _ -> Set.member pid subs) currentMetrics- -- Send updates for changed processors+ lifecycle <- getLifecycleSnapshotIO master+ subscription <- readTVarIO connState.subscriptions+ lastSent <- readTVarIO connState.lastMetrics+ let filteredMetrics = filterMetrics subscription currentMetrics _ <- Map.traverseWithKey (sendIfChanged lastSent conn) filteredMetrics- -- Update last sent+ let removed = Map.keysSet lastSent `Set.difference` Map.keysSet currentMetrics+ mapM_ (sendTerminal lifecycle conn) $ Set.toList removed atomically $ writeTVar connState.lastMetrics filteredMetrics++filterMetrics :: Subscription -> MetricsMap -> MetricsMap+filterMetrics subscription =+ Map.filterWithKey $ \pid _ -> case subscription of+ AllProcessors excluded -> Set.notMember pid excluded+ SelectedProcessors selected -> Set.member pid selected++sendTerminal :: Map.Map ProcessorId ProcessorLifecycle -> WS.Connection -> ProcessorId -> IO ()+sendTerminal lifecycle conn pid =+ case Map.lookup pid lifecycle >>= terminalStatus of+ Nothing -> pure ()+ Just status -> WS.sendTextData conn $ encode $ ProcessorTerminal pid status++terminalStatus :: ProcessorLifecycle -> Maybe ProcessorTerminalStatus+terminalStatus = \case+ LifecycleStopped -> Just TerminalStopped+ LifecycleFailed failure messageId ->+ Just $ TerminalFailed failure (messageIdText <$> messageId)+ LifecycleRunning -> Nothing+ LifecycleDraining -> Nothing++messageIdText :: MessageId -> Text+messageIdText (MessageId value) = value -- | Send update if metrics have changed. sendIfChanged ::
+ test/Main.hs view
@@ -0,0 +1,18 @@+module Main (main) where++import Shibuya.Metrics.HealthSpec qualified+import Shibuya.Metrics.JSONSpec qualified+import Shibuya.Metrics.PrometheusSpec qualified+import Shibuya.Metrics.ServerSpec qualified+import Shibuya.Metrics.TypesSpec qualified+import Shibuya.Metrics.WebSocketSpec qualified+import Test.Hspec (hspec)++main :: IO ()+main = hspec $ do+ Shibuya.Metrics.ServerSpec.spec+ Shibuya.Metrics.JSONSpec.spec+ Shibuya.Metrics.PrometheusSpec.spec+ Shibuya.Metrics.TypesSpec.spec+ Shibuya.Metrics.WebSocketSpec.spec+ Shibuya.Metrics.HealthSpec.spec
+ test/Shibuya/Metrics/HealthSpec.hs view
@@ -0,0 +1,236 @@+module Shibuya.Metrics.HealthSpec (spec) where++import Control.Concurrent (newEmptyMVar, putMVar, takeMVar, threadDelay)+import Control.Concurrent.Async (async, cancel, waitCatch)+import Control.Concurrent.NQE.Supervisor (Strategy (IgnoreAll))+import Control.Exception (bracket, throwIO)+import Data.Atomics.Counter (readCounter)+import Data.Either (isLeft)+import Data.IORef (newIORef, readIORef, writeIORef)+import Data.Maybe (isJust)+import Data.Time.Clock (addUTCTime, getCurrentTime)+import Data.Word (Word64)+import Effectful (runEff)+import Shibuya.App (Master)+import Shibuya.Core.Metrics+ ( AckDecisionMetric (CountProcessed),+ HotCounters (..),+ InFlightInfo (..),+ MetricsHandle (..),+ ProcessorId (..),+ ProcessorMetrics (..),+ ProcessorState (..),+ beginProcessing,+ finishProcessing,+ newMetricsHandle,+ newMetricsHandleWithClock,+ sampleMetrics,+ )+import Shibuya.Internal.Runner.Master+ ( markMasterDraining,+ markProcessorFailed,+ registerProcessor,+ startMaster,+ stopMaster,+ unregisterProcessor,+ )+import Shibuya.Metrics.Health+ ( ApplicationStatus (..),+ DependencyStatus (..),+ HealthConfig (..),+ LivenessStatus (..),+ ProcessorHealth (..),+ ReadinessStatus (..),+ checkDetailedHealth,+ checkLiveness,+ checkReadiness,+ defaultHealthConfig,+ )+import Shibuya.Metrics.TestSupport (registerFailedProcessor, withMaster)+import System.Timeout (timeout)+import Test.Hspec (Spec, around, describe, it, shouldBe)++spec :: Spec+spec = do+ around withMaster $ describe "health characterization" $ do+ it "reports a running, intentionally empty master live and ready" $ \master -> do+ checkLiveness defaultHealthConfig master `shouldReturn` LivenessStatus {alive = True}+ checkReadiness defaultHealthConfig master []+ `shouldReturn` ReadinessStatus+ { ready = True,+ application = ConfiguredEmpty,+ processors = ProcessorHealth {total = 0, healthy = 0, failed = 0, stuck = 0},+ dependencies = []+ }++ it "reports a failed processor while it remains registered" $ \master -> do+ _ <- registerFailedProcessor master (ProcessorId "failed")+ readiness <- checkReadiness defaultHealthConfig master []+ readiness.ready `shouldBe` False+ readiness.application `shouldBe` Running+ readiness.processors `shouldBe` ProcessorHealth {total = 1, healthy = 0, failed = 1, stuck = 0}++ (detailed, metrics) <- checkDetailedHealth defaultHealthConfig master []+ detailed `shouldBe` readiness+ length metrics `shouldBe` 1++ it "reports an unhealthy dependency unready with its diagnostic fields" $ \master -> do+ let dependency =+ DependencyStatus+ { name = "database",+ healthy = False,+ latencyMs = Just 7,+ errorMsg = Just "unavailable"+ }+ readiness <- checkReadiness defaultHealthConfig master [pure dependency]+ readiness.ready `shouldBe` False+ readiness.dependencies `shouldBe` [dependency]++ it "restamps separated bursts and reports progress independently" $ \master -> do+ now <- getCurrentTime+ clock <- newIORef 0+ handle <- newMetricsHandleWithClock (readIORef clock) (addUTCTime (-120) now)+ runEff $ registerProcessor master (ProcessorId "bursts") handle++ _ <- beginProcessing handle 1+ writeIORef clock (seconds 1)+ finishProcessing handle (Right CountProcessed)+ readIORef handle.stateActiveRef `shouldReturn` False++ writeIORef clock (seconds 119)+ _ <- beginProcessing handle 1+ metrics <- sampleMetrics handle+ case metrics.state of+ Processing (InFlightInfo 1 1) burstStarted lastProgress -> do+ burstStarted `shouldBe` lastProgress+ lastProgress `shouldBe` addUTCTime (-1) now+ other -> fail $ "expected processing metrics, got " <> show other++ readiness <- checkReadiness defaultHealthConfig master []+ readiness.ready `shouldBe` True++ it "keeps sustained concurrent progress ready after the burst threshold" $ \master -> do+ now <- getCurrentTime+ clock <- newIORef 0+ handle <- newMetricsHandleWithClock (readIORef clock) (addUTCTime (-120) now)+ runEff $ registerProcessor master (ProcessorId "sustained") handle++ _ <- beginProcessing handle 2+ writeIORef clock (seconds 10)+ _ <- beginProcessing handle 2+ writeIORef clock (seconds 30)+ finishProcessing handle (Right CountProcessed)+ writeIORef clock (seconds 119)+ _ <- beginProcessing handle 2++ readiness <- checkReadiness defaultHealthConfig master []+ readiness.ready `shouldBe` True+ readiness.processors.stuck `shouldBe` 0++ it "reports a genuinely non-progressing handler stuck" $ \master -> do+ now <- getCurrentTime+ clock <- newIORef 0+ handle <- newMetricsHandleWithClock (readIORef clock) (addUTCTime (-120) now)+ runEff $ registerProcessor master (ProcessorId "stuck") handle+ _ <- beginProcessing handle 1+ _ <- sampleMetrics handle+ writeIORef clock (seconds 120)++ readiness <- checkReadiness defaultHealthConfig master []+ readiness.ready `shouldBe` False+ readiness.processors.stuck `shouldBe` 1++ it "never lets duplicate completion drive in-flight below zero" $ \master -> do+ handle <- registerTestHandle master (ProcessorId "floor")+ _ <- beginProcessing handle 1+ finishProcessing handle (Right CountProcessed)+ finishProcessing handle (Right CountProcessed)+ readCounter handle.hot.inFlight `shouldReturn` 0++ it "retains a configured processor failure after metrics unregister" $ \master -> do+ _ <- registerTestHandle master (ProcessorId "failed-and-gone")+ runEff $ do+ markProcessorFailed master (ProcessorId "failed-and-gone") "boom" Nothing+ unregisterProcessor master (ProcessorId "failed-and-gone")+ readiness <- checkReadiness defaultHealthConfig master []+ readiness.ready `shouldBe` False+ readiness.application `shouldBe` ApplicationFailed+ readiness.processors `shouldBe` ProcessorHealth {total = 1, healthy = 0, failed = 1, stuck = 0}++ it "reports draining and stopped masters unavailable" $ \master -> do+ runEff $ markMasterDraining master+ draining <- checkReadiness defaultHealthConfig master []+ draining.ready `shouldBe` False+ draining.application `shouldBe` Draining++ runEff $ stopMaster master+ stopped <- checkReadiness defaultHealthConfig master []+ stopped.ready `shouldBe` False+ stopped.application `shouldBe` ApplicationStopped+ checkLiveness defaultHealthConfig master `shouldReturn` LivenessStatus {alive = False}++ it "bounds each hung dependency check" $ \master -> do+ let config = defaultHealthConfig {dependencyTimeoutMicros = 10_000}+ result <-+ timeout 100_000 $+ checkReadiness config master [threadDelay 5_000_000 >> pure healthyDependency]+ result `shouldSatisfy` isJust+ case result of+ Just readiness -> do+ readiness.ready `shouldBe` False+ readiness.dependencies+ `shouldBe` [DependencyStatus "unknown" False Nothing (Just "Dependency check timed out after 10000 microseconds")]+ Nothing -> fail "health check exceeded its dependency timeout"++ it "normalizes a synchronous dependency exception" $ \master -> do+ readiness <-+ checkReadiness defaultHealthConfig master [throwIO $ userError "database exploded"]+ readiness.ready `shouldBe` False+ readiness.dependencies+ `shouldBe` [DependencyStatus "unknown" False Nothing (Just "user error (database exploded)")]++ it "preserves asynchronous cancellation of a dependency check" $ \master -> do+ started <- newEmptyMVar+ worker <- async $ checkReadiness defaultHealthConfig master [putMVar started () >> threadDelay 5_000_000 >> pure healthyDependency]+ takeMVar started+ cancel worker+ waitCatch worker >>= (`shouldSatisfy` isLeft)++ it "keeps repeated master stop observable and idempotent" $ \master -> do+ runEff $ stopMaster master+ runEff $ stopMaster master+ checkLiveness defaultHealthConfig master `shouldReturn` LivenessStatus {alive = False}+ readiness <- checkReadiness defaultHealthConfig master []+ readiness.application `shouldBe` ApplicationStopped+ readiness.ready `shouldBe` False++ describe "starting health" $+ it "distinguishes a starting master from a configured-empty running master" $+ bracket+ (runEff $ startMaster IgnoreAll)+ (\master -> runEff $ stopMaster master)+ ( \master -> do+ readiness <- checkReadiness defaultHealthConfig master []+ readiness.ready `shouldBe` False+ readiness.application `shouldBe` Starting+ )++shouldReturn :: (Eq a, Show a) => IO a -> a -> IO ()+shouldReturn action expected = action >>= (`shouldBe` expected)++registerTestHandle :: Master -> ProcessorId -> IO MetricsHandle+registerTestHandle master pid = do+ now <- getCurrentTime+ handle <- newMetricsHandle now+ runEff $ registerProcessor master pid handle+ pure handle++seconds :: Word64 -> Word64+seconds value = value * 1_000_000_000++healthyDependency :: DependencyStatus+healthyDependency = DependencyStatus "hung" True Nothing Nothing++shouldSatisfy :: (Show a) => a -> (a -> Bool) -> IO ()+shouldSatisfy actual predicate =+ if predicate actual then pure () else fail $ "predicate failed for " <> show actual
+ test/Shibuya/Metrics/JSONSpec.hs view
@@ -0,0 +1,11 @@+module Shibuya.Metrics.JSONSpec (spec) where++import Data.Aeson (encode)+import Shibuya.Metrics.TestSupport (assertGolden, fixtureMetrics)+import Test.Hspec (Spec, describe, it)++spec :: Spec+spec =+ describe "JSON wire contract" $+ it "matches the golden encoding for all four processor states" $+ assertGolden "processor-metrics.json.golden" (encode fixtureMetrics <> "\n")
+ test/Shibuya/Metrics/PrometheusSpec.hs view
@@ -0,0 +1,24 @@+module Shibuya.Metrics.PrometheusSpec (spec) where++import Network.HTTP.Types (status200)+import Network.Wai.Test (SResponse (..))+import Shibuya.Metrics.Config (MetricsServerConfig (..), defaultConfig)+import Shibuya.Metrics.Server (combinedApp)+import Shibuya.Metrics.TestSupport+ ( assertGolden,+ getResponse,+ registerPrometheusFixtures,+ withMaster,+ )+import Shibuya.Metrics.WebSocket (newWebSocketState)+import Test.Hspec (Spec, around, describe, it, shouldBe)++spec :: Spec+spec = around withMaster $+ describe "Prometheus wire contract" $+ it "matches the golden series names, labels, state values, and counters" $ \master -> do+ registerPrometheusFixtures master+ wsState <- newWebSocketState defaultConfig.wsMaxConnections+ response <- getResponse (combinedApp defaultConfig master wsState []) "/metrics/prometheus"+ response.simpleStatus `shouldBe` status200+ assertGolden "prometheus.golden" response.simpleBody
+ test/Shibuya/Metrics/ServerSpec.hs view
@@ -0,0 +1,116 @@+module Shibuya.Metrics.ServerSpec (spec) where++import Control.Exception (SomeException, throwIO)+import Data.Aeson (decode, object, (.=))+import Data.ByteString (ByteString)+import Data.ByteString.Lazy qualified as LBS+import Network.HTTP.Types (Status, hContentType, status200, status404, status503)+import Network.Wai (Application)+import Network.Wai.Test (SResponse (..))+import Shibuya.App (Master)+import Shibuya.Core.Metrics (ProcessorId (..))+import Shibuya.Metrics.Config (MetricsServerConfig (..), defaultConfig)+import Shibuya.Metrics.Health (DependencyCheck, DependencyStatus (..))+import Shibuya.Metrics.Server (combinedApp, startMetricsServer)+import Shibuya.Metrics.TestSupport+ ( getResponse,+ registerIdleProcessor,+ withMaster,+ )+import Shibuya.Metrics.WebSocket (newWebSocketState)+import Test.Hspec (Spec, anyException, around, describe, it, shouldBe, shouldSatisfy, shouldThrow)++spec :: Spec+spec = around withMaster $ do+ describe "built-in server configuration" $ do+ it "defaults to a loopback-only listener" $ \_ ->+ defaultConfig.host `shouldBe` "127.0.0.1"++ it "rejects nonpositive WebSocket resource limits" $ \master ->+ startMetricsServer defaultConfig {wsMaxSubscriptions = 0} master+ `shouldThrow` (anyException :: SomeException -> Bool)++ describe "combinedApp HTTP routes" $ do+ it "serves every enabled JSON, health, Prometheus, and WebSocket path" $ \master -> do+ _ <- registerIdleProcessor master (ProcessorId "known")+ app <- appFor defaultConfig master []++ assertResponse app "/metrics" status200 (Just "application/json")+ assertResponse app "/metrics/known" status200 (Just "application/json")+ assertResponse app "/health" status200 (Just "application/json")+ assertResponse app "/health/live" status200 (Just "application/json")+ assertResponse app "/health/ready" status200 (Just "application/json")+ assertResponse app "/metrics/prometheus" status200 (Just "text/plain; version=0.0.4; charset=utf-8")++ wsResponse <- getResponse app "/ws"+ wsResponse.simpleStatus `shouldBe` status404+ decode wsResponse.simpleBody+ `shouldBe` Just (object ["error" .= ("WebSocket endpoint - use ws:// protocol" :: String)])++ it "returns the published JSON error for an unknown processor" $ \master -> do+ app <- appFor defaultConfig master []+ response <- getResponse app "/metrics/missing"+ response.simpleStatus `shouldBe` status404+ decode response.simpleBody+ `shouldBe` Just+ ( object+ [ "error" .= ("Processor not found" :: String),+ "processor" .= ("missing" :: String)+ ]+ )++ it "returns the published JSON error for an unknown path" $ \master -> do+ app <- appFor defaultConfig master []+ response <- getResponse app "/unknown"+ response.simpleStatus `shouldBe` status404+ decode response.simpleBody+ `shouldBe` Just (object ["error" .= ("Not found" :: String)])++ it "returns 404 for every JSON route when JSON endpoints are disabled" $ \master -> do+ app <- appFor defaultConfig {enableJSON = False} master []+ mapM_+ (\path -> assertResponse app path status404 (Just "application/json"))+ ["/metrics", "/metrics/known", "/health", "/health/live", "/health/ready"]++ it "returns 404 when Prometheus is disabled" $ \master -> do+ app <- appFor defaultConfig {enablePrometheus = False} master []+ assertResponse app "/metrics/prometheus" status404 (Just "application/json")++ it "uses the generic 404 for plain HTTP when WebSockets are disabled" $ \master -> do+ app <- appFor defaultConfig {enableWebSocket = False} master []+ response <- getResponse app "/ws"+ response.simpleStatus `shouldBe` status404+ decode response.simpleBody+ `shouldBe` Just (object ["error" .= ("Not found" :: String)])++ it "returns 503 from readiness and detailed health for an unhealthy dependency" $ \master -> do+ app <- appFor defaultConfig master [failingDependency]+ assertResponse app "/health" status503 (Just "application/json")+ assertResponse app "/health/ready" status503 (Just "application/json")++ it "returns 503 when a dependency check throws synchronously" $ \master -> do+ app <- appFor defaultConfig master [throwIO $ userError "database exploded"]+ assertResponse app "/health" status503 (Just "application/json")+ assertResponse app "/health/ready" status503 (Just "application/json")++appFor :: MetricsServerConfig -> Master -> [DependencyCheck] -> IO Application+appFor config master dependencies = do+ wsState <- newWebSocketState config.wsMaxConnections+ pure $ combinedApp config master wsState dependencies++assertResponse :: Application -> ByteString -> Status -> Maybe ByteString -> IO ()+assertResponse app path expectedStatus expectedContentType = do+ response <- getResponse app path+ response.simpleStatus `shouldBe` expectedStatus+ lookup hContentType response.simpleHeaders `shouldBe` expectedContentType+ response.simpleBody `shouldSatisfy` (not . LBS.null)++failingDependency :: DependencyCheck+failingDependency =+ pure+ DependencyStatus+ { name = "database",+ healthy = False,+ latencyMs = Just 7,+ errorMsg = Just "unavailable"+ }
+ test/Shibuya/Metrics/TestSupport.hs view
@@ -0,0 +1,122 @@+module Shibuya.Metrics.TestSupport+ ( fixedTime,+ fixtureMetrics,+ fixtureProcessor,+ withMaster,+ registerIdleProcessor,+ registerFailedProcessor,+ registerPrometheusFixtures,+ getResponse,+ assertGolden,+ )+where++import Control.Concurrent.NQE.Supervisor (Strategy (IgnoreAll))+import Control.Concurrent.STM (atomically, modifyTVar')+import Control.Exception (bracket)+import Control.Monad (replicateM_)+import Data.ByteString (ByteString)+import Data.ByteString.Lazy qualified as LBS+import Data.Map.Strict qualified as Map+import Data.Time.Calendar (fromGregorian)+import Data.Time.Clock (UTCTime (..))+import Effectful (runEff)+import Network.Wai (Application)+import Network.Wai.Test (SResponse, defaultRequest, request, runSession, setPath)+import Paths_shibuya_metrics (getDataFileName)+import Shibuya.App (Master)+import Shibuya.Core.Metrics+ ( AckDecisionMetric (CountProcessed),+ BatchStats (..),+ InFlightInfo (..),+ MetricsHandle (..),+ MetricsMap,+ ProcessorId (..),+ ProcessorMetrics (..),+ ProcessorState (..),+ StreamStats (..),+ beginProcessing,+ finishProcessing,+ incrementReceived,+ newMetricsHandle,+ )+import Shibuya.Internal.Runner.Master+ ( markMasterRunning,+ registerProcessor,+ startMaster,+ stopMaster,+ )+import Test.Hspec (Expectation, shouldBe)++fixedTime :: UTCTime+fixedTime = UTCTime (fromGregorian 2026 9 20) 12_345++fixtureMetrics :: MetricsMap+fixtureMetrics =+ Map.fromList+ [ (ProcessorId "failed", fixtureProcessor (Failed "boom" (at 90)) 20),+ (ProcessorId "idle", fixtureProcessor Idle 0),+ (ProcessorId "processing", fixtureProcessor (Processing (InFlightInfo 2 4) (at 60) (at 75)) 10),+ (ProcessorId "stopped", fixtureProcessor Stopped 30)+ ]+ where+ at seconds = fixedTime {utctDayTime = seconds}++fixtureProcessor :: ProcessorState -> Int -> ProcessorMetrics+fixtureProcessor state offset =+ ProcessorMetrics+ { state,+ stats = StreamStats (offset + 1) (offset + 2) (offset + 3),+ batch = BatchStats (offset + 4) (offset + 5) (offset + 6) (offset + 7) (offset + 8) (offset + 9),+ startedAt = fixedTime+ }++withMaster :: (Master -> IO a) -> IO a+withMaster = bracket acquire release+ where+ acquire = runEff $ do+ master <- startMaster IgnoreAll+ markMasterRunning master+ pure master+ release master = runEff $ stopMaster master++registerIdleProcessor :: Master -> ProcessorId -> IO MetricsHandle+registerIdleProcessor master pid = do+ handle <- newMetricsHandle fixedTime+ runEff $ registerProcessor master pid handle+ pure handle++registerFailedProcessor :: Master -> ProcessorId -> IO MetricsHandle+registerFailedProcessor master pid = do+ handle <- registerIdleProcessor master pid+ _ <- beginProcessing handle 1+ finishProcessing handle (Left "fixture failure")+ pure handle++registerPrometheusFixtures :: Master -> IO ()+registerPrometheusFixtures master = do+ _ <- registerIdleProcessor master (ProcessorId "idle")++ processing <- registerIdleProcessor master (ProcessorId "processing")+ replicateM_ 2 $ incrementReceived processing+ _ <- beginProcessing processing 4++ failed <- registerIdleProcessor master (ProcessorId "failed")+ incrementReceived failed+ _ <- beginProcessing failed 1+ finishProcessing failed (Left "fixture failure")++ stopped <- registerIdleProcessor master (ProcessorId "stopped")+ replicateM_ 4 $ incrementReceived stopped+ _ <- beginProcessing stopped 1+ finishProcessing stopped (Right CountProcessed)+ atomically $ modifyTVar' stopped.cold $ \metrics -> metrics {state = Stopped}++getResponse :: Application -> ByteString -> IO SResponse+getResponse app path = runSession (request $ setPath defaultRequest path) app++assertGolden :: FilePath -> LBS.ByteString -> Expectation+assertGolden name actual = do+ path <- getDataFileName $ "test/golden/" <> name+ expected <- LBS.readFile path+ actual `shouldBe` expected
+ test/Shibuya/Metrics/TypesSpec.hs view
@@ -0,0 +1,67 @@+module Shibuya.Metrics.TypesSpec (spec) where++import Data.Aeson (FromJSON, ToJSON, Value, eitherDecode, encode, object, toJSON, (.=))+import Data.Map.Strict qualified as Map+import Shibuya.Core.Metrics (ProcessorId (..))+import Shibuya.Metrics.TestSupport (fixtureMetrics)+import Shibuya.Metrics.Types+ ( ClientMessage (..),+ ProcessorTerminalStatus (..),+ ServerMessage (..),+ )+import Test.Hspec (Spec, describe, it, shouldBe, shouldSatisfy)++spec :: Spec+spec = do+ describe "ClientMessage JSON contract" $ do+ messageCase SubscribeAll $ object ["type" .= ("subscribe_all" :: String)]+ messageCase (Subscribe [ProcessorId "alpha"]) $+ object ["type" .= ("subscribe" :: String), "processors" .= ["alpha" :: String]]+ messageCase (Unsubscribe [ProcessorId "alpha"]) $+ object ["type" .= ("unsubscribe" :: String), "processors" .= ["alpha" :: String]]+ messageCase Ping $ object ["type" .= ("ping" :: String)]++ it "rejects an unknown client message tag" $+ (eitherDecode "{\"type\":\"unknown\"}" :: Either String ClientMessage)+ `shouldSatisfy` isLeft++ describe "ServerMessage JSON contract" $ do+ let processing = fixtureMetrics Map.! ProcessorId "processing"+ messageCase (MetricsSnapshot fixtureMetrics) $+ object ["type" .= ("snapshot" :: String), "metrics" .= fixtureMetrics]+ messageCase (ProcessorUpdate (ProcessorId "processing") processing) $+ object+ [ "type" .= ("update" :: String),+ "processor" .= ("processing" :: String),+ "metrics" .= processing+ ]+ messageCase Pong $ object ["type" .= ("pong" :: String)]+ messageCase (ProcessorTerminal (ProcessorId "alpha") TerminalStopped) $+ object+ [ "type" .= ("terminal" :: String),+ "processor" .= ("alpha" :: String),+ "status" .= ("stopped" :: String)+ ]+ messageCase+ (ProcessorTerminal (ProcessorId "alpha") (TerminalFailed "boom" (Just "message-1")))+ $ object+ [ "type" .= ("terminal" :: String),+ "processor" .= ("alpha" :: String),+ "status" .= ("failed" :: String),+ "error" .= ("boom" :: String),+ "messageId" .= (Just "message-1" :: Maybe String)+ ]+ messageCase Goodbye $ object ["type" .= ("goodbye" :: String)]++ it "rejects an unknown server message tag" $+ (eitherDecode "{\"type\":\"unknown\"}" :: Either String ServerMessage)+ `shouldSatisfy` isLeft++messageCase :: (Eq a, Show a, ToJSON a, FromJSON a) => a -> Value -> Spec+messageCase message expected =+ it (show message) $ do+ toJSON message `shouldBe` expected+ eitherDecode (encode message) `shouldBe` Right message++isLeft :: Either a b -> Bool+isLeft = \case Left _ -> True; Right _ -> False
+ test/Shibuya/Metrics/WebSocketSpec.hs view
@@ -0,0 +1,237 @@+module Shibuya.Metrics.WebSocketSpec (spec) where++import Control.Concurrent (MVar, newEmptyMVar, putMVar, takeMVar)+import Control.Concurrent.Async (async, wait)+import Control.Concurrent.STM (atomically, check, readTVar)+import Control.Exception (SomeException, throwIO, try)+import Data.Aeson (eitherDecode, encode)+import Data.ByteString.Lazy (ByteString)+import Data.IORef (atomicModifyIORef', newIORef)+import Data.Map.Strict qualified as Map+import Effectful (runEff)+import Network.Wai.Handler.Warp qualified as Warp+import Network.WebSockets qualified as WS+import Shibuya.App (Master, getAllMetricsIO)+import Shibuya.Core.Metrics+ ( MetricsMap,+ ProcessorId (..),+ ProcessorMetrics (..),+ StreamStats (..),+ beginProcessing,+ incrementReceived,+ newMetricsHandleWithClock,+ )+import Shibuya.Internal.Runner.Master (markProcessorFailedIO, registerProcessor, unregisterProcessor)+import Shibuya.Metrics.Config (MetricsServerConfig (..), defaultConfig)+import Shibuya.Metrics.Server (combinedApp)+import Shibuya.Metrics.TestSupport+ ( fixedTime,+ registerIdleProcessor,+ withMaster,+ )+import Shibuya.Metrics.Types+ ( ClientMessage (..),+ ProcessorTerminalStatus (..),+ ServerMessage (..),+ )+import Shibuya.Metrics.WebSocket+ ( WebSocketState (..),+ newWebSocketState,+ shutdownWebSockets,+ )+import System.Timeout (timeout)+import Test.Hspec+ ( Spec,+ anyException,+ around,+ describe,+ expectationFailure,+ it,+ shouldBe,+ shouldThrow,+ )++spec :: Spec+spec = around withMaster $ do+ describe "WebSocket wire protocol" $ do+ it "sends an initial snapshot" $ \master -> do+ _ <- registerIdleProcessor master (ProcessorId "alpha")+ withServer defaultConfig master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ expected <- MetricsSnapshot <$> currentSnapshot master+ receiveServer conn `shouldReturn` expected++ it "answers subscribe_all and selective subscribe with snapshots" $ \master -> do+ _ <- registerIdleProcessor master (ProcessorId "alpha")+ _ <- registerIdleProcessor master (ProcessorId "beta")+ withServer defaultConfig master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn++ WS.sendTextData conn $ encode SubscribeAll+ allSnapshot <- receiveServer conn+ expected <- MetricsSnapshot <$> currentSnapshot master+ allSnapshot `shouldBe` expected++ WS.sendTextData conn $ encode $ Subscribe [ProcessorId "alpha"]+ selective <- receiveServer conn+ case selective of+ MetricsSnapshot metrics -> Map.keys metrics `shouldBe` [ProcessorId "alpha"]+ other -> expectationFailure $ "expected selective snapshot, got " <> show other++ it "answers ping with pong" $ \master ->+ withServer defaultConfig master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ WS.sendTextData conn $ encode Ping+ receiveServer conn `shouldReturn` Pong++ it "pushes an update after metrics change" $ \master -> do+ handle <- registerIdleProcessor master (ProcessorId "alpha")+ withServer fastConfig master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ incrementReceived handle+ receiveServer conn >>= \case+ ProcessorUpdate (ProcessorId "alpha") ProcessorMetrics {stats = StreamStats {received}} -> received `shouldBe` 1+ other -> expectationFailure $ "expected processor update, got " <> show other++ it "does not push an update when metrics have not changed" $ \master -> do+ _ <- registerIdleProcessor master (ProcessorId "alpha")+ withServer fastConfig master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ timeout 100_000 (receiveServer conn) `shouldReturn` Nothing++ it "rejects a connection beyond the configured limit" $ \master -> do+ let config = fastConfig {wsMaxConnections = 1}+ withServer config master $ \port -> do+ ready <- newEmptyMVar+ release <- newEmptyMVar+ first <- async $ holdConnection port ready release+ takeMVar ready+ (WS.runClient "127.0.0.1" port "/ws" $ \conn -> receiveServer conn)+ `shouldThrow` (anyException :: SomeException -> Bool)+ putMVar release ()+ wait first++ it "rejects WebSocket upgrades when disabled" $ \master ->+ withServer defaultConfig {enableWebSocket = False} master $ \port ->+ (WS.runClient "127.0.0.1" port "/ws" $ \conn -> receiveServer conn)+ `shouldThrow` (anyException :: SomeException -> Bool)++ it "closes clients that exceed the retained processor subscription limit" $ \master ->+ withServer fastConfig {wsMaxSubscriptions = 2} master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ WS.sendTextData conn $ encode $ Subscribe [ProcessorId "one", ProcessorId "two", ProcessorId "three"]+ (WS.receiveData conn :: IO ByteString)+ `shouldThrow` (anyException :: SomeException -> Bool)++ it "also bounds subscribe-all exclusion state" $ \master ->+ withServer fastConfig {wsMaxSubscriptions = 1} master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ WS.sendTextData conn $ encode $ Unsubscribe [ProcessorId "one", ProcessorId "two"]+ (WS.receiveData conn :: IO ByteString)+ `shouldThrow` (anyException :: SomeException -> Bool)++ it "excludes processors unsubscribed from subscribe-all" $ \master -> do+ alpha <- registerIdleProcessor master (ProcessorId "alpha")+ beta <- registerIdleProcessor master (ProcessorId "beta")+ withServer fastConfig master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ WS.sendTextData conn $ encode $ Unsubscribe [ProcessorId "alpha"]+ incrementReceived alpha+ incrementReceived beta+ receiveServer conn >>= \case+ ProcessorUpdate pid _ -> pid `shouldBe` ProcessorId "beta"+ other -> expectationFailure $ "expected beta update, got " <> show other++ it "restores a slot after a peer disconnects" $ \master -> do+ wsState <- newWebSocketState 1+ let app = combinedApp fastConfig master wsState []+ Warp.testWithApplication (pure app) $ \port -> do+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ pure ()+ released <- timeout 1_000_000 $ atomically $ do+ count <- readTVar wsState.connectionCount+ check $ count == 0+ released `shouldBe` Just ()++ it "restores a slot when initial snapshot generation fails" $ \master -> do+ clockCalls <- newIORef (0 :: Int)+ let failingClock = do+ call <- atomicModifyIORef' clockCalls $ \count -> (count + 1, count)+ if call == 0 then pure 0 else throwIO $ userError "snapshot failed"+ handle <- newMetricsHandleWithClock failingClock fixedTime+ _ <- beginProcessing handle 1+ runEff $ registerProcessor master (ProcessorId "broken") handle+ wsState <- newWebSocketState 1+ let app = combinedApp fastConfig master wsState []+ quietSettings = Warp.setOnException (\_ _ -> pure ()) Warp.defaultSettings+ Warp.withApplicationSettings quietSettings (pure app) $ \port -> do+ _ <-+ try (WS.runClient "127.0.0.1" port "/ws" receiveServer) ::+ IO (Either SomeException ServerMessage)+ released <- timeout 1_000_000 $ atomically $ do+ count <- readTVar wsState.connectionCount+ check $ count == 0+ released `shouldBe` Just ()++ it "sends goodbye when WebSocket shutdown is requested" $ \master -> do+ wsState <- newWebSocketState 1+ let app = combinedApp fastConfig master wsState []+ Warp.testWithApplication (pure app) $ \port -> do+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ shutdownWebSockets wsState+ shutdownWebSockets wsState+ receiveServer conn `shouldReturn` Goodbye+ released <- timeout 1_000_000 $ atomically $ do+ count <- readTVar wsState.connectionCount+ check $ count == 0+ released `shouldBe` Just ()++ it "reports a retained terminal failure once when a processor disappears" $ \master -> do+ _ <- registerIdleProcessor master (ProcessorId "alpha")+ withServer fastConfig master $ \port ->+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ markProcessorFailedIO master (ProcessorId "alpha") "boom" (Just "message-1")+ runEff $ unregisterProcessor master (ProcessorId "alpha")+ receiveServer conn+ `shouldReturn` ProcessorTerminal+ (ProcessorId "alpha")+ (TerminalFailed "boom" (Just "message-1"))+ timeout 50_000 (receiveServer conn) `shouldReturn` Nothing++fastConfig :: MetricsServerConfig+fastConfig = defaultConfig {wsPushIntervalUs = 10_000}++withServer :: MetricsServerConfig -> Master -> (Int -> IO a) -> IO a+withServer config master action = do+ wsState <- newWebSocketState config.wsMaxConnections+ Warp.testWithApplication (pure $ combinedApp config master wsState []) action++currentSnapshot :: Master -> IO MetricsMap+currentSnapshot = getAllMetricsIO++receiveServer :: WS.Connection -> IO ServerMessage+receiveServer conn = do+ payload <- WS.receiveData conn :: IO ByteString+ case eitherDecode payload of+ Left err -> expectationFailure err >> fail err+ Right message -> pure message++holdConnection :: Int -> MVar () -> MVar () -> IO ()+holdConnection port ready release =+ WS.runClient "127.0.0.1" port "/ws" $ \conn -> do+ _ <- receiveServer conn+ putMVar ready ()+ takeMVar release++shouldReturn :: (Eq a, Show a) => IO a -> a -> IO ()+shouldReturn action expected = action >>= (`shouldBe` expected)
+ test/golden/processor-metrics.json.golden view
@@ -0,0 +1,1 @@+{"failed":{"batch":{"batchedMessages":25,"batchesEmitted":24,"flushTriggered":29,"partialFailures":26,"sizeTriggered":27,"timeoutTriggered":28},"startedAt":"2026-09-20T03:25:45Z","state":{"error":"boom","status":"failed","timestamp":"2026-09-20T00:01:30Z"},"stats":{"failed":23,"processed":22,"received":21}},"idle":{"batch":{"batchedMessages":5,"batchesEmitted":4,"flushTriggered":9,"partialFailures":6,"sizeTriggered":7,"timeoutTriggered":8},"startedAt":"2026-09-20T03:25:45Z","state":{"status":"idle"},"stats":{"failed":3,"processed":2,"received":1}},"processing":{"batch":{"batchedMessages":15,"batchesEmitted":14,"flushTriggered":19,"partialFailures":16,"sizeTriggered":17,"timeoutTriggered":18},"startedAt":"2026-09-20T03:25:45Z","state":{"inFlight":2,"lastActivity":"2026-09-20T00:01:00Z","lastProgress":"2026-09-20T00:01:15Z","maxConcurrency":4,"status":"processing"},"stats":{"failed":13,"processed":12,"received":11}},"stopped":{"batch":{"batchedMessages":35,"batchesEmitted":34,"flushTriggered":39,"partialFailures":36,"sizeTriggered":37,"timeoutTriggered":38},"startedAt":"2026-09-20T03:25:45Z","state":{"status":"stopped"},"stats":{"failed":33,"processed":32,"received":31}}}
+ test/golden/prometheus.golden view
@@ -0,0 +1,30 @@+# HELP shibuya_messages_received_total Total messages received by processor+# TYPE shibuya_messages_received_total counter+shibuya_messages_received_total{processor="failed"} 1.0+shibuya_messages_received_total{processor="idle"} 0.0+shibuya_messages_received_total{processor="processing"} 2.0+shibuya_messages_received_total{processor="stopped"} 4.0+# HELP shibuya_messages_processed_total Total messages successfully processed+# TYPE shibuya_messages_processed_total counter+shibuya_messages_processed_total{processor="failed"} 0.0+shibuya_messages_processed_total{processor="idle"} 0.0+shibuya_messages_processed_total{processor="processing"} 0.0+shibuya_messages_processed_total{processor="stopped"} 1.0+# HELP shibuya_messages_failed_total Total messages that failed processing+# TYPE shibuya_messages_failed_total counter+shibuya_messages_failed_total{processor="failed"} 1.0+shibuya_messages_failed_total{processor="idle"} 0.0+shibuya_messages_failed_total{processor="processing"} 0.0+shibuya_messages_failed_total{processor="stopped"} 0.0+# HELP shibuya_processor_state Current processor state (1=idle, 2=processing, 3=failed, 4=stopped)+# TYPE shibuya_processor_state gauge+shibuya_processor_state{processor="failed"} 3.0+shibuya_processor_state{processor="idle"} 1.0+shibuya_processor_state{processor="processing"} 2.0+shibuya_processor_state{processor="stopped"} 4.0+# HELP shibuya_processor_in_flight Number of messages currently being processed+# TYPE shibuya_processor_in_flight gauge+shibuya_processor_in_flight{processor="failed"} 0.0+shibuya_processor_in_flight{processor="idle"} 0.0+shibuya_processor_in_flight{processor="processing"} 1.0+shibuya_processor_in_flight{processor="stopped"} 0.0