Simplify channelHassEval now that the writer owns the flag decision
This commit is contained in:
@@ -1,5 +1,6 @@
|
|||||||
{-# LANGUAGE LambdaCase #-}
|
{-# LANGUAGE LambdaCase #-}
|
||||||
{-# LANGUAGE OverloadedStrings #-}
|
{-# LANGUAGE OverloadedStrings #-}
|
||||||
|
{-# LANGUAGE ScopedTypeVariables #-}
|
||||||
|
|
||||||
module HomeAssistant.Runtime.Bus
|
module HomeAssistant.Runtime.Bus
|
||||||
( Bus(..)
|
( Bus(..)
|
||||||
@@ -31,10 +32,8 @@ import System.IO (stdout)
|
|||||||
import System.Metrics.Counter (inc)
|
import System.Metrics.Counter (inc)
|
||||||
import Data.UUID (toText)
|
import Data.UUID (toText)
|
||||||
import AFRP (Request(..), Event(..))
|
import AFRP (Request(..), Event(..))
|
||||||
import Control.Monad (when)
|
|
||||||
import Control.Monad.IO.Class (MonadIO, liftIO)
|
import Control.Monad.IO.Class (MonadIO, liftIO)
|
||||||
import HomeAssistant.Runtime.Flags (Flags, isEnabled)
|
import HomeAssistant.Runtime.Flags (Flags)
|
||||||
import Data.Foldable (forM_)
|
|
||||||
|
|
||||||
-- | Shared runtime state: inbound is a broadcast channel (controllers
|
-- | Shared runtime state: inbound is a broadcast channel (controllers
|
||||||
-- read from 'dupTChan' copies), outbound queues service calls for the
|
-- read from 'dupTChan' copies), outbound queues service calls for the
|
||||||
@@ -71,19 +70,16 @@ recordInbound bus = inc (amTriggersIn (busMetrics bus))
|
|||||||
recordOutbound :: Bus -> IO ()
|
recordOutbound :: Bus -> IO ()
|
||||||
recordOutbound bus = inc (amServicesOut (busMetrics bus))
|
recordOutbound bus = inc (amServicesOut (busMetrics bus))
|
||||||
|
|
||||||
channelHassEval :: (MonadIO m, KatipContext m) => Bus -> HASSEff a -> m a
|
channelHassEval :: forall m a. (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 -> emit req svc
|
||||||
logFM DebugS (ls $ show svc)
|
CallServices req svcs -> mapM_ (emit req) svcs
|
||||||
enabled <- isEnabled (busFlags bus) (requestHandler req)
|
|
||||||
when enabled $ liftIO $ atomically $ writeTChan (busOutbound bus) (req, svc)
|
|
||||||
CallServices req svcs -> katipAddContext (sl "traceId" (toText (requestTraceId req))) $ forM_ svcs $ \svc -> do
|
|
||||||
logFM DebugS (ls $ show svc)
|
|
||||||
enabled <- isEnabled (busFlags bus) (requestHandler req)
|
|
||||||
when enabled $ liftIO $ atomically $ writeTChan (busOutbound bus) (req, svc)
|
|
||||||
Debug x -> logFM DebugS (ls $ show x)
|
Debug x -> logFM DebugS (ls $ show x)
|
||||||
Trace req x -> katipAddContext (sl "traceId" (toText (requestTraceId req))) $
|
Trace req x -> katipAddContext (sl "traceId" (toText (requestTraceId req))) $
|
||||||
logFM InfoS (ls $ show x)
|
logFM InfoS (ls $ show x)
|
||||||
|
where
|
||||||
|
emit :: Request -> Service -> m ()
|
||||||
|
emit req svc = liftIO $ atomically $ writeTChan (busOutbound bus) (req, svc)
|
||||||
|
|
||||||
newtype CallIdGen = CallIdGen { generateCallId :: IO Int }
|
newtype CallIdGen = CallIdGen { generateCallId :: IO Int }
|
||||||
|
|
||||||
|
|||||||
+3
-2
@@ -43,14 +43,15 @@ spec = describe "Bus" $ do
|
|||||||
(_, svc') <- atomically $ readTChan (busOutbound bus)
|
(_, svc') <- atomically $ readTChan (busOutbound bus)
|
||||||
svc' `shouldBe` svc
|
svc' `shouldBe` svc
|
||||||
|
|
||||||
it "channelHassEval drops CallService when the handler is disabled" $
|
it "channelHassEval writes CallService to the outbound channel even when the handler is disabled" $
|
||||||
withTestFlagsBus $ \bus -> do
|
withTestFlagsBus $ \bus -> do
|
||||||
setEnabled (busFlags bus) "test-handler" False
|
setEnabled (busFlags bus) "test-handler" False
|
||||||
let req = Request (UTCTime (toEnum 0) 0) utc nil "test-handler"
|
let req = Request (UTCTime (toEnum 0) 0) utc nil "test-handler"
|
||||||
svc = Service "light" "turn_on" Nothing [EntityId "light.test"]
|
svc = Service "light" "turn_on" Nothing [EntityId "light.test"]
|
||||||
runKatipContextT (busLogEnv bus) () (Namespace ["test"]) $
|
runKatipContextT (busLogEnv bus) () (Namespace ["test"]) $
|
||||||
channelHassEval bus (CallService req svc)
|
channelHassEval bus (CallService req svc)
|
||||||
atomically (isEmptyTChan (busOutbound bus)) `shouldReturn` True
|
(_, svc') <- atomically $ readTChan (busOutbound bus)
|
||||||
|
svc' `shouldBe` svc
|
||||||
|
|
||||||
|
|
||||||
it "generates unique sequential call ids" $ do
|
it "generates unique sequential call ids" $ do
|
||||||
|
|||||||
Reference in New Issue
Block a user