From c2425e8fb0ac1f951926639280a2482e07532bc8 Mon Sep 17 00:00:00 2001 From: Mats Rauhala Date: Wed, 30 Sep 2026 08:06:59 +0300 Subject: [PATCH] Move handler flag check from bus to writer with skip logging --- src/HomeAssistant/Runtime/Connection.hs | 23 ++++++++++++++++++----- test/ConnectionSpec.hs | 25 ++++++++++++++++++++++++- 2 files changed, 42 insertions(+), 6 deletions(-) diff --git a/src/HomeAssistant/Runtime/Connection.hs b/src/HomeAssistant/Runtime/Connection.hs index 616b30d..5bfe287 100644 --- a/src/HomeAssistant/Runtime/Connection.hs +++ b/src/HomeAssistant/Runtime/Connection.hs @@ -4,6 +4,7 @@ module HomeAssistant.Runtime.Connection ( readerAction , writerAction + , dispatchService , encodeService , isTriggerEvent ) where @@ -29,9 +30,10 @@ import qualified Data.Text as T import Data.Void (Void) import HomeAssistant.Controller (Service (..), Target (..)) import HomeAssistant.Runtime.Bus +import HomeAssistant.Runtime.Flags (isEnabled) import HomeAssistant.Runtime.Supervisor (Fatal (..)) import qualified Network.WebSockets as WS -import Katip (sl, logFM, Severity (..), ls, KatipContext, katipAddNamespace, katipAddContext) +import Katip (sl, logFM, Severity (..), ls, Namespace (Namespace), KatipContext, katipAddNamespace, katipAddContext) import Data.UUID (toText) import AFRP (Request(..), Event(..)) import Control.Monad.IO.Class (liftIO, MonadIO) @@ -122,10 +124,21 @@ sendWithId bus conn svc = katipAddNamespace "connection" $ do writerAction :: (MonadIO m, KatipContext m, MonadCatch m) => RateLimiter -> Bus -> m Void writerAction rateLimiter bus = forever $ do (request, svc) <- liftIO $ atomically $ readTChan (busOutbound bus) - katipAddContext (sl "traceId" (toText (requestTraceId request))) $ do - conn <- liftIO $ atomically $ readTVar (busConn bus) >>= maybe retry pure - -- runRateLimited rateLimiter $ sendWithId bus conn request svc - runRateLimited rateLimiter $ sendWithId bus conn svc + katipAddContext (sl "traceId" (toText (requestTraceId request))) $ + katipAddNamespace (Namespace [requestHandler request]) $ + dispatchService rateLimiter bus request svc + +-- | Sends the service call, or logs and drops it when its handler is disabled. +dispatchService + :: (MonadIO m, KatipContext m, MonadCatch m) + => RateLimiter -> Bus -> Request -> Service -> m () +dispatchService rateLimiter bus request svc = do + enabled <- isEnabled (busFlags bus) (requestHandler request) + if enabled + then do + conn <- liftIO $ atomically $ readTVar (busConn bus) >>= maybe retry pure + runRateLimited rateLimiter $ sendWithId bus conn svc + else logFM InfoS (ls $ "skipped disabled handler: " <> show svc) encodeService :: Int -> Service -> Value encodeService callId Service{..} = object $ diff --git a/test/ConnectionSpec.hs b/test/ConnectionSpec.hs index f4c30a0..ae81d41 100644 --- a/test/ConnectionSpec.hs +++ b/test/ConnectionSpec.hs @@ -3,13 +3,22 @@ module ConnectionSpec (spec) where import AFRP (Request(..)) +import Control.Monad.IO.Class (liftIO) import Data.Aeson (object, (.=)) +import Data.Foldable (forM_) import Data.Maybe (fromJust) import Data.Text (Text) import Data.Time (UTCTime (..), utc) import Data.UUID (fromString) import HomeAssistant.Controller (Service (..), Target(..)) -import HomeAssistant.Runtime.Connection (encodeService, isTriggerEvent) +import HomeAssistant.Runtime.Bus +import HomeAssistant.Runtime.Flags (loadFlags, setEnabled) +import HomeAssistant.Runtime.Metrics (registerAppMetrics) +import HomeAssistant.Runtime.RateLimit (registerRateLimitMetrics, slidingWindowLimiter) +import HomeAssistant.Runtime.Connection (dispatchService, encodeService, isTriggerEvent) +import Katip (Namespace (Namespace), Severity (InfoS), runKatipContextT) +import System.IO.Temp (withTempDirectory) +import qualified System.Metrics as Metrics import Test.Hspec spec :: Spec @@ -44,6 +53,20 @@ spec = do it "is false when there is no type" $ isTriggerEvent (object ["id" .= (1 :: Int)]) `shouldBe` False + describe "dispatchService" $ do + it "skips disabled handlers without consuming a call id" $ do + store <- Metrics.newStore + m <- registerAppMetrics store + rlMetrics <- registerRateLimitMetrics store + limiter <- slidingWindowLimiter rlMetrics 10 20 + withTempDirectory "/tmp" "connection-spec" $ \dir -> do + flags <- loadFlags dir ["test"] + setEnabled flags "test" False + withBus InfoS m flags $ \bus -> do + runKatipContextT (busLogEnv bus) () (Namespace ["connection"]) $ + dispatchService limiter bus (req 1) (lightOn [EntityId "light.test"]) + generateCallId (busGen bus) `shouldReturn` 1 + req :: Int -> Request req n = Request (UTCTime (toEnum 0) (fromIntegral (0 :: Int))) utc