From 39bc2c0ec13d726669474e6bda1b367a643ec9ac Mon Sep 17 00:00:00 2001 From: Mats Rauhala Date: Thu, 20 Aug 2026 19:15:23 +0300 Subject: [PATCH] Add concurrent runtime implementation plan --- .../plans/2026-08-20-concurrent-runtime.md | 869 ++++++++++++++++++ 1 file changed, 869 insertions(+) create mode 100644 docs/superpowers/plans/2026-08-20-concurrent-runtime.md diff --git a/docs/superpowers/plans/2026-08-20-concurrent-runtime.md b/docs/superpowers/plans/2026-08-20-concurrent-runtime.md new file mode 100644 index 0000000..4b16e06 --- /dev/null +++ b/docs/superpowers/plans/2026-08-20-concurrent-runtime.md @@ -0,0 +1,869 @@ +# Concurrent Runtime Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Support multiple controllers over a single websocket connection: one reader thread broadcasting to per-controller channels, a writer thread sending queued service calls, supervised worker threads with backoff restarts, and reconnect-on-network-failure. + +**Architecture:** A `Bus` (broadcast `TChan` inbound, `TChan` outbound, `TVar (Maybe Connection)`, existing `CallIdGen`) connects a reader thread (connect/auth/subscribe/receive-broadcast; restart = reconnect) and a writer thread to N controller threads, each stepping a `Mealy HASSEff` machine over its own dup'd channel. Slice 1 wires them with `mapConcurrently_` (any crash exits). Slice 2 wraps each worker in `supervised` (annotated-exception, exponential backoff, `Fatal` rethrow). + +**Tech Stack:** Haskell (GHC 9.10, GHC2024), stm, async, annotated-exception, websockets, aeson, hspec, hedgehog. Nix/cabal build. + +**Spec:** `docs/superpowers/specs/2026-08-20-concurrent-runtime-design.md` + +## Global Constraints + +- Run cabal only through the devShell: `nix develop -c cabal build` / `nix develop -c cabal test`. +- After changing the .cabal file, regenerate the derivation: `nix run nixpkgs#cabal2nix -- ./. > default.nix`. Never edit it by hand. +- `-Wall` (common `warnings` stanza) must produce zero warnings. +- Comments are minimal, contract-style: describe what, not how. +- Never read, print, or commit `secrets.yaml`. +- Commit after every task. Message style: short imperative, no prefix (see `git log`). +- Unit tests: hspec. Property tests: hedgehog. +- Modules are vertical (split by feature/concept, not layer). + +## File Structure + +| File | Responsibility | +|---|---| +| `src/HomeAssistant/Runtime/Bus.hs` (create) | Shared channels + connection cell + call-id gen + channel interpreter | +| `src/HomeAssistant/Runtime/Connection.hs` (create) | Reader/writer thread actions, pure `encodeService` | +| `src/HomeAssistant/Runtime/Supervisor.hs` (create, slice 2) | Generic restart-with-backoff combinator, `Backoff`, `Fatal` | +| `src/HomeAssistant/Runtime.hs` (rewrite) | Glue: `Controller`, `runController`, `controllers`, `defaultMain`, `step`, `dryRunHassEval` | +| `src/HomeAssistant/Controller.hs` (1-word change) | Derive `Eq` on `Service` (tests need it) | +| `test/Main.hs` (rewrite) | hspec runner | +| `test/BusSpec.hs`, `test/ConnectionSpec.hs`, `test/RuntimeSpec.hs`, `test/SupervisorSpec.hs`, `test/BackoffProp.hs` (create) | Specs | +| `home-assistant-controller.cabal` | Deps + module lists per task | + +Slices: Tasks 1–3 = slice 1 (concurrency). Tasks 4–5 = slice 2 (robustness). + +--- + +### Task 1: Bus module + +**Files:** +- Create: `src/HomeAssistant/Runtime/Bus.hs` +- Modify: `src/HomeAssistant/Runtime.hs` (move `CallIdGen` out, re-export from Bus) +- Modify: `src/HomeAssistant/Controller.hs:43` (derive `Eq`) +- Modify: `home-assistant-controller.cabal` (exposed module, `stm` dep, test deps) +- Rewrite: `test/Main.hs` +- Create: `test/BusSpec.hs` + +**Interfaces:** +- Consumes: `HASSEff(..)`, `Service` from `HomeAssistant.Controller`. +- Produces (Bus exports): `Bus(..)` with fields `busInbound :: TChan Value`, `busOutbound :: TChan Service`, `busConn :: TVar (Maybe WS.Connection)`, `busGen :: CallIdGen`; `CallIdGen(..)` (record field `generateCallId :: IO Int`); `mkCallIdGen :: Int -> IO CallIdGen`; `newBus :: Int -> IO Bus`; `channelHassEval :: Bus -> HASSEff a -> IO a`. + +- [ ] **Step 1: Write the failing test** + +`test/BusSpec.hs`: + +```haskell +module BusSpec (spec) where + +import Control.Concurrent.STM + ( atomically + , dupTChanIO + , readTChan + , writeTChan + ) +import Data.Aeson (Value (..)) +import HomeAssistant.Controller (HASSEff (..), Service (..)) +import HomeAssistant.Runtime.Bus +import Test.Hspec + +spec :: Spec +spec = describe "Bus" $ do + it "broadcasts inbound messages to every dup'd channel in order" $ do + bus <- newBus 0 + p1 <- dupTChanIO (busInbound bus) + p2 <- dupTChanIO (busInbound bus) + atomically $ writeTChan (busInbound bus) (Number 1) + atomically $ writeTChan (busInbound bus) (Number 2) + r1 <- atomically $ (,) <$> readTChan p1 <*> readTChan p1 + r2 <- atomically $ (,) <$> readTChan p2 <*> readTChan p2 + r1 `shouldBe` (Number 1, Number 2) + r2 `shouldBe` (Number 1, Number 2) + + it "channelHassEval writes CallService to the outbound channel" $ do + bus <- newBus 0 + let svc = Service "light" "turn_on" Nothing "light.bedroom_masse" + channelHassEval bus (CallService svc) + atomically (readTChan (busOutbound bus)) `shouldReturn` svc + + it "channelHassEval leaves Pure untouched" $ do + bus <- newBus 0 + channelHassEval bus (Pure 42) `shouldReturn` (42 :: Int) + + it "generates unique sequential call ids" $ do + gen <- mkCallIdGen 0 + a <- generateCallId gen + b <- generateCallId gen + (a, b) `shouldBe` (1, 2) +``` + +`test/Main.hs`: + +```haskell +module Main (main) where + +import Test.Hspec (hspec) +import qualified BusSpec + +main :: IO () +main = hspec BusSpec.spec +``` + +In `home-assistant-controller.cabal`: + +Library section: add to `exposed-modules`: `HomeAssistant.Runtime.Bus`; add to `build-depends`: `stm`. + +Test suite section: add `other-modules: BusSpec` and set: + +``` + build-depends: + base ^>=4.20.2.0, + home-assistant-controller, + hspec, + stm, + aeson, + text +``` + +In `src/HomeAssistant/Controller.hs` change line 43 from `deriving Show` to `deriving (Show, Eq)`. + +- [ ] **Step 2: Run test to verify it fails** + +Run: `nix develop -c cabal test` +Expected: FAIL — compile error, `Could not find module 'HomeAssistant.Runtime.Bus'`. + +- [ ] **Step 3: Write the implementation** + +`src/HomeAssistant/Runtime/Bus.hs`: + +```haskell +{-# LANGUAGE LambdaCase #-} + +module HomeAssistant.Runtime.Bus + ( Bus(..) + , CallIdGen(..) + , mkCallIdGen + , newBus + , channelHassEval + ) where + +import Control.Concurrent.STM + ( TChan + , TVar + , atomically + , newBroadcastTChanIO + , newTChanIO + , newTVarIO + , writeTChan + ) +import Data.Aeson (Value) +import Data.IORef (IORef, atomicModifyIORef', newIORef) +import HomeAssistant.Controller (HASSEff (..), Service) +import Network.WebSockets (Connection) + +-- | Shared runtime state: inbound is a broadcast channel (controllers +-- read from 'dupTChanIO' copies), outbound queues service calls for the +-- writer, conn holds the current websocket (Nothing before first connect). +data Bus = Bus + { busInbound :: TChan Value + , busOutbound :: TChan Service + , busConn :: TVar (Maybe Connection) + , busGen :: CallIdGen + } + +newBus :: Int -> IO Bus +newBus start = Bus + <$> newBroadcastTChanIO + <*> newTChanIO + <*> newTVarIO Nothing + <*> mkCallIdGen start + +channelHassEval :: Bus -> HASSEff a -> IO a +channelHassEval bus = \case + CallService svc -> atomically $ writeTChan (busOutbound bus) svc + Pure a -> pure a + +newtype CallIdGen = CallIdGen { generateCallId :: IO Int } + +mkCallIdGen :: Int -> IO CallIdGen +mkCallIdGen start = do + gen <- newIORef start + pure $ CallIdGen $ atomicModifyIORef' gen (\old -> let new = old + 1 in new `seq` (new, new)) +``` + +Then in `src/HomeAssistant/Runtime.hs`: +- Delete the local `newtype CallIdGen`/`mkCallIdGen` definitions (lines 105–110). +- Add import: `import HomeAssistant.Runtime.Bus (Bus(..), CallIdGen(..), mkCallIdGen)`. +- Export list stays the same (`CallIdGen`, `mkCallIdGen` now re-exported from Bus). + +- [ ] **Step 4: Regenerate the nix derivation** + +Run: `nix run nixpkgs#cabal2nix -- ./. > default.nix` +Expected: no output; `git diff default.nix` shows the new `stm`/test dependencies. + +- [ ] **Step 5: Run tests to verify they pass** + +Run: `nix develop -c cabal build && nix develop -c cabal test` +Expected: build with zero warnings; 4 examples, 0 failures. + +- [ ] **Step 6: Commit** + +```bash +git add src test home-assistant-controller.cabal default.nix +git commit -m "Add Runtime.Bus with channel-based effect interpreter" +``` + +--- + +### Task 2: Connection module (reader/writer actions) + +**Files:** +- Create: `src/HomeAssistant/Runtime/Connection.hs` +- Create: `test/ConnectionSpec.hs` +- Modify: `test/Main.hs`, `home-assistant-controller.cabal` + +**Interfaces:** +- Consumes: `Bus(..)`, `CallIdGen(..)` from Task 1. +- Produces (Connection exports): `readerAction :: String -> Int -> String -> Bus -> IO Void` (host, port, token; restart = reconnect); `writerAction :: Bus -> IO Void`; `encodeService :: Int -> Service -> Data.Aeson.Value`. + +**Background for the implementer (verified facts, do not re-verify):** +- `WS.runClient` runs the `ClientApp` under `bracket` and closes the socket/stream when the app throws — reconnects do not leak fds. +- Exceptions from `WS.receiveData` (e.g. `ConnectionClosed`) are what make the reader restartable; connect failures throw `IOException`. +- HA sends a `result` ack for `subscribe_events`; we deliberately do not read it — it flows to controllers and is filtered out by their entity-id lenses (same as current behavior). + +- [ ] **Step 1: Write the failing test** + +`test/ConnectionSpec.hs`: + +```haskell +module ConnectionSpec (spec) where + +import Data.Aeson (Value, object, (.=)) +import Data.Text (Text) +import HomeAssistant.Controller (Service (..)) +import HomeAssistant.Runtime.Connection (encodeService) +import Test.Hspec + +spec :: Spec +spec = describe "encodeService" $ do + it "encodes a call_service message" $ + encodeService 7 (Service "light" "turn_on" Nothing "light.bedroom_masse") + `shouldBe` object + [ "id" .= (7 :: Int) + , "type" .= ("call_service" :: Text) + , "domain" .= ("light" :: Text) + , "service" .= ("turn_on" :: Text) + , "target" .= object ["entity_id" .= ("light.bedroom_masse" :: Text)] + ] + + it "includes service_data when present" $ + encodeService 8 (Service "light" "turn_on" (Just (object ["brightness" .= (200 :: Int)])) "light.bedroom_masse") + `shouldBe` object + [ "id" .= (8 :: Int) + , "type" .= ("call_service" :: Text) + , "domain" .= ("light" :: Text) + , "service" .= ("turn_on" :: Text) + , "target" .= object ["entity_id" .= ("light.bedroom_masse" :: Text)] + , "service_data" .= object ["brightness" .= (200 :: Int)] + ] +``` + +`test/Main.hs`: add `import qualified ConnectionSpec` and change `main` to: + +```haskell +main :: IO () +main = hspec $ do + BusSpec.spec + ConnectionSpec.spec +``` + +Cabal: library `exposed-modules` += `HomeAssistant.Runtime.Connection`; test `other-modules` += `ConnectionSpec`. + +- [ ] **Step 2: Run test to verify it fails** + +Run: `nix develop -c cabal test` +Expected: FAIL — `Could not find module 'HomeAssistant.Runtime.Connection'`. + +- [ ] **Step 3: Write the implementation** + +`src/HomeAssistant/Runtime/Connection.hs`: + +```haskell +{-# LANGUAGE OverloadedStrings #-} +{-# LANGUAGE RecordWildCards #-} + +module HomeAssistant.Runtime.Connection + ( readerAction + , writerAction + , encodeService + ) where + +import Control.Concurrent.STM + ( atomically + , readTChan + , readTVar + , retry + , writeTChan + , writeTVar + ) +import Control.Lens ((^?)) +import Control.Monad (forever) +import Data.Aeson (Value, eitherDecode, encode, object, (.=)) +import Data.Aeson.Lens (key, _String) +import qualified Data.ByteString.Lazy as BL +import qualified Data.Text as T +import Data.Void (Void) +import HomeAssistant.Controller (Service (..)) +import HomeAssistant.Runtime.Bus +import qualified Network.WebSockets as WS + +-- | 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 = + WS.runClient host port "/api/websocket" $ \conn -> do + handshake conn token + subscribe bus conn + atomically $ writeTVar (busConn bus) (Just conn) + putStrLn "[reader] connected" + receiveLoop bus conn + +handshake :: WS.Connection -> String -> IO () +handshake conn token = do + required <- receiveJSON conn + expectType "auth_required" required + WS.sendTextData conn $ encode $ object + [ "type" .= ("auth" :: T.Text) + , "access_token" .= token + ] + ok <- receiveJSON conn + expectType "auth_ok" ok + +expectType :: T.Text -> Value -> IO () +expectType expected msg = + case msg ^? key "type" . _String of + Just t | t == expected -> pure () + _ -> fail $ "expected " <> T.unpack expected <> ", got: " <> 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) + ] + +-- | Undecodable messages are skipped: reconnecting cannot fix a decode +-- problem, so crashing here would only produce a hot restart loop. +receiveLoop :: Bus -> WS.Connection -> IO Void +receiveLoop bus conn = forever $ do + msg <- WS.receiveData conn :: IO BL.ByteString + case eitherDecode msg of + Left err -> putStrLn $ "[reader] skipping undecodable message: " <> err + Right v -> atomically $ writeTChan (busInbound bus) v + +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 + +writerAction :: Bus -> IO Void +writerAction bus = forever $ do + svc <- atomically $ readTChan (busOutbound bus) + conn <- atomically $ readTVar (busConn bus) >>= maybe retry pure + callId <- generateCallId (busGen bus) + WS.sendTextData conn $ encode $ encodeService callId svc + +encodeService :: Int -> Service -> Value +encodeService callId Service{..} = object $ + [ "id" .= callId + , "type" .= ("call_service" :: T.Text) + , "domain" .= serviceDomain + , "service" .= serviceName + , "target" .= object ["entity_id" .= serviceTarget] + ] <> maybe [] (\d -> ["service_data" .= d]) serviceData +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `nix develop -c cabal build && nix develop -c cabal test` +Expected: zero warnings; 6 examples, 0 failures. + +- [ ] **Step 5: Commit** + +```bash +git add src test home-assistant-controller.cabal +git commit -m "Add Runtime.Connection reader and writer actions" +``` + +--- + +### Task 3: Runtime rewrite — controller threads and wiring (completes slice 1) + +**Files:** +- Rewrite: `src/HomeAssistant/Runtime.hs` +- Create: `test/RuntimeSpec.hs` +- Modify: `test/Main.hs`, `home-assistant-controller.cabal` + +**Interfaces:** +- Consumes: `Bus(..)`, `newBus`, `channelHassEval`, `CallIdGen(..)`, `mkCallIdGen` (Task 1); `readerAction`, `writerAction` (Task 2). +- Produces (Runtime exports): `defaultMain :: IO ()`, `step`, `dryRunHassEval`, `CallIdGen`, `mkCallIdGen`, `Controller(..)` with `data Controller = forall b. Controller T.Text (HASS (Event Value) b)`, `runController :: Bus -> Controller -> IO Void`. +- Removes exports: `app`, `hassEval`, `wsCallService`, `receiveJSON` (deleted or internal to Connection). The `get_states` debug dump goes away. + +**Background:** `mapConcurrently_` rethrows the first worker exception and cancels the rest — a crash exits the process, same as today. That is intended for slice 1. + +- [ ] **Step 1: Write the failing test** + +`test/RuntimeSpec.hs`: + +```haskell +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] + ] + ] + ] +``` + +Note: the first event only seeds the machine (`changes` does not fire on the initial value), hence two expected service calls, not three. + +`test/Main.hs`: add `import qualified RuntimeSpec`; run all three specs. Cabal: test `other-modules` += `RuntimeSpec`, test `build-depends` += `async`. + +- [ ] **Step 2: Run test to verify it fails** + +Run: `nix develop -c cabal test` +Expected: FAIL — `Controller` / `runController` not in scope. + +- [ ] **Step 3: Rewrite the Runtime module** + +Replace `src/HomeAssistant/Runtime.hs` entirely with: + +```haskell +{-# LANGUAGE ExistentialQuantification #-} +{-# LANGUAGE GADTs #-} +{-# LANGUAGE LambdaCase #-} +{-# LANGUAGE OverloadedStrings #-} + +module HomeAssistant.Runtime + ( defaultMain + , step + , CallIdGen + , mkCallIdGen + , dryRunHassEval + , Controller(..) + , runController + ) where + +import AFRP (Event (..), Mealy (..)) +import Control.Concurrent.Async (mapConcurrently_) +import Control.Concurrent.STM (atomically, dupTChanIO, readTChan) +import Data.Aeson (Value) +import qualified Data.Text as T +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) + +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 <- dupTChanIO (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" + 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 + CallService x -> do + callId <- generateCallId gen + print (callId, x) + Pure a -> pure a +``` + +- [ ] **Step 4: Regenerate derivation and run tests** + +Run: `nix run nixpkgs#cabal2nix -- ./. > default.nix && nix develop -c cabal build && nix develop -c cabal test` +Expected: zero warnings; 7 examples, 0 failures. + +- [ ] **Step 5: Commit** + +```bash +git add src test home-assistant-controller.cabal default.nix +git commit -m "Run controllers on their own threads over the bus" +``` + +--- + +### Task 4: Supervisor module (starts slice 2) + +**Files:** +- Create: `src/HomeAssistant/Runtime/Supervisor.hs` +- Create: `test/SupervisorSpec.hs`, `test/BackoffProp.hs` +- Modify: `test/Main.hs`, `home-assistant-controller.cabal` + +**Interfaces:** +- Consumes: nothing from other runtime modules (generic). +- Produces: `supervised :: Text -> Backoff -> IO Void -> IO Void`; `Backoff(..)` (`backoffBase`, `backoffCap`, `backoffQuiet`, all `NominalDiffTime`); `defaultBackoff` (= `Backoff 0.1 30 30`); `backoffDelay :: Backoff -> Int -> NominalDiffTime`; `nextAttempt :: Backoff -> NominalDiffTime -> Int -> Int`; `Fatal (..)` (`newtype Fatal = Fatal Text`). + +**Background (verified, do not re-verify):** `Control.Exception.Annotated` is built on safe-exceptions: its `catch`/`catches`/`try` only catch *synchronous* exceptions, so async exceptions propagate without any manual `SomeAsyncException` filtering. `catch @(AnnotatedException e)` sees through a single `AnnotatedException` wrapper (bare `e` gets an empty annotation set), which is how the `Fatal` handler below works. + +**Invariant (document in code):** actions passed to `supervised` must not be wrapped in `checkpoint`s around `Fatal`-throwing code — a double-wrapped `Fatal` is invisible to the handler and would restart instead of crashing. The reader honors this (Task 5). + +- [ ] **Step 1: Write the failing tests** + +`test/SupervisorSpec.hs`: + +```haskell +module SupervisorSpec (spec) where + +import Control.Concurrent (newEmptyMVar, putMVar, readMVar, threadDelay) +import Control.Concurrent.Async (async, cancel, poll, waitCatch) +import Control.Exception (SomeException, fromException) +import Control.Exception.Annotated (AnnotatedException (..), throw) +import Control.Monad (forever) +import Data.IORef (atomicModifyIORef', newIORef, readIORef) +import Data.Maybe (isJust, isNothing) +import HomeAssistant.Runtime.Supervisor +import System.IO.Error (ioError, userError) +import Test.Hspec + +tinyBackoff :: Backoff +tinyBackoff = Backoff 0.001 0.002 0.001 + +spec :: Spec +spec = describe "supervised" $ do + it "restarts a crashing action until it stays up" $ do + counter <- newIORef (0 :: Int) + up <- newEmptyMVar + let action = do + n <- atomicModifyIORef' counter (\c -> (c + 1, c + 1)) + if n < 3 + then ioError (userError "boom") + else do putMVar up (); forever (threadDelay 1000000) + sup <- async (supervised "test" tinyBackoff action) + readMVar up + threadDelay 50000 + status <- poll sup + isNothing status `shouldBe` True + readIORef counter `shouldReturn` 3 + cancel sup + + it "rethrows Fatal instead of restarting" $ do + counter <- newIORef (0 :: Int) + let action = do + atomicModifyIORef' counter (\c -> (c + 1, c + 1)) + throw (Fatal "auth_invalid") + sup <- async (supervised "test" tinyBackoff action) + res <- waitCatch sup + case res of + Left se -> case fromException se :: Maybe (AnnotatedException Fatal) of + Just _ -> pure () + Nothing -> expectationFailure "expected Fatal to propagate" + Right _ -> expectationFailure "supervised returned" + threadDelay 50000 + readIORef counter `shouldReturn` 1 + + it "does not restart on async exceptions" $ do + counter <- newIORef (0 :: Int) + let action = do + atomicModifyIORef' counter (\c -> (c + 1, c + 1)) + forever (threadDelay 1000000) + sup <- async (supervised "test" tinyBackoff action) + threadDelay 100000 + cancel sup + threadDelay 100000 + readIORef counter `shouldReturn` 1 + status <- poll sup + isJust status `shouldBe` True + + describe "nextAttempt" $ do + it "resets after a quiet period" $ + nextAttempt tinyBackoff 0.001 5 `shouldBe` 1 + it "increments otherwise" $ + nextAttempt tinyBackoff 0.0005 5 `shouldBe` 6 + + describe "backoffDelay" $ do + it "starts at base" $ backoffDelay tinyBackoff 1 `shouldBe` 0.001 + it "doubles" $ backoffDelay tinyBackoff 2 `shouldBe` 0.002 + it "clamps at cap" $ backoffDelay tinyBackoff 3 `shouldBe` 0.002 +``` + +`test/BackoffProp.hs`: + +```haskell +module BackoffProp (spec) where + +import Data.Time (NominalDiffTime) +import Hedgehog +import qualified Hedgehog.Gen as Gen +import qualified Hedgehog.Range as Range +import HomeAssistant.Runtime.Supervisor (Backoff (..), backoffDelay) +import Test.Hspec (Spec, describe) + +spec :: Spec +spec = describe "backoffDelay" $ + prop "doubles from base, clamped at cap" $ property $ do + baseD <- forAll $ Gen.double (Range.constant 0.0001 10) + ratio <- forAll $ Gen.double (Range.constant 1 100) + let base = realToFrac baseD :: NominalDiffTime + cap = realToFrac (baseD * ratio) :: NominalDiffTime + backoff = Backoff base cap 1 + delays = map (backoffDelay backoff) [1 .. 100 :: Int] + head delays === min cap base + mapM_ (\(a, b) -> b === min cap (a * 2)) (zip delays (drop 1 delays)) +``` + +`test/Main.hs`: add imports and run all five specs. Cabal: library `exposed-modules` += `HomeAssistant.Runtime.Supervisor`, library `build-depends` += `annotated-exception`; test `other-modules` += `SupervisorSpec, BackoffProp`, test `build-depends` += `hedgehog, annotated-exception, time`. + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `nix develop -c cabal test` +Expected: FAIL — `Could not find module 'HomeAssistant.Runtime.Supervisor'`. + +- [ ] **Step 3: Write the implementation** + +`src/HomeAssistant/Runtime/Supervisor.hs`: + +```haskell +{-# LANGUAGE ScopedTypeVariables #-} + +module HomeAssistant.Runtime.Supervisor + ( supervised + , Backoff(..) + , defaultBackoff + , backoffDelay + , nextAttempt + , Fatal(..) + ) where + +import Control.Concurrent (threadDelay) +import Control.Exception.Annotated + ( Exception + , Handler (..) + , SomeException + , catches + , displayException + , throw + ) +import Data.Text (Text) +import qualified Data.Text as T +import Data.Time (NominalDiffTime, diffUTCTime, getCurrentTime) +import Data.Void (Void) + +-- | A failure that cannot be fixed by restarting; propagates out of +-- 'supervised' and terminates the process. +newtype Fatal = Fatal Text + deriving (Show, Exception) + +data Backoff = Backoff + { backoffBase :: NominalDiffTime -- ^ Delay before the first restart + , backoffCap :: NominalDiffTime -- ^ Maximum delay between restarts + , backoffQuiet :: NominalDiffTime -- ^ Uptime after which the delay resets + } + deriving (Eq, Show) + +defaultBackoff :: Backoff +defaultBackoff = Backoff 0.1 30 30 + +-- | Delay before the @attempt@-th restart: doubles from base, clamped at cap. +backoffDelay :: Backoff -> Int -> NominalDiffTime +backoffDelay (Backoff base cap _) attempt = go (attempt - 1) base + where + go 0 d = d + go n d = go (n - 1) (min cap (d * 2)) + +-- | Attempt number to use after a crash that ran for the given uptime. +nextAttempt :: Backoff -> NominalDiffTime -> Int -> Int +nextAttempt (Backoff _ _ quiet) uptime attempt + | uptime >= quiet = 1 + | otherwise = attempt + 1 + +-- | Runs the action forever, restarting it with backoff after synchronous +-- exceptions; async exceptions propagate. 'Fatal' is rethrown (crashing the +-- caller) rather than restarted. The action must never return normally and +-- must not be wrapped in checkpoints around 'Fatal'-throwing code: a +-- doubly-wrapped 'Fatal' is indistinguishable from a crash and would be +-- restarted instead of escalated. +supervised :: Text -> Backoff -> IO Void -> IO Void +supervised name backoff action = go 1 + where + go attempt = do + start <- getCurrentTime + outcome <- (action >> pure (Nothing :: Maybe (Either Fatal SomeException))) `catches` + [ Handler $ \(f :: Fatal) -> pure (Just (Left f)) + , Handler $ \(e :: SomeException) -> pure (Just (Right e)) + ] + case outcome of + Nothing -> error "unreachable: supervised action returned" + Just (Left f) -> throw f + Just (Right e) -> do + putStrLn $ "[" <> T.unpack name <> "] attempt " <> show attempt <> " crashed: " <> displayException e + let delay = backoffDelay backoff attempt + putStrLn $ "[" <> T.unpack name <> "] restarting in " <> show delay <> "s" + threadDelay (round (realToFrac delay * 1000000)) + end <- getCurrentTime + go (nextAttempt backoff (diffUTCTime end start) attempt) +``` + +- [ ] **Step 4: Regenerate derivation and run tests** + +Run: `nix run nixpkgs#cabal2nix -- ./. > default.nix && nix develop -c cabal build && nix develop -c cabal test` +Expected: zero warnings; 15 examples and 1 property, 0 failures. + +- [ ] **Step 5: Commit** + +```bash +git add src test home-assistant-controller.cabal default.nix +git commit -m "Add supervisor with backoff restarts" +``` + +--- + +### Task 5: Supervised wiring and fatal auth (completes slice 2) + +**Files:** +- Modify: `src/HomeAssistant/Runtime.hs` (defaultMain only) +- Modify: `src/HomeAssistant/Runtime/Connection.hs` (Fatal in handshake) + +**Interfaces:** +- Consumes: `supervised`, `defaultBackoff`, `Fatal` (Task 4). +- Produces: no new exports; `defaultMain` behavior changes: workers are supervised, auth failure exits the process with `Fatal`. + +**Background:** `waitAny` rethrows the exception of the first completed async. Supervised workers only complete by rethrowing `Fatal`, so `waitAny` blocks forever in normal operation and propagates `Fatal` otherwise. + +- [ ] **Step 1: Change `defaultMain` in `src/HomeAssistant/Runtime.hs`** + +Update imports: replace `mapConcurrently_` with `async, waitAny` from `Control.Concurrent.Async`; add `absurd` to the `Data.Void` import; add `import HomeAssistant.Runtime.Supervisor (defaultBackoff, supervised)`. Replace `defaultMain` with: + +```haskell +defaultMain :: IO () +defaultMain = withSocketsDo $ do + token <- getEnv "HA_TOKEN" + bus <- newBus 0 + let workers = + [ ("reader", readerAction "last-resort-redux" 8123 token bus) + , ("writer", writerAction bus) + ] ++ [ (name, runController bus c) | c@(Controller name _) <- controllers ] + as <- mapM (\(name, act) -> async (supervised name defaultBackoff act)) workers + (_, v) <- waitAny as + absurd v +``` + +- [ ] **Step 2: Make handshake failures fatal in `src/HomeAssistant/Runtime/Connection.hs`** + +Add imports: `import HomeAssistant.Runtime.Supervisor (Fatal)` and `throw` from `Control.Exception.Annotated`. Change `expectType` and `receiveJSON`: + +```haskell +expectType :: T.Text -> Value -> IO () +expectType expected msg = + case msg ^? key "type" . _String of + Just t | t == expected -> pure () + _ -> throw (Fatal $ "expected " <> expected <> ", got: " <> T.pack (show msg)) +``` + +```haskell +receiveJSON :: WS.Connection -> IO Value +receiveJSON conn = do + msg <- WS.receiveData conn + case eitherDecode msg of + Left err -> throw (Fatal $ "Invalid JSON from Home Assistant: " <> T.pack err) + Right x -> pure x +``` + +The `fail`-based behavior (crash the process) is preserved for slice 1's semantics but now carries a `Fatal` marker the supervisor escalates. Do **not** add checkpoints around the handshake — see the `supervised` contract. + +- [ ] **Step 3: Build and test** + +Run: `nix develop -c cabal build && nix develop -c cabal test` +Expected: zero warnings, zero test failures. + +- [ ] **Step 4: Manual smoke check (optional, needs real HA)** + +Run: `HA_TOKEN=... nix develop -c cabal run home-assistant-controller` (only if a Home Assistant instance is reachable; otherwise skip — unit tests cover the wiring logic). + +- [ ] **Step 5: Commit** + +```bash +git add src +git commit -m "Supervise workers and make auth failures fatal" +``` + +--- + +## Verification (all tasks) + +- `nix develop -c cabal build` — zero warnings under `-Wall`. +- `nix develop -c cabal test` — all specs green. +- `git status` clean after each commit. + +## Non-goals (from spec) + +- Request/response correlation for service calls; dynamic controller registration; structured logging; env-based host/port config; state pre-seeding via `get_states`; multiple websocket connections. + +## Gaps accepted by spec + +- Reader's connect/auth/subscribe loop has no integration test (thin IO glue over `WS.runClient`; localhost fake-server scaffolding judged not worth the complexity). +- `Fatal` detection relies on the no-checkpoints-around-handshake invariant (documented in `supervised`'s contract and Task 5).