Bedroom spec and partially fix writing
LLM had done "rate limiting" via deduplicating outbound calls. But the ordering matters..
This commit is contained in:
@@ -44,10 +44,10 @@ bedroomPresence = presence "binary_sensor.presence_sensor_bedroom_occupancy"
|
|||||||
|
|
||||||
bedroomPresenceController :: HASS (Event Value) ()
|
bedroomPresenceController :: HASS (Event Value) ()
|
||||||
bedroomPresenceController = proc x -> do
|
bedroomPresenceController = proc x -> do
|
||||||
p <- bedroomPresence -< x
|
p <- bedroomPresence >>> traceEvent -< x
|
||||||
case p of
|
case p of
|
||||||
Event Unoccupied -> callService createBedroomScene >>> callService (light bedroomLights Off) -< ()
|
Event Unoccupied -> callService createBedroomScene >>> callService (light bedroomLights Off) -< ()
|
||||||
Event Occupied -> callService (activateScene "makuuhuone_lights_snapshot") -< ()
|
Event Occupied -> callService (activateScene "scene.makuuhuone_lights_snapshot") -< ()
|
||||||
_ -> returnA -< ()
|
_ -> returnA -< ()
|
||||||
|
|
||||||
data IkeaGesture
|
data IkeaGesture
|
||||||
|
|||||||
@@ -39,6 +39,7 @@ import qualified Network.WebSockets as WS
|
|||||||
import Katip (runKatipContextT, sl, logFM, Severity (..), ls)
|
import Katip (runKatipContextT, sl, logFM, Severity (..), ls)
|
||||||
import Data.UUID (toText)
|
import Data.UUID (toText)
|
||||||
import AFRP (Request(..), Event(..))
|
import AFRP (Request(..), Event(..))
|
||||||
|
import Control.Monad.IO.Class (liftIO)
|
||||||
|
|
||||||
-- | Connect, authenticate, subscribe, then receive and broadcast forever.
|
-- | Connect, authenticate, subscribe, then receive and broadcast forever.
|
||||||
-- Restarting this action reconnects. All setup sends happen before the
|
-- 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)
|
Left err -> throw (Fatal $ "Invalid JSON from Home Assistant: " <> T.pack err)
|
||||||
Right x -> pure x
|
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 -> WS.Connection -> Request -> Service -> IO ()
|
||||||
sendWithId bus conn request svc = do
|
sendWithId bus conn request svc =
|
||||||
callId <- generateCallId (busGen bus)
|
runKatipContextT (busLogEnv bus) (sl "traceId" (toText (requestTraceId request))) "connection" $ do
|
||||||
let textData = encode $ encodeService callId svc
|
callId <- liftIO $ generateCallId (busGen bus)
|
||||||
runKatipContextT (busLogEnv bus) (sl "traceId" (toText (requestTraceId request))) "connection" $
|
let textData = encode $ encodeService callId svc
|
||||||
logFM DebugS (ls textData)
|
logFM DebugS (ls textData)
|
||||||
recordOutbound bus
|
liftIO $ recordOutbound bus
|
||||||
WS.sendTextData conn textData
|
liftIO $ WS.sendTextData conn textData
|
||||||
|
|
||||||
writerAction :: Bus -> IO Void
|
writerAction :: Bus -> IO Void
|
||||||
writerAction bus = forever $ do
|
writerAction bus = forever $ do
|
||||||
first <- atomically $ readTChan (busOutbound bus)
|
(request, svc) <- atomically $ readTChan (busOutbound bus)
|
||||||
rest <- drainTry (busOutbound bus)
|
|
||||||
let deduped = dedupeBatch (first : rest)
|
|
||||||
conn <- atomically $ readTVar (busConn bus) >>= maybe retry pure
|
conn <- atomically $ readTVar (busConn bus) >>= maybe retry pure
|
||||||
forM_ deduped $ \(request, svc) -> do
|
sendWithId bus conn request svc
|
||||||
sendWithId bus conn request svc
|
|
||||||
threadDelay minInterval
|
|
||||||
|
|
||||||
encodeService :: Int -> Service -> Value
|
encodeService :: Int -> Service -> Value
|
||||||
encodeService callId Service{..} = object $
|
encodeService callId Service{..} = object $
|
||||||
|
|||||||
Reference in New Issue
Block a user