Move handler flag check from bus to writer with skip logging
This commit is contained in:
@@ -4,6 +4,7 @@
|
|||||||
module HomeAssistant.Runtime.Connection
|
module HomeAssistant.Runtime.Connection
|
||||||
( readerAction
|
( readerAction
|
||||||
, writerAction
|
, writerAction
|
||||||
|
, dispatchService
|
||||||
, encodeService
|
, encodeService
|
||||||
, isTriggerEvent
|
, isTriggerEvent
|
||||||
) where
|
) where
|
||||||
@@ -29,9 +30,10 @@ import qualified Data.Text as T
|
|||||||
import Data.Void (Void)
|
import Data.Void (Void)
|
||||||
import HomeAssistant.Controller (Service (..), Target (..))
|
import HomeAssistant.Controller (Service (..), Target (..))
|
||||||
import HomeAssistant.Runtime.Bus
|
import HomeAssistant.Runtime.Bus
|
||||||
|
import HomeAssistant.Runtime.Flags (isEnabled)
|
||||||
import HomeAssistant.Runtime.Supervisor (Fatal (..))
|
import HomeAssistant.Runtime.Supervisor (Fatal (..))
|
||||||
import qualified Network.WebSockets as WS
|
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 Data.UUID (toText)
|
||||||
import AFRP (Request(..), Event(..))
|
import AFRP (Request(..), Event(..))
|
||||||
import Control.Monad.IO.Class (liftIO, MonadIO)
|
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 :: (MonadIO m, KatipContext m, MonadCatch m) => RateLimiter -> Bus -> m Void
|
||||||
writerAction rateLimiter bus = forever $ do
|
writerAction rateLimiter bus = forever $ do
|
||||||
(request, svc) <- liftIO $ atomically $ readTChan (busOutbound bus)
|
(request, svc) <- liftIO $ atomically $ readTChan (busOutbound bus)
|
||||||
katipAddContext (sl "traceId" (toText (requestTraceId request))) $ do
|
katipAddContext (sl "traceId" (toText (requestTraceId request))) $
|
||||||
conn <- liftIO $ atomically $ readTVar (busConn bus) >>= maybe retry pure
|
katipAddNamespace (Namespace [requestHandler request]) $
|
||||||
-- runRateLimited rateLimiter $ sendWithId bus conn request svc
|
dispatchService rateLimiter bus request svc
|
||||||
runRateLimited rateLimiter $ sendWithId bus conn 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 :: Int -> Service -> Value
|
||||||
encodeService callId Service{..} = object $
|
encodeService callId Service{..} = object $
|
||||||
|
|||||||
+24
-1
@@ -3,13 +3,22 @@
|
|||||||
module ConnectionSpec (spec) where
|
module ConnectionSpec (spec) where
|
||||||
|
|
||||||
import AFRP (Request(..))
|
import AFRP (Request(..))
|
||||||
|
import Control.Monad.IO.Class (liftIO)
|
||||||
import Data.Aeson (object, (.=))
|
import Data.Aeson (object, (.=))
|
||||||
|
import Data.Foldable (forM_)
|
||||||
import Data.Maybe (fromJust)
|
import Data.Maybe (fromJust)
|
||||||
import Data.Text (Text)
|
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, 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
|
import Test.Hspec
|
||||||
|
|
||||||
spec :: Spec
|
spec :: Spec
|
||||||
@@ -44,6 +53,20 @@ spec = do
|
|||||||
it "is false when there is no type" $
|
it "is false when there is no type" $
|
||||||
isTriggerEvent (object ["id" .= (1 :: Int)]) `shouldBe` False
|
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 :: Int -> Request
|
||||||
req n = Request (UTCTime (toEnum 0) (fromIntegral (0 :: Int))) utc
|
req n = Request (UTCTime (toEnum 0) (fromIntegral (0 :: Int))) utc
|
||||||
|
|||||||
Reference in New Issue
Block a user