Compare commits
28
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0f7c6ff98a | ||
|
|
aaba917a9b | ||
|
|
a3dd63f26c | ||
|
|
b20320c779 | ||
|
|
e16010bf2f | ||
|
|
1f786815e6 | ||
|
|
54ef6b5cc3 | ||
|
|
1edb2ac5a2 | ||
|
|
d3518e8ff2 | ||
|
|
e30a28c9ef | ||
|
|
f7849ac889 | ||
|
|
0a4b7dbb01 | ||
|
|
88af6a3139 | ||
|
|
055c02d0f2 | ||
|
|
7b6b529120 | ||
|
|
ef4811f89b | ||
|
|
9fadbae777 | ||
|
|
3f5aa4735c | ||
|
|
6bb96833e1 | ||
|
|
2730657952 | ||
|
|
771d702abd | ||
|
|
7051eaf244 | ||
|
|
7a9f30e1be | ||
|
|
aa8ea07249 | ||
|
|
2664ca67c0 | ||
|
|
d341270ed1 | ||
|
|
dc55834507 | ||
|
|
50c82393ed |
@@ -5,3 +5,8 @@ dist-newstyle
|
||||
.worktrees/
|
||||
|
||||
docs/superpowers
|
||||
|
||||
*.hp
|
||||
*.eventlog
|
||||
*.eventlog.html
|
||||
*.rrd
|
||||
|
||||
+11
-6
@@ -1,6 +1,8 @@
|
||||
{ mkDerivation, aeson, annotated-exception, async, base, bytestring
|
||||
, hedgehog, hspec, hspec-hedgehog, katip, lens, lens-aeson, lib
|
||||
, network, stm, text, time, uuid, websockets
|
||||
, cereal, cereal-conduit, conduit, containers, directory, ekg-core
|
||||
, filepath, hedgehog, hspec, hspec-hedgehog, katip, lens
|
||||
, lens-aeson, lib, network, process, stm, text, time
|
||||
, unordered-containers, uuid, websockets
|
||||
}:
|
||||
mkDerivation {
|
||||
pname = "home-assistant-controller";
|
||||
@@ -9,13 +11,16 @@ mkDerivation {
|
||||
isLibrary = true;
|
||||
isExecutable = true;
|
||||
libraryHaskellDepends = [
|
||||
aeson annotated-exception async base bytestring katip lens
|
||||
lens-aeson network stm text time uuid websockets
|
||||
aeson annotated-exception async base bytestring cereal
|
||||
cereal-conduit conduit containers directory ekg-core filepath katip
|
||||
lens lens-aeson network process stm text time unordered-containers
|
||||
uuid websockets
|
||||
];
|
||||
executableHaskellDepends = [ base ];
|
||||
testHaskellDepends = [
|
||||
aeson annotated-exception async base hedgehog hspec hspec-hedgehog
|
||||
stm text time
|
||||
aeson annotated-exception async base cereal containers directory
|
||||
ekg-core hedgehog hspec hspec-hedgehog katip process stm text time
|
||||
unordered-containers uuid
|
||||
];
|
||||
license = lib.meta.getLicenseFromSpdxId "BSD-3-Clause";
|
||||
mainProgram = "home-assistant-controller";
|
||||
|
||||
@@ -20,7 +20,15 @@
|
||||
});
|
||||
});
|
||||
in rec {
|
||||
packages.home-assistant-controller = pkgs.haskell.lib.justStaticExecutables hp.home-assistant-controller;
|
||||
packages.home-assistant-controller = pkgs.symlinkJoin {
|
||||
name = "home-assistant-controller";
|
||||
paths = [ (pkgs.haskell.lib.justStaticExecutables hp.home-assistant-controller) ];
|
||||
nativeBuildInputs = [ pkgs.makeWrapper ];
|
||||
postBuild = ''
|
||||
wrapProgram $out/bin/home-assistant-controller \
|
||||
--set HA_RRDTOOL ${pkgs.lib.getBin pkgs.rrdtool}/bin/rrdtool
|
||||
'';
|
||||
};
|
||||
defaultPackage = packages.home-assistant-controller;
|
||||
devShell = hp.shellFor {
|
||||
packages = h: [h.home-assistant-controller];
|
||||
@@ -34,6 +42,8 @@
|
||||
hp.graphmod
|
||||
|
||||
hp.haskell-language-server
|
||||
|
||||
rrdtool
|
||||
];
|
||||
};
|
||||
}
|
||||
|
||||
@@ -62,9 +62,13 @@ library
|
||||
exposed-modules: AFRP
|
||||
, HomeAssistant.Controller
|
||||
, HomeAssistant.Controller.Bedroom
|
||||
, HomeAssistant.Controller.Kitchen
|
||||
, HomeAssistant.Controller.Children
|
||||
, HomeAssistant.Controller.Ruuvi
|
||||
, HomeAssistant.Runtime
|
||||
, HomeAssistant.Runtime.Bus
|
||||
, HomeAssistant.Runtime.Connection
|
||||
, HomeAssistant.Runtime.Metrics
|
||||
, HomeAssistant.Runtime.Supervisor
|
||||
|
||||
-- Modules included in this library but not exported.
|
||||
@@ -88,6 +92,16 @@ library
|
||||
, annotated-exception
|
||||
, uuid
|
||||
, katip
|
||||
, containers
|
||||
, ekg-core
|
||||
, unordered-containers
|
||||
, process
|
||||
, directory
|
||||
, cereal
|
||||
, containers
|
||||
, filepath
|
||||
, conduit
|
||||
, cereal-conduit
|
||||
|
||||
-- Directories containing source files.
|
||||
hs-source-dirs: src
|
||||
@@ -118,7 +132,7 @@ executable home-assistant-controller
|
||||
|
||||
-- Base language which the package is written in.
|
||||
default-language: GHC2024
|
||||
ghc-options: -threaded
|
||||
ghc-options: -threaded -with-rtsopts=-T
|
||||
|
||||
test-suite home-assistant-controller-test
|
||||
-- Import common warning flags.
|
||||
@@ -128,11 +142,16 @@ test-suite home-assistant-controller-test
|
||||
default-language: GHC2024
|
||||
|
||||
-- Modules included in this executable, other than Main.
|
||||
other-modules: BusSpec
|
||||
other-modules: AFRPLawsSpec
|
||||
, AFRPSpec
|
||||
, BackoffProp
|
||||
, BedroomSpec
|
||||
, BusSpec
|
||||
, ConnectionSpec
|
||||
, MetricsSpec
|
||||
, RuntimeSpec
|
||||
, SupervisorSpec
|
||||
, BackoffProp
|
||||
, Support
|
||||
|
||||
-- LANGUAGE extensions used by modules in this package.
|
||||
-- other-extensions:
|
||||
@@ -153,6 +172,7 @@ test-suite home-assistant-controller-test
|
||||
hspec,
|
||||
stm,
|
||||
aeson,
|
||||
cereal,
|
||||
text,
|
||||
async,
|
||||
hedgehog,
|
||||
@@ -160,4 +180,9 @@ test-suite home-assistant-controller-test
|
||||
annotated-exception,
|
||||
time,
|
||||
uuid,
|
||||
katip
|
||||
katip,
|
||||
containers,
|
||||
ekg-core,
|
||||
unordered-containers,
|
||||
process,
|
||||
directory
|
||||
|
||||
+476
-87
@@ -1,12 +1,16 @@
|
||||
{-# LANGUAGE LambdaCase #-}
|
||||
{-# LANGUAGE Arrows #-}
|
||||
|
||||
module AFRP
|
||||
( Mealy(..)
|
||||
, Auto(..)
|
||||
, DecodedAuto(..)
|
||||
, eff
|
||||
, withEntities
|
||||
, Event(..)
|
||||
, hold
|
||||
, events
|
||||
, switch
|
||||
-- , switch
|
||||
, preMapAccum
|
||||
, preMapAccumRequest
|
||||
, mapAccum
|
||||
@@ -18,83 +22,308 @@ module AFRP
|
||||
, (>>|)
|
||||
, toEvent
|
||||
, lMerge
|
||||
, Pair(..)
|
||||
, Request(..)
|
||||
, SerializeUTCTime(..)
|
||||
, SerializeLocalTime(..)
|
||||
, edge
|
||||
, waitFor
|
||||
, duration
|
||||
, tag
|
||||
, isEvent
|
||||
, delayEvent
|
||||
, sample
|
||||
, rollup
|
||||
, sliding
|
||||
, debounce
|
||||
, currentTime
|
||||
, onEvent
|
||||
, save
|
||||
, load
|
||||
, stepAuto
|
||||
, stepAutoSerializing
|
||||
) where
|
||||
|
||||
import Control.Category (Category(..), (>>>))
|
||||
import Prelude hiding ((.), id)
|
||||
import Control.Arrow (Arrow(..), ArrowChoice(..), ArrowLoop(..))
|
||||
import Data.Time (UTCTime, NominalDiffTime, diffUTCTime, addUTCTime)
|
||||
import Control.Monad.Fix (MonadFix (mfix))
|
||||
import Control.Arrow (Arrow(..), ArrowChoice(..))
|
||||
import Data.Time (UTCTime (UTCTime), NominalDiffTime, diffUTCTime, addUTCTime, TimeZone, LocalTime (LocalTime), utcToLocalTime, Day (..), diffTimeToPicoseconds, picosecondsToDiffTime, TimeOfDay (TimeOfDay), diffLocalTime)
|
||||
import Data.Either (fromLeft)
|
||||
import Data.Bool (bool)
|
||||
import Data.UUID (UUID)
|
||||
import qualified Data.Set as S
|
||||
import qualified Data.Text as T
|
||||
import Data.Serialize (Get, Putter, Serialize (put), runGet, get)
|
||||
import qualified Data.ByteString as B
|
||||
import Control.Exception (IOException, handle, throwIO)
|
||||
import System.IO.Error (isDoesNotExistError)
|
||||
import GHC.Generics (Generic)
|
||||
import Data.Sequence (Seq, (|>))
|
||||
import qualified Data.Foldable as F
|
||||
import Control.Monad.IO.Class (MonadIO, liftIO)
|
||||
import Conduit (ConduitT, (.|))
|
||||
import qualified Data.Conduit.Cereal as CC
|
||||
import qualified Conduit as C
|
||||
|
||||
|
||||
data Codec s = Codec { getter :: !(Get s), putter :: !(Putter s) }
|
||||
|
||||
data State s = State {state :: !s, dirty :: !Bool}
|
||||
deriving Functor
|
||||
|
||||
|
||||
instance Semigroup s => Semigroup (State s) where
|
||||
s1 <> s2 = State (state s1 <> state s2) (dirty s1 || dirty s2)
|
||||
|
||||
instance Monoid s => Monoid (State s) where
|
||||
mempty = State mempty False
|
||||
|
||||
instance Applicative State where
|
||||
pure a = State a False
|
||||
s1 <*> s2 = State
|
||||
{ state =
|
||||
let a = state s2
|
||||
f = state s1
|
||||
in f a
|
||||
, dirty = dirty s1 || dirty s2
|
||||
}
|
||||
|
||||
mergeState :: State s1 -> State s2 -> State (s1, s2)
|
||||
mergeState s1 s2 = (,) <$> s1 <*> s2
|
||||
|
||||
|
||||
mergeCodec :: Codec s -> Codec s1 -> Codec (s, s1)
|
||||
mergeCodec (Codec agetter aputter) (Codec bgetter bputter) = Codec (mergeGet agetter bgetter) (mergePut aputter bputter)
|
||||
where
|
||||
mergePut :: Putter s -> Putter s1 -> Putter (s, s1)
|
||||
mergePut p1 p2 (s, s1) = p1 s >> p2 s1
|
||||
mergeGet :: Get s -> Get s' -> Get (s, s')
|
||||
mergeGet g1 g2 = (,) <$> g1 <*> g2
|
||||
|
||||
data Pair a b = Pair !a !b
|
||||
|
||||
data Request = Request
|
||||
{ requestTime :: !UTCTime
|
||||
, requestTimeZone :: !TimeZone
|
||||
, requestTraceId :: !UUID
|
||||
}
|
||||
deriving Show
|
||||
} deriving (Show, Eq)
|
||||
|
||||
newtype Mealy eff a b = Mealy
|
||||
{ runMealy :: forall m. MonadFix m => (forall x. eff x -> m x) -> Request -> a -> m (b, Mealy eff a b) }
|
||||
|
||||
data Auto m a b
|
||||
= Fun (Request -> a -> b) -- Stateless variant, needed at least for 'id'
|
||||
| forall s. Stateful !(Codec s) !(State s) !(State s -> Request -> a -> m (b, State s)) -- State is explicitly part of it
|
||||
|
||||
|
||||
instance Monad m => Functor (Auto m a) where
|
||||
fmap f = \case
|
||||
Fun x -> Fun $ \req -> f . x req
|
||||
Stateful codec s x -> Stateful codec s $ \s' req a -> do
|
||||
(a',s'') <- x s' req a
|
||||
pure (f a', s'')
|
||||
|
||||
instance Monad m => Applicative (Auto m a) where
|
||||
pure a = Fun (\_req -> const a)
|
||||
fa <*> fb =
|
||||
case (fa,fb) of
|
||||
(Fun af, Fun bf) -> Fun $ \req -> (af req <*> bf req)
|
||||
(Stateful codec s af, Fun bf) -> Stateful codec s
|
||||
(\s' req x -> do
|
||||
let a = bf req x
|
||||
(h, s'') <- af s' req x
|
||||
pure (h a, s'')
|
||||
)
|
||||
(Fun af, Stateful codec s bf) -> Stateful codec s
|
||||
(\s' req x -> do
|
||||
(a, s'') <- bf s' req x
|
||||
let h = af req x
|
||||
pure (h a, s'')
|
||||
)
|
||||
(Stateful acodec as af, Stateful bcodec bs bf) -> Stateful (mergeCodec acodec bcodec) (mergeState as bs)
|
||||
(\s' req x -> do
|
||||
(a, as') <- bf (snd <$> s') req x
|
||||
(h, bs') <- af (fst <$> s') req x
|
||||
pure (h a, mergeState bs' as')
|
||||
)
|
||||
|
||||
|
||||
instance (Monad m, Semigroup b) => Semigroup (Auto m a b) where
|
||||
fa <> fb =
|
||||
case (fa,fb) of
|
||||
(Fun af, Fun bf) -> Fun (af <> bf)
|
||||
(Stateful codec s af, Fun bf) -> Stateful codec s
|
||||
(\s' req a -> do
|
||||
(ab, s'') <- af s' req a
|
||||
let bb = bf req a
|
||||
pure (ab <> bb, s'')
|
||||
)
|
||||
(Fun af, Stateful codec s bf) -> Stateful codec s
|
||||
(\s' req a -> do
|
||||
let ab = af req a
|
||||
(bb, s'') <- bf s' req a
|
||||
pure (ab <> bb, s'')
|
||||
)
|
||||
(Stateful acodec as af , Stateful bcodec bs bf) -> Stateful (mergeCodec acodec bcodec) (mergeState as bs)
|
||||
(\s req a -> do
|
||||
(ab, as'') <- af (fst <$> s) req a
|
||||
(bb, bs'') <- bf (snd <$> s) req a
|
||||
pure (ab <> bb, mergeState as'' bs'')
|
||||
)
|
||||
|
||||
|
||||
instance (Monad m, Monoid b) => Monoid (Auto m a b) where
|
||||
mempty = Fun $ \_req _ -> mempty
|
||||
|
||||
|
||||
instance Monad m => Category (Auto m) where
|
||||
id = Fun $ \_ -> id
|
||||
af . ag =
|
||||
case (af, ag) of
|
||||
(Fun f, Fun g) -> Fun (\req -> f req . g req)
|
||||
(Stateful codec s f, Fun g) -> Stateful codec s (\s' req -> f s' req . g req)
|
||||
(Fun f, Stateful codec s g) -> Stateful codec s (\s' req -> fmap (first (f req)) . g s' req)
|
||||
(Stateful fcodec fs f , Stateful gcodec gs g) ->
|
||||
Stateful (mergeCodec fcodec gcodec) (mergeState fs gs) (\s req a -> do
|
||||
(b, s') <- g (snd <$> s) req a
|
||||
(c, s'') <- f (fst <$> s) req b
|
||||
pure (c, mergeState s'' s'))
|
||||
|
||||
|
||||
|
||||
instance Monad m => Arrow (Auto m) where
|
||||
arr f = Fun $ const f
|
||||
first = \case
|
||||
Fun f -> Fun $ \req -> first (f req)
|
||||
Stateful codec s f -> Stateful codec s $ \s' req (b,d) -> do
|
||||
(c, s'') <- f s' req b
|
||||
pure ((c,d), s'')
|
||||
|
||||
instance Monad m => ArrowChoice (Auto m) where
|
||||
left = \case
|
||||
Fun f -> Fun $ \req ->
|
||||
\case
|
||||
Left b -> Left $ f req b
|
||||
Right d -> Right d
|
||||
Stateful codec s f -> Stateful codec s $ \s' req -> \case
|
||||
Right d -> pure (Right d, s')
|
||||
Left b -> do
|
||||
(c, s'') <- f s' req b
|
||||
pure (Left c, s'')
|
||||
|
||||
|
||||
serialize :: Monad m => Auto eff a b -> ConduitT i B.ByteString m ()
|
||||
serialize = \case
|
||||
Fun _ -> CC.sourcePut (put ())
|
||||
Stateful Codec{putter} s _ -> CC.sourcePut (putter (state s))
|
||||
|
||||
data DecodedAuto m a b
|
||||
= Decoded (Auto m a b) -- decoded from serialized state
|
||||
| FailDecode String (Auto m a b) -- gives back the original + errmsg
|
||||
|
||||
deserialize :: B.ByteString -> Auto m a b -> DecodedAuto m a b
|
||||
deserialize bs = \case
|
||||
Fun f -> Decoded (Fun f) -- no state to decode, success by default
|
||||
Stateful codec s f ->
|
||||
either
|
||||
(\err -> FailDecode err (Stateful codec s f))
|
||||
(\s' -> Decoded $ Stateful codec (pure s') f)
|
||||
$ runGet (getter codec) bs
|
||||
|
||||
|
||||
save :: FilePath -> Auto m a b -> IO (Auto m a b)
|
||||
save path s
|
||||
| isDirty s = do
|
||||
-- Using conduit machinery as it handles the exception handling for me
|
||||
() <- C.runResourceT $ C.runConduit (serialize s .| C.sinkFileCautious path)
|
||||
pure $ cleanDirty s
|
||||
| otherwise = pure s
|
||||
where
|
||||
cleanDirty :: Auto m a b -> Auto m a b
|
||||
cleanDirty (Stateful codec s' f) = Stateful codec s'{dirty=False} f
|
||||
cleanDirty a = a
|
||||
isDirty :: Auto m a b -> Bool
|
||||
isDirty (Stateful _ s' _) = dirty s'
|
||||
isDirty _ = False
|
||||
|
||||
|
||||
load :: forall m a b. FilePath -> Auto m a b -> IO (DecodedAuto m a b)
|
||||
load path a = handle defaultOnMissingFile (flip deserialize a <$> B.readFile path)
|
||||
where
|
||||
defaultOnMissingFile :: IOException -> IO (DecodedAuto m a b)
|
||||
defaultOnMissingFile e
|
||||
| isDoesNotExistError e = pure $ FailDecode "State doesn't exist yet" a
|
||||
| otherwise = throwIO e
|
||||
|
||||
-- | The set of entity ids an arrow subscribes to. Static: it does not
|
||||
-- change as the machine steps, so the runtime can read it once to build
|
||||
-- trigger subscriptions.
|
||||
data Mealy eff a b = Mealy
|
||||
{ entities :: S.Set T.Text
|
||||
, runMealy :: forall m. Monad m => (forall x. eff x -> m x) -> Auto m a b
|
||||
}
|
||||
|
||||
|
||||
instance Semigroup b => Semigroup (Mealy eff a b) where
|
||||
Mealy ast af <> Mealy bst bf = Mealy (ast <> bst) $ \nt -> af nt <> bf nt
|
||||
|
||||
|
||||
instance Monoid b => Monoid (Mealy eff a b) where
|
||||
mempty = Mealy mempty $ \_nt -> mempty
|
||||
-- where
|
||||
-- m = Mealy mempty $ \_ _ _ -> pure (Pair mempty m)
|
||||
|
||||
eff :: (Request -> a -> eff b) -> Mealy eff a b
|
||||
eff f = Mealy $ \nt req x ->
|
||||
nt (f req x) >>= \b -> pure (b, eff f)
|
||||
eff f = Mealy mempty $ \nt -> do
|
||||
Stateful (Codec get put) (State () False) $ \s req a -> do
|
||||
b <- nt (f req a)
|
||||
pure (b, s)
|
||||
|
||||
-- | Override the static entity set of an arrow. Use when a combinator
|
||||
-- (e.g. 'switch') hides continuation entities from the runtime's
|
||||
-- startup subscription scan.
|
||||
withEntities :: S.Set T.Text -> Mealy eff a b -> Mealy eff a b
|
||||
withEntities es (Mealy _ f) = Mealy es f
|
||||
|
||||
instance Category (Mealy eff) where
|
||||
id = Mealy (\_ _ x -> pure (x, id))
|
||||
(Mealy f) . (Mealy g) = Mealy $ \nt t a -> do
|
||||
(b, g') <- g nt t a
|
||||
(c, f') <- f nt t b
|
||||
pure (c, f' . g')
|
||||
id = Mealy mempty (\_ -> id)
|
||||
(Mealy ast f) . (Mealy bst g) = Mealy (ast <> bst) $ \nt -> do
|
||||
f nt . g nt
|
||||
|
||||
instance Arrow (Mealy eff) where
|
||||
arr f = Mealy $ \_ _ b -> pure (f b, arr f)
|
||||
first (Mealy f) = Mealy $ \nt t (b,d) -> do
|
||||
(c, f') <- f nt t b
|
||||
pure ((c, d), first f')
|
||||
arr f = Mealy mempty $ \_nt -> arr f
|
||||
first (Mealy st f) = Mealy st $ \nt -> first (f nt)
|
||||
|
||||
instance ArrowChoice (Mealy eff) where
|
||||
left (Mealy f) = Mealy $ \nt t -> \case
|
||||
Left b -> do
|
||||
(c, f') <- f nt t b
|
||||
pure (Left c, left f')
|
||||
Right d -> pure (Right d, left (Mealy f))
|
||||
left (Mealy st f) = Mealy st $ \nt -> left (f nt)
|
||||
|
||||
instance ArrowLoop (Mealy eff) where
|
||||
loop (Mealy f) = Mealy $ \nt t b -> do
|
||||
((c,_), f') <- mfix $ \((_,d), _) -> f nt t (b,d)
|
||||
pure (c, loop f')
|
||||
|
||||
instance Functor (Mealy eff a) where
|
||||
fmap f (Mealy g) = Mealy $ \nt t a -> do
|
||||
(b, g') <- g nt t a
|
||||
pure (f b, fmap f g')
|
||||
fmap f (Mealy st g) = Mealy st $ \nt -> fmap f (g nt)
|
||||
|
||||
instance Applicative (Mealy eff a) where
|
||||
pure b = Mealy $ \_ _ _ -> pure (b, pure b)
|
||||
Mealy f <*> Mealy x = Mealy $ \nt t a -> do
|
||||
(f', fNext) <- f nt t a
|
||||
(x', xNext) <- x nt t a
|
||||
pure (f' x', fNext <*> xNext)
|
||||
pure b = Mealy mempty $ \_ -> pure b
|
||||
Mealy ast f <*> Mealy bst x = Mealy (ast <> bst) $ \nt ->
|
||||
f nt <*> x nt
|
||||
|
||||
data Event a
|
||||
= Tick
|
||||
| Event a
|
||||
deriving (Show, Functor, Foldable, Traversable)
|
||||
deriving (Show, Eq, Functor, Foldable, Traversable, Generic)
|
||||
|
||||
instance Serialize a => Serialize (Event a)
|
||||
|
||||
instance Semigroup (Event a) where
|
||||
(<>) = lMerge
|
||||
|
||||
instance Monoid (Event a) where
|
||||
mempty = Tick
|
||||
|
||||
hold :: (Serialize a, Eq a) => a -> Mealy m (Event a) a
|
||||
hold def = mapAccum step def id
|
||||
where
|
||||
step prev = \case
|
||||
Tick -> prev
|
||||
Event new -> new
|
||||
|
||||
hold :: a -> Mealy eff (Event a) a
|
||||
hold a = Mealy $ \_ _ -> \case
|
||||
Tick -> pure (a, hold a)
|
||||
Event a' -> pure (a', hold a')
|
||||
|
||||
events :: Mealy eff (Event a) (Either () a)
|
||||
events = arr $ \case
|
||||
@@ -108,51 +337,89 @@ isEvent _ = True
|
||||
tag :: b -> Event a -> Event b
|
||||
tag b ev = b <$ ev
|
||||
|
||||
switch :: Mealy eff a (b, Event c) -> (c -> Mealy eff a b) -> Mealy eff a b
|
||||
switch (Mealy f) s = Mealy $ \nt t a -> do
|
||||
((b, ev), f') <- f nt t a
|
||||
case ev of
|
||||
Tick -> pure (b, switch f' s)
|
||||
Event x -> runMealy (s x) nt t a
|
||||
-- switch :: Mealy eff a (b, Event c) -> (c -> Mealy eff a b) -> Mealy eff a b
|
||||
-- switch (Mealy st f) s = Mealy st $ \nt t a -> do
|
||||
-- Pair (b, ev) f' <- f nt t a
|
||||
-- case ev of
|
||||
-- Tick -> pure (Pair b (switch f' s))
|
||||
-- Event x -> runMealy (s x) nt t a
|
||||
|
||||
sample :: Mealy eff (a, Event b) (Event a)
|
||||
sample = arr (uncurry tag)
|
||||
|
||||
preMapAccum :: (x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
preMapAccum f x extract = go x
|
||||
where
|
||||
go b = Mealy $ \_ _ a ->
|
||||
let next = f b a
|
||||
in pure (extract b, go next)
|
||||
|
||||
preMapAccumRequest :: (Request -> x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
preMapAccumRequest f x extract = go x
|
||||
preMapAccum' :: forall m x a b. (Monad m, Eq x, Serialize x) => (x -> a -> x) -> x -> (x -> b) -> Auto m a b
|
||||
preMapAccum' f x extract = Stateful (Codec get put) (pure x) (\s _req a -> pure $ step s a)
|
||||
where
|
||||
go b = Mealy $ \_ t a ->
|
||||
let next = f t b a
|
||||
in pure (extract b, go next)
|
||||
step :: State x -> a -> (b, State x)
|
||||
step s a = let s' = f (state s) a in (extract (state s), State s' (dirty s || state s /= s'))
|
||||
|
||||
mapAccum :: (x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
mapAccum f x extract = go x
|
||||
|
||||
mapAccum :: (Eq x, Serialize x) => (x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
mapAccum step x extract = Mealy mempty $ \_nt -> mapAccum' step x extract
|
||||
|
||||
preMapAccum :: (Eq x, Serialize x) => (x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
preMapAccum step x extract = Mealy mempty $ \_nt -> preMapAccum' step x extract
|
||||
|
||||
mapAccum' :: forall m x a b. (Monad m, Eq x, Serialize x) => (x -> a -> x) -> x -> (x -> b) -> Auto m a b
|
||||
mapAccum' f x extract = Stateful (Codec get put) (pure x) (\s _req a -> pure $ step s a)
|
||||
where
|
||||
go b = Mealy $ \_ _ a ->
|
||||
let next = f b a
|
||||
in pure (extract next, go next)
|
||||
step :: State x -> a -> (b, State x)
|
||||
step s a = let s' = f (state s) a in (extract s', State s' (dirty s || state s /= s'))
|
||||
|
||||
mapAccumRequest :: (Request -> x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
mapAccumRequest f x extract = go x
|
||||
preMapAccumRequest :: (Serialize x, Eq x) => (Request -> x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
preMapAccumRequest step x extract = Mealy mempty $ \_ -> preMapAccumRequest' step x extract
|
||||
|
||||
preMapAccumRequest' :: forall m x a b. (Serialize x, Eq x, Monad m) => (Request -> x -> a -> x) -> x -> (x -> b) -> Auto m a b
|
||||
preMapAccumRequest' f x extract = Stateful (Codec get put) (State x False) (\s req a -> pure $ step s req a)
|
||||
where
|
||||
go b = Mealy $ \_ t a ->
|
||||
let next = f t b a
|
||||
in pure (extract next, go next)
|
||||
step :: State x -> Request -> a -> (b, State x)
|
||||
step s req a = let s' = f req (state s) a in (extract (state s), State s' (dirty s || state s /= s'))
|
||||
|
||||
data DelayState a = DelayState
|
||||
{ pending :: [(UTCTime, a)]
|
||||
, output :: Event a
|
||||
mapAccumRequest :: (Serialize x, Eq x) => (Request -> x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
|
||||
mapAccumRequest step x extract = Mealy mempty $ \_ -> mapAccumRequest' step x extract
|
||||
|
||||
mapAccumRequest' :: forall m x a b. (Monad m, Serialize x, Eq x) => (Request -> x -> a -> x) -> x -> (x -> b) -> Auto m a b
|
||||
mapAccumRequest' f x extract = Stateful (Codec get put) (State x False) (\s req a -> pure $ step s req a)
|
||||
where
|
||||
step :: State x -> Request -> a -> (b, State x)
|
||||
step s req a = let s' = f req (state s) a in (extract s', State s' (dirty s || state s /= s'))
|
||||
|
||||
data DelayState x a = DelayState
|
||||
{ pending :: x
|
||||
, output :: !(Event a)
|
||||
}
|
||||
deriving (Generic, Eq)
|
||||
|
||||
instance (Serialize x, Serialize a) => Serialize (DelayState x a)
|
||||
|
||||
delayEvent :: NominalDiffTime -> Mealy eff (Event a) (Event a)
|
||||
newtype SerializeUTCTime = SerializeUTCTime UTCTime
|
||||
deriving (Eq, Show)
|
||||
|
||||
instance Serialize SerializeUTCTime where
|
||||
put (SerializeUTCTime (UTCTime day time)) = do
|
||||
put (toModifiedJulianDay day)
|
||||
put (diffTimeToPicoseconds time)
|
||||
|
||||
get = do
|
||||
day <- ModifiedJulianDay <$> get
|
||||
time <- picosecondsToDiffTime <$> get
|
||||
pure $ SerializeUTCTime (UTCTime day time)
|
||||
|
||||
newtype SerializeLocalTime = SerializeLocalTime LocalTime
|
||||
deriving (Eq, Show)
|
||||
|
||||
instance Serialize SerializeLocalTime where
|
||||
put (SerializeLocalTime (LocalTime day time)) = do
|
||||
put (toModifiedJulianDay day)
|
||||
let TimeOfDay h m s = time
|
||||
put (h,m, toRational s)
|
||||
get = do
|
||||
day <- ModifiedJulianDay <$> get
|
||||
(h,m,s) <- get
|
||||
pure $ SerializeLocalTime (LocalTime day (TimeOfDay h m (fromRational s)))
|
||||
|
||||
delayEvent :: (Eq a, Serialize a) => NominalDiffTime -> Mealy eff (Event a) (Event a)
|
||||
delayEvent delay =
|
||||
mapAccumRequest step initial output
|
||||
where
|
||||
@@ -164,17 +431,32 @@ delayEvent delay =
|
||||
queued =
|
||||
case input of
|
||||
Tick -> pending st
|
||||
Event x -> pending st ++ [(delay `addUTCTime` now, x)]
|
||||
Event x -> pending st ++ [(SerializeUTCTime $ delay `addUTCTime` now, x)]
|
||||
|
||||
in case queued of
|
||||
(due, x) : rest
|
||||
(SerializeUTCTime due, x) : rest
|
||||
| due <= now ->
|
||||
DelayState rest (Event x)
|
||||
|
||||
_ ->
|
||||
DelayState queued Tick
|
||||
|
||||
changes :: Eq a => Mealy eff a (Event a)
|
||||
debounce :: (Serialize a, Eq a) => NominalDiffTime -> Mealy eff (Event a) (Event a)
|
||||
debounce delay =
|
||||
mapAccumRequest step initial output
|
||||
where
|
||||
initial = DelayState Nothing Tick
|
||||
step req st input =
|
||||
let now = requestTime req
|
||||
held = case input of
|
||||
Tick -> pending st
|
||||
Event x -> Just (SerializeUTCTime $ delay `addUTCTime` now, x)
|
||||
in case held of
|
||||
Just (SerializeUTCTime due, x)
|
||||
| due <= now -> DelayState Nothing (Event x)
|
||||
_ -> DelayState held Tick
|
||||
|
||||
changes :: (Serialize a, Eq a) => Mealy eff a (Event a)
|
||||
changes = mapAccum go Nothing (maybe Tick snd)
|
||||
where
|
||||
go :: Eq a => Maybe (a, Event a) -> a -> Maybe (a, Event a)
|
||||
@@ -207,19 +489,126 @@ lMerge (Event a) _ = Event a
|
||||
lMerge Tick (Event a) = Event a
|
||||
|
||||
edge :: Mealy eff Bool (Event ())
|
||||
edge = go False
|
||||
where
|
||||
go True = Mealy $ \_ _ -> \case
|
||||
True -> pure (Tick, go True)
|
||||
False -> pure (Tick, go False)
|
||||
go False = Mealy $ \_ _ -> \case
|
||||
True -> pure (Event (), go True)
|
||||
False -> pure (Tick, go False)
|
||||
edge =
|
||||
mapAccum
|
||||
(\(_, current) new -> (current, new))
|
||||
(False, False)
|
||||
(\(old, current) ->
|
||||
if not old && current
|
||||
then Event ()
|
||||
else Tick)
|
||||
|
||||
|
||||
|
||||
data WaitingFor
|
||||
= Waiting
|
||||
| Pending { waitingForStart :: SerializeLocalTime, waitingForCurrent :: SerializeLocalTime }
|
||||
deriving (Show, Eq, Generic)
|
||||
|
||||
instance Serialize WaitingFor
|
||||
|
||||
waitFor :: NominalDiffTime -> Mealy eff Bool (Event ())
|
||||
waitFor delta =
|
||||
mapAccumRequest step Waiting extract >>> edge
|
||||
where
|
||||
extract :: WaitingFor -> Bool
|
||||
extract Waiting = False
|
||||
extract Pending{waitingForStart=SerializeLocalTime s, waitingForCurrent=SerializeLocalTime e} =
|
||||
e `diffLocalTime` s >= delta
|
||||
step :: Request -> WaitingFor -> Bool -> WaitingFor
|
||||
step _req _prev False = Waiting
|
||||
step req prev True =
|
||||
let now = SerializeLocalTime $ requestLocalTime req
|
||||
in case prev of
|
||||
Waiting -> Pending now now
|
||||
pending -> pending{waitingForCurrent = now}
|
||||
|
||||
duration :: forall eff a. Mealy eff a NominalDiffTime
|
||||
duration = mapAccumRequest go (Nothing @(UTCTime, NominalDiffTime)) (maybe 0 snd)
|
||||
duration = mapAccumRequest go (Nothing @(SerializeUTCTime, SerializeUTCTime)) (maybe 0 delta)
|
||||
where
|
||||
go :: Request -> Maybe (UTCTime, NominalDiffTime) -> a -> Maybe (UTCTime, NominalDiffTime)
|
||||
go req Nothing _ = Just $ (requestTime req, requestTime req `diffUTCTime` requestTime req)
|
||||
go req (Just (startTime, _)) _ = Just (startTime, requestTime req `diffUTCTime` startTime)
|
||||
delta :: (SerializeUTCTime, SerializeUTCTime) -> NominalDiffTime
|
||||
delta (SerializeUTCTime start, SerializeUTCTime end) = end `diffUTCTime` start
|
||||
go :: Request -> Maybe (SerializeUTCTime, SerializeUTCTime) -> a -> Maybe (SerializeUTCTime, SerializeUTCTime)
|
||||
go req Nothing _ = Just (SerializeUTCTime $ requestTime req, SerializeUTCTime $ requestTime req)
|
||||
go req (Just (startTime, _)) _ = Just (startTime, SerializeUTCTime $ requestTime req)
|
||||
|
||||
|
||||
-- | Rollup, hold back bursty messages
|
||||
--
|
||||
-- Consider a case where you have a bursty set of data. You care to get an immediate response,
|
||||
-- but don't want to spam the output.
|
||||
rollup
|
||||
:: (Serialize a, Eq a)
|
||||
=> Int -- ^ How many items to pass through before burst protection
|
||||
-> Int -- ^ How many seconds to collect the bursty data
|
||||
-> Mealy eff (Event a) (Event [a])
|
||||
rollup limit seconds =
|
||||
mapAccumRequest
|
||||
go
|
||||
(Left Tick)
|
||||
(either id (\(_, _, _, ev) -> ev))
|
||||
where
|
||||
go
|
||||
:: Request
|
||||
-> Either (Event [a]) (SerializeUTCTime, Int, Seq a, Event [a])
|
||||
-> Event a
|
||||
-> Either (Event [a]) (SerializeUTCTime, Int, Seq a, Event [a])
|
||||
|
||||
go _ (Left _) Tick =
|
||||
Left Tick
|
||||
|
||||
go req (Left _) (Event a) =
|
||||
Right
|
||||
( SerializeUTCTime $ addUTCTime (fromIntegral seconds) (requestTime req)
|
||||
, 1
|
||||
, mempty
|
||||
, Event [a]
|
||||
)
|
||||
|
||||
go req (Right (SerializeUTCTime end, n, acc, _)) Tick
|
||||
| requestTime req >= end =
|
||||
Left (Event $ F.toList acc)
|
||||
| otherwise =
|
||||
Right (SerializeUTCTime end, n, acc, Tick)
|
||||
|
||||
go req (Right (SerializeUTCTime end, n, acc, _)) (Event a)
|
||||
| requestTime req >= end =
|
||||
Left (Event $ F.toList (acc |> a))
|
||||
| n < limit =
|
||||
Right (SerializeUTCTime end, n + 1, acc, Event [a])
|
||||
| otherwise =
|
||||
Right (SerializeUTCTime end, n + 1, acc |> a, Tick)
|
||||
|
||||
|
||||
-- Sliding window into the events
|
||||
sliding :: (Serialize a, Eq a) => Int -> Mealy eff (Event a) [a]
|
||||
sliding size = mapAccum go [] id
|
||||
where
|
||||
go :: [a] -> Event a -> [a]
|
||||
go acc Tick = acc
|
||||
go acc (Event a) = let xs = acc ++ [a] in drop (max 0 (length xs - size)) xs
|
||||
|
||||
|
||||
requestLocalTime :: Request -> LocalTime
|
||||
requestLocalTime Request{requestTime, requestTimeZone} = utcToLocalTime requestTimeZone requestTime
|
||||
|
||||
currentTime :: Mealy eff a LocalTime
|
||||
currentTime = Mealy mempty $ \_nt -> Fun $ \req _ ->
|
||||
requestLocalTime req
|
||||
|
||||
|
||||
stepAuto :: Monad m => Auto m a b -> Request -> a -> m (b, Auto m a b)
|
||||
stepAuto (Fun f) req a = pure (f req a, Fun f)
|
||||
stepAuto (Stateful codec s f) req a = do
|
||||
(b, s') <- f s req a
|
||||
pure (b, Stateful codec s' f)
|
||||
|
||||
|
||||
stepAutoSerializing :: MonadIO m => FilePath -> Auto m a b -> Request -> a -> m (b, Auto m a b)
|
||||
stepAutoSerializing path f req a = do
|
||||
(b,x) <- stepAuto f req a
|
||||
y <- liftIO $ save path x
|
||||
pure (b,y)
|
||||
|
||||
onEvent :: Mealy eff a () -> Mealy eff (Event a) ()
|
||||
onEvent f = events >>> (arr (const ()) ||| f)
|
||||
|
||||
@@ -8,16 +8,13 @@ module HomeAssistant.Controller
|
||||
, HASSEff(..)
|
||||
, HASS
|
||||
, callService
|
||||
, callServiceDyn
|
||||
, entityChangeEvent
|
||||
, entityChangeEvent'
|
||||
, entityRead
|
||||
, entityRead'
|
||||
, entityBool
|
||||
, entityBool'
|
||||
, Ruuvi(..)
|
||||
, ruuvi
|
||||
, ruuviTemperatures
|
||||
, ruuviPressures
|
||||
, DoorState(..)
|
||||
, light
|
||||
, Presence(..)
|
||||
@@ -27,20 +24,25 @@ module HomeAssistant.Controller
|
||||
, traceValue
|
||||
, switch
|
||||
, Target(..)
|
||||
, brightness
|
||||
, Light(..)
|
||||
) where
|
||||
|
||||
import AFRP (Mealy (..), eff, Event(..), hold, events, changes, filterA, (>>|), toEvent, Request)
|
||||
import AFRP (Mealy (..), eff, Event(..), events, filterA, (>>|), toEvent, Request)
|
||||
import Control.Arrow (Arrow(..), returnA)
|
||||
import Control.Category ((>>>))
|
||||
import Data.Aeson (Value)
|
||||
import Data.Aeson (Value, object, (.=))
|
||||
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 Data.Aeson.Lens (key, _String, _Integral)
|
||||
import qualified Data.Text.Lens as TL
|
||||
import Data.Bool (bool)
|
||||
import Data.Serialize (Serialize)
|
||||
import GHC.Generics (Generic)
|
||||
|
||||
data Target = EntityId !T.Text | AreaId !T.Text
|
||||
deriving (Show,Eq)
|
||||
deriving (Show,Eq,Ord)
|
||||
|
||||
data Service = Service
|
||||
{ serviceDomain :: T.Text
|
||||
@@ -60,38 +62,35 @@ type HASS a b = Mealy HASSEff a b
|
||||
callService :: Service -> HASS a ()
|
||||
callService service = eff (\req _ -> CallService req service)
|
||||
|
||||
callServiceDyn :: (a -> Service) -> HASS a ()
|
||||
callServiceDyn mkService = eff (\req a -> CallService req (mkService a))
|
||||
|
||||
debug :: Show a => HASS a a
|
||||
debug = proc x -> do
|
||||
eff (const Debug) -< x
|
||||
returnA -< x
|
||||
|
||||
traceEvent :: Show a => HASS (Event a) (Event a)
|
||||
traceEvent = Mealy $ \nt req -> \case
|
||||
Event a -> nt (Trace req a) >>= \() -> pure (Event a, traceEvent)
|
||||
Tick -> pure (Tick, traceEvent)
|
||||
traceEvent = proc ev -> do
|
||||
case ev of
|
||||
Event a -> eff Trace -< a
|
||||
Tick -> returnA -< ()
|
||||
returnA -< ev
|
||||
|
||||
traceValue :: Show a => HASS a a
|
||||
traceValue = proc x -> do
|
||||
eff Trace -< x
|
||||
returnA -< x
|
||||
|
||||
ruuviTemperatures :: Mealy eff (Event Value) Double
|
||||
ruuviTemperatures = entityRead @Double "sensor.ruuvitag_b168_temperature" >>> hold 0
|
||||
|
||||
ruuviPressures :: Mealy eff (Event Value) Double
|
||||
ruuviPressures = entityRead "sensor.ruuvitag_b168_pressure" >>> hold 0
|
||||
|
||||
data Ruuvi = Ruuvi { ruuviTemperature :: Double, ruuviPressure :: Double }
|
||||
deriving (Show, Eq)
|
||||
|
||||
ruuvi :: Mealy eff (Event Value) (Event Ruuvi)
|
||||
ruuvi = (Ruuvi <$> ruuviTemperatures <*> ruuviPressures) >>> changes
|
||||
|
||||
data DoorState = Open | Closed
|
||||
deriving (Show, Eq)
|
||||
deriving (Show, Eq, Generic)
|
||||
|
||||
instance Serialize DoorState
|
||||
|
||||
data Presence = Occupied | Unoccupied
|
||||
deriving (Show, Eq)
|
||||
deriving (Show, Eq, Generic)
|
||||
|
||||
instance Serialize Presence
|
||||
|
||||
presence :: T.Text -> HASS (Event Value) (Event Presence)
|
||||
presence entityId =entityBool entityId
|
||||
@@ -99,12 +98,21 @@ presence entityId =entityBool entityId
|
||||
|
||||
|
||||
|
||||
data Light
|
||||
= Off
|
||||
| On { brightnessPercentage :: Maybe Double }
|
||||
|
||||
-- Turn off lights when door is closed
|
||||
light :: [Target] -> Bool -> Service
|
||||
light targets b = Service
|
||||
light :: [Target] -> Light -> Service
|
||||
light targets (On {brightnessPercentage}) = Service
|
||||
{ serviceDomain="light"
|
||||
, serviceName= bool "turn_off" "turn_on" b
|
||||
, serviceName= "turn_on"
|
||||
, serviceData=fmap (\pct -> object ["brightness_pct" .= pct]) brightnessPercentage
|
||||
, serviceTarget=targets
|
||||
}
|
||||
light targets Off = Service
|
||||
{ serviceDomain="light"
|
||||
, serviceName= "turn_off"
|
||||
, serviceData=Nothing
|
||||
, serviceTarget=targets
|
||||
}
|
||||
@@ -123,20 +131,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
|
||||
@@ -147,3 +155,11 @@ entityRead entityId = entityRead' entityId >>> toEvent
|
||||
|
||||
entityBool :: T.Text -> Mealy eff (Event Value) (Event Bool)
|
||||
entityBool entityId = entityBool' entityId >>> toEvent
|
||||
|
||||
brightness :: T.Text -> HASS (Event Value) (Event Int)
|
||||
brightness entityId =
|
||||
entityChangeEvent' entityId
|
||||
>>| arr (maybe (Left ()) Right . eventBrightness)
|
||||
>>> AFRP.toEvent
|
||||
where
|
||||
eventBrightness v = v ^? key "event" . key "variables" . key "trigger" . key "to_state" . key "attributes" . key "brightness" . _Integral
|
||||
|
||||
@@ -46,7 +46,7 @@ bedroomPresenceController :: HASS (Event Value) ()
|
||||
bedroomPresenceController = proc x -> do
|
||||
p <- bedroomPresence -< x
|
||||
case p of
|
||||
Event Unoccupied -> callService createBedroomScene >>> callService (light bedroomLights False) -< ()
|
||||
Event Unoccupied -> callService createBedroomScene >>> callService (light bedroomLights Off) -< ()
|
||||
Event Occupied -> callService (activateScene "makuuhuone_lights_snapshot") -< ()
|
||||
_ -> returnA -< ()
|
||||
|
||||
@@ -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
|
||||
@@ -98,14 +98,14 @@ bedroomButtonController :: HASS (Event Value) ()
|
||||
bedroomButtonController = proc x -> do
|
||||
ev <- bedroomButton >>> traceEvent -< x
|
||||
case ev of
|
||||
Event (Masse (OnButton ShortRelease)) -> callService (activateScene "scene.makuuhuone_masse") -< () -- this should be on release
|
||||
Event (Masse (OnButton ShortRelease)) -> callService (activateScene "scene.makuuhuone_masse") -< ()
|
||||
Event (Masse (OnButton DoubleClick)) -> callService (activateScene "scene.makuuhuone_keski") -< ()
|
||||
Event (Masse (OnButton LongClick)) -> callService (activateScene "scene.makuuhuone_kirkas") -< ()
|
||||
Event (Masse (OffButton _)) -> callService (light [AreaId "makuuhuone"] False) -< ()
|
||||
Event (Masse (OffButton _)) -> callService (light [AreaId "makuuhuone"] Off) -< ()
|
||||
Event (Enishen (OnButton ShortRelease)) -> callService (activateScene "scene.makuuhuone_jemina") -< ()
|
||||
Event (Enishen (OnButton DoubleClick)) -> callService (activateScene "scene.makuuhuone_keski") -< ()
|
||||
Event (Enishen (OnButton LongClick)) -> callService (activateScene "scene.makuuhuone_kirkas") -< ()
|
||||
Event (Enishen (OffButton _)) -> callService (light [AreaId "makuuhuone"] False) -< ()
|
||||
Event (Enishen (OffButton _)) -> callService (light [AreaId "makuuhuone"] Off) -< ()
|
||||
_ -> returnA -< ()
|
||||
|
||||
|
||||
@@ -131,11 +131,12 @@ door = entityBool "binary_sensor.makuuhuone_ovi_contact"
|
||||
waitFor :: NominalDiffTime -> HASS a (Event ())
|
||||
waitFor n = duration >>> arr (> n) >>> edge
|
||||
|
||||
|
||||
delayedDoor :: HASS (Event Value) (Event DoorState)
|
||||
delayedDoor = door
|
||||
>>> (AFRP.hold Open &&& AFRP.delayEvent 15)
|
||||
>>> AFRP.sample
|
||||
>>> traceEvent
|
||||
>>> AFRP.debounce 15
|
||||
>>> AFRP.hold Open
|
||||
>>> AFRP.changes
|
||||
|
||||
humidifierController :: HASS (Event Value) ()
|
||||
humidifierController = proc x -> do
|
||||
|
||||
@@ -0,0 +1,82 @@
|
||||
{-# LANGUAGE Arrows #-}
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
module HomeAssistant.Controller.Children where
|
||||
import HomeAssistant.Controller (HASS, callService, Target (AreaId), light, Light(..))
|
||||
import AFRP (Event)
|
||||
import qualified AFRP
|
||||
import Control.Arrow (Arrow(..), (>>>))
|
||||
import Data.Time (Day, TimeOfDay (..), localDay, LocalTime (..))
|
||||
import Data.Time.Calendar.OrdinalDate (WeekOfYear, mondayStartWeek)
|
||||
import Data.Functor.Contravariant (Predicate (..), (>$<))
|
||||
|
||||
-- Let's see building some reasonable interface for utctime
|
||||
dow :: Day -> (WeekOfYear, Int)
|
||||
dow = mondayStartWeek
|
||||
|
||||
|
||||
weekday :: Predicate Day
|
||||
weekday = Predicate (betweenInclusive 1 5 . snd . dow)
|
||||
where
|
||||
betweenInclusive a b c = c >= a && c <= b
|
||||
|
||||
time :: (Int, Int) -> Predicate TimeOfDay
|
||||
time (h,m) = mconcat
|
||||
[ Predicate (equals h . todHour)
|
||||
, Predicate (equals m . todMin)
|
||||
]
|
||||
where
|
||||
equals a b = a == b
|
||||
|
||||
|
||||
atTime :: Predicate LocalTime -> HASS a (Event ())
|
||||
atTime p = AFRP.currentTime
|
||||
>>> arr (getPredicate p)
|
||||
>>> AFRP.edge
|
||||
|
||||
-- I don't have any proper presence sensors in their bedroom
|
||||
-- and they are notoriously bad at changing clothes in complete darkness
|
||||
-- So I have set up an automation that attempts to turn on the lights sometime
|
||||
-- before they leave for school and turns them off a bit later
|
||||
|
||||
-- Don't mconcat these predicates they have && behavior
|
||||
-- if you mconcat the actual arrows, they combine the behaviors of the separate branches
|
||||
-- essentially becoming || behavior
|
||||
timersOff :: [Predicate LocalTime]
|
||||
timersOff =
|
||||
[ day 1 <> at (08,15)
|
||||
, day 2 <> at (09,15)
|
||||
, day 3 <> at (08,15)
|
||||
, day 4 <> at (08,15)
|
||||
, day 5 <> at (08,15)
|
||||
, at (18,57) -- debug
|
||||
]
|
||||
where
|
||||
dayOfWeek = snd . mondayStartWeek . localDay
|
||||
at (h,m) = localTimeOfDay >$< Predicate (\TimeOfDay{todHour, todMin} -> todHour == h && todMin == m)
|
||||
day n = dayOfWeek >$< Predicate (== n)
|
||||
|
||||
timersOn :: [Predicate LocalTime]
|
||||
timersOn =
|
||||
[ day 1 <> at (07,30)
|
||||
, day 2 <> at (08,30)
|
||||
, day 3 <> at (07,30)
|
||||
, day 4 <> at (07,30)
|
||||
, day 5 <> at (07,30)
|
||||
, at (18,55) -- debug
|
||||
]
|
||||
where
|
||||
dayOfWeek = snd . mondayStartWeek . localDay
|
||||
at (h,m) = localTimeOfDay >$< Predicate (\TimeOfDay{todHour, todMin} -> todHour == h && todMin == m)
|
||||
day n = dayOfWeek >$< Predicate (== n)
|
||||
|
||||
schoolLightController :: HASS a ()
|
||||
schoolLightController = lightsOn <> lightsOff
|
||||
|
||||
|
||||
lightsOn :: HASS a ()
|
||||
lightsOn = foldMap atTime timersOn
|
||||
>>> AFRP.onEvent (callService (light [AreaId "lasten_makuuhuone"] On{brightnessPercentage = Just 100}))
|
||||
|
||||
lightsOff :: HASS a ()
|
||||
lightsOff = foldMap atTime timersOff
|
||||
>>> AFRP.onEvent (callService (light [AreaId "lasten_makuuhuone"] Off))
|
||||
@@ -0,0 +1,59 @@
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
{-# LANGUAGE Arrows #-}
|
||||
module HomeAssistant.Controller.Kitchen where
|
||||
import HomeAssistant.Controller
|
||||
import AFRP (Event (..))
|
||||
import qualified AFRP
|
||||
import Data.Aeson (Value)
|
||||
import Control.Arrow ((>>>), Arrow (..), returnA)
|
||||
import Data.Bool (bool)
|
||||
import GHC.Generics (Generic)
|
||||
import Data.Serialize (Serialize)
|
||||
import Prelude hiding (id)
|
||||
import Data.Time (LocalTime(..), TimeOfDay (..))
|
||||
|
||||
|
||||
-- Kitchen has two "presence" sensors. One IKEA motion sensor and one SwitchBot presence sensor
|
||||
data Motion = MotionDetected | MotionNotDetected | MotionUnknown
|
||||
deriving (Show, Eq, Generic)
|
||||
|
||||
instance Serialize Motion
|
||||
|
||||
data Lights = LightsOn | LightsOff
|
||||
deriving (Show)
|
||||
|
||||
kitchenMotion :: HASS (Event Value) Motion
|
||||
kitchenMotion = entityBool "binary_sensor.kitchen_movement_occupancy"
|
||||
>>> arr (fmap (bool MotionNotDetected MotionDetected))
|
||||
>>> traceEvent
|
||||
>>> AFRP.hold MotionUnknown
|
||||
|
||||
|
||||
kitchenPresence :: HASS Motion Presence
|
||||
kitchenPresence = (eventOccupied &&& eventUnoccupied)
|
||||
>>> arr (uncurry AFRP.lMerge) >>> traceEvent
|
||||
>>> AFRP.hold Unoccupied
|
||||
where
|
||||
eventOccupied :: HASS Motion (Event Presence)
|
||||
eventOccupied = arr (== MotionDetected) >>> AFRP.edge >>> arr (fmap (const Occupied))
|
||||
eventUnoccupied :: HASS Motion (Event Presence)
|
||||
eventUnoccupied = arr (== MotionNotDetected) >>> AFRP.waitFor 300 >>> arr (AFRP.tag Unoccupied)
|
||||
|
||||
eventLights :: HASS Presence (Event Lights)
|
||||
eventLights = AFRP.changes >>> arr (fmap presenceLights)
|
||||
where
|
||||
presenceLights Occupied = LightsOn
|
||||
presenceLights Unoccupied = LightsOff
|
||||
|
||||
kitchenMotionController :: HASS (Event Value) ()
|
||||
kitchenMotionController = proc x -> do
|
||||
now <- AFRP.currentTime -< ()
|
||||
p <- kitchenMotion >>> kitchenPresence -< x
|
||||
ev <- eventLights -< p
|
||||
traceEvent -< ev
|
||||
case ev of
|
||||
Event LightsOn | lightsAllowed now -> callServiceDyn (light [EntityId "light.kitchen_ceiling"]) -< On Nothing
|
||||
Event LightsOff -> callServiceDyn (light [EntityId "light.kitchen_ceiling"]) -< Off
|
||||
_ -> returnA -< ()
|
||||
where
|
||||
lightsAllowed (LocalTime _ tod) = not (tod > TimeOfDay 1 45 0 && tod < TimeOfDay 5 0 0)
|
||||
@@ -0,0 +1,35 @@
|
||||
{-# LANGUAGE Arrows #-}
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
module HomeAssistant.Controller.Ruuvi
|
||||
( Ruuvi(..)
|
||||
, ruuvi
|
||||
, ruuviTemperatures
|
||||
, ruuviPressures
|
||||
, ruuviController
|
||||
) where
|
||||
|
||||
import AFRP (Mealy, Event, hold, changes, rollup)
|
||||
import Control.Category ((>>>))
|
||||
import Data.Aeson (Value)
|
||||
import HomeAssistant.Controller (entityRead, traceEvent, HASS)
|
||||
import Control.Arrow (Arrow(..))
|
||||
import GHC.Generics (Generic)
|
||||
import Data.Serialize (Serialize)
|
||||
|
||||
ruuviTemperatures :: Mealy eff (Event Value) Double
|
||||
ruuviTemperatures = entityRead @Double "sensor.ruuvitag_b168_temperature" >>> hold 0
|
||||
|
||||
ruuviPressures :: Mealy eff (Event Value) Double
|
||||
ruuviPressures = entityRead "sensor.ruuvitag_b168_pressure" >>> hold 0
|
||||
|
||||
data Ruuvi = Ruuvi { ruuviTemperature :: Double, ruuviPressure :: Double }
|
||||
deriving (Show, Eq, Generic)
|
||||
|
||||
instance Serialize Ruuvi
|
||||
|
||||
ruuvi :: Mealy eff (Event Value) (Event Ruuvi)
|
||||
ruuvi = (Ruuvi <$> ruuviTemperatures <*> ruuviPressures) >>> changes
|
||||
|
||||
|
||||
ruuviController :: HASS (Event Value) ()
|
||||
ruuviController = ruuvi >>> rollup 1 30 >>> traceEvent >>> arr (const ())
|
||||
@@ -13,12 +13,12 @@ module HomeAssistant.Runtime
|
||||
, runController
|
||||
) where
|
||||
|
||||
import AFRP (Event (..), Mealy (..), Request (..))
|
||||
import AFRP (Event (..), Mealy (..), Request (..), Auto, stepAutoSerializing, load, DecodedAuto (..))
|
||||
import Control.Concurrent.Async (async, waitAny)
|
||||
import Control.Concurrent.STM (atomically, dupTChan, readTChan)
|
||||
import Data.Aeson (Value)
|
||||
import qualified Data.Text as T
|
||||
import Data.Time (getCurrentTime)
|
||||
import Data.Time (getCurrentTime, getCurrentTimeZone)
|
||||
import Data.Void (Void, absurd)
|
||||
import HomeAssistant.Controller (HASS, HASSEff (..))
|
||||
import HomeAssistant.Runtime.Bus
|
||||
@@ -31,12 +31,20 @@ import Data.UUID (UUID, toText)
|
||||
import qualified Data.UUID.V4 as UUID.V4
|
||||
import Katip (runKatipT, logF, sl, Severity (..), ls, Namespace (Namespace), runKatipContextT)
|
||||
import Control.Monad.IO.Class (liftIO, MonadIO)
|
||||
import Control.Monad.Fix (MonadFix)
|
||||
import HomeAssistant.Controller.Ruuvi (ruuviController)
|
||||
import HomeAssistant.Controller.Children (schoolLightController)
|
||||
import Data.Maybe (fromMaybe)
|
||||
import qualified System.Metrics
|
||||
import qualified HomeAssistant.Runtime.Metrics
|
||||
import System.FilePath ((</>))
|
||||
import HomeAssistant.Controller.Kitchen (kitchenMotionController)
|
||||
|
||||
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 :: (MonadIO m) => FilePath -> UUID -> Auto m a b -> a -> m (b, Auto m a b)
|
||||
step path trace st a = do
|
||||
now <- liftIO getCurrentTime
|
||||
f nt (Request now trace) a
|
||||
tz <- liftIO getCurrentTimeZone
|
||||
let req = Request now tz trace
|
||||
stepAutoSerializing path st req a
|
||||
|
||||
data Controller = forall b. Controller T.Text (HASS (Event Value) b) Bool
|
||||
|
||||
@@ -46,22 +54,30 @@ controllers =
|
||||
, Controller "bedroom-button" bedroomButtonController False -- This works but leaving for vacation
|
||||
, Controller "bedroom-drawer" bedroomDrawerController True
|
||||
, Controller "bedroom-humidifier" humidifierController True
|
||||
, Controller "ruuvi-controller" ruuviController False
|
||||
, Controller "school-light-controller" schoolLightController True
|
||||
, Controller "kitchen-motion-controller" kitchenMotionController True
|
||||
]
|
||||
|
||||
-- | 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 _enabled) = do
|
||||
runController :: FilePath -> Bus -> Controller -> IO Void
|
||||
runController rootDir bus (Controller name machine _enabled) = do
|
||||
inbound <- atomically (dupTChan (busInbound bus))
|
||||
go inbound machine
|
||||
let ns = Namespace [name]
|
||||
let workerDefinition = runMealy machine (runKatipContextT (busLogEnv bus) () ns . channelHassEval bus)
|
||||
let path = rootDir </> T.unpack name
|
||||
worker <- load path workerDefinition >>= \case
|
||||
Decoded a -> pure a
|
||||
FailDecode err a -> a <$ putStrLn ("Failed to load (" <> T.unpack name <> "): " <> err)
|
||||
go path inbound worker
|
||||
where
|
||||
go inbound f = do
|
||||
go path inbound f = do
|
||||
msg <- atomically (readTChan inbound)
|
||||
uuid <- UUID.V4.nextRandom
|
||||
let ns = Namespace [name]
|
||||
(_, f') <- step (runKatipContextT (busLogEnv bus) () ns . channelHassEval bus) uuid f (Event msg)
|
||||
go inbound f'
|
||||
(_, next) <- step path uuid f msg
|
||||
go path inbound next
|
||||
|
||||
defaultMain :: IO ()
|
||||
defaultMain = withSocketsDo $ do
|
||||
@@ -69,10 +85,18 @@ defaultMain = withSocketsDo $ do
|
||||
withBus severity $ \bus -> do
|
||||
token <- getEnv "HA_TOKEN"
|
||||
host <- getEnv "HA_HOST"
|
||||
let workers =
|
||||
[ ("reader", readerAction host 8123 token bus)
|
||||
rootPath <- fromMaybe "/tmp/" <$> lookupEnv "HA_LIB_DIR"
|
||||
store <- System.Metrics.newStore
|
||||
System.Metrics.registerGcMetrics store
|
||||
rrdPath <- fromMaybe "hass-controller.rrd" <$> lookupEnv "HA_RRD_PATH"
|
||||
rrdtool <- fromMaybe "rrdtool" <$> lookupEnv "HA_RRDTOOL"
|
||||
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 ]
|
||||
] ++ [ (name, runController rootPath bus c) | c@(Controller name _ True) <- controllers ]
|
||||
++ [("metrics", HomeAssistant.Runtime.Metrics.metricsAction store rrdPath rrdtool)]
|
||||
as <- mapM (\(name, act) -> async (supervised name defaultBackoff act)) workers
|
||||
(_, v) <- waitAny as
|
||||
absurd v
|
||||
|
||||
@@ -26,14 +26,14 @@ import Katip (LogEnv, closeScribes, mkHandleScribe, ColorStrategy (..), permitIt
|
||||
import Control.Exception (bracket)
|
||||
import System.IO (stdout)
|
||||
import Data.UUID (toText)
|
||||
import AFRP (Request(..))
|
||||
import AFRP (Request(..), Event(..))
|
||||
import Control.Monad.IO.Class (MonadIO, liftIO)
|
||||
|
||||
-- | Shared runtime state: inbound is a broadcast channel (controllers
|
||||
-- read from 'dupTChan' copies), outbound queues service calls for the
|
||||
-- writer, conn holds the current websocket (Nothing before first connect).
|
||||
data Bus = Bus
|
||||
{ busInbound :: TChan Value
|
||||
{ busInbound :: TChan (Event Value)
|
||||
, busOutbound :: TChan (Request, Service)
|
||||
, busConn :: TVar (Maybe Connection)
|
||||
, busGen :: CallIdGen
|
||||
|
||||
@@ -5,23 +5,30 @@ module HomeAssistant.Runtime.Connection
|
||||
( readerAction
|
||||
, writerAction
|
||||
, encodeService
|
||||
, dedupeBatch
|
||||
) where
|
||||
|
||||
import Control.Concurrent.STM
|
||||
( atomically
|
||||
( TChan
|
||||
, atomically
|
||||
, readTChan
|
||||
, readTVar
|
||||
, retry
|
||||
, tryReadTChan
|
||||
, writeTChan
|
||||
, writeTVar
|
||||
)
|
||||
import Control.Concurrent.Async (race)
|
||||
import Control.Concurrent (threadDelay)
|
||||
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 Data.List (sort)
|
||||
import qualified Data.Map.Strict as M
|
||||
import qualified Data.Set as S
|
||||
import qualified Data.Text as T
|
||||
import Data.Void (Void)
|
||||
import HomeAssistant.Controller (Service (..), Target (..))
|
||||
@@ -30,16 +37,16 @@ import HomeAssistant.Runtime.Supervisor (Fatal (..))
|
||||
import qualified Network.WebSockets as WS
|
||||
import Katip (runKatipContextT, sl, logFM, Severity (..), ls)
|
||||
import Data.UUID (toText)
|
||||
import AFRP (Request(..))
|
||||
import AFRP (Request(..), Event(..))
|
||||
|
||||
-- | 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,23 +69,34 @@ 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.
|
||||
--
|
||||
-- Each read races a one-second timeout: a timeout broadcasts 'Tick' so
|
||||
-- time-based primitives (debounce, rollup, fixed, ...) keep advancing
|
||||
-- even when no state changes arrive.
|
||||
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
|
||||
winner <- race (threadDelay 1_000_000) (WS.receiveData conn)
|
||||
case winner of
|
||||
Left () -> atomically $ writeTChan (busInbound bus) Tick
|
||||
Right msg -> case eitherDecode msg of
|
||||
Left err -> putStrLn $ "[reader] skipping undecodable message: " <> err
|
||||
Right v -> atomically $ writeTChan (busInbound bus) (Event v)
|
||||
|
||||
receiveJSON :: WS.Connection -> IO Value
|
||||
receiveJSON conn = do
|
||||
@@ -87,16 +105,41 @@ receiveJSON conn = do
|
||||
Left err -> throw (Fatal $ "Invalid JSON from Home Assistant: " <> T.pack err)
|
||||
Right x -> pure x
|
||||
|
||||
writerAction :: Bus -> IO Void
|
||||
writerAction bus = forever $ do
|
||||
(request, svc) <- atomically $ readTChan (busOutbound bus)
|
||||
conn <- atomically $ readTVar (busConn bus) >>= maybe retry pure
|
||||
-- | 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" $
|
||||
logFM DebugS (ls textData)
|
||||
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)
|
||||
conn <- atomically $ readTVar (busConn bus) >>= maybe retry pure
|
||||
forM_ deduped $ \(request, svc) -> do
|
||||
sendWithId bus conn request svc
|
||||
threadDelay minInterval
|
||||
|
||||
encodeService :: Int -> Service -> Value
|
||||
encodeService callId Service{..} = object $
|
||||
[ "id" .= callId
|
||||
@@ -106,6 +149,18 @@ encodeService callId Service{..} = object $
|
||||
, "target" .= targetObject serviceTarget
|
||||
] <> maybe [] (\d -> ["service_data" .= d]) serviceData
|
||||
|
||||
-- | Collapse a drained batch of outbound calls: the newest call per
|
||||
-- `(domain, service, sorted-targets)` survives; older duplicates are
|
||||
-- dropped. `serviceData` is not part of the key, so a newer `turn_on`
|
||||
-- with different brightness supersedes an older one to the same target.
|
||||
dedupeBatch :: [(Request, Service)] -> [(Request, Service)]
|
||||
dedupeBatch = M.elems . foldl' ins M.empty
|
||||
where
|
||||
ins m (req, svc) = M.insert (dedupeKey svc) (req, svc) m
|
||||
|
||||
dedupeKey :: Service -> (T.Text, T.Text, [Target])
|
||||
dedupeKey Service{..} = (serviceDomain, serviceName, sort serviceTarget)
|
||||
|
||||
-- | A single target encodes as a scalar; multiple encode as a list. Empty
|
||||
-- lists are omitted so Home Assistant receives only populated keys.
|
||||
targetObject :: [Target] -> Value
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
|
||||
module HomeAssistant.Runtime.Metrics
|
||||
( DsType (..)
|
||||
, DsSpec (..)
|
||||
, dsTypeOf
|
||||
, sanitizeName
|
||||
, buildSchema
|
||||
, buildCreateArgs
|
||||
, buildUpdateArgs
|
||||
, ensureRrd
|
||||
, sampleAndUpdate
|
||||
, metricsAction
|
||||
) where
|
||||
|
||||
import Control.Concurrent (threadDelay)
|
||||
import Control.Monad (forever, unless)
|
||||
import Data.Char (ord)
|
||||
import Data.Int (Int64)
|
||||
import Data.List (intercalate, sortBy)
|
||||
import Data.Ord (comparing)
|
||||
import Data.Text (Text)
|
||||
import Data.Void (Void)
|
||||
import Numeric (showHex)
|
||||
import qualified Data.Text as T
|
||||
import qualified Data.HashMap.Strict as HM
|
||||
import qualified System.Metrics as M (Value (..), Sample, Store, sampleAll)
|
||||
import System.Directory (doesFileExist)
|
||||
import System.Process (callProcess)
|
||||
|
||||
data DsType = Derive | Gauge
|
||||
deriving (Eq, Show)
|
||||
|
||||
data DsSpec = DsSpec
|
||||
{ dsEkgName :: Text
|
||||
, dsName :: String
|
||||
, dsType :: DsType
|
||||
}
|
||||
deriving (Eq, Show)
|
||||
|
||||
dsTypeOf :: M.Value -> Maybe DsType
|
||||
dsTypeOf (M.Counter _) = Just Derive
|
||||
dsTypeOf (M.Gauge _) = Just Gauge
|
||||
dsTypeOf _ = Nothing
|
||||
|
||||
-- | Maps an ekg metric label to a valid rrd DS name (≤19 chars, [A-Za-z0-9_]).
|
||||
-- Long names keep the first 15 chars plus a 3-hex hash suffix for uniqueness.
|
||||
sanitizeName :: Text -> String
|
||||
sanitizeName name
|
||||
| length sanitized <= 19 = sanitized
|
||||
| otherwise = take 15 sanitized ++ "_" ++ paddedHash
|
||||
where
|
||||
sanitized = T.unpack (T.replace "." "_" name)
|
||||
paddedHash = let h = showHex (sum (map ord sanitized) `mod` 4096) ""
|
||||
in replicate (3 - length h) '0' ++ h
|
||||
|
||||
buildSchema :: M.Sample -> [DsSpec]
|
||||
buildSchema sample =
|
||||
sortBy (comparing dsName)
|
||||
[ DsSpec ekgName (sanitizeName ekgName) dt
|
||||
| (ekgName, val) <- HM.toList sample
|
||||
, Just dt <- [dsTypeOf val]
|
||||
]
|
||||
|
||||
buildCreateArgs :: FilePath -> Int -> [(String, DsType)] -> [String]
|
||||
buildCreateArgs path step specs =
|
||||
["create", path, "--step", show step]
|
||||
++ concatMap dsArg specs
|
||||
++ rras
|
||||
where
|
||||
dsArg (name, Derive) = ["DS:" ++ name ++ ":DERIVE:20:0:U"]
|
||||
dsArg (name, Gauge) = ["DS:" ++ name ++ ":GAUGE:20:0:U"]
|
||||
rras =
|
||||
[ "RRA:AVERAGE:0.5:1:6000"
|
||||
, "RRA:MAX:0.5:1:6000"
|
||||
, "RRA:AVERAGE:0.5:360:1680"
|
||||
, "RRA:MAX:0.5:360:1680"
|
||||
]
|
||||
|
||||
buildUpdateArgs :: FilePath -> [String] -> [Maybe Int64] -> [String]
|
||||
buildUpdateArgs path names values =
|
||||
[ "update"
|
||||
, path
|
||||
, "--template"
|
||||
, intercalate ":" names
|
||||
, "N:" ++ intercalate ":" (map renderValue values)
|
||||
]
|
||||
where
|
||||
renderValue Nothing = "U"
|
||||
renderValue (Just n) = show n
|
||||
|
||||
lookupValue :: M.Sample -> Text -> Maybe Int64
|
||||
lookupValue sample name = case HM.lookup name sample of
|
||||
Just (M.Counter n) -> Just n
|
||||
Just (M.Gauge n) -> Just n
|
||||
_ -> Nothing
|
||||
|
||||
ensureRrd :: FilePath -> FilePath -> [DsSpec] -> IO ()
|
||||
ensureRrd rrdtool rrdPath schema = do
|
||||
exists <- doesFileExist rrdPath
|
||||
unless exists $
|
||||
callProcess rrdtool (buildCreateArgs rrdPath 10 (map toPair schema))
|
||||
where
|
||||
toPair s = (dsName s, dsType s)
|
||||
|
||||
sampleAndUpdate :: M.Store -> FilePath -> FilePath -> [DsSpec] -> IO ()
|
||||
sampleAndUpdate store rrdtool rrdPath schema = do
|
||||
sample <- M.sampleAll store
|
||||
let names = map dsName schema
|
||||
values = map (lookupValue sample . dsEkgName) schema
|
||||
callProcess rrdtool (buildUpdateArgs rrdPath names values)
|
||||
|
||||
metricsAction :: M.Store -> FilePath -> FilePath -> IO Void
|
||||
metricsAction store rrdPath rrdtool = do
|
||||
schema <- buildSchema <$> M.sampleAll store
|
||||
ensureRrd rrdtool rrdPath schema
|
||||
forever $ do
|
||||
sampleAndUpdate store rrdtool rrdPath schema
|
||||
threadDelay 10000000
|
||||
@@ -0,0 +1,256 @@
|
||||
module AFRPLawsSpec (spec) where
|
||||
|
||||
import AFRP
|
||||
import Control.Arrow (arr, first, left, (***), (+++))
|
||||
import Control.Category ((>>>))
|
||||
import qualified Control.Category as Cat (id)
|
||||
import Data.Functor.Identity (Identity (..))
|
||||
import Hedgehog (Gen, PropertyT)
|
||||
import qualified Hedgehog.Gen as Gen
|
||||
import qualified Hedgehog.Range as Range
|
||||
import Support (fakeRequest)
|
||||
import Test.Hspec (Spec, describe, it)
|
||||
import Test.Hspec.Hedgehog (forAll, forAllWith, hedgehog, (===))
|
||||
|
||||
-- | Law tests for the Mealy instances. Two machines count as equal when
|
||||
-- they emit equal outputs on every input sequence, so each law runs both
|
||||
-- sides on generated inputs.
|
||||
|
||||
runPure :: Mealy Identity a b -> [a] -> [b]
|
||||
runPure m = go (runMealy m id)
|
||||
where
|
||||
go _ [] = []
|
||||
go w (a : as) = case runIdentity (stepAuto w fakeRequest a) of
|
||||
(b, w') -> b : go w' as
|
||||
|
||||
-- | Machines wrap functions and have no Show; name them for forAll instead.
|
||||
forAllMealy :: Gen (Mealy Identity a b) -> PropertyT IO (Mealy Identity a b)
|
||||
forAllMealy = forAllWith (const "<mealy>")
|
||||
|
||||
intGen :: Gen Int
|
||||
intGen = Gen.int (Range.linear (-5) 5)
|
||||
|
||||
ints :: Gen [Int]
|
||||
ints = Gen.list (Range.linear 0 30) intGen
|
||||
|
||||
intPairs :: Gen [(Int, Int)]
|
||||
intPairs = Gen.list (Range.linear 0 30) ((,) <$> intGen <*> intGen)
|
||||
|
||||
intEithers :: Gen [Either Int Int]
|
||||
intEithers = Gen.list (Range.linear 0 30) $
|
||||
Gen.choice [Left <$> intGen, Right <$> intGen]
|
||||
|
||||
nestedPairs :: Gen [((Int, Int), Int)]
|
||||
nestedPairs = Gen.list (Range.linear 0 30) ((,) <$> ((,) <$> intGen <*> intGen) <*> intGen)
|
||||
|
||||
nestedEithers :: Gen [Either (Either Int Int) Int]
|
||||
nestedEithers = Gen.list (Range.linear 0 30) $
|
||||
Gen.choice
|
||||
[ Left <$> Gen.choice [Left <$> intGen, Right <$> intGen]
|
||||
, Right <$> intGen
|
||||
]
|
||||
|
||||
-- | Stateful Int machines: the arrow variables of the laws.
|
||||
statefulGen :: Gen (Mealy Identity Int Int)
|
||||
statefulGen = Gen.choice
|
||||
[ (\k -> mapAccum (+) k id) <$> intGen
|
||||
, (\k -> preMapAccum (+) k id) <$> intGen
|
||||
, (\k -> mapAccum (*) 1 (+ k)) <$> intGen
|
||||
]
|
||||
|
||||
eventArrowGen :: Gen (Mealy Identity Int (Event Int))
|
||||
eventArrowGen = Gen.choice
|
||||
[ pure changes
|
||||
, (\k -> mapAccum (+) k Event) <$> intGen
|
||||
, (\k -> preMapAccum (+) k (Event . (* 2))) <$> intGen
|
||||
]
|
||||
|
||||
funArrowGen :: Gen (Mealy Identity Int (Int -> Int))
|
||||
funArrowGen = Gen.choice
|
||||
[ (\k -> mapAccum (+) k (*)) <$> intGen
|
||||
, pure (preMapAccum (*) 1 (+))
|
||||
]
|
||||
|
||||
spec :: Spec
|
||||
spec = describe "Mealy laws" $ do
|
||||
semigroupSpec
|
||||
monoidSpec
|
||||
categorySpec
|
||||
arrowSpec
|
||||
arrowChoiceSpec
|
||||
functorSpec
|
||||
applicativeSpec
|
||||
|
||||
semigroupSpec :: Spec
|
||||
semigroupSpec = describe "Semigroup (<>)" $ do
|
||||
it "(a <> b) <> c = a <> (b <> c)" $ hedgehog $ do
|
||||
a <- forAllMealy eventArrowGen
|
||||
b <- forAllMealy eventArrowGen
|
||||
c <- forAllMealy eventArrowGen
|
||||
xs <- forAll ints
|
||||
runPure ((a <> b) <> c) xs === runPure (a <> (b <> c)) xs
|
||||
|
||||
monoidSpec :: Spec
|
||||
monoidSpec = describe "Monoid" $ do
|
||||
it "mempty <> a = a" $ hedgehog $ do
|
||||
a <- forAllMealy eventArrowGen
|
||||
xs <- forAll ints
|
||||
runPure (mempty <> a) xs === runPure a xs
|
||||
|
||||
it "a <> mempty = a" $ hedgehog $ do
|
||||
a <- forAllMealy eventArrowGen
|
||||
xs <- forAll ints
|
||||
runPure (a <> mempty) xs === runPure a xs
|
||||
|
||||
categorySpec :: Spec
|
||||
categorySpec = describe "Category" $ do
|
||||
it "id >>> f = f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
xs <- forAll ints
|
||||
runPure (Cat.id >>> f) xs === runPure f xs
|
||||
|
||||
it "f >>> id = f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
xs <- forAll ints
|
||||
runPure (f >>> Cat.id) xs === runPure f xs
|
||||
|
||||
it "(f >>> g) >>> h = f >>> (g >>> h)" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
g <- forAllMealy statefulGen
|
||||
h <- forAllMealy statefulGen
|
||||
xs <- forAll ints
|
||||
runPure ((f >>> g) >>> h) xs === runPure (f >>> (g >>> h)) xs
|
||||
|
||||
arrowSpec :: Spec
|
||||
arrowSpec = describe "Arrow" $ do
|
||||
it "arr id = id" $ hedgehog $ do
|
||||
xs <- forAll ints
|
||||
runPure (arr id :: Mealy Identity Int Int) xs === runPure Cat.id xs
|
||||
|
||||
it "arr (f >>> g) = arr f >>> arr g" $ hedgehog $ do
|
||||
p <- forAll intGen
|
||||
q <- forAll intGen
|
||||
xs <- forAll ints
|
||||
let f = (+ p)
|
||||
g = (* q)
|
||||
runPure (arr (f >>> g)) xs === runPure (arr f >>> arr g) xs
|
||||
|
||||
it "first (arr f) = arr (first f)" $ hedgehog $ do
|
||||
p <- forAll intGen
|
||||
ps <- forAll intPairs
|
||||
let f = (+ p)
|
||||
runPure (first (arr f)) ps === runPure (arr (first f)) ps
|
||||
|
||||
it "first (f >>> g) = first f >>> first g" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
g <- forAllMealy statefulGen
|
||||
ps <- forAll intPairs
|
||||
runPure (first (f >>> g)) ps === runPure (first f >>> first g) ps
|
||||
|
||||
it "first f >>> arr fst = arr fst >>> f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
ps <- forAll intPairs
|
||||
runPure (first f >>> arr fst) ps === runPure (arr fst >>> f) ps
|
||||
|
||||
it "first f >>> arr (id *** g) = arr (id *** g) >>> first f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
p <- forAll intGen
|
||||
ps <- forAll intPairs
|
||||
let g = (* p)
|
||||
runPure (first f >>> arr (id *** g)) ps
|
||||
=== runPure (arr (id *** g) >>> first f) ps
|
||||
|
||||
it "first (first f) >>> arr assoc = arr assoc >>> first f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
ts <- forAll nestedPairs
|
||||
let assoc ((a, b), c) = (a, (b, c))
|
||||
runPure (first (first f) >>> arr assoc) ts
|
||||
=== runPure (arr assoc >>> first f) ts
|
||||
|
||||
arrowChoiceSpec :: Spec
|
||||
arrowChoiceSpec = describe "ArrowChoice" $ do
|
||||
it "left (arr f) = arr (left f)" $ hedgehog $ do
|
||||
p <- forAll intGen
|
||||
es <- forAll intEithers
|
||||
let f = (+ p)
|
||||
runPure (left (arr f)) es === runPure (arr (left f)) es
|
||||
|
||||
it "left (f >>> g) = left f >>> left g" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
g <- forAllMealy statefulGen
|
||||
es <- forAll intEithers
|
||||
runPure (left (f >>> g)) es === runPure (left f >>> left g) es
|
||||
|
||||
it "f >>> arr Left = arr Left >>> left f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
xs <- forAll ints
|
||||
runPure (f >>> arr (Left @Int @Int)) xs
|
||||
=== runPure (arr (Left @Int @Int) >>> left f) xs
|
||||
|
||||
it "left f >>> arr (id +++ g) = arr (id +++ g) >>> left f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
p <- forAll intGen
|
||||
es <- forAll intEithers
|
||||
let g = (* p)
|
||||
runPure (left f >>> arr (id +++ g)) es
|
||||
=== runPure (arr (id +++ g) >>> left f) es
|
||||
|
||||
it "left (left f) >>> arr assocsum = arr assocsum >>> left f" $ hedgehog $ do
|
||||
f <- forAllMealy statefulGen
|
||||
es <- forAll nestedEithers
|
||||
let assocsum (Left (Left x)) = Left x
|
||||
assocsum (Left (Right y)) = Right (Left y)
|
||||
assocsum (Right z) = Right (Right z)
|
||||
runPure (left (left f) >>> arr assocsum) es
|
||||
=== runPure (arr assocsum >>> left f) es
|
||||
|
||||
functorSpec :: Spec
|
||||
functorSpec = describe "Functor" $ do
|
||||
it "fmap id = id" $ hedgehog $ do
|
||||
m <- forAllMealy statefulGen
|
||||
xs <- forAll ints
|
||||
runPure (fmap id m) xs === runPure m xs
|
||||
|
||||
it "fmap (f . g) = fmap f . fmap g" $ hedgehog $ do
|
||||
m <- forAllMealy statefulGen
|
||||
p <- forAll intGen
|
||||
q <- forAll intGen
|
||||
xs <- forAll ints
|
||||
let f = (+ p)
|
||||
g = (* q)
|
||||
runPure (fmap (f . g) m) xs === runPure (fmap f (fmap g m)) xs
|
||||
|
||||
applicativeSpec :: Spec
|
||||
applicativeSpec = describe "Applicative" $ do
|
||||
it "pure id <*> v = v" $ hedgehog $ do
|
||||
v <- forAllMealy statefulGen
|
||||
xs <- forAll ints
|
||||
runPure (pure id <*> v) xs === runPure v xs
|
||||
|
||||
it "pure f <*> pure x = pure (f x)" $ hedgehog $ do
|
||||
p <- forAll intGen
|
||||
x <- forAll intGen
|
||||
xs <- forAll ints
|
||||
let f = (+ p)
|
||||
runPure (pure f <*> pure x) xs === runPure (pure (f x)) xs
|
||||
|
||||
it "u <*> pure y = pure ($ y) <*> u" $ hedgehog $ do
|
||||
u <- forAllMealy funArrowGen
|
||||
y <- forAll intGen
|
||||
xs <- forAll ints
|
||||
runPure (u <*> pure y) xs === runPure (pure ($ y) <*> u) xs
|
||||
|
||||
it "pure (.) <*> u <*> v <*> w = u <*> (v <*> w)" $ hedgehog $ do
|
||||
u <- forAllMealy funArrowGen
|
||||
v <- forAllMealy funArrowGen
|
||||
w <- forAllMealy statefulGen
|
||||
xs <- forAll ints
|
||||
runPure (pure (.) <*> u <*> v <*> w) xs
|
||||
=== runPure (u <*> (v <*> w)) xs
|
||||
|
||||
it "fmap f x = pure f <*> x" $ hedgehog $ do
|
||||
x <- forAllMealy statefulGen
|
||||
p <- forAll intGen
|
||||
xs <- forAll ints
|
||||
let f = (+ p)
|
||||
runPure (fmap f x) xs === runPure (pure f <*> x) xs
|
||||
@@ -0,0 +1,587 @@
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
|
||||
module AFRPSpec (spec) where
|
||||
|
||||
import Control.Arrow (arr, (&&&), first, left)
|
||||
import Control.Category ((>>>))
|
||||
import Control.Monad.Fix (MonadFix (..))
|
||||
import Data.Foldable (for_)
|
||||
import AFRP
|
||||
import Data.Functor.Identity (Identity (..))
|
||||
import Data.List (sort)
|
||||
import Data.Serialize (get, put, runGet, runPut)
|
||||
import qualified Data.Set as S
|
||||
import qualified Data.Text as T
|
||||
import Data.Time (Day (..), NominalDiffTime, LocalTime(..), TimeOfDay(..), UTCTime (..), picosecondsToDiffTime, utc)
|
||||
import Data.UUID (nil)
|
||||
import Hedgehog
|
||||
import qualified Hedgehog.Gen as Gen
|
||||
import qualified Hedgehog.Range as Range
|
||||
import Test.Hspec
|
||||
import Test.Hspec.Hedgehog
|
||||
|
||||
fakeRequest :: Request
|
||||
fakeRequest = Request (sec 0) utc nil
|
||||
|
||||
sec :: Integer -> UTCTime
|
||||
sec n = UTCTime (toEnum 0) (fromIntegral n)
|
||||
|
||||
untime :: SerializeUTCTime -> UTCTime
|
||||
untime (SerializeUTCTime t) = t
|
||||
|
||||
runPure :: Mealy Identity a b -> [a] -> [b]
|
||||
runPure m = go (runMealy m id)
|
||||
where
|
||||
go _ [] = []
|
||||
go w (a : as) = case runIdentity (stepAuto w fakeRequest a) of
|
||||
(b, w') -> b : go w' as
|
||||
|
||||
-- | Run a Mealy with a per-step wall clock (seconds since the day-0 epoch).
|
||||
runTimed :: Mealy Identity a b -> [(Integer, a)] -> [b]
|
||||
runTimed m = go (runMealy m id)
|
||||
where
|
||||
go _ [] = []
|
||||
go w ((s, a) : as) =
|
||||
case runIdentity (stepAuto w (Request (sec s) utc nil) a) of
|
||||
(b, w') -> b : go w' as
|
||||
|
||||
-- | A minimal State monad for observing effectful arrows (e.g. whenA gating).
|
||||
newtype St a = St { unSt :: Int -> (a, Int) }
|
||||
|
||||
instance Functor St where
|
||||
fmap f (St g) = St $ \s -> let (a, s') = g s in (f a, s')
|
||||
|
||||
instance Applicative St where
|
||||
pure a = St (\s -> (a, s))
|
||||
St f <*> St x = St $ \s -> let (f', s') = f s; (a, s'') = x s' in (f' a, s'')
|
||||
|
||||
instance Monad St where
|
||||
St m >>= k = St $ \s -> let (a, s') = m s; (b, s'') = unSt (k a) s' in (b, s'')
|
||||
|
||||
instance MonadFix St where
|
||||
mfix f = St $ \s -> let (a, s') = unSt (f a) s in (a, s')
|
||||
|
||||
runStEff :: Mealy St a b -> Int -> [a] -> ([b], Int)
|
||||
runStEff m = go (runMealy m id)
|
||||
where
|
||||
go _ s [] = ([], s)
|
||||
go w s (a : rest) = case unSt (stepAuto w fakeRequest a) s of
|
||||
((b, w'), s') -> let (bs, s'') = go w' s' rest in (b : bs, s'')
|
||||
|
||||
spec :: Spec
|
||||
spec = describe "AFRP" $ do
|
||||
entitiesSpec
|
||||
holdSpec
|
||||
eventsSpec
|
||||
isEventSpec
|
||||
tagSpec
|
||||
toEventSpec
|
||||
lMergeSpec
|
||||
changesSpec
|
||||
edgeSpec
|
||||
filterASpec
|
||||
slidingSpec
|
||||
mapAccumSpec
|
||||
preMapAccumSpec
|
||||
durationSpec
|
||||
delayEventSpec
|
||||
debounceSpec
|
||||
rollupSpec
|
||||
serializeSpec
|
||||
effSpec
|
||||
mapAccumRequestSpec
|
||||
preMapAccumRequestSpec
|
||||
whenASpec
|
||||
thenASpec
|
||||
sampleSpec
|
||||
|
||||
holdSpec :: Spec
|
||||
holdSpec = describe "hold" $ do
|
||||
it "holds initial value until an Event arrives" $
|
||||
runPure (hold 'a') [Tick, Event 'b', Tick, Event 'c']
|
||||
`shouldBe` ['a', 'b', 'b', 'c']
|
||||
|
||||
it "never changes on Tick" $
|
||||
runPure (hold (0 :: Int)) (replicate 5 Tick) `shouldBe` replicate 5 (0 :: Int)
|
||||
|
||||
eventsSpec :: Spec
|
||||
eventsSpec = describe "events" $ do
|
||||
it "converts Tick to Left () and Event a to Right a" $
|
||||
runPure events [Tick, Event 'a', Tick, Event 'b']
|
||||
`shouldBe` [Left (), Right 'a', Left (), Right 'b']
|
||||
|
||||
isEventSpec :: Spec
|
||||
isEventSpec = describe "isEvent" $ do
|
||||
it "returns False for Tick" $
|
||||
isEvent Tick `shouldBe` False
|
||||
|
||||
it "returns True for Event x" $
|
||||
isEvent (Event ()) `shouldBe` True
|
||||
|
||||
tagSpec :: Spec
|
||||
tagSpec = describe "tag" $ do
|
||||
it "replaces value preserving structure" $ do
|
||||
tag 'b' Tick `shouldBe` Tick
|
||||
tag 'b' (Event 'a') `shouldBe` Event 'b'
|
||||
|
||||
toEventSpec :: Spec
|
||||
toEventSpec = describe "toEvent" $ do
|
||||
it "round-trips through events" $
|
||||
runPure toEvent [Left (), Right 'a', Left ()]
|
||||
`shouldBe` [Tick, Event 'a', Tick]
|
||||
|
||||
it "is inverse of events modulo Event/Either" $
|
||||
runPure (events >>> toEvent) [Tick, Event 'a', Event 'b']
|
||||
`shouldBe` [Tick, Event 'a', Event 'b']
|
||||
|
||||
lMergeSpec :: Spec
|
||||
lMergeSpec = describe "lMerge" $ do
|
||||
it "both Tick gives Tick" $
|
||||
lMerge (Tick :: Event Int) (Tick :: Event Int) `shouldBe` Tick
|
||||
|
||||
it "prefers left Event" $
|
||||
lMerge (Event (1 :: Int)) (Event (2 :: Int)) `shouldBe` Event (1 :: Int)
|
||||
|
||||
it "prefers right Event if left is Tick" $
|
||||
lMerge Tick (Event (2 :: Int)) `shouldBe` Event (2 :: Int)
|
||||
|
||||
it "Tick is a left identity" $
|
||||
hedgehog $ do
|
||||
e <- forAll eventGen
|
||||
lMerge Tick e === e
|
||||
|
||||
it "Tick is a right identity" $
|
||||
hedgehog $ do
|
||||
e <- forAll eventGen
|
||||
lMerge e Tick === e
|
||||
|
||||
it "is associative" $
|
||||
hedgehog $ do
|
||||
a <- forAll eventGen
|
||||
b <- forAll eventGen
|
||||
c <- forAll eventGen
|
||||
lMerge a (lMerge b c) === lMerge (lMerge a b) c
|
||||
|
||||
changesSpec :: Spec
|
||||
changesSpec = describe "changes" $ do
|
||||
it "first output is always Tick" $
|
||||
runPure changes "hello" !! 0 `shouldBe` Tick
|
||||
|
||||
it "outputs Event only on value change" $
|
||||
runPure changes "aaaabbbcca"
|
||||
`shouldBe` [Tick, Tick, Tick, Tick
|
||||
, Event 'b', Tick, Tick
|
||||
, Event 'c', Tick
|
||||
, Event 'a'
|
||||
]
|
||||
|
||||
it "first output is Tick, subsequent outputs are Event iff value changed" $
|
||||
hedgehog $ do
|
||||
xs <- forAll $ Gen.list (Range.linear 0 100) Gen.alpha
|
||||
let out = runPure changes xs
|
||||
length out === length xs
|
||||
case out of
|
||||
[] -> pure ()
|
||||
(Tick : rest) -> do
|
||||
let triples = zip3 xs (drop 1 xs) rest
|
||||
for_ triples $ \(prev, curr, o) ->
|
||||
if prev /= curr
|
||||
then o === Event curr
|
||||
else o === Tick
|
||||
_ -> failure
|
||||
|
||||
edgeSpec :: Spec
|
||||
edgeSpec = describe "edge" $ do
|
||||
it "emits Event () only on rising edge" $
|
||||
runPure edge [False, True, True, False, True]
|
||||
`shouldBe` [Tick, Event (), Tick, Tick, Event ()]
|
||||
|
||||
it "starts from False, so first True is a rising edge" $
|
||||
runPure edge [True, False, True]
|
||||
`shouldBe` [Event (), Tick, Event ()]
|
||||
|
||||
it "Event () only on False -> True transition" $
|
||||
hedgehog $ do
|
||||
bs <- forAll $ Gen.list (Range.linear 0 50) Gen.bool
|
||||
let out = runPure edge bs
|
||||
length out === length bs
|
||||
case (bs, out) of
|
||||
([], []) -> pure ()
|
||||
(b : _, o : _) -> do
|
||||
if b then o === Event () else o === Tick
|
||||
let triples = zip3 bs (drop 1 bs) (drop 1 out)
|
||||
for_ triples $ \(prev, curr, o') ->
|
||||
if not prev && curr
|
||||
then o' === Event ()
|
||||
else o' === Tick
|
||||
_ -> failure
|
||||
|
||||
|
||||
filterASpec :: Spec
|
||||
filterASpec = describe "filterA" $ do
|
||||
it "lets through values matching predicate" $
|
||||
runPure (filterA (even @Int)) [1, 2, 3, 4]
|
||||
`shouldBe` [Left (), Right 2, Left (), Right 4]
|
||||
|
||||
it "output is Right a iff predicate holds" $
|
||||
hedgehog $ do
|
||||
threshold <- forAll $ Gen.int (Range.linear (-10) 10)
|
||||
let p = (> threshold)
|
||||
xs <- forAll $ Gen.list (Range.linear 0 50) (Gen.int (Range.linear (-20) 20))
|
||||
let out = runPure (filterA p) xs
|
||||
length out === length xs
|
||||
for_ (zip xs out) $ \(x, o) ->
|
||||
if p x
|
||||
then o === Right x
|
||||
else o === Left ()
|
||||
|
||||
slidingSpec :: Spec
|
||||
slidingSpec = describe "sliding" $ do
|
||||
it "accumulates events up to the window size" $
|
||||
runPure (sliding (3 :: Int)) [Event (1 :: Int), Event 2, Event 3, Event 4]
|
||||
`shouldBe` [[1], [1, 2], [1, 2, 3], [2, 3, 4]]
|
||||
|
||||
it "Ticks don't change the accumulator" $
|
||||
runPure (sliding 2) [Event (1 :: Int), Tick, Event 2]
|
||||
`shouldBe` [[1], [1], [1, 2]]
|
||||
|
||||
it "empty list stays empty" $
|
||||
runPure (sliding (5 :: Int)) ([] :: [Event Int]) `shouldBe` []
|
||||
|
||||
it "output length never exceeds window size" $
|
||||
hedgehog $ do
|
||||
n <- forAll $ Gen.int (Range.constant 1 10)
|
||||
evs <- forAll $ Gen.list (Range.linear 0 20) (Gen.frequency
|
||||
[(3, Event <$> Gen.alpha), (1, pure Tick)])
|
||||
let out = runPure (sliding n) evs
|
||||
for_ out $ \xs -> assert (length xs <= n)
|
||||
|
||||
mapAccumSpec :: Spec
|
||||
mapAccumSpec = describe "mapAccum" $ do
|
||||
it "running sum" $
|
||||
runPure (mapAccum (+) (0 :: Int) id) [1, 2, 3]
|
||||
`shouldBe` [1, 3, 6]
|
||||
|
||||
it "post-state extraction: output uses state after applying f" $
|
||||
runPure (mapAccum (\s x -> s ++ [x]) ([] :: [Int]) id) [1, 2, 3]
|
||||
`shouldBe` [[1], [1, 2], [1, 2, 3]]
|
||||
|
||||
it "output equals running sum of all inputs so far" $
|
||||
hedgehog $ do
|
||||
xs <- forAll $ Gen.list (Range.linear 0 50) (Gen.int (Range.linear (-100) 100))
|
||||
let out = runPure (mapAccum (+) (0 :: Int) id) xs
|
||||
length out === length xs
|
||||
for_ (zip3 [0 ..] xs out) $ \(i, _x, cur) ->
|
||||
cur === sum (take (i + 1) xs)
|
||||
|
||||
preMapAccumSpec :: Spec
|
||||
preMapAccumSpec = describe "preMapAccum" $ do
|
||||
it "running sum with pre-state extraction" $
|
||||
runPure (preMapAccum (+) (0 :: Int) id) [1, 2, 3]
|
||||
`shouldBe` [0, 1, 3]
|
||||
|
||||
it "pre-state extraction: output uses state before applying f" $
|
||||
runPure (preMapAccum (\s x -> s ++ [x]) ([] :: [Int]) id) [1, 2, 3]
|
||||
`shouldBe` [[], [1], [1, 2]]
|
||||
|
||||
it "output equals running sum before current input" $
|
||||
hedgehog $ do
|
||||
xs <- forAll $ Gen.list (Range.linear 0 50) (Gen.int (Range.linear (-100) 100))
|
||||
let out = runPure (preMapAccum (+) (0 :: Int) id) xs
|
||||
length out === length xs
|
||||
for_ (zip3 [0 ..] xs out) $ \(i, _x, cur) ->
|
||||
cur === sum (take i xs)
|
||||
|
||||
durationSpec :: Spec
|
||||
durationSpec = describe "duration" $ do
|
||||
it "first sample is 0, then elapsed time since first sample" $
|
||||
runTimed duration [(0, 'a'), (5, 'b'), (10, 'c')]
|
||||
`shouldBe` [0, 5, 10]
|
||||
|
||||
it "measures from the first observation, not the most recent" $
|
||||
runTimed duration [(2, 'a'), (3, 'b'), (7, 'c')]
|
||||
`shouldBe` [0, 1, 5]
|
||||
|
||||
it "output i equals times[i] - times[0]" $
|
||||
hedgehog $ do
|
||||
secs <- forAll $ Gen.list (Range.linear 0 50) (Gen.int (Range.linear 0 1000))
|
||||
let out = runTimed duration [(fromIntegral s, ()) | s <- secs]
|
||||
length out === length secs
|
||||
case secs of
|
||||
[] -> pure ()
|
||||
(t0 : _) -> for_ (zip secs out) $ \(s, d) ->
|
||||
d === fromIntegral (s - t0)
|
||||
|
||||
delayEventSpec :: Spec
|
||||
delayEventSpec = describe "delayEvent" $ do
|
||||
let delay = 5 :: NominalDiffTime
|
||||
|
||||
it "emits a queued event once the delay has elapsed" $
|
||||
runTimed (delayEvent delay)
|
||||
[(0, Event 'a'), (3, Tick), (6, Tick)]
|
||||
`shouldBe` [Tick, Tick, Event 'a']
|
||||
|
||||
it "preserves order when multiple events are queued" $
|
||||
runTimed (delayEvent delay)
|
||||
[(0, Event 'a'), (1, Event 'b'), (10, Tick), (12, Tick)]
|
||||
`shouldBe` [Tick, Tick, Event 'a', Event 'b']
|
||||
|
||||
it "emits nothing on pure Tick input" $
|
||||
runTimed (delayEvent delay) [(0, Tick :: Event Char), (10, Tick)]
|
||||
`shouldBe` [Tick, Tick]
|
||||
|
||||
it "emits exactly one output Event per input Event" $
|
||||
hedgehog $ do
|
||||
evs <- forAll $ Gen.list (Range.linear 0 30) eventGen
|
||||
let evs' = evs ++ replicate 3 Tick
|
||||
n = length [() | Event _ <- evs]
|
||||
out = runTimed (delayEvent (1 :: NominalDiffTime))
|
||||
(zip [0, 2 ..] evs')
|
||||
length [() | Event _ <- out] === n
|
||||
|
||||
debounceSpec :: Spec
|
||||
debounceSpec = describe "debounce" $ do
|
||||
let delay = 5 :: NominalDiffTime
|
||||
|
||||
it "fires the last event after the quiet period" $
|
||||
runTimed (debounce delay)
|
||||
[(0, Event 'a'), (3, Tick), (6, Tick)]
|
||||
`shouldBe` [Tick, Tick, Event 'a']
|
||||
|
||||
it "a newer event before firing resets the timer" $
|
||||
runTimed (debounce delay)
|
||||
[(0, Event 'a'), (3, Event 'b'), (6, Tick), (8, Tick)]
|
||||
`shouldBe` [Tick, Tick, Tick, Event 'b']
|
||||
|
||||
it "collapses a burst into a single emission" $
|
||||
runTimed (debounce delay)
|
||||
[(0, Event 'a'), (1, Event 'b'), (2, Event 'c'), (10, Tick)]
|
||||
`shouldBe` [Tick, Tick, Tick, Event 'c']
|
||||
|
||||
it "emits no more output Events than input Events" $
|
||||
hedgehog $ do
|
||||
evs <- forAll $ Gen.list (Range.linear 0 30) eventGen
|
||||
let nIn = length [() | Event _ <- evs]
|
||||
out = runTimed (debounce (1 :: NominalDiffTime))
|
||||
(zip [0, 2 ..] evs)
|
||||
nOut = length [() | Event _ <- out]
|
||||
assert (nOut <= nIn)
|
||||
|
||||
rollupSpec :: Spec
|
||||
rollupSpec = describe "rollup" $ do
|
||||
it "passes the first `limit` events through immediately, then bursts" $
|
||||
runTimed (rollup 2 10)
|
||||
[ (0, Event 'a'), (1, Event 'b')
|
||||
, (2, Event 'c'), (3, Event 'd')
|
||||
, (12, Tick)
|
||||
]
|
||||
`shouldBe` [ Event ['a'], Event ['b']
|
||||
, Tick, Tick
|
||||
, Event ['c', 'd']
|
||||
]
|
||||
|
||||
it "flushes the accumulator when the window ends on a Tick" $
|
||||
runTimed (rollup 1 10)
|
||||
[(0, Event 'a'), (1, Event 'b'), (2, Tick), (12, Tick)]
|
||||
`shouldBe` [Event ['a'], Tick, Tick, Event ['b']]
|
||||
|
||||
it "is idle (Tick) until the first Event" $
|
||||
runTimed (rollup 2 10) [(0, Tick :: Event Char), (1, Tick)]
|
||||
`shouldBe` [Tick, Tick]
|
||||
|
||||
it "every input Event appears exactly once across the outputs" $
|
||||
hedgehog $ do
|
||||
limit <- forAll $ Gen.int (Range.constant 1 5)
|
||||
evs <- forAll $ Gen.list (Range.linear 0 30) eventGen
|
||||
let evs' = evs ++ replicate 10 Tick
|
||||
times = [0 ..]
|
||||
out = runTimed (rollup limit 5) (zip times evs')
|
||||
emitted = concat [xs | Event xs <- out]
|
||||
sort emitted === sort [x | Event x <- evs]
|
||||
|
||||
|
||||
eventGen :: Gen (Event Char)
|
||||
eventGen = Gen.frequency
|
||||
[ (3, Event <$> Gen.alpha)
|
||||
, (1, pure Tick)
|
||||
]
|
||||
|
||||
timeGen :: Gen SerializeUTCTime
|
||||
timeGen = do
|
||||
day <- ModifiedJulianDay . fromIntegral <$> Gen.int (Range.linear 0 100000)
|
||||
pico <- picosecondsToDiffTime . fromIntegral
|
||||
<$> Gen.int (Range.linear 0 (86400 * 10 ^ (12 :: Int) - 1))
|
||||
pure $ SerializeUTCTime (UTCTime day pico)
|
||||
|
||||
localTimeGen :: Gen SerializeLocalTime
|
||||
localTimeGen = do
|
||||
day <- ModifiedJulianDay . fromIntegral <$> genInt 0 100000
|
||||
tod <- TimeOfDay <$> genInt 0 23 <*> genInt 0 59 <*> (fromIntegral <$> genInt 0 60)
|
||||
pure $ SerializeLocalTime (LocalTime day tod)
|
||||
where
|
||||
genInt a b = Gen.int (Range.linear a b)
|
||||
|
||||
serializeSpec :: Spec
|
||||
serializeSpec = do
|
||||
describe "SerializeUTCTime" $ do
|
||||
it "get (put x) == pure x" $
|
||||
hedgehog $ do
|
||||
x <- forAll timeGen
|
||||
tripping x (runPut . put) (runGet get)
|
||||
describe "SerializeLocalTime" $ do
|
||||
it "get (put x) == pure x" $
|
||||
hedgehog $ do
|
||||
x <- forAll localTimeGen
|
||||
tripping x (runPut . put) (runGet get)
|
||||
|
||||
effSpec :: Spec
|
||||
effSpec = describe "eff" $ do
|
||||
it "lifts a pure effect function into a stateless Mealy" $
|
||||
runPure (eff (\_ x -> Identity (x + 1))) [1, 2, 3]
|
||||
`shouldBe` [2 :: Int, 3, 4]
|
||||
|
||||
it "output equals f(input) for every step" $
|
||||
hedgehog $ do
|
||||
xs <- forAll $ Gen.list (Range.linear 0 50) (Gen.int (Range.linear (-100) 100))
|
||||
let out = runPure (eff (\_ x -> Identity (x * 2))) xs
|
||||
out === map (* 2) xs
|
||||
|
||||
mapAccumRequestSpec :: Spec
|
||||
mapAccumRequestSpec = describe "mapAccumRequest" $ do
|
||||
it "accumulates request times, post-state extraction" $
|
||||
runTimed (mapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
|
||||
[(0, 'a'), (5, 'b'), (10, 'c')]
|
||||
`shouldBe` [ [sec 0], [sec 0, sec 5], [sec 0, sec 5, sec 10] ]
|
||||
|
||||
it "output i is every request time seen so far" $
|
||||
hedgehog $ do
|
||||
secs' <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear 0 100))
|
||||
let out = runTimed (mapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
|
||||
[(fromIntegral s, ()) | s <- secs']
|
||||
expected = [ map (sec . fromIntegral) (take (i + 1) secs') | i <- [0 .. length secs' - 1] ]
|
||||
out === expected
|
||||
|
||||
preMapAccumRequestSpec :: Spec
|
||||
preMapAccumRequestSpec = describe "preMapAccumRequest" $ do
|
||||
it "accumulates request times, pre-state extraction" $
|
||||
runTimed (preMapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
|
||||
[(0, 'a'), (5, 'b'), (10, 'c')]
|
||||
`shouldBe` [ [], [sec 0], [sec 0, sec 5] ]
|
||||
|
||||
it "output i is every request time before the current step" $
|
||||
hedgehog $ do
|
||||
secs' <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear 0 100))
|
||||
let out = runTimed (preMapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
|
||||
[(fromIntegral s, ()) | s <- secs']
|
||||
expected = [ map (sec . fromIntegral) (take i secs') | i <- [0 .. length secs' - 1] ]
|
||||
out === expected
|
||||
|
||||
whenASpec :: Spec
|
||||
whenASpec = describe "whenA" $ do
|
||||
let counter = eff (\_ (_ :: Int) -> St (\s -> ((), s + 1)))
|
||||
|
||||
it "output is always () regardless of the predicate" $
|
||||
fst (runStEff (whenA (> 5) counter) 0 [1, 6, 2, 7])
|
||||
`shouldBe` [(), (), (), ()]
|
||||
|
||||
it "runs the inner arrow only when the predicate holds" $
|
||||
snd (runStEff (whenA (> 5) counter) 0 [1, 6, 2, 7])
|
||||
`shouldBe` 2
|
||||
|
||||
it "never runs the inner arrow when the predicate is always false" $
|
||||
snd (runStEff (whenA (const False) counter) 0 [1, 6, 2, 7])
|
||||
`shouldBe` 0
|
||||
|
||||
it "runs the inner arrow on every input when the predicate is always true" $
|
||||
snd (runStEff (whenA (const True) counter) 0 [1, 6, 2, 7])
|
||||
`shouldBe` 4
|
||||
|
||||
thenASpec :: Spec
|
||||
thenASpec = describe "thenA" $ do
|
||||
it "short-circuits on Left and continues on Right" $
|
||||
runPure (filterA (even @Int) `thenA` filterA (> 3)) [1..6]
|
||||
`shouldBe` [Left (), Left (), Left (), Right 4, Left (), Right 6]
|
||||
|
||||
it "(>>|) is an infix alias for thenA" $
|
||||
runPure (filterA (even @Int) >>| filterA (> 3)) [1..6]
|
||||
`shouldBe` [Left (), Left (), Left (), Right 4, Left (), Right 6]
|
||||
|
||||
it "second stage runs only when the first produces Right" $
|
||||
hedgehog $ do
|
||||
xs <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear (-20) 20))
|
||||
let out = runPure (filterA (even @Int) >>| filterA (> 0)) xs
|
||||
expected =
|
||||
[ if odd x then Left ()
|
||||
else if x > 0 then Right x
|
||||
else Left ()
|
||||
| x <- xs ]
|
||||
out === expected
|
||||
|
||||
sampleSpec :: Spec
|
||||
sampleSpec = describe "sample" $ do
|
||||
it "tags the current value onto the Event structure" $
|
||||
runPure sample [(1, Tick), (2, Event 'a'), (3, Tick)]
|
||||
`shouldBe` [Tick, Event @Int 2, Tick]
|
||||
|
||||
it "output is Event a iff the input event is present" $
|
||||
hedgehog $ do
|
||||
vals <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear 0 100))
|
||||
evs <- forAll $ Gen.list (Range.linear 0 30) eventGen
|
||||
let n = min (length vals) (length evs)
|
||||
ps = take n $ zip vals evs
|
||||
out = runPure sample ps
|
||||
for_ (zip ps out) $ \((v, ev), o) ->
|
||||
o === tag v ev
|
||||
|
||||
-- | A stateless arrow carrying a fixed entity set, for testing propagation.
|
||||
subscribed :: S.Set T.Text -> Mealy Identity Int Int
|
||||
subscribed ents = Mealy ents $ \_ -> Fun $ \_ a -> a
|
||||
|
||||
-- | Same as 'subscribed' but yields a function, for testing '<*>'.
|
||||
subscribedF :: S.Set T.Text -> Mealy Identity Int (Int -> Int)
|
||||
subscribedF ents = Mealy ents $ \_ -> Fun $ \_ a -> (a +)
|
||||
|
||||
entitiesSpec :: Spec
|
||||
entitiesSpec = describe "entities" $ do
|
||||
it "id carries no entities" $
|
||||
entities (arr id :: Mealy Identity Int Int) `shouldBe` S.empty
|
||||
|
||||
it "arr carries no entities" $
|
||||
entities (arr (+ 1) :: Mealy Identity Int Int) `shouldBe` S.empty
|
||||
|
||||
it "eff carries no entities" $
|
||||
entities (eff (\_ x -> Identity (x + 1 :: Int))) `shouldBe` S.empty
|
||||
|
||||
it "primitive combinators carry no entities" $ do
|
||||
entities (hold 'a') `shouldBe` S.empty
|
||||
entities (changes @Int) `shouldBe` S.empty
|
||||
entities edge `shouldBe` S.empty
|
||||
entities (sliding (3 :: Int) :: Mealy Identity (Event Int) [Int]) `shouldBe` S.empty
|
||||
|
||||
it "Category (.) unions entity sets" $
|
||||
entities (subscribed (S.singleton "a") >>> subscribed (S.singleton "b"))
|
||||
`shouldBe` S.fromList ["a", "b"]
|
||||
|
||||
it "Applicative (<*>) unions entity sets" $
|
||||
entities (subscribedF (S.singleton "a") <*> subscribed (S.singleton "b"))
|
||||
`shouldBe` S.fromList ["a", "b"]
|
||||
|
||||
it "Arrow (&&&) unions entity sets" $
|
||||
entities (subscribed (S.singleton "a") &&& subscribed (S.singleton "b"))
|
||||
`shouldBe` S.fromList ["a", "b"]
|
||||
|
||||
it "(>>>) unions entity sets" $
|
||||
entities (subscribed (S.singleton "a") >>> arr id >>> subscribed (S.singleton "b"))
|
||||
`shouldBe` S.fromList ["a", "b"]
|
||||
|
||||
it "left preserves the entity set" $
|
||||
entities (left (subscribed (S.singleton "a")) :: Mealy Identity (Either Int Int) (Either Int Int))
|
||||
`shouldBe` S.singleton "a"
|
||||
|
||||
it "first preserves the entity set" $
|
||||
entities (first (subscribed (S.singleton "a")) :: Mealy Identity (Int, Int) (Int, Int))
|
||||
`shouldBe` S.singleton "a"
|
||||
|
||||
it "fmap preserves the entity set" $
|
||||
entities (fmap (+ 1) (subscribed (S.singleton "a")))
|
||||
`shouldBe` S.singleton "a"
|
||||
@@ -0,0 +1,134 @@
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
|
||||
module BedroomSpec (spec) where
|
||||
|
||||
import AFRP (Event (..), entities)
|
||||
import qualified Data.Set as S
|
||||
import qualified Data.Text as T
|
||||
import Data.Aeson (Value)
|
||||
import HomeAssistant.Controller
|
||||
import HomeAssistant.Controller.Bedroom
|
||||
import Support
|
||||
import Test.Hspec
|
||||
|
||||
masseButton :: T.Text -> Event Value
|
||||
masseButton = Event . buttonEvent "event.bedroom_quick_remote_masse_action"
|
||||
|
||||
enishenButton :: T.Text -> Event Value
|
||||
enishenButton = Event . buttonEvent "event.bedroom_quick_jemina_action"
|
||||
|
||||
drawerState :: T.Text -> Event Value
|
||||
drawerState = Event . stateEvent "binary_sensor.bedroom_nightstand_drawer_sensor_masse_contact"
|
||||
|
||||
spec :: Spec
|
||||
spec = describe "Bedroom" $ do
|
||||
entitySpec
|
||||
drawerSpec
|
||||
buttonSpec
|
||||
|
||||
entitySpec :: Spec
|
||||
entitySpec = describe "entities" $ do
|
||||
it "entityChangeEvent' registers its entity id" $
|
||||
entities (entityChangeEvent' "sensor.foo")
|
||||
`shouldBe` S.singleton "sensor.foo"
|
||||
|
||||
it "entityRead / entityBool inherit the entity id" $ do
|
||||
entities (entityRead @Double "sensor.bar") `shouldBe` S.singleton "sensor.bar"
|
||||
entities (entityBool "binary_sensor.baz") `shouldBe` S.singleton "binary_sensor.baz"
|
||||
|
||||
it "bedroomDrawerController subscribes to the drawer sensor" $
|
||||
entities bedroomDrawerController
|
||||
`shouldBe` S.singleton "binary_sensor.bedroom_nightstand_drawer_sensor_masse_contact"
|
||||
|
||||
it "bedroomButtonController subscribes to both remote event entities" $
|
||||
entities bedroomButtonController
|
||||
`shouldBe` S.fromList
|
||||
[ "event.bedroom_quick_remote_masse_action"
|
||||
, "event.bedroom_quick_jemina_action"
|
||||
]
|
||||
|
||||
it "bedroomPresenceController subscribes to the presence sensor" $
|
||||
entities bedroomPresenceController
|
||||
`shouldBe` S.singleton "binary_sensor.presence_sensor_bedroom_occupancy"
|
||||
|
||||
it "humidifierController subscribes to the door sensor" $
|
||||
entities humidifierController
|
||||
`shouldBe` S.singleton "binary_sensor.makuuhuone_ovi_contact"
|
||||
|
||||
drawerSpec :: Spec
|
||||
drawerSpec = describe "bedroomDrawerController" $ do
|
||||
let entity = EntityId "switch.bedroom_drawer_light_masse"
|
||||
|
||||
it "turns the drawer light on when the drawer opens" $
|
||||
services (runHASS bedroomDrawerController [drawerState "on"])
|
||||
`shouldBe` [[switch [entity] True]]
|
||||
|
||||
it "turns the drawer light off when the drawer closes" $
|
||||
services (runHASS bedroomDrawerController [drawerState "off"])
|
||||
`shouldBe` [[switch [entity] False]]
|
||||
|
||||
it "toggles the light as the drawer opens and closes" $
|
||||
services (runHASS bedroomDrawerController
|
||||
[ drawerState "on", drawerState "off", drawerState "on" ])
|
||||
`shouldBe` [ [switch [entity] True]
|
||||
, [switch [entity] False]
|
||||
, [switch [entity] True]
|
||||
]
|
||||
|
||||
it "ignores state events for other entities" $
|
||||
services (runHASS bedroomDrawerController
|
||||
[Event (stateEvent "binary_sensor.some_other_contact" "on")])
|
||||
`shouldBe` [[]]
|
||||
|
||||
it "does nothing on Tick" $
|
||||
services (runHASS bedroomDrawerController [Tick])
|
||||
`shouldBe` [[]]
|
||||
|
||||
buttonSpec :: Spec
|
||||
buttonSpec = describe "bedroomButtonController" $ do
|
||||
let area = AreaId "makuuhuone"
|
||||
|
||||
it "Masse single click turns on his nightstand scene (lowest)" $
|
||||
services (runHASS bedroomButtonController [masseButton "1_short_release"])
|
||||
`shouldBe` [[activateScene "scene.makuuhuone_masse"]]
|
||||
|
||||
it "Masse double click turns on the middle scene" $
|
||||
services (runHASS bedroomButtonController [masseButton "1_double_press"])
|
||||
`shouldBe` [[activateScene "scene.makuuhuone_keski"]]
|
||||
|
||||
it "Masse long click turns on the bright scene" $
|
||||
services (runHASS bedroomButtonController [masseButton "1_long_press"])
|
||||
`shouldBe` [[activateScene "scene.makuuhuone_kirkas"]]
|
||||
|
||||
it "Masse off click turns the bedroom lights off" $
|
||||
services (runHASS bedroomButtonController [masseButton "2_short_release"])
|
||||
`shouldBe` [[light [area] Off]]
|
||||
|
||||
it "Enishen single click turns on her nightstand scene (lowest)" $
|
||||
services (runHASS bedroomButtonController [enishenButton "1_short_release"])
|
||||
`shouldBe` [[activateScene "scene.makuuhuone_jemina"]]
|
||||
|
||||
it "Enishen double click turns on the middle scene" $
|
||||
services (runHASS bedroomButtonController [enishenButton "1_double_press"])
|
||||
`shouldBe` [[activateScene "scene.makuuhuone_keski"]]
|
||||
|
||||
it "Enishen long click turns on the bright scene" $
|
||||
services (runHASS bedroomButtonController [enishenButton "1_long_press"])
|
||||
`shouldBe` [[activateScene "scene.makuuhuone_kirkas"]]
|
||||
|
||||
it "Enishen off click turns the bedroom lights off" $
|
||||
services (runHASS bedroomButtonController [enishenButton "2_short_release"])
|
||||
`shouldBe` [[light [area] Off]]
|
||||
|
||||
it "ignores the initial press (scene only fires on release)" $
|
||||
services (runHASS bedroomButtonController [masseButton "1_initial_press"])
|
||||
`shouldBe` [[]]
|
||||
|
||||
it "ignores button events for other entities" $
|
||||
services (runHASS bedroomButtonController
|
||||
[Event (buttonEvent "event.some_other_action" "1_short_release")])
|
||||
`shouldBe` [[]]
|
||||
|
||||
it "does nothing on Tick" $
|
||||
services (runHASS bedroomButtonController [Tick])
|
||||
`shouldBe` [[]]
|
||||
+7
-7
@@ -2,7 +2,7 @@
|
||||
|
||||
module BusSpec (spec) where
|
||||
|
||||
import AFRP (Request (..))
|
||||
import AFRP (Request (..), Event (..))
|
||||
import Control.Concurrent.STM
|
||||
( atomically
|
||||
, dupTChan
|
||||
@@ -10,7 +10,7 @@ import Control.Concurrent.STM
|
||||
, writeTChan
|
||||
)
|
||||
import Data.Aeson (Value (..))
|
||||
import Data.Time (UTCTime (..))
|
||||
import Data.Time (UTCTime (..), utc)
|
||||
import Data.UUID (nil)
|
||||
import HomeAssistant.Controller (HASSEff (..), Service (..), Target(..))
|
||||
import HomeAssistant.Runtime.Bus
|
||||
@@ -22,15 +22,15 @@ spec = describe "Bus" $ do
|
||||
it "broadcasts inbound messages to every dup'd channel in order" $ withBus InfoS $ \bus -> do
|
||||
p1 <- atomically $ dupTChan (busInbound bus)
|
||||
p2 <- atomically $ dupTChan (busInbound bus)
|
||||
atomically $ writeTChan (busInbound bus) (Number 1)
|
||||
atomically $ writeTChan (busInbound bus) (Number 2)
|
||||
atomically $ writeTChan (busInbound bus) (Event (Number 1))
|
||||
atomically $ writeTChan (busInbound bus) (Event (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)
|
||||
r1 `shouldBe` (Event (Number 1), Event (Number 2))
|
||||
r2 `shouldBe` (Event (Number 1), Event (Number 2))
|
||||
|
||||
it "channelHassEval writes CallService to the outbound channel" $ withBus InfoS $ \bus -> do
|
||||
let req = Request (UTCTime (toEnum 0) 0) nil
|
||||
let req = Request (UTCTime (toEnum 0) 0) utc nil
|
||||
svc = Service "light" "turn_on" Nothing [EntityId "light.bedroom_masse"]
|
||||
runKatipContextT (busLogEnv bus) () (Namespace ["test"]) $
|
||||
channelHassEval bus (CallService req svc)
|
||||
|
||||
+68
-21
@@ -2,31 +2,78 @@
|
||||
|
||||
module ConnectionSpec (spec) where
|
||||
|
||||
import AFRP (Request(..))
|
||||
import Data.Aeson (object, (.=))
|
||||
import Data.Maybe (fromJust)
|
||||
import Data.Text (Text)
|
||||
import Data.Time (UTCTime (..), utc)
|
||||
import Data.UUID (fromString)
|
||||
import HomeAssistant.Controller (Service (..), Target(..))
|
||||
import HomeAssistant.Runtime.Connection (encodeService)
|
||||
import HomeAssistant.Runtime.Connection (encodeService, dedupeBatch)
|
||||
import Test.Hspec
|
||||
|
||||
spec :: Spec
|
||||
spec = describe "encodeService" $ do
|
||||
it "encodes a call_service message" $
|
||||
encodeService 7 (Service "light" "turn_on" Nothing [EntityId "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)]
|
||||
]
|
||||
spec = do
|
||||
describe "encodeService" $ do
|
||||
it "encodes a call_service message" $
|
||||
encodeService 7 (Service "light" "turn_on" Nothing [EntityId "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)])) [EntityId "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)]
|
||||
]
|
||||
it "includes service_data when present" $
|
||||
encodeService 8 (Service "light" "turn_on" (Just (object ["brightness" .= (200 :: Int)])) [EntityId "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)]
|
||||
]
|
||||
|
||||
describe "dedupeBatch" $ do
|
||||
it "collapses identical calls to one" $
|
||||
let batch = [ (req 1, lightOn [AreaId "x"])
|
||||
, (req 2, lightOn [AreaId "x"])
|
||||
, (req 3, lightOn [AreaId "x"])
|
||||
]
|
||||
in dedupeBatch batch `shouldBe` [(req 3, lightOn [AreaId "x"])]
|
||||
|
||||
it "keeps same-target different-service calls separate" $
|
||||
let batch = [ (req 1, lightOn [AreaId "x"])
|
||||
, (req 2, lightOff [AreaId "x"])
|
||||
]
|
||||
result = dedupeBatch batch
|
||||
in length result `shouldBe` 2
|
||||
|
||||
it "newest call wins for the same key" $
|
||||
let batch = [ (req 1, lightOn [AreaId "x"])
|
||||
, (req 2, lightOn [AreaId "x"])
|
||||
, (req 3, lightOn [AreaId "x"])
|
||||
]
|
||||
in map requestTraceId (map fst (dedupeBatch batch)) `shouldBe`
|
||||
[fromJust (fromString "00000000-0000-0000-0000-000000000003")]
|
||||
|
||||
it "treats target lists in different order as the same key" $
|
||||
let batch = [ (req 1, lightOn [EntityId "a", EntityId "b"])
|
||||
, (req 2, lightOn [EntityId "b", EntityId "a"])
|
||||
]
|
||||
in length (dedupeBatch batch) `shouldBe` 1
|
||||
|
||||
req :: Int -> Request
|
||||
req n = Request (UTCTime (toEnum 0) (fromIntegral (0 :: Int))) utc
|
||||
(fromJust (fromString uuid))
|
||||
where
|
||||
pad i = replicate (12 - length (show i)) '0' <> show i
|
||||
uuid = "00000000-0000-0000-0000-" <> pad n
|
||||
|
||||
lightOn :: [Target] -> Service
|
||||
lightOn targets = Service "light" "turn_on" Nothing targets
|
||||
|
||||
lightOff :: [Target] -> Service
|
||||
lightOff targets = Service "light" "turn_off" Nothing targets
|
||||
|
||||
@@ -1,16 +1,24 @@
|
||||
module Main (main) where
|
||||
|
||||
import Test.Hspec (hspec)
|
||||
import qualified AFRPLawsSpec
|
||||
import qualified AFRPSpec
|
||||
import qualified BackoffProp
|
||||
import qualified BedroomSpec
|
||||
import qualified BusSpec
|
||||
import qualified ConnectionSpec
|
||||
import qualified MetricsSpec
|
||||
import qualified RuntimeSpec
|
||||
import qualified SupervisorSpec
|
||||
|
||||
main :: IO ()
|
||||
main = hspec $ do
|
||||
AFRPLawsSpec.spec
|
||||
AFRPSpec.spec
|
||||
BedroomSpec.spec
|
||||
BusSpec.spec
|
||||
ConnectionSpec.spec
|
||||
MetricsSpec.spec
|
||||
RuntimeSpec.spec
|
||||
SupervisorSpec.spec
|
||||
BackoffProp.spec
|
||||
|
||||
@@ -0,0 +1,150 @@
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
|
||||
module MetricsSpec (spec) where
|
||||
|
||||
import Data.HashMap.Strict (HashMap)
|
||||
import qualified Data.HashMap.Strict as HM
|
||||
import Data.Int (Int64)
|
||||
import Data.Text (Text)
|
||||
import HomeAssistant.Runtime.Metrics
|
||||
( DsType (..)
|
||||
, DsSpec (..)
|
||||
, dsTypeOf
|
||||
, sanitizeName
|
||||
, buildSchema
|
||||
, buildCreateArgs
|
||||
, buildUpdateArgs
|
||||
, ensureRrd
|
||||
, sampleAndUpdate
|
||||
, metricsAction
|
||||
)
|
||||
import qualified System.Metrics as M (Value (..))
|
||||
import Test.Hspec
|
||||
import Control.Exception (try, SomeException)
|
||||
import System.Directory (findExecutable, getTemporaryDirectory, removeFile)
|
||||
import System.Exit (ExitCode (ExitSuccess))
|
||||
import System.Process (readProcessWithExitCode)
|
||||
import qualified System.Metrics as Metrics
|
||||
import qualified System.Metrics.Counter as Counter
|
||||
import qualified System.Metrics.Gauge as Gauge
|
||||
|
||||
spec :: Spec
|
||||
spec = do
|
||||
describe "dsTypeOf" $ do
|
||||
it "maps Counter to Derive" $
|
||||
dsTypeOf (M.Counter 1000) `shouldBe` Just Derive
|
||||
|
||||
it "maps Gauge to Gauge" $
|
||||
dsTypeOf (M.Gauge 500) `shouldBe` Just Gauge
|
||||
|
||||
it "maps Label to Nothing" $
|
||||
dsTypeOf (M.Label "hello") `shouldBe` Nothing
|
||||
|
||||
describe "sanitizeName" $ do
|
||||
it "replaces dots with underscores for short names" $
|
||||
sanitizeName ("rts.gc.cpu_ms" :: Text) `shouldBe` "rts_gc_cpu_ms"
|
||||
|
||||
it "shortens names longer than 19 chars to 15 chars + _ + 3 hex" $ do
|
||||
let result = sanitizeName ("rts.gc.par_balanced_bytes_copied" :: Text)
|
||||
length result `shouldBe` 19
|
||||
take 15 result `shouldBe` "rts_gc_par_bala"
|
||||
drop 15 result `shouldBe` "_" ++ drop 16 result
|
||||
|
||||
it "is deterministic (same input -> same output)" $
|
||||
sanitizeName ("rts.gc.peak_megabytes_allocated" :: Text)
|
||||
`shouldBe` sanitizeName ("rts.gc.peak_megabytes_allocated" :: Text)
|
||||
|
||||
describe "buildSchema" $ do
|
||||
it "builds sorted DsSpecs from counters and gauges, skipping labels" $
|
||||
let sample :: HashMap Text M.Value
|
||||
sample = HM.fromList
|
||||
[ ("x.allocated", M.Counter 1000)
|
||||
, ("a.bytes_used", M.Gauge 500)
|
||||
, ("c.label_thing", M.Label "irrelevant")
|
||||
]
|
||||
in buildSchema sample `shouldBe`
|
||||
[ DsSpec "a.bytes_used" "a_bytes_used" Gauge
|
||||
, DsSpec "x.allocated" "x_allocated" Derive
|
||||
]
|
||||
|
||||
it "produces dsName <= 19 chars for long ekg GC metric names" $
|
||||
let sample :: HashMap Text M.Value
|
||||
sample = HM.fromList
|
||||
[ ("rts.gc.par_balanced_bytes_copied", M.Gauge 1)
|
||||
, ("rts.gc.peak_megabytes_allocated", M.Gauge 2)
|
||||
, ("rts.gc.cumulative_bytes_used", M.Counter 3)
|
||||
]
|
||||
in map (length . dsName) (buildSchema sample) `shouldSatisfy` all (<= 19)
|
||||
|
||||
describe "buildCreateArgs" $ do
|
||||
it "builds create argv with mixed DERIVE and GAUGE DSes and RRAs" $
|
||||
buildCreateArgs "test.rrd" 10
|
||||
[ ("ds1", Derive)
|
||||
, ("ds2", Gauge)
|
||||
]
|
||||
`shouldBe`
|
||||
[ "create"
|
||||
, "test.rrd"
|
||||
, "--step"
|
||||
, "10"
|
||||
, "DS:ds1:DERIVE:20:0:U"
|
||||
, "DS:ds2:GAUGE:20:0:U"
|
||||
, "RRA:AVERAGE:0.5:1:6000"
|
||||
, "RRA:MAX:0.5:1:6000"
|
||||
, "RRA:AVERAGE:0.5:360:1680"
|
||||
, "RRA:MAX:0.5:360:1680"
|
||||
]
|
||||
|
||||
describe "buildUpdateArgs" $ do
|
||||
it "builds update argv with numeric values" $
|
||||
buildUpdateArgs "test.rrd" ["ds1", "ds2"] [Just 100, Just 200]
|
||||
`shouldBe`
|
||||
[ "update"
|
||||
, "test.rrd"
|
||||
, "--template"
|
||||
, "ds1:ds2"
|
||||
, "N:100:200"
|
||||
]
|
||||
|
||||
it "renders Nothing as U (unknown)" $
|
||||
buildUpdateArgs "test.rrd" ["ds1", "ds2"] [Just 100, Nothing]
|
||||
`shouldBe`
|
||||
[ "update"
|
||||
, "test.rrd"
|
||||
, "--template"
|
||||
, "ds1:ds2"
|
||||
, "N:100:U"
|
||||
]
|
||||
|
||||
it "renders all-Nothing as all-U" $
|
||||
buildUpdateArgs "test.rrd" ["ds1"] [Nothing]
|
||||
`shouldBe`
|
||||
[ "update"
|
||||
, "test.rrd"
|
||||
, "--template"
|
||||
, "ds1"
|
||||
, "N:U"
|
||||
]
|
||||
|
||||
describe "end-to-end (rrdtool-gated)" $ do
|
||||
it "creates an rrd, samples, and updates it" $ do
|
||||
mRrdtool <- findExecutable "rrdtool"
|
||||
case mRrdtool of
|
||||
Nothing -> pendingWith "rrdtool not on PATH"
|
||||
Just rrdtool -> do
|
||||
store <- Metrics.newStore
|
||||
c <- Metrics.createCounter "test.counter" store
|
||||
g <- Metrics.createGauge "test.gauge" store
|
||||
Counter.inc c
|
||||
Gauge.set g 42
|
||||
tmp <- getTemporaryDirectory
|
||||
let rrdPath = tmp ++ "/hass-controller-metrics-test.rrd"
|
||||
_ <- try (removeFile rrdPath) :: IO (Either SomeException ())
|
||||
schema <- buildSchema <$> Metrics.sampleAll store
|
||||
ensureRrd rrdtool rrdPath schema
|
||||
sampleAndUpdate store rrdtool rrdPath schema
|
||||
(rc, out, _) <- readProcessWithExitCode rrdtool ["fetch", rrdPath, "AVERAGE"] ""
|
||||
rc `shouldBe` ExitSuccess
|
||||
length out `shouldSatisfy` (> 0)
|
||||
_ <- try (removeFile rrdPath) :: IO (Either SomeException ())
|
||||
return ()
|
||||
+36
-37
@@ -1,43 +1,42 @@
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
module RuntimeSpec (spec) where
|
||||
|
||||
import Control.Concurrent (threadDelay)
|
||||
import Control.Concurrent.Async (async, race)
|
||||
import Control.Concurrent.STM (atomically, isEmptyTChan, readTChan, writeTChan)
|
||||
import Data.Aeson (Value, object, (.=))
|
||||
import Data.Text (Text)
|
||||
import HomeAssistant.Controller (light, lightController, Target (..))
|
||||
import HomeAssistant.Runtime (Controller (..), runController)
|
||||
import HomeAssistant.Runtime.Bus
|
||||
import Test.Hspec
|
||||
import Katip (Severity(..))
|
||||
|
||||
spec :: Spec
|
||||
spec = describe "runController" $ do
|
||||
it "feeds inbound events through the machine and forwards service calls" $ withBus InfoS $ \bus -> do
|
||||
_ <- async (runController bus (Controller "test" lightController True))
|
||||
putStrLn "Before the delay"
|
||||
threadDelay 100000 -- let the controller dup its inbound channel
|
||||
putStrLn "After the delay"
|
||||
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
|
||||
putStrLn "After the writes"
|
||||
Right (_, svc1) <- boundedRead (busOutbound bus)
|
||||
Right (_, svc2) <- boundedRead (busOutbound bus)
|
||||
putStrLn "After the reads"
|
||||
svc1 `shouldBe` light [EntityId "light.bedroom_masse"] True
|
||||
svc2 `shouldBe` light [EntityId "light.bedroom_masse"] False
|
||||
atomically (isEmptyTChan (busOutbound bus)) `shouldReturn` True
|
||||
where
|
||||
boundedRead channel = race (threadDelay 10000) (atomically (readTChan channel))
|
||||
spec = pure ()
|
||||
|
||||
doorEvent :: Text -> Value
|
||||
doorEvent state = object
|
||||
[ "event" .= object
|
||||
[ "data" .= object
|
||||
[ "entity_id" .= ("binary_sensor.makuuhuone_ovi_contact" :: Text)
|
||||
, "new_state" .= object ["state" .= state]
|
||||
]
|
||||
]
|
||||
]
|
||||
-- I'm commenting out the test as that specific controller
|
||||
-- was just a dummy for seeing how the data flows (dead code)
|
||||
-- But leaving it commented as it retains some code sample on how and what to test
|
||||
--
|
||||
-- spec :: Spec
|
||||
-- spec = describe "runController" $ do
|
||||
-- it "feeds inbound events through the machine and forwards service calls" $ withBus InfoS $ \bus -> do
|
||||
-- _ <- async (runController bus (Controller "test" lightController True))
|
||||
-- putStrLn "Before the delay"
|
||||
-- threadDelay 100000 -- let the controller dup its inbound channel
|
||||
-- putStrLn "After the delay"
|
||||
-- atomically $ writeTChan (busInbound bus) (Event (doorEvent "on")) -- initial value: no change event
|
||||
-- atomically $ writeTChan (busInbound bus) (Event (doorEvent "off")) -- door closes: lights on
|
||||
-- atomically $ writeTChan (busInbound bus) (Event (doorEvent "on")) -- door opens: lights off
|
||||
-- putStrLn "After the writes"
|
||||
-- Right (_, svc1) <- boundedRead (busOutbound bus)
|
||||
-- Right (_, svc2) <- boundedRead (busOutbound bus)
|
||||
-- putStrLn "After the reads"
|
||||
-- svc1 `shouldBe` light [EntityId "light.bedroom_masse"] True
|
||||
-- svc2 `shouldBe` light [EntityId "light.bedroom_masse"] False
|
||||
-- atomically (isEmptyTChan (busOutbound bus)) `shouldReturn` True
|
||||
-- where
|
||||
-- boundedRead channel = race (threadDelay 10000) (atomically (readTChan channel))
|
||||
--
|
||||
-- doorEvent :: Text -> Value
|
||||
-- doorEvent state = object
|
||||
-- [ "event" .= object
|
||||
-- [ "variables" .= object
|
||||
-- [ "trigger" .= object
|
||||
-- [ "entity_id" .= ("binary_sensor.makuuhuone_ovi_contact" :: Text)
|
||||
-- , "to_state" .= object ["state" .= state]
|
||||
-- ]
|
||||
-- ]
|
||||
-- ]
|
||||
-- ]
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
{-# LANGUAGE OverloadedStrings #-}
|
||||
|
||||
module Support
|
||||
( Acc(..)
|
||||
, interp
|
||||
, runHASS
|
||||
, services
|
||||
, fakeRequest
|
||||
, sec
|
||||
, stateEvent
|
||||
, buttonEvent
|
||||
) where
|
||||
|
||||
import Control.Monad.Fix (MonadFix (..))
|
||||
import Data.Aeson (Value, object, (.=))
|
||||
import qualified Data.Text as T
|
||||
import Data.Time (UTCTime (..), utc)
|
||||
import Data.UUID (nil)
|
||||
import AFRP (Mealy (..), Request (..), stepAuto)
|
||||
import HomeAssistant.Controller (HASSEff (..), Service)
|
||||
|
||||
fakeRequest :: Request
|
||||
fakeRequest = Request (sec 0) utc nil
|
||||
|
||||
sec :: Integer -> UTCTime
|
||||
sec n = UTCTime (toEnum 0) (fromIntegral n)
|
||||
|
||||
-- | A pure Writer-like monad accumulating `Service` calls per step.
|
||||
newtype Acc a = Acc { runAcc :: [Service] -> (a, [Service]) }
|
||||
|
||||
instance Functor Acc where
|
||||
fmap f (Acc g) = Acc $ \s -> let (a, s') = g s in (f a, s')
|
||||
|
||||
instance Applicative Acc where
|
||||
pure a = Acc (\s -> (a, s))
|
||||
Acc f <*> Acc x = Acc $ \s -> let (f', s') = f s; (a, s'') = x s' in (f' a, s'')
|
||||
|
||||
instance Monad Acc where
|
||||
Acc m >>= k = Acc $ \s -> let (a, s') = m s; (b, s'') = runAcc (k a) s' in (b, s'')
|
||||
|
||||
instance MonadFix Acc where
|
||||
mfix f = Acc $ \s -> let (a, s') = runAcc (f a) s in (a, s')
|
||||
|
||||
-- | Interpret `HASSEff` in `Acc`: record `CallService`, drop tracing/debug.
|
||||
interp :: HASSEff a -> Acc a
|
||||
interp (CallService _ svc) = Acc $ \s -> ((), s ++ [svc])
|
||||
interp (Debug _) = pure ()
|
||||
interp (Trace _ _) = pure ()
|
||||
|
||||
-- | Run a HASS arrow over a list of inputs, collecting per-step emitted services.
|
||||
runHASS :: Mealy HASSEff a b -> [a] -> [(b, [Service])]
|
||||
runHASS m = go (runMealy m interp)
|
||||
where
|
||||
go _ [] = []
|
||||
go w (a : as) =
|
||||
case runAcc (stepAuto w fakeRequest a) [] of
|
||||
((b, w'), svcs) -> (b, svcs) : go w' as
|
||||
|
||||
services :: [(b, [Service])] -> [[Service]]
|
||||
services = map snd
|
||||
|
||||
-- | Build a state-trigger payload matching `entityChangeEvent'` / `entityBool'`
|
||||
-- lenses. The subscribe_trigger websocket event wraps the trigger datum under
|
||||
-- `event.variables.trigger`, with `entity_id` and `to_state.state` fields.
|
||||
stateEvent :: T.Text -> T.Text -> Value
|
||||
stateEvent entityId state = object
|
||||
[ "event" .= object
|
||||
[ "variables" .= object
|
||||
[ "trigger" .= object
|
||||
[ "entity_id" .= entityId
|
||||
, "to_state" .= object [ "state" .= state ]
|
||||
]
|
||||
]
|
||||
]
|
||||
]
|
||||
|
||||
-- | Build an Ikea button trigger payload matching `ikeaQuickButton` lenses.
|
||||
buttonEvent :: T.Text -> T.Text -> Value
|
||||
buttonEvent entityId eventType = object
|
||||
[ "event" .= object
|
||||
[ "variables" .= object
|
||||
[ "trigger" .= object
|
||||
[ "entity_id" .= entityId
|
||||
, "to_state" .= object
|
||||
[ "attributes" .= object [ "event_type" .= eventType ] ]
|
||||
]
|
||||
]
|
||||
]
|
||||
]
|
||||
Reference in New Issue
Block a user