Merge branch 'feature/trigger-service-metrics'
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
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ module HomeAssistant.Runtime.Connection
|
|||||||
, writerAction
|
, writerAction
|
||||||
, encodeService
|
, encodeService
|
||||||
, dedupeBatch
|
, dedupeBatch
|
||||||
|
, isTriggerEvent
|
||||||
) where
|
) where
|
||||||
|
|
||||||
import Control.Concurrent.STM
|
import Control.Concurrent.STM
|
||||||
@@ -23,7 +24,7 @@ import Control.Concurrent (threadDelay)
|
|||||||
import Control.Exception (onException)
|
import Control.Exception (onException)
|
||||||
import Control.Exception.Annotated (throw)
|
import Control.Exception.Annotated (throw)
|
||||||
import Control.Lens ((^?))
|
import Control.Lens ((^?))
|
||||||
import Control.Monad (forever, forM_)
|
import Control.Monad (forever, forM_, when)
|
||||||
import Data.Aeson (Value, eitherDecode, encode, object, (.=))
|
import Data.Aeson (Value, eitherDecode, encode, object, (.=))
|
||||||
import Data.Aeson.Lens (key, _String)
|
import Data.Aeson.Lens (key, _String)
|
||||||
import Data.List (sort)
|
import Data.List (sort)
|
||||||
@@ -69,6 +70,9 @@ expectType expected msg =
|
|||||||
Just t | t == expected -> pure ()
|
Just t | t == expected -> pure ()
|
||||||
_ -> throw (Fatal $ "expected " <> expected <> ", got: " <> T.pack (show msg))
|
_ -> throw (Fatal $ "expected " <> expected <> ", got: " <> T.pack (show msg))
|
||||||
|
|
||||||
|
isTriggerEvent :: Value -> Bool
|
||||||
|
isTriggerEvent v = v ^? key "type" . _String == Just "event"
|
||||||
|
|
||||||
subscribe :: Bus -> WS.Connection -> S.Set T.Text -> IO ()
|
subscribe :: Bus -> WS.Connection -> S.Set T.Text -> IO ()
|
||||||
subscribe bus conn ents =
|
subscribe bus conn ents =
|
||||||
forM_ (S.toList ents) $ \entityId -> do
|
forM_ (S.toList ents) $ \entityId -> do
|
||||||
@@ -96,7 +100,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
|
||||||
|
when (isTriggerEvent v) (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 +134,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
|
||||||
|
|||||||
@@ -8,6 +8,10 @@ module HomeAssistant.Runtime.Metrics
|
|||||||
, buildSchema
|
, buildSchema
|
||||||
, buildCreateArgs
|
, buildCreateArgs
|
||||||
, buildUpdateArgs
|
, buildUpdateArgs
|
||||||
|
, parseInfoDs
|
||||||
|
, schemaMatches
|
||||||
|
, AppMetrics (..)
|
||||||
|
, registerAppMetrics
|
||||||
, ensureRrd
|
, ensureRrd
|
||||||
, sampleAndUpdate
|
, sampleAndUpdate
|
||||||
, metricsAction
|
, metricsAction
|
||||||
@@ -17,19 +21,23 @@ import Control.Concurrent (threadDelay)
|
|||||||
import Control.Monad (forever, unless)
|
import Control.Monad (forever, unless)
|
||||||
import Data.Char (ord)
|
import Data.Char (ord)
|
||||||
import Data.Int (Int64)
|
import Data.Int (Int64)
|
||||||
import Data.List (intercalate, sortBy)
|
import Data.List (intercalate, sort, sortBy, stripPrefix)
|
||||||
|
import Data.Maybe (mapMaybe)
|
||||||
import Data.Ord (comparing)
|
import Data.Ord (comparing)
|
||||||
import Data.Text (Text)
|
import Data.Text (Text)
|
||||||
|
import Data.Time (defaultTimeLocale, formatTime, getCurrentTime)
|
||||||
import Data.Void (Void)
|
import Data.Void (Void)
|
||||||
import Numeric (showHex)
|
import Numeric (showHex)
|
||||||
import qualified Data.Text as T
|
import qualified Data.Text as T
|
||||||
import qualified Data.HashMap.Strict as HM
|
import qualified Data.HashMap.Strict as HM
|
||||||
import qualified System.Metrics as M (Value (..), Sample, Store, sampleAll)
|
import qualified System.Metrics as M (Value (..), Sample, Store, createCounter, sampleAll)
|
||||||
import System.Directory (doesFileExist)
|
import System.Metrics.Counter (Counter)
|
||||||
import System.Process (callProcess)
|
import System.Directory (doesFileExist, renameFile)
|
||||||
|
import System.Exit (ExitCode (..))
|
||||||
|
import System.Process (callProcess, readProcessWithExitCode)
|
||||||
|
|
||||||
data DsType = Derive | Gauge
|
data DsType = Derive | Gauge
|
||||||
deriving (Eq, Show)
|
deriving (Eq, Ord, Show)
|
||||||
|
|
||||||
data DsSpec = DsSpec
|
data DsSpec = DsSpec
|
||||||
{ dsEkgName :: Text
|
{ dsEkgName :: Text
|
||||||
@@ -62,6 +70,17 @@ buildSchema sample =
|
|||||||
, Just dt <- [dsTypeOf val]
|
, Just dt <- [dsTypeOf val]
|
||||||
]
|
]
|
||||||
|
|
||||||
|
data AppMetrics = AppMetrics
|
||||||
|
{ amTriggersIn :: Counter
|
||||||
|
, amServicesOut :: Counter
|
||||||
|
}
|
||||||
|
|
||||||
|
registerAppMetrics :: M.Store -> IO AppMetrics
|
||||||
|
registerAppMetrics store =
|
||||||
|
AppMetrics
|
||||||
|
<$> M.createCounter "hass.trigger.in" store
|
||||||
|
<*> M.createCounter "hass.service.out" store
|
||||||
|
|
||||||
buildCreateArgs :: FilePath -> Int -> [(String, DsType)] -> [String]
|
buildCreateArgs :: FilePath -> Int -> [(String, DsType)] -> [String]
|
||||||
buildCreateArgs path step specs =
|
buildCreateArgs path step specs =
|
||||||
["create", path, "--step", show step]
|
["create", path, "--step", show step]
|
||||||
@@ -89,18 +108,57 @@ buildUpdateArgs path names values =
|
|||||||
renderValue Nothing = "U"
|
renderValue Nothing = "U"
|
||||||
renderValue (Just n) = show n
|
renderValue (Just n) = show n
|
||||||
|
|
||||||
|
dsTypeFromString :: String -> Maybe DsType
|
||||||
|
dsTypeFromString "DERIVE" = Just Derive
|
||||||
|
dsTypeFromString "GAUGE" = Just Gauge
|
||||||
|
dsTypeFromString _ = Nothing
|
||||||
|
|
||||||
|
parseInfoDs :: String -> [(String, DsType)]
|
||||||
|
parseInfoDs = mapMaybe parseLine . lines
|
||||||
|
where
|
||||||
|
parseLine line = do
|
||||||
|
rest0 <- stripPrefix "ds[" line
|
||||||
|
let (name, rest1) = break (== ']') rest0
|
||||||
|
rest2 <- stripPrefix "].type = \"" rest1
|
||||||
|
let (ty, _) = break (== '"') rest2
|
||||||
|
dt <- dsTypeFromString ty
|
||||||
|
pure (name, dt)
|
||||||
|
|
||||||
|
schemaMatches :: [DsSpec] -> [(String, DsType)] -> Bool
|
||||||
|
schemaMatches schema existing =
|
||||||
|
sort (map (\s -> (dsName s, dsType s)) schema) == sort existing
|
||||||
|
|
||||||
lookupValue :: M.Sample -> Text -> Maybe Int64
|
lookupValue :: M.Sample -> Text -> Maybe Int64
|
||||||
lookupValue sample name = case HM.lookup name sample of
|
lookupValue sample name = case HM.lookup name sample of
|
||||||
Just (M.Counter n) -> Just n
|
Just (M.Counter n) -> Just n
|
||||||
Just (M.Gauge n) -> Just n
|
Just (M.Gauge n) -> Just n
|
||||||
_ -> Nothing
|
_ -> Nothing
|
||||||
|
|
||||||
|
readExistingDs :: FilePath -> FilePath -> IO [(String, DsType)]
|
||||||
|
readExistingDs rrdtool rrdPath = do
|
||||||
|
(code, out, _) <- readProcessWithExitCode rrdtool ["info", rrdPath] ""
|
||||||
|
pure $ case code of
|
||||||
|
ExitSuccess -> parseInfoDs out
|
||||||
|
_ -> []
|
||||||
|
|
||||||
|
backupPath :: FilePath -> IO FilePath
|
||||||
|
backupPath path = do
|
||||||
|
now <- getCurrentTime
|
||||||
|
pure $ path ++ ".bak-" ++ formatTime defaultTimeLocale "%Y%m%dT%H%M%S" now
|
||||||
|
|
||||||
ensureRrd :: FilePath -> FilePath -> [DsSpec] -> IO ()
|
ensureRrd :: FilePath -> FilePath -> [DsSpec] -> IO ()
|
||||||
ensureRrd rrdtool rrdPath schema = do
|
ensureRrd rrdtool rrdPath schema = do
|
||||||
exists <- doesFileExist rrdPath
|
exists <- doesFileExist rrdPath
|
||||||
unless exists $
|
if not exists
|
||||||
callProcess rrdtool (buildCreateArgs rrdPath 10 (map toPair schema))
|
then create
|
||||||
|
else do
|
||||||
|
current <- readExistingDs rrdtool rrdPath
|
||||||
|
unless (schemaMatches schema current) $ do
|
||||||
|
backup <- backupPath rrdPath
|
||||||
|
renameFile rrdPath backup
|
||||||
|
create
|
||||||
where
|
where
|
||||||
|
create = callProcess rrdtool (buildCreateArgs rrdPath 10 (map toPair schema))
|
||||||
toPair s = (dsName s, dsType s)
|
toPair s = (dsName s, dsType s)
|
||||||
|
|
||||||
sampleAndUpdate :: M.Store -> FilePath -> FilePath -> [DsSpec] -> IO ()
|
sampleAndUpdate :: M.Store -> FilePath -> FilePath -> [DsSpec] -> IO ()
|
||||||
|
|||||||
+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
|
||||||
|
|||||||
@@ -9,7 +9,7 @@ import Data.Text (Text)
|
|||||||
import Data.Time (UTCTime (..), utc)
|
import Data.Time (UTCTime (..), utc)
|
||||||
import Data.UUID (fromString)
|
import Data.UUID (fromString)
|
||||||
import HomeAssistant.Controller (Service (..), Target(..))
|
import HomeAssistant.Controller (Service (..), Target(..))
|
||||||
import HomeAssistant.Runtime.Connection (encodeService, dedupeBatch)
|
import HomeAssistant.Runtime.Connection (encodeService, dedupeBatch, isTriggerEvent)
|
||||||
import Test.Hspec
|
import Test.Hspec
|
||||||
|
|
||||||
spec :: Spec
|
spec :: Spec
|
||||||
@@ -36,6 +36,14 @@ spec = do
|
|||||||
, "service_data" .= object ["brightness" .= (200 :: Int)]
|
, "service_data" .= object ["brightness" .= (200 :: Int)]
|
||||||
]
|
]
|
||||||
|
|
||||||
|
describe "isTriggerEvent" $ do
|
||||||
|
it "is true for event frames" $
|
||||||
|
isTriggerEvent (object ["type" .= ("event" :: Text)]) `shouldBe` True
|
||||||
|
it "is false for result frames" $
|
||||||
|
isTriggerEvent (object ["type" .= ("result" :: Text)]) `shouldBe` False
|
||||||
|
it "is false when there is no type" $
|
||||||
|
isTriggerEvent (object ["id" .= (1 :: Int)]) `shouldBe` False
|
||||||
|
|
||||||
describe "dedupeBatch" $ do
|
describe "dedupeBatch" $ do
|
||||||
it "collapses identical calls to one" $
|
it "collapses identical calls to one" $
|
||||||
let batch = [ (req 1, lightOn [AreaId "x"])
|
let batch = [ (req 1, lightOn [AreaId "x"])
|
||||||
|
|||||||
+80
-4
@@ -4,8 +4,9 @@ module MetricsSpec (spec) where
|
|||||||
|
|
||||||
import Data.HashMap.Strict (HashMap)
|
import Data.HashMap.Strict (HashMap)
|
||||||
import qualified Data.HashMap.Strict as HM
|
import qualified Data.HashMap.Strict as HM
|
||||||
import Data.Int (Int64)
|
import Data.List (isPrefixOf)
|
||||||
import Data.Text (Text)
|
import Data.Text (Text)
|
||||||
|
import qualified Data.Text as T
|
||||||
import HomeAssistant.Runtime.Metrics
|
import HomeAssistant.Runtime.Metrics
|
||||||
( DsType (..)
|
( DsType (..)
|
||||||
, DsSpec (..)
|
, DsSpec (..)
|
||||||
@@ -16,12 +17,16 @@ import HomeAssistant.Runtime.Metrics
|
|||||||
, buildUpdateArgs
|
, buildUpdateArgs
|
||||||
, ensureRrd
|
, ensureRrd
|
||||||
, sampleAndUpdate
|
, sampleAndUpdate
|
||||||
, metricsAction
|
, parseInfoDs
|
||||||
|
, schemaMatches
|
||||||
|
, AppMetrics (..)
|
||||||
|
, registerAppMetrics
|
||||||
)
|
)
|
||||||
import qualified System.Metrics as M (Value (..))
|
import qualified System.Metrics as M (Value (..))
|
||||||
import Test.Hspec
|
import Test.Hspec
|
||||||
import Control.Exception (try, SomeException)
|
import Control.Exception (try, SomeException, finally)
|
||||||
import System.Directory (findExecutable, getTemporaryDirectory, removeFile)
|
import Control.Monad (forM_)
|
||||||
|
import System.Directory (findExecutable, getTemporaryDirectory, listDirectory, removeFile)
|
||||||
import System.Exit (ExitCode (ExitSuccess))
|
import System.Exit (ExitCode (ExitSuccess))
|
||||||
import System.Process (readProcessWithExitCode)
|
import System.Process (readProcessWithExitCode)
|
||||||
import qualified System.Metrics as Metrics
|
import qualified System.Metrics as Metrics
|
||||||
@@ -126,6 +131,49 @@ spec = do
|
|||||||
, "N:U"
|
, "N:U"
|
||||||
]
|
]
|
||||||
|
|
||||||
|
describe "parseInfoDs" $ do
|
||||||
|
it "extracts DS names and types from rrdtool info output" $
|
||||||
|
let info = unlines
|
||||||
|
[ "filename = \"test.rrd\""
|
||||||
|
, "step = 10"
|
||||||
|
, "ds[foo].index = 0"
|
||||||
|
, "ds[foo].type = \"DERIVE\""
|
||||||
|
, "ds[bar].type = \"GAUGE\""
|
||||||
|
, "rra[0].cf = \"AVERAGE\""
|
||||||
|
]
|
||||||
|
in parseInfoDs info `shouldBe` [("foo", Derive), ("bar", Gauge)]
|
||||||
|
|
||||||
|
it "ignores lines that are not ds type declarations" $
|
||||||
|
parseInfoDs "step = 10\nrra[0].cf = \"AVERAGE\"\n" `shouldBe` []
|
||||||
|
|
||||||
|
describe "schemaMatches" $ do
|
||||||
|
let ds name ty = DsSpec (T.pack name) name ty
|
||||||
|
it "is true for identical schemas regardless of order" $
|
||||||
|
schemaMatches [ds "b" Gauge, ds "a" Derive] [("a", Derive), ("b", Gauge)]
|
||||||
|
`shouldBe` True
|
||||||
|
|
||||||
|
it "is false when the derived schema has an added DS" $
|
||||||
|
schemaMatches [ds "a" Derive, ds "b" Gauge] [("a", Derive)]
|
||||||
|
`shouldBe` False
|
||||||
|
|
||||||
|
it "is false when the derived schema has a removed DS" $
|
||||||
|
schemaMatches [ds "a" Derive] [("a", Derive), ("b", Gauge)]
|
||||||
|
`shouldBe` False
|
||||||
|
|
||||||
|
it "is false when a DS type changed" $
|
||||||
|
schemaMatches [ds "a" Derive] [("a", Gauge)]
|
||||||
|
`shouldBe` False
|
||||||
|
|
||||||
|
describe "registerAppMetrics" $ do
|
||||||
|
it "registers both counters with the store" $ do
|
||||||
|
store <- Metrics.newStore
|
||||||
|
m <- registerAppMetrics store
|
||||||
|
Counter.inc (amTriggersIn m)
|
||||||
|
Counter.inc (amServicesOut m)
|
||||||
|
sample <- Metrics.sampleAll store
|
||||||
|
HM.lookup "hass.trigger.in" sample `shouldBe` Just (M.Counter 1)
|
||||||
|
HM.lookup "hass.service.out" sample `shouldBe` Just (M.Counter 1)
|
||||||
|
|
||||||
describe "end-to-end (rrdtool-gated)" $ do
|
describe "end-to-end (rrdtool-gated)" $ do
|
||||||
it "creates an rrd, samples, and updates it" $ do
|
it "creates an rrd, samples, and updates it" $ do
|
||||||
mRrdtool <- findExecutable "rrdtool"
|
mRrdtool <- findExecutable "rrdtool"
|
||||||
@@ -148,3 +196,31 @@ spec = do
|
|||||||
length out `shouldSatisfy` (> 0)
|
length out `shouldSatisfy` (> 0)
|
||||||
_ <- try (removeFile rrdPath) :: IO (Either SomeException ())
|
_ <- try (removeFile rrdPath) :: IO (Either SomeException ())
|
||||||
return ()
|
return ()
|
||||||
|
|
||||||
|
it "backs up and recreates the rrd when the schema changes" $ do
|
||||||
|
mRrdtool <- findExecutable "rrdtool"
|
||||||
|
case mRrdtool of
|
||||||
|
Nothing -> pendingWith "rrdtool not on PATH"
|
||||||
|
Just rrdtool -> do
|
||||||
|
tmp <- getTemporaryDirectory
|
||||||
|
let rrdPath = tmp ++ "/hass-controller-schema-change-test.rrd"
|
||||||
|
bakPrefix = "hass-controller-schema-change-test.rrd.bak-"
|
||||||
|
base = [ DsSpec "a" "a" Derive, DsSpec "b" "b" Gauge ]
|
||||||
|
expanded = base ++ [ DsSpec "c" "c" Derive ]
|
||||||
|
clean = do
|
||||||
|
stale <- filter (isPrefixOf bakPrefix) <$> listDirectory tmp
|
||||||
|
forM_ (rrdPath : map (\f -> tmp ++ "/" ++ f) stale) $ \f ->
|
||||||
|
try (removeFile f) :: IO (Either SomeException ())
|
||||||
|
clean
|
||||||
|
( do
|
||||||
|
ensureRrd rrdtool rrdPath base
|
||||||
|
ensureRrd rrdtool rrdPath base
|
||||||
|
noBackups <- filter (isPrefixOf bakPrefix) <$> listDirectory tmp
|
||||||
|
noBackups `shouldBe` []
|
||||||
|
ensureRrd rrdtool rrdPath expanded
|
||||||
|
backups <- filter (isPrefixOf bakPrefix) <$> listDirectory tmp
|
||||||
|
length backups `shouldBe` 1
|
||||||
|
(rc, out, _) <- readProcessWithExitCode rrdtool ["info", rrdPath] ""
|
||||||
|
rc `shouldBe` ExitSuccess
|
||||||
|
out `shouldContain` "ds[c].type"
|
||||||
|
) `finally` clean
|
||||||
|
|||||||
+1
-1
@@ -11,7 +11,7 @@ spec = pure ()
|
|||||||
--
|
--
|
||||||
-- spec :: Spec
|
-- spec :: Spec
|
||||||
-- spec = describe "runController" $ do
|
-- spec = describe "runController" $ do
|
||||||
-- it "feeds inbound events through the machine and forwards service calls" $ withBus InfoS $ \bus -> do
|
-- it "feeds inbound events through the machine and forwards service calls" $ withBus severity appMetrics $ \bus -> do
|
||||||
-- _ <- async (runController bus (Controller "test" lightController True))
|
-- _ <- async (runController bus (Controller "test" lightController True))
|
||||||
-- putStrLn "Before the delay"
|
-- putStrLn "Before the delay"
|
||||||
-- threadDelay 100000 -- let the controller dup its inbound channel
|
-- threadDelay 100000 -- let the controller dup its inbound channel
|
||||||
|
|||||||
Reference in New Issue
Block a user