Triggers instead of full state events

This commit is contained in:
2026-08-25 13:13:50 +03:00
parent 771d702abd
commit 2730657952
11 changed files with 180 additions and 66 deletions
+6 -5
View File
@@ -30,6 +30,7 @@ import Control.Arrow (Arrow(..), returnA)
import Control.Category ((>>>))
import Data.Aeson (Value)
import qualified Data.Text as T
import qualified Data.Set as S
import Control.Lens (has, only, (^?), to)
import Data.Aeson.Lens (key, _String)
import qualified Data.Text.Lens as TL
@@ -62,7 +63,7 @@ debug = proc x -> do
returnA -< x
traceEvent :: Show a => HASS (Event a) (Event a)
traceEvent = Mealy $ \nt req -> \case
traceEvent = Mealy mempty $ \nt req -> \case
Event a -> nt (Trace req a) >>= \() -> pure (Event a, traceEvent)
Tick -> pure (Tick, traceEvent)
@@ -107,20 +108,20 @@ entityChangeEvent :: T.Text -> Mealy eff (Event Value) (Event Value)
entityChangeEvent entityId = entityChangeEvent' entityId >>> toEvent
entityChangeEvent' :: T.Text -> Mealy eff (Event Value) (Either () Value)
entityChangeEvent' entityId = events >>| filterA isEntity
entityChangeEvent' entityId = Mealy (S.singleton entityId) $ runMealy (events >>| filterA isEntity)
where
isEntity :: Value -> Bool
isEntity = has (key "event" . key "data" . key "entity_id" . _String . only entityId)
isEntity = has (key "event" . key "variables" . key "trigger" . key "entity_id" . _String . only entityId)
entityRead' :: (Read a) => T.Text -> Mealy eff (Event Value) (Either () a)
entityRead' entityId = entityChangeEvent' entityId >>| (arr state >>> arr (maybe (Left ()) Right))
where
state v = v ^? key "event" . key "data" . key "new_state" . key "state" . _String . TL.unpacked . to read
state v = v ^? key "event" . key "variables" . key "trigger" . key "to_state" . key "state" . _String . TL.unpacked . to read
entityBool' :: T.Text -> Mealy eff (Event Value) (Either () Bool)
entityBool' entityId = entityChangeEvent' entityId >>| (arr state >>> arr (maybe (Left ()) Right))
where
state v = v ^? key "event" . key "data" . key "new_state" . key "state" . _String . TL.unpacked . to toBool
state v = v ^? key "event" . key "variables" . key "trigger" . key "to_state" . key "state" . _String . TL.unpacked . to toBool
toBool = \case
"on" -> True
"off" -> False
+1 -1
View File
@@ -79,7 +79,7 @@ ikeaQuickButton entityId =
>>| arr (maybe (Left ()) Right . eventType)
>>> toEvent
where
eventType v = v ^? key "event" . key "data" . key "new_state" . key "attributes" . key "event_type" . _String . to toIkeaQuickButton . traversed
eventType v = v ^? key "event" . key "variables" . key "trigger" . key "to_state" . key "attributes" . key "event_type" . _String . to toIkeaQuickButton . traversed
data BedroomControls
+6 -4
View File
@@ -35,7 +35,7 @@ import Control.Monad.Fix (MonadFix)
import HomeAssistant.Controller.Ruuvi (ruuviController)
step :: (MonadFix m, MonadIO m) => (forall x. eff x -> m x) -> UUID -> Mealy eff a b -> a -> m (b, Mealy eff a b)
step nt trace (Mealy f) a = do
step nt trace (Mealy _ f) a = do
now <- liftIO getCurrentTime
f nt (Request now trace) a
@@ -47,7 +47,7 @@ controllers =
, Controller "bedroom-button" bedroomButtonController False -- This works but leaving for vacation
, Controller "bedroom-drawer" bedroomDrawerController True
, Controller "bedroom-humidifier" humidifierController False
, Controller "ruuvi-controller" ruuviController False
, Controller "ruuvi-controller" ruuviController True
]
-- | Steps the machine for every inbound message; service calls go to the
@@ -71,8 +71,10 @@ defaultMain = withSocketsDo $ do
withBus severity $ \bus -> do
token <- getEnv "HA_TOKEN"
host <- getEnv "HA_HOST"
let workers =
[ ("reader", readerAction host 8123 token bus)
let active = [c | c@(Controller _ _ True) <- controllers]
ents = foldMap (\(Controller _ m _) -> entities m) active
workers =
[ ("reader", readerAction host 8123 token ents bus)
, ("writer", writerAction bus)
] ++ [ (name, runController bus c) | c@(Controller name _ True) <- controllers ]
as <- mapM (\(name, act) -> async (supervised name defaultBackoff act)) workers
+18 -12
View File
@@ -18,10 +18,11 @@ import Control.Concurrent.STM
import Control.Exception (onException)
import Control.Exception.Annotated (throw)
import Control.Lens ((^?))
import Control.Monad (forever)
import Control.Monad (forever, forM_)
import Data.Aeson (Value, eitherDecode, encode, object, (.=))
import Data.Aeson.Lens (key, _String)
import qualified Data.ByteString.Lazy as BL
import qualified Data.Set as S
import qualified Data.Text as T
import Data.Void (Void)
import HomeAssistant.Controller (Service (..), Target (..))
@@ -35,11 +36,11 @@ import AFRP (Request(..))
-- | Connect, authenticate, subscribe, then receive and broadcast forever.
-- Restarting this action reconnects. All setup sends happen before the
-- connection is published in the bus, so only the writer sends afterwards.
readerAction :: String -> Int -> String -> Bus -> IO Void
readerAction host port token bus =
readerAction :: String -> Int -> String -> S.Set T.Text -> Bus -> IO Void
readerAction host port token ents bus =
WS.runClient host port "/api/websocket" $ \conn -> do
handshake conn token
subscribe bus conn
subscribe bus conn ents
atomically $ writeTVar (busConn bus) (Just conn)
putStrLn "[reader] connected"
-- Unpublish on exit so the writer blocks and the backlog survives the outage.
@@ -62,14 +63,19 @@ expectType expected msg =
Just t | t == expected -> pure ()
_ -> throw (Fatal $ "expected " <> expected <> ", got: " <> T.pack (show msg))
subscribe :: Bus -> WS.Connection -> IO ()
subscribe bus conn = do
sid <- generateCallId (busGen bus)
WS.sendTextData conn $ encode $ object
[ "id" .= sid
, "type" .= ("subscribe_events" :: T.Text)
, "event_type" .= ("state_changed" :: T.Text)
]
subscribe :: Bus -> WS.Connection -> S.Set T.Text -> IO ()
subscribe bus conn ents =
forM_ (S.toList ents) $ \entityId -> do
print entityId
sid <- generateCallId (busGen bus)
WS.sendTextData conn $ encode $ object
[ "id" .= sid
, "type" .= ("subscribe_trigger" :: T.Text)
, "trigger" .= object
[ "platform" .= ("state" :: T.Text)
, "entity_id" .= entityId
]
]
-- | Undecodable messages are skipped: reconnecting cannot fix a decode
-- problem, so crashing here would only produce a hot restart loop.