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..d77a1cd 100644 --- a/src/HomeAssistant/Runtime/Connection.hs +++ b/src/HomeAssistant/Runtime/Connection.hs @@ -6,6 +6,7 @@ module HomeAssistant.Runtime.Connection , writerAction , encodeService , dedupeBatch + , isTriggerEvent ) where import Control.Concurrent.STM @@ -23,7 +24,7 @@ import Control.Concurrent (threadDelay) import Control.Exception (onException) import Control.Exception.Annotated (throw) import Control.Lens ((^?)) -import Control.Monad (forever, forM_) +import Control.Monad (forever, forM_, when) import Data.Aeson (Value, eitherDecode, encode, object, (.=)) import Data.Aeson.Lens (key, _String) import Data.List (sort) @@ -69,6 +70,9 @@ expectType expected msg = Just t | t == expected -> pure () _ -> 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 conn ents = forM_ (S.toList ents) $ \entityId -> do @@ -96,7 +100,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 + when (isTriggerEvent v) (recordInbound bus) + atomically $ writeTChan (busInbound bus) (Event v) receiveJSON :: WS.Connection -> IO Value receiveJSON conn = do @@ -128,6 +134,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/src/HomeAssistant/Runtime/Metrics.hs b/src/HomeAssistant/Runtime/Metrics.hs index 0e07597..5d9d92f 100644 --- a/src/HomeAssistant/Runtime/Metrics.hs +++ b/src/HomeAssistant/Runtime/Metrics.hs @@ -8,6 +8,10 @@ module HomeAssistant.Runtime.Metrics , buildSchema , buildCreateArgs , buildUpdateArgs + , parseInfoDs + , schemaMatches + , AppMetrics (..) + , registerAppMetrics , ensureRrd , sampleAndUpdate , metricsAction @@ -17,19 +21,23 @@ import Control.Concurrent (threadDelay) import Control.Monad (forever, unless) import Data.Char (ord) 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.Text (Text) +import Data.Time (defaultTimeLocale, formatTime, getCurrentTime) import Data.Void (Void) import Numeric (showHex) import qualified Data.Text as T import qualified Data.HashMap.Strict as HM -import qualified System.Metrics as M (Value (..), Sample, Store, sampleAll) -import System.Directory (doesFileExist) -import System.Process (callProcess) +import qualified System.Metrics as M (Value (..), Sample, Store, createCounter, sampleAll) +import System.Metrics.Counter (Counter) +import System.Directory (doesFileExist, renameFile) +import System.Exit (ExitCode (..)) +import System.Process (callProcess, readProcessWithExitCode) data DsType = Derive | Gauge - deriving (Eq, Show) + deriving (Eq, Ord, Show) data DsSpec = DsSpec { dsEkgName :: Text @@ -62,6 +70,17 @@ buildSchema sample = , 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 path step specs = ["create", path, "--step", show step] @@ -89,18 +108,57 @@ buildUpdateArgs path names values = renderValue Nothing = "U" 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 sample name = case HM.lookup name sample of Just (M.Counter n) -> Just n Just (M.Gauge n) -> Just n _ -> 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 rrdtool rrdPath schema = do exists <- doesFileExist rrdPath - unless exists $ - callProcess rrdtool (buildCreateArgs rrdPath 10 (map toPair schema)) + if not exists + then create + else do + current <- readExistingDs rrdtool rrdPath + unless (schemaMatches schema current) $ do + backup <- backupPath rrdPath + renameFile rrdPath backup + create where + create = callProcess rrdtool (buildCreateArgs rrdPath 10 (map toPair schema)) toPair s = (dsName s, dsType s) sampleAndUpdate :: M.Store -> FilePath -> FilePath -> [DsSpec] -> IO () 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 diff --git a/test/ConnectionSpec.hs b/test/ConnectionSpec.hs index 90726ab..b279380 100644 --- a/test/ConnectionSpec.hs +++ b/test/ConnectionSpec.hs @@ -9,7 +9,7 @@ import Data.Text (Text) import Data.Time (UTCTime (..), utc) import Data.UUID (fromString) import HomeAssistant.Controller (Service (..), Target(..)) -import HomeAssistant.Runtime.Connection (encodeService, dedupeBatch) +import HomeAssistant.Runtime.Connection (encodeService, dedupeBatch, isTriggerEvent) import Test.Hspec spec :: Spec @@ -36,6 +36,14 @@ spec = do , "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 it "collapses identical calls to one" $ let batch = [ (req 1, lightOn [AreaId "x"]) diff --git a/test/MetricsSpec.hs b/test/MetricsSpec.hs index 2b3c6b2..e648be5 100644 --- a/test/MetricsSpec.hs +++ b/test/MetricsSpec.hs @@ -4,8 +4,9 @@ module MetricsSpec (spec) where import Data.HashMap.Strict (HashMap) import qualified Data.HashMap.Strict as HM -import Data.Int (Int64) +import Data.List (isPrefixOf) import Data.Text (Text) +import qualified Data.Text as T import HomeAssistant.Runtime.Metrics ( DsType (..) , DsSpec (..) @@ -16,12 +17,16 @@ import HomeAssistant.Runtime.Metrics , buildUpdateArgs , ensureRrd , sampleAndUpdate - , metricsAction + , parseInfoDs + , schemaMatches + , AppMetrics (..) + , registerAppMetrics ) import qualified System.Metrics as M (Value (..)) import Test.Hspec -import Control.Exception (try, SomeException) -import System.Directory (findExecutable, getTemporaryDirectory, removeFile) +import Control.Exception (try, SomeException, finally) +import Control.Monad (forM_) +import System.Directory (findExecutable, getTemporaryDirectory, listDirectory, removeFile) import System.Exit (ExitCode (ExitSuccess)) import System.Process (readProcessWithExitCode) import qualified System.Metrics as Metrics @@ -126,6 +131,49 @@ spec = do , "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 it "creates an rrd, samples, and updates it" $ do mRrdtool <- findExecutable "rrdtool" @@ -148,3 +196,31 @@ spec = do length out `shouldSatisfy` (> 0) _ <- try (removeFile rrdPath) :: IO (Either SomeException ()) 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 diff --git a/test/RuntimeSpec.hs b/test/RuntimeSpec.hs index 2f85fdd..53dcdbc 100644 --- a/test/RuntimeSpec.hs +++ b/test/RuntimeSpec.hs @@ -11,7 +11,7 @@ spec = pure () -- -- spec :: Spec -- 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)) -- putStrLn "Before the delay" -- threadDelay 100000 -- let the controller dup its inbound channel