From 64561718a5fd34c1eb689dba869087939bbfb642 Mon Sep 17 00:00:00 2001 From: Mats Rauhala Date: Wed, 16 Sep 2026 10:52:10 +0300 Subject: [PATCH] Bedroom spec and partially fix writing LLM had done "rate limiting" via deduplicating outbound calls. But the ordering matters.. --- src/HomeAssistant/Controller/Bedroom.hs | 4 +-- src/HomeAssistant/Runtime/Connection.hs | 36 +++++++------------------ 2 files changed, 11 insertions(+), 29 deletions(-) diff --git a/src/HomeAssistant/Controller/Bedroom.hs b/src/HomeAssistant/Controller/Bedroom.hs index 0748fbe..fd80c0b 100644 --- a/src/HomeAssistant/Controller/Bedroom.hs +++ b/src/HomeAssistant/Controller/Bedroom.hs @@ -44,10 +44,10 @@ bedroomPresence = presence "binary_sensor.presence_sensor_bedroom_occupancy" bedroomPresenceController :: HASS (Event Value) () bedroomPresenceController = proc x -> do - p <- bedroomPresence -< x + p <- bedroomPresence >>> traceEvent -< x case p of Event Unoccupied -> callService createBedroomScene >>> callService (light bedroomLights Off) -< () - Event Occupied -> callService (activateScene "makuuhuone_lights_snapshot") -< () + Event Occupied -> callService (activateScene "scene.makuuhuone_lights_snapshot") -< () _ -> returnA -< () data IkeaGesture diff --git a/src/HomeAssistant/Runtime/Connection.hs b/src/HomeAssistant/Runtime/Connection.hs index d77a1cd..ecad52e 100644 --- a/src/HomeAssistant/Runtime/Connection.hs +++ b/src/HomeAssistant/Runtime/Connection.hs @@ -39,6 +39,7 @@ import qualified Network.WebSockets as WS import Katip (runKatipContextT, sl, logFM, Severity (..), ls) import Data.UUID (toText) import AFRP (Request(..), Event(..)) +import Control.Monad.IO.Class (liftIO) -- | Connect, authenticate, subscribe, then receive and broadcast forever. -- Restarting this action reconnects. All setup sends happen before the @@ -111,41 +112,22 @@ receiveJSON conn = do Left err -> throw (Fatal $ "Invalid JSON from Home Assistant: " <> T.pack err) Right x -> pure x --- | Floor between sends within a batch: 100ms, so a many-distinct-target --- flood still caps at ~10 sends/sec even after dedupe. -minInterval :: Int -minInterval = 100000 --- | Non-blocking drain of everything queued on a channel. Returns items --- oldest-first (FIFO from the channel), so prepending the blocking --- `readTChan` item keeps the whole batch oldest-first for `dedupeBatch`. -drainTry :: TChan a -> IO [a] -drainTry chan = go [] - where - go acc = do - m <- atomically $ tryReadTChan chan - case m of - Nothing -> pure (reverse acc) - Just x -> go (x : acc) sendWithId :: Bus -> WS.Connection -> Request -> Service -> IO () -sendWithId bus conn request svc = do - callId <- generateCallId (busGen bus) - let textData = encode $ encodeService callId svc - runKatipContextT (busLogEnv bus) (sl "traceId" (toText (requestTraceId request))) "connection" $ +sendWithId bus conn request svc = + runKatipContextT (busLogEnv bus) (sl "traceId" (toText (requestTraceId request))) "connection" $ do + callId <- liftIO $ generateCallId (busGen bus) + let textData = encode $ encodeService callId svc logFM DebugS (ls textData) - recordOutbound bus - WS.sendTextData conn textData + liftIO $ recordOutbound bus + liftIO $ WS.sendTextData conn textData writerAction :: Bus -> IO Void writerAction bus = forever $ do - first <- atomically $ readTChan (busOutbound bus) - rest <- drainTry (busOutbound bus) - let deduped = dedupeBatch (first : rest) + (request, svc) <- atomically $ readTChan (busOutbound bus) conn <- atomically $ readTVar (busConn bus) >>= maybe retry pure - forM_ deduped $ \(request, svc) -> do - sendWithId bus conn request svc - threadDelay minInterval + sendWithId bus conn request svc encodeService :: Int -> Service -> Value encodeService callId Service{..} = object $