From 2a477b4977602dc296448045740245fa4ec0ce67 Mon Sep 17 00:00:00 2001 From: Mats Rauhala Date: Tue, 15 Sep 2026 16:34:35 +0300 Subject: [PATCH] runtime: increment trigger/service counters on the bus --- src/HomeAssistant/Runtime.hs | 7 ++++--- src/HomeAssistant/Runtime/Bus.hs | 16 ++++++++++++++-- src/HomeAssistant/Runtime/Connection.hs | 5 ++++- test/BusSpec.hs | 25 +++++++++++++++++++++++-- 4 files changed, 45 insertions(+), 8 deletions(-) diff --git a/src/HomeAssistant/Runtime.hs b/src/HomeAssistant/Runtime.hs index 7960bfa..bc99bc4 100644 --- a/src/HomeAssistant/Runtime.hs +++ b/src/HomeAssistant/Runtime.hs @@ -82,12 +82,13 @@ runController rootDir bus (Controller name machine _enabled) = do defaultMain :: IO () defaultMain = withSocketsDo $ do severity <- maybe InfoS (const DebugS) <$> lookupEnv "HA_DEBUG" - withBus severity $ \bus -> do + store <- System.Metrics.newStore + System.Metrics.registerGcMetrics store + appMetrics <- HomeAssistant.Runtime.Metrics.registerAppMetrics store + withBus severity appMetrics $ \bus -> do token <- getEnv "HA_TOKEN" host <- getEnv "HA_HOST" rootPath <- fromMaybe "/tmp/" <$> lookupEnv "HA_LIB_DIR" - store <- System.Metrics.newStore - System.Metrics.registerGcMetrics store rrdPath <- fromMaybe "hass-controller.rrd" <$> lookupEnv "HA_RRD_PATH" rrdtool <- fromMaybe "rrdtool" <$> lookupEnv "HA_RRDTOOL" let active = [c | c@(Controller _ _ True) <- controllers] diff --git a/src/HomeAssistant/Runtime/Bus.hs b/src/HomeAssistant/Runtime/Bus.hs index c02c2e1..23509f1 100644 --- a/src/HomeAssistant/Runtime/Bus.hs +++ b/src/HomeAssistant/Runtime/Bus.hs @@ -6,6 +6,8 @@ module HomeAssistant.Runtime.Bus , CallIdGen(..) , mkCallIdGen , withBus + , recordInbound + , recordOutbound , channelHassEval ) where @@ -21,10 +23,12 @@ import Control.Concurrent.STM import Data.Aeson (Value) import Data.IORef (atomicModifyIORef', newIORef) import HomeAssistant.Controller (HASSEff (..), Service) +import HomeAssistant.Runtime.Metrics (AppMetrics (..)) import Network.WebSockets (Connection) import Katip (LogEnv, closeScribes, mkHandleScribe, ColorStrategy (..), permitItem, Severity (..), Verbosity (V2), registerScribe, defaultScribeSettings, initLogEnv, ls, sl, logFM, katipAddContext, KatipContext) import Control.Exception (bracket) import System.IO (stdout) +import System.Metrics.Counter (inc) import Data.UUID (toText) import AFRP (Request(..), Event(..)) import Control.Monad.IO.Class (MonadIO, liftIO) @@ -38,10 +42,11 @@ data Bus = Bus , busConn :: TVar (Maybe Connection) , busGen :: CallIdGen , busLogEnv :: LogEnv + , busMetrics :: AppMetrics } -withBus :: Severity -> (Bus -> IO a) -> IO a -withBus severity callback = do +withBus :: Severity -> AppMetrics -> (Bus -> IO a) -> IO a +withBus severity metrics callback = do handleScribe <- mkHandleScribe ColorIfTerminal stdout (permitItem severity) V2 let makeLogEnv = registerScribe "stdout" handleScribe defaultScribeSettings =<< initLogEnv "hass-controller" "production" -- closeScribes will stop accepting new logs, flush existing ones and clean up resources @@ -52,8 +57,15 @@ withBus severity callback = do <*> newTVarIO Nothing <*> mkCallIdGen 0 <*> pure le + <*> pure metrics callback bus +recordInbound :: Bus -> IO () +recordInbound bus = inc (amTriggersIn (busMetrics bus)) + +recordOutbound :: Bus -> IO () +recordOutbound bus = inc (amServicesOut (busMetrics bus)) + channelHassEval :: (MonadIO m, KatipContext m) => Bus -> HASSEff a -> m a channelHassEval bus = \case CallService req svc -> katipAddContext (sl "traceId" (toText (requestTraceId req))) $ do diff --git a/src/HomeAssistant/Runtime/Connection.hs b/src/HomeAssistant/Runtime/Connection.hs index a4076dc..6ed8e78 100644 --- a/src/HomeAssistant/Runtime/Connection.hs +++ b/src/HomeAssistant/Runtime/Connection.hs @@ -96,7 +96,9 @@ receiveLoop bus conn = forever $ do Left () -> atomically $ writeTChan (busInbound bus) Tick Right msg -> case eitherDecode msg of Left err -> putStrLn $ "[reader] skipping undecodable message: " <> err - Right v -> atomically $ writeTChan (busInbound bus) (Event v) + Right v -> do + recordInbound bus + atomically $ writeTChan (busInbound bus) (Event v) receiveJSON :: WS.Connection -> IO Value receiveJSON conn = do @@ -128,6 +130,7 @@ sendWithId bus conn request svc = do let textData = encode $ encodeService callId svc runKatipContextT (busLogEnv bus) (sl "traceId" (toText (requestTraceId request))) "connection" $ logFM DebugS (ls textData) + recordOutbound bus WS.sendTextData conn textData writerAction :: Bus -> IO Void diff --git a/test/BusSpec.hs b/test/BusSpec.hs index f7cf2f2..a8c37d8 100644 --- a/test/BusSpec.hs +++ b/test/BusSpec.hs @@ -10,16 +10,19 @@ import Control.Concurrent.STM , writeTChan ) import Data.Aeson (Value (..)) +import qualified Data.HashMap.Strict as HM import Data.Time (UTCTime (..), utc) import Data.UUID (nil) import HomeAssistant.Controller (HASSEff (..), Service (..), Target(..)) import HomeAssistant.Runtime.Bus +import HomeAssistant.Runtime.Metrics (registerAppMetrics) import Katip (Namespace (Namespace), runKatipContextT, Severity (..)) +import qualified System.Metrics as Metrics import Test.Hspec spec :: Spec spec = describe "Bus" $ do - it "broadcasts inbound messages to every dup'd channel in order" $ withBus InfoS $ \bus -> do + it "broadcasts inbound messages to every dup'd channel in order" $ withTestBus $ \bus -> do p1 <- atomically $ dupTChan (busInbound bus) p2 <- atomically $ dupTChan (busInbound bus) atomically $ writeTChan (busInbound bus) (Event (Number 1)) @@ -29,7 +32,7 @@ spec = describe "Bus" $ do r1 `shouldBe` (Event (Number 1), Event (Number 2)) r2 `shouldBe` (Event (Number 1), Event (Number 2)) - it "channelHassEval writes CallService to the outbound channel" $ withBus InfoS $ \bus -> do + it "channelHassEval writes CallService to the outbound channel" $ withTestBus $ \bus -> do let req = Request (UTCTime (toEnum 0) 0) utc nil svc = Service "light" "turn_on" Nothing [EntityId "light.bedroom_masse"] runKatipContextT (busLogEnv bus) () (Namespace ["test"]) $ @@ -43,3 +46,21 @@ spec = describe "Bus" $ do a <- generateCallId gen b <- generateCallId gen (a, b) `shouldBe` (1, 2) + + describe "recordInbound / recordOutbound" $ do + it "increments the trigger and service counters" $ do + store <- Metrics.newStore + m <- registerAppMetrics store + withBus InfoS m $ \bus -> do + recordInbound bus + recordInbound bus + recordOutbound bus + sample <- Metrics.sampleAll store + HM.lookup "hass.trigger.in" sample `shouldBe` Just (Metrics.Counter 2) + HM.lookup "hass.service.out" sample `shouldBe` Just (Metrics.Counter 1) + +withTestBus :: (Bus -> IO a) -> IO a +withTestBus action = do + store <- Metrics.newStore + m <- registerAppMetrics store + withBus InfoS m action