From 9c355adf47d7122e88c153e864e39604636c150d Mon Sep 17 00:00:00 2001 From: Mats Rauhala Date: Thu, 20 Aug 2026 19:40:01 +0300 Subject: [PATCH] Run controllers on their own threads over the bus --- default.nix | 8 +- home-assistant-controller.cabal | 5 +- src/HomeAssistant/Runtime.hs | 127 ++++++++++---------------------- test/Main.hs | 2 + test/RuntimeSpec.hs | 36 +++++++++ 5 files changed, 84 insertions(+), 94 deletions(-) create mode 100644 test/RuntimeSpec.hs diff --git a/default.nix b/default.nix index 5061456..9227cea 100644 --- a/default.nix +++ b/default.nix @@ -1,5 +1,5 @@ -{ mkDerivation, aeson, base, bytestring, hspec, lens, lens-aeson -, lib, network, stm, text, time, websockets +{ mkDerivation, aeson, async, base, bytestring, hspec, lens +, lens-aeson, lib, network, stm, text, time, websockets }: mkDerivation { pname = "home-assistant-controller"; @@ -8,11 +8,11 @@ mkDerivation { isLibrary = true; isExecutable = true; libraryHaskellDepends = [ - aeson base bytestring lens lens-aeson network stm text time + aeson async base bytestring lens lens-aeson network stm text time websockets ]; executableHaskellDepends = [ base ]; - testHaskellDepends = [ aeson base hspec stm text ]; + testHaskellDepends = [ aeson async base hspec stm text ]; license = lib.meta.getLicenseFromSpdxId "BSD-3-Clause"; mainProgram = "home-assistant-controller"; } diff --git a/home-assistant-controller.cabal b/home-assistant-controller.cabal index adda1d6..7336ecc 100644 --- a/home-assistant-controller.cabal +++ b/home-assistant-controller.cabal @@ -82,6 +82,7 @@ library , bytestring , time , stm + , async -- Directories containing source files. hs-source-dirs: src @@ -123,6 +124,7 @@ test-suite home-assistant-controller-test -- Modules included in this executable, other than Main. other-modules: BusSpec , ConnectionSpec + , RuntimeSpec -- LANGUAGE extensions used by modules in this package. -- other-extensions: @@ -143,4 +145,5 @@ test-suite home-assistant-controller-test hspec, stm, aeson, - text + text, + async diff --git a/src/HomeAssistant/Runtime.hs b/src/HomeAssistant/Runtime.hs index d198174..c0f7dc4 100644 --- a/src/HomeAssistant/Runtime.hs +++ b/src/HomeAssistant/Runtime.hs @@ -1,113 +1,62 @@ -{-# LANGUAGE OverloadedStrings #-} -{-# LANGUAGE LambdaCase #-} +{-# LANGUAGE ExistentialQuantification #-} {-# LANGUAGE GADTs #-} +{-# LANGUAGE LambdaCase #-} +{-# LANGUAGE OverloadedStrings #-} module HomeAssistant.Runtime ( defaultMain - , app , step , CallIdGen , mkCallIdGen - , hassEval , dryRunHassEval - , receiveJSON - , wsCallService + , Controller(..) + , runController ) where -import AFRP (Mealy(..), Event(..)) -import HomeAssistant.Controller (HASSEff(..), lightController, Service(..)) -import HomeAssistant.Runtime.Bus (CallIdGen(..), mkCallIdGen) -import Data.Aeson ((.=), Value (Null), encode, eitherDecode, object) -import qualified Data.ByteString.Lazy as BL +import AFRP (Event (..), Mealy (..)) +import Control.Concurrent.Async (mapConcurrently_) +import Control.Concurrent.STM (atomically, dupTChan, readTChan) +import Data.Aeson (Value) import qualified Data.Text as T -import qualified Network.WebSockets as WS +import Data.Time (getCurrentTime) +import Data.Void (Void) +import HomeAssistant.Controller (HASS, HASSEff (..), lightController) +import HomeAssistant.Runtime.Bus +import HomeAssistant.Runtime.Connection (readerAction, writerAction) import Network.Socket (withSocketsDo) import System.Environment (getEnv) -import Data.Time (getCurrentTime) step :: (forall x. eff x -> IO x) -> Mealy eff a b -> a -> IO (b, Mealy eff a b) step nt (Mealy f) a = do now <- getCurrentTime f nt now a +data Controller = forall b. Controller T.Text (HASS (Event Value) b) + +controllers :: [Controller] +controllers = [Controller "light" lightController] + +-- | Steps the machine for every inbound message; service calls go to the +-- bus. A restart re-dups the inbound channel and starts from the machine's +-- initial state; messages broadcast during the restart window are lost. +runController :: Bus -> Controller -> IO Void +runController bus (Controller _name machine) = do + inbound <- atomically (dupTChan (busInbound bus)) + go inbound machine + where + go inbound f = do + msg <- atomically (readTChan inbound) + (_, f') <- step (channelHassEval bus) f (Event msg) + go inbound f' + defaultMain :: IO () defaultMain = withSocketsDo $ do - token <- getEnv "HA_TOKEN" - gen <- mkCallIdGen 0 - WS.runClient "last-resort-redux" 8123 "/api/websocket" (app gen token) - -app :: CallIdGen -> String -> WS.ClientApp () -app gen token conn = do - -- HA speaks first: {"type":"auth_required", ...} - authRequired <- receiveJSON conn - print authRequired - - WS.sendTextData conn $ encode $ object - [ "type" .= ("auth" :: T.Text) - , "access_token" .= token - ] - - -- Expect {"type":"auth_ok", ...} - authResult <- receiveJSON conn - print authResult - - getStateId <- generateCallId gen - WS.sendTextData conn $ encode $ object - [ "id" .= getStateId - , "type" .= ("get_states" :: T.Text) - ] - msg <- WS.receiveData conn :: IO BL.ByteString - BL.writeFile "/tmp/states.json" msg - - subscribeId <- generateCallId gen - -- Subscription 1: all entity state changes - WS.sendTextData conn $ encode $ object - [ "id" .= subscribeId - , "type" .= ("subscribe_events" :: T.Text) - , "event_type" .= ("state_changed" :: T.Text) - ] - - go lightController - - where - go f = do - msg <- WS.receiveData conn :: IO BL.ByteString - let decoded = Event $ either (const Null) id $ eitherDecode @Value msg - (x, f') <- step (dryRunHassEval gen) f decoded - mapM_ print x - go f' - -receiveJSON :: WS.Connection -> IO Value -receiveJSON conn = do - msg <- WS.receiveData conn - case eitherDecode msg of - Left err -> fail $ "Invalid JSON from Home Assistant: " ++ err - Right x -> pure x - -wsCallService - :: WS.Connection - -> Int -- ^ request id - -> T.Text -- ^ domain - -> T.Text -- ^ service - -> T.Text -- ^ entity id - -> IO () -wsCallService conn requestId domain service entityId = - WS.sendTextData conn $ encode $ object - [ "id" .= requestId - , "type" .= ("call_service" :: T.Text) - , "domain" .= domain - , "service" .= service - , "target" .= object - [ "entity_id" .= entityId - ] - ] - -hassEval :: CallIdGen -> WS.Connection -> HASSEff a -> IO a -hassEval gen conn = \case - CallService x -> do - callId <- generateCallId gen - wsCallService conn callId (serviceDomain x) (serviceName x) (serviceTarget x) - Pure a -> pure a + token <- getEnv "HA_TOKEN" + bus <- newBus 0 + mapConcurrently_ id $ + [ readerAction "last-resort-redux" 8123 token bus + , writerAction bus + ] ++ map (runController bus) controllers dryRunHassEval :: CallIdGen -> HASSEff a -> IO a dryRunHassEval gen = \case diff --git a/test/Main.hs b/test/Main.hs index 0c2739e..f4dd601 100644 --- a/test/Main.hs +++ b/test/Main.hs @@ -3,8 +3,10 @@ module Main (main) where import Test.Hspec (hspec) import qualified BusSpec import qualified ConnectionSpec +import qualified RuntimeSpec main :: IO () main = hspec $ do BusSpec.spec ConnectionSpec.spec + RuntimeSpec.spec diff --git a/test/RuntimeSpec.hs b/test/RuntimeSpec.hs new file mode 100644 index 0000000..03ccbc5 --- /dev/null +++ b/test/RuntimeSpec.hs @@ -0,0 +1,36 @@ +{-# LANGUAGE OverloadedStrings #-} +module RuntimeSpec (spec) where + +import Control.Concurrent (threadDelay) +import Control.Concurrent.Async (async) +import Control.Concurrent.STM (atomically, readTChan, writeTChan) +import Data.Aeson (Value, object, (.=)) +import Data.Text (Text) +import HomeAssistant.Controller (light, lightController) +import HomeAssistant.Runtime (Controller (..), runController) +import HomeAssistant.Runtime.Bus +import Test.Hspec + +spec :: Spec +spec = describe "runController" $ do + it "feeds inbound events through the machine and forwards service calls" $ do + bus <- newBus 0 + _ <- async (runController bus (Controller "test" lightController)) + threadDelay 100000 -- let the controller dup its inbound channel + atomically $ writeTChan (busInbound bus) (doorEvent "on") -- initial value: no change event + atomically $ writeTChan (busInbound bus) (doorEvent "off") -- door closes: lights on + atomically $ writeTChan (busInbound bus) (doorEvent "on") -- door opens: lights off + svc1 <- atomically (readTChan (busOutbound bus)) + svc2 <- atomically (readTChan (busOutbound bus)) + svc1 `shouldBe` light True + svc2 `shouldBe` light False + +doorEvent :: Text -> Value +doorEvent state = object + [ "event" .= object + [ "data" .= object + [ "entity_id" .= ("binary_sensor.makuuhuone_ovi_contact" :: Text) + , "new_state" .= object ["state" .= state] + ] + ] + ]