runtime: increment trigger/service counters on the bus
This commit is contained in:
@@ -82,12 +82,13 @@ runController rootDir bus (Controller name machine _enabled) = do
|
|||||||
defaultMain :: IO ()
|
defaultMain :: IO ()
|
||||||
defaultMain = withSocketsDo $ do
|
defaultMain = withSocketsDo $ do
|
||||||
severity <- maybe InfoS (const DebugS) <$> lookupEnv "HA_DEBUG"
|
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"
|
token <- getEnv "HA_TOKEN"
|
||||||
host <- getEnv "HA_HOST"
|
host <- getEnv "HA_HOST"
|
||||||
rootPath <- fromMaybe "/tmp/" <$> lookupEnv "HA_LIB_DIR"
|
rootPath <- fromMaybe "/tmp/" <$> lookupEnv "HA_LIB_DIR"
|
||||||
store <- System.Metrics.newStore
|
|
||||||
System.Metrics.registerGcMetrics store
|
|
||||||
rrdPath <- fromMaybe "hass-controller.rrd" <$> lookupEnv "HA_RRD_PATH"
|
rrdPath <- fromMaybe "hass-controller.rrd" <$> lookupEnv "HA_RRD_PATH"
|
||||||
rrdtool <- fromMaybe "rrdtool" <$> lookupEnv "HA_RRDTOOL"
|
rrdtool <- fromMaybe "rrdtool" <$> lookupEnv "HA_RRDTOOL"
|
||||||
let active = [c | c@(Controller _ _ True) <- controllers]
|
let active = [c | c@(Controller _ _ True) <- controllers]
|
||||||
|
|||||||
@@ -6,6 +6,8 @@ module HomeAssistant.Runtime.Bus
|
|||||||
, CallIdGen(..)
|
, CallIdGen(..)
|
||||||
, mkCallIdGen
|
, mkCallIdGen
|
||||||
, withBus
|
, withBus
|
||||||
|
, recordInbound
|
||||||
|
, recordOutbound
|
||||||
, channelHassEval
|
, channelHassEval
|
||||||
) where
|
) where
|
||||||
|
|
||||||
@@ -21,10 +23,12 @@ import Control.Concurrent.STM
|
|||||||
import Data.Aeson (Value)
|
import Data.Aeson (Value)
|
||||||
import Data.IORef (atomicModifyIORef', newIORef)
|
import Data.IORef (atomicModifyIORef', newIORef)
|
||||||
import HomeAssistant.Controller (HASSEff (..), Service)
|
import HomeAssistant.Controller (HASSEff (..), Service)
|
||||||
|
import HomeAssistant.Runtime.Metrics (AppMetrics (..))
|
||||||
import Network.WebSockets (Connection)
|
import Network.WebSockets (Connection)
|
||||||
import Katip (LogEnv, closeScribes, mkHandleScribe, ColorStrategy (..), permitItem, Severity (..), Verbosity (V2), registerScribe, defaultScribeSettings, initLogEnv, ls, sl, logFM, katipAddContext, KatipContext)
|
import Katip (LogEnv, closeScribes, mkHandleScribe, ColorStrategy (..), permitItem, Severity (..), Verbosity (V2), registerScribe, defaultScribeSettings, initLogEnv, ls, sl, logFM, katipAddContext, KatipContext)
|
||||||
import Control.Exception (bracket)
|
import Control.Exception (bracket)
|
||||||
import System.IO (stdout)
|
import System.IO (stdout)
|
||||||
|
import System.Metrics.Counter (inc)
|
||||||
import Data.UUID (toText)
|
import Data.UUID (toText)
|
||||||
import AFRP (Request(..), Event(..))
|
import AFRP (Request(..), Event(..))
|
||||||
import Control.Monad.IO.Class (MonadIO, liftIO)
|
import Control.Monad.IO.Class (MonadIO, liftIO)
|
||||||
@@ -38,10 +42,11 @@ data Bus = Bus
|
|||||||
, busConn :: TVar (Maybe Connection)
|
, busConn :: TVar (Maybe Connection)
|
||||||
, busGen :: CallIdGen
|
, busGen :: CallIdGen
|
||||||
, busLogEnv :: LogEnv
|
, busLogEnv :: LogEnv
|
||||||
|
, busMetrics :: AppMetrics
|
||||||
}
|
}
|
||||||
|
|
||||||
withBus :: Severity -> (Bus -> IO a) -> IO a
|
withBus :: Severity -> AppMetrics -> (Bus -> IO a) -> IO a
|
||||||
withBus severity callback = do
|
withBus severity metrics callback = do
|
||||||
handleScribe <- mkHandleScribe ColorIfTerminal stdout (permitItem severity) V2
|
handleScribe <- mkHandleScribe ColorIfTerminal stdout (permitItem severity) V2
|
||||||
let makeLogEnv = registerScribe "stdout" handleScribe defaultScribeSettings =<< initLogEnv "hass-controller" "production"
|
let makeLogEnv = registerScribe "stdout" handleScribe defaultScribeSettings =<< initLogEnv "hass-controller" "production"
|
||||||
-- closeScribes will stop accepting new logs, flush existing ones and clean up resources
|
-- closeScribes will stop accepting new logs, flush existing ones and clean up resources
|
||||||
@@ -52,8 +57,15 @@ withBus severity callback = do
|
|||||||
<*> newTVarIO Nothing
|
<*> newTVarIO Nothing
|
||||||
<*> mkCallIdGen 0
|
<*> mkCallIdGen 0
|
||||||
<*> pure le
|
<*> pure le
|
||||||
|
<*> pure metrics
|
||||||
callback bus
|
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 :: (MonadIO m, KatipContext m) => Bus -> HASSEff a -> m a
|
||||||
channelHassEval bus = \case
|
channelHassEval bus = \case
|
||||||
CallService req svc -> katipAddContext (sl "traceId" (toText (requestTraceId req))) $ do
|
CallService req svc -> katipAddContext (sl "traceId" (toText (requestTraceId req))) $ do
|
||||||
|
|||||||
@@ -96,7 +96,9 @@ receiveLoop bus conn = forever $ do
|
|||||||
Left () -> atomically $ writeTChan (busInbound bus) Tick
|
Left () -> atomically $ writeTChan (busInbound bus) Tick
|
||||||
Right msg -> case eitherDecode msg of
|
Right msg -> case eitherDecode msg of
|
||||||
Left err -> putStrLn $ "[reader] skipping undecodable message: " <> err
|
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 :: WS.Connection -> IO Value
|
||||||
receiveJSON conn = do
|
receiveJSON conn = do
|
||||||
@@ -128,6 +130,7 @@ sendWithId bus conn request svc = do
|
|||||||
let textData = encode $ encodeService callId svc
|
let textData = encode $ encodeService callId svc
|
||||||
runKatipContextT (busLogEnv bus) (sl "traceId" (toText (requestTraceId request))) "connection" $
|
runKatipContextT (busLogEnv bus) (sl "traceId" (toText (requestTraceId request))) "connection" $
|
||||||
logFM DebugS (ls textData)
|
logFM DebugS (ls textData)
|
||||||
|
recordOutbound bus
|
||||||
WS.sendTextData conn textData
|
WS.sendTextData conn textData
|
||||||
|
|
||||||
writerAction :: Bus -> IO Void
|
writerAction :: Bus -> IO Void
|
||||||
|
|||||||
+23
-2
@@ -10,16 +10,19 @@ import Control.Concurrent.STM
|
|||||||
, writeTChan
|
, writeTChan
|
||||||
)
|
)
|
||||||
import Data.Aeson (Value (..))
|
import Data.Aeson (Value (..))
|
||||||
|
import qualified Data.HashMap.Strict as HM
|
||||||
import Data.Time (UTCTime (..), utc)
|
import Data.Time (UTCTime (..), utc)
|
||||||
import Data.UUID (nil)
|
import Data.UUID (nil)
|
||||||
import HomeAssistant.Controller (HASSEff (..), Service (..), Target(..))
|
import HomeAssistant.Controller (HASSEff (..), Service (..), Target(..))
|
||||||
import HomeAssistant.Runtime.Bus
|
import HomeAssistant.Runtime.Bus
|
||||||
|
import HomeAssistant.Runtime.Metrics (registerAppMetrics)
|
||||||
import Katip (Namespace (Namespace), runKatipContextT, Severity (..))
|
import Katip (Namespace (Namespace), runKatipContextT, Severity (..))
|
||||||
|
import qualified System.Metrics as Metrics
|
||||||
import Test.Hspec
|
import Test.Hspec
|
||||||
|
|
||||||
spec :: Spec
|
spec :: Spec
|
||||||
spec = describe "Bus" $ do
|
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)
|
p1 <- atomically $ dupTChan (busInbound bus)
|
||||||
p2 <- atomically $ dupTChan (busInbound bus)
|
p2 <- atomically $ dupTChan (busInbound bus)
|
||||||
atomically $ writeTChan (busInbound bus) (Event (Number 1))
|
atomically $ writeTChan (busInbound bus) (Event (Number 1))
|
||||||
@@ -29,7 +32,7 @@ spec = describe "Bus" $ do
|
|||||||
r1 `shouldBe` (Event (Number 1), Event (Number 2))
|
r1 `shouldBe` (Event (Number 1), Event (Number 2))
|
||||||
r2 `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
|
let req = Request (UTCTime (toEnum 0) 0) utc nil
|
||||||
svc = Service "light" "turn_on" Nothing [EntityId "light.bedroom_masse"]
|
svc = Service "light" "turn_on" Nothing [EntityId "light.bedroom_masse"]
|
||||||
runKatipContextT (busLogEnv bus) () (Namespace ["test"]) $
|
runKatipContextT (busLogEnv bus) () (Namespace ["test"]) $
|
||||||
@@ -43,3 +46,21 @@ spec = describe "Bus" $ do
|
|||||||
a <- generateCallId gen
|
a <- generateCallId gen
|
||||||
b <- generateCallId gen
|
b <- generateCallId gen
|
||||||
(a, b) `shouldBe` (1, 2)
|
(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
|
||||||
|
|||||||
Reference in New Issue
Block a user