4 Commits
Author SHA1 Message Date
MasseR aaba917a9b Test the SerializeUTCTime property 2026-09-12 20:25:24 +03:00
MasseR a3dd63f26c Update tests 2026-09-12 20:18:13 +03:00
MasseR b20320c779 Explicit state 2026-09-12 20:09:42 +03:00
MasseR e16010bf2f Backport serializable arrow 2026-09-12 16:27:12 +03:00
9 changed files with 458 additions and 273 deletions
+8 -8
View File
@@ -1,7 +1,7 @@
{ mkDerivation, aeson, annotated-exception, async, base, bytestring { mkDerivation, aeson, annotated-exception, async, base, bytestring
, containers, directory, ekg-core, hedgehog, hspec, hspec-hedgehog , cereal, containers, directory, ekg-core, filepath, hedgehog
, katip, lens, lens-aeson, lib, network, process, stm, text, time , hspec, hspec-hedgehog, katip, lens, lens-aeson, lib, network
, unordered-containers, uuid, websockets , process, stm, text, time, unordered-containers, uuid, websockets
}: }:
mkDerivation { mkDerivation {
pname = "home-assistant-controller"; pname = "home-assistant-controller";
@@ -10,14 +10,14 @@ mkDerivation {
isLibrary = true; isLibrary = true;
isExecutable = true; isExecutable = true;
libraryHaskellDepends = [ libraryHaskellDepends = [
aeson annotated-exception async base bytestring containers aeson annotated-exception async base bytestring cereal containers
directory ekg-core katip lens lens-aeson network process stm text directory ekg-core filepath katip lens lens-aeson network process
time unordered-containers uuid websockets stm text time unordered-containers uuid websockets
]; ];
executableHaskellDepends = [ base ]; executableHaskellDepends = [ base ];
testHaskellDepends = [ testHaskellDepends = [
aeson annotated-exception async base containers directory ekg-core aeson annotated-exception async base cereal containers directory
hedgehog hspec hspec-hedgehog katip process stm text time ekg-core hedgehog hspec hspec-hedgehog katip process stm text time
unordered-containers uuid unordered-containers uuid
]; ];
license = lib.meta.getLicenseFromSpdxId "BSD-3-Clause"; license = lib.meta.getLicenseFromSpdxId "BSD-3-Clause";
+4
View File
@@ -96,6 +96,9 @@ library
, unordered-containers , unordered-containers
, process , process
, directory , directory
, cereal
, containers
, filepath
-- Directories containing source files. -- Directories containing source files.
hs-source-dirs: src hs-source-dirs: src
@@ -166,6 +169,7 @@ test-suite home-assistant-controller-test
hspec, hspec,
stm, stm,
aeson, aeson,
cereal,
text, text,
async, async,
hedgehog, hedgehog,
+359 -146
View File
@@ -3,12 +3,13 @@
module AFRP module AFRP
( Mealy(..) ( Mealy(..)
, Auto(..)
, eff , eff
, withEntities , withEntities
, Event(..) , Event(..)
, hold , hold
, events , events
, switch -- , switch
, preMapAccum , preMapAccum
, preMapAccumRequest , preMapAccumRequest
, mapAccum , mapAccum
@@ -22,8 +23,8 @@ module AFRP
, lMerge , lMerge
, Pair(..) , Pair(..)
, Request(..) , Request(..)
, SerializeUTCTime(..)
, edge , edge
, dropFirst
, duration , duration
, tag , tag
, isEvent , isEvent
@@ -31,23 +32,67 @@ module AFRP
, sample , sample
, rollup , rollup
, sliding , sliding
, fixed
, debounce , debounce
, currentTime , currentTime
, onEvent , onEvent
, save
, load
, stepAuto
, stepAutoSerializing
) where ) where
import Control.Category (Category(..), (>>>)) import Control.Category (Category(..), (>>>))
import Prelude hiding ((.), id) import Prelude hiding ((.), id)
import Control.Arrow (Arrow(..), ArrowChoice(..), ArrowLoop(..)) import Control.Arrow (Arrow(..), ArrowChoice(..))
import Data.Time (UTCTime, NominalDiffTime, diffUTCTime, addUTCTime, TimeZone, LocalTime, utcToLocalTime) import Data.Time (UTCTime (UTCTime), NominalDiffTime, diffUTCTime, addUTCTime, TimeZone, LocalTime, utcToLocalTime, Day (..), diffTimeToPicoseconds, picosecondsToDiffTime)
import Control.Monad.Fix (MonadFix (mfix))
import Data.Either (fromLeft) import Data.Either (fromLeft)
import Data.Bool (bool) import Data.Bool (bool)
import Data.Monoid (Endo(..))
import Data.UUID (UUID) import Data.UUID (UUID)
import qualified Data.Set as S import qualified Data.Set as S
import qualified Data.Text as T import qualified Data.Text as T
import Data.Serialize (Get, Putter, runPut, 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)
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 Pair a b = Pair !a !b
@@ -57,32 +102,171 @@ data Request = Request
, requestTraceId :: !UUID , requestTraceId :: !UUID
} deriving (Show, Eq) } deriving (Show, Eq)
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 :: Auto m a b -> B.ByteString
serialize = \case
Fun _ -> runPut $ put ()
Stateful Codec{putter} s _ -> runPut $ 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
_ <- B.writeFile path $ serialize s
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 eyt" a
| otherwise = throwIO e
-- | The set of entity ids an arrow subscribes to. Static: it does not -- | 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 -- change as the machine steps, so the runtime can read it once to build
-- trigger subscriptions. -- trigger subscriptions.
data Mealy eff a b = Mealy data Mealy eff a b = Mealy
{ entities :: S.Set T.Text { entities :: S.Set T.Text
, runMealy :: forall m. MonadFix m => (forall x. eff x -> m x) -> Request -> a -> m (Pair b (Mealy eff a b)) , runMealy :: forall m. Monad m => (forall x. eff x -> m x) -> Auto m a b
} }
instance Semigroup b => Semigroup (Mealy eff a b) where instance Semigroup b => Semigroup (Mealy eff a b) where
Mealy ast af <> Mealy bst bf = Mealy (ast <> bst) $ \nt r a -> do Mealy ast af <> Mealy bst bf = Mealy (ast <> bst) $ \nt -> af nt <> bf nt
Pair x af' <- af nt r a
Pair x' bf' <- bf nt r a
pure (Pair (x <> x') (af' <> bf'))
instance Monoid b => Monoid (Mealy eff a b) where instance Monoid b => Monoid (Mealy eff a b) where
mempty = m mempty = Mealy mempty $ \_nt -> mempty
where -- where
m = Mealy mempty $ \_ _ _ -> pure (Pair mempty m) -- m = Mealy mempty $ \_ _ _ -> pure (Pair mempty m)
eff :: (Request -> a -> eff b) -> Mealy eff a b eff :: (Request -> a -> eff b) -> Mealy eff a b
eff f = m eff f = Mealy mempty $ \nt -> do
where Stateful (Codec get put) (State () False) $ \s req a -> do
m = Mealy mempty $ \nt req x -> b <- nt (f req a)
nt (f req x) >>= \b -> pure (Pair b m) pure (b, s)
-- | Override the static entity set of an arrow. Use when a combinator -- | Override the static entity set of an arrow. Use when a combinator
-- (e.g. 'switch') hides continuation entities from the runtime's -- (e.g. 'switch') hides continuation entities from the runtime's
@@ -91,53 +275,32 @@ withEntities :: S.Set T.Text -> Mealy eff a b -> Mealy eff a b
withEntities es (Mealy _ f) = Mealy es f withEntities es (Mealy _ f) = Mealy es f
instance Category (Mealy eff) where instance Category (Mealy eff) where
id = Mealy mempty (\_ _ x -> pure (Pair x id)) id = Mealy mempty (\_ -> id)
(Mealy ast f) . (Mealy bst g) = Mealy (ast <> bst) $ \nt t a -> do (Mealy ast f) . (Mealy bst g) = Mealy (ast <> bst) $ \nt -> do
Pair b g' <- g nt t a f nt . g nt
Pair c f' <- f nt t b
pure (Pair c (f' . g'))
instance Arrow (Mealy eff) where instance Arrow (Mealy eff) where
arr f = mealy arr f = Mealy mempty $ \_nt -> arr f
where first (Mealy st f) = Mealy st $ \nt -> first (f nt)
mealy = Mealy mempty $ \_ _ b -> pure (Pair (f b) mealy)
first (Mealy st f) = Mealy st $ \nt t (b,d) -> do
Pair c f' <- f nt t b
pure (Pair (c, d) (first f'))
instance ArrowChoice (Mealy eff) where instance ArrowChoice (Mealy eff) where
left (Mealy st f) = lm left (Mealy st f) = Mealy st $ \nt -> left (f nt)
where
lm = Mealy st $ \nt t -> \case
Left b -> do
Pair c f' <- f nt t b
pure (Pair (Left c) (left f'))
Right d -> pure (Pair (Right d) lm)
-- ArrowLoop is incompatible with strict Pair (strict fields prevent
-- the lazy knot-tying that mfix requires with loop).
-- instance ArrowLoop (Mealy eff) where
-- loop (Mealy st f) = Mealy st $ \nt t b -> do
-- Pair (c,_) f' <- mfix $ \(Pair (_,d) _) -> f nt t (b,d)
-- pure (Pair c (loop f'))
instance Functor (Mealy eff a) where instance Functor (Mealy eff a) where
fmap f (Mealy st g) = Mealy st $ \nt t a -> do fmap f (Mealy st g) = Mealy st $ \nt -> fmap f (g nt)
Pair b g' <- g nt t a
pure (Pair (f b) (fmap f g'))
instance Applicative (Mealy eff a) where instance Applicative (Mealy eff a) where
pure b = Mealy mempty $ \_ _ _ -> pure (Pair b (pure b)) pure b = Mealy mempty $ \_ -> pure b
Mealy ast f <*> Mealy bst x = Mealy (ast <> bst) $ \nt t a -> do Mealy ast f <*> Mealy bst x = Mealy (ast <> bst) $ \nt ->
b' <- x nt t a f nt <*> x nt
let Pair x' xNext = b'
Pair f' fNext <- f nt t a
pure (Pair (f' x') (fNext <*> xNext))
data Event a data Event a
= Tick = Tick
| Event a | Event a
deriving (Show, Eq, Functor, Foldable, Traversable) deriving (Show, Eq, Functor, Foldable, Traversable, Generic)
instance Serialize a => Serialize (Event a)
instance Semigroup (Event a) where instance Semigroup (Event a) where
(<>) = lMerge (<>) = lMerge
@@ -145,10 +308,18 @@ instance Semigroup (Event a) where
instance Monoid (Event a) where instance Monoid (Event a) where
mempty = Tick mempty = Tick
hold :: a -> Mealy eff (Event a) a hold :: (Serialize a, Eq a) => a -> Mealy m (Event a) a
hold a = Mealy mempty $ \_ _ -> \case hold def = mapAccum step def id
Tick -> pure (Pair a (hold a)) where
Event a' -> pure (Pair a' (hold a')) step prev = \case
Tick -> prev
Event new -> new
-- hold :: a -> Mealy eff (Event a) a
-- hold a = Mealy mempty $ \_ _ -> \case
-- Tick -> pure (Pair a (hold a))
-- Event a' -> pure (Pair a' (hold a'))
events :: Mealy eff (Event a) (Either () a) events :: Mealy eff (Event a) (Either () a)
events = arr $ \case events = arr $ \case
@@ -162,51 +333,76 @@ isEvent _ = True
tag :: b -> Event a -> Event b tag :: b -> Event a -> Event b
tag b ev = b <$ ev tag b ev = b <$ ev
switch :: Mealy eff a (b, Event c) -> (c -> Mealy eff a b) -> Mealy eff a b -- 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 -- switch (Mealy st f) s = Mealy st $ \nt t a -> do
Pair (b, ev) f' <- f nt t a -- Pair (b, ev) f' <- f nt t a
case ev of -- case ev of
Tick -> pure (Pair b (switch f' s)) -- Tick -> pure (Pair b (switch f' s))
Event x -> runMealy (s x) nt t a -- Event x -> runMealy (s x) nt t a
sample :: Mealy eff (a, Event b) (Event a) sample :: Mealy eff (a, Event b) (Event a)
sample = arr (uncurry tag) 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 mempty $ \_ _ a ->
let next = f b a
in pure (Pair (extract b) (go next))
preMapAccumRequest :: (Request -> x -> a -> x) -> x -> (x -> b) -> Mealy eff a b preMapAccum' :: forall m x a b. (Monad m, Eq x, Serialize x) => (x -> a -> x) -> x -> (x -> b) -> Auto m a b
preMapAccumRequest f x extract = go x preMapAccum' f x extract = Stateful (Codec get put) (State x False) (\s _req a -> pure $ step s a)
where where
go b = Mealy mempty $ \_ t a -> step :: State x -> a -> (b, State x)
let next = f t b a step s a = let s' = f (state s) a in (extract (state s), State s' (dirty s || state s /= s'))
in pure (Pair (extract b) (go next))
mapAccum :: (x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
mapAccum f x extract = go x
where
go b = Mealy mempty $ \_ _ a ->
let next = f b a
in pure (Pair (extract next) (go next))
mapAccumRequest :: (Request -> x -> a -> x) -> x -> (x -> b) -> Mealy eff a b mapAccum :: (Eq x, Serialize x) => (x -> a -> x) -> x -> (x -> b) -> Mealy eff a b
mapAccumRequest f x extract = go x 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) (State x False) (\s _req a -> pure $ step s a)
where where
go b = Mealy mempty $ \_ t a -> step :: State x -> a -> (b, State x)
let next = f t b a step s a = let s' = f (state s) a in (extract s', State s' (dirty s || state s /= s'))
in pure (Pair (extract next) (go next))
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
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'))
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 data DelayState x a = DelayState
{ pending :: x { pending :: x
, output :: !(Event a) , 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)
delayEvent :: (Eq a, Serialize a) => NominalDiffTime -> Mealy eff (Event a) (Event a)
delayEvent delay = delayEvent delay =
mapAccumRequest step initial output mapAccumRequest step initial output
where where
@@ -218,17 +414,17 @@ delayEvent delay =
queued = queued =
case input of case input of
Tick -> pending st Tick -> pending st
Event x -> pending st ++ [(delay `addUTCTime` now, x)] Event x -> pending st ++ [(SerializeUTCTime $ delay `addUTCTime` now, x)]
in case queued of in case queued of
(due, x) : rest (SerializeUTCTime due, x) : rest
| due <= now -> | due <= now ->
DelayState rest (Event x) DelayState rest (Event x)
_ -> _ ->
DelayState queued Tick DelayState queued Tick
debounce :: NominalDiffTime -> Mealy eff (Event a) (Event a) debounce :: (Serialize a, Eq a) => NominalDiffTime -> Mealy eff (Event a) (Event a)
debounce delay = debounce delay =
mapAccumRequest step initial output mapAccumRequest step initial output
where where
@@ -237,13 +433,13 @@ debounce delay =
let now = requestTime req let now = requestTime req
held = case input of held = case input of
Tick -> pending st Tick -> pending st
Event x -> Just (delay `addUTCTime` now, x) Event x -> Just (SerializeUTCTime $ delay `addUTCTime` now, x)
in case held of in case held of
Just (due, x) Just (SerializeUTCTime due, x)
| due <= now -> DelayState Nothing (Event x) | due <= now -> DelayState Nothing (Event x)
_ -> DelayState held Tick _ -> DelayState held Tick
changes :: Eq a => Mealy eff a (Event a) changes :: (Serialize a, Eq a) => Mealy eff a (Event a)
changes = mapAccum go Nothing (maybe Tick snd) changes = mapAccum go Nothing (maybe Tick snd)
where where
go :: Eq a => Maybe (a, Event a) -> a -> Maybe (a, Event a) go :: Eq a => Maybe (a, Event a) -> a -> Maybe (a, Event a)
@@ -276,86 +472,103 @@ lMerge (Event a) _ = Event a
lMerge Tick (Event a) = Event a lMerge Tick (Event a) = Event a
edge :: Mealy eff Bool (Event ()) edge :: Mealy eff Bool (Event ())
edge = go False edge =
where mapAccum
go True = Mealy mempty $ \_ _ -> \case (\(_, current) new -> (current, new))
True -> pure (Pair Tick (go True)) (False, False)
False -> pure (Pair Tick (go False)) (\(old, current) ->
go False = Mealy mempty $ \_ _ -> \case if not old && current
True -> pure (Pair (Event ()) (go True)) then Event ()
False -> pure (Pair Tick (go False)) else Tick)
-- | Drop the first 'Event' and pass through everything after. Useful for
-- ignoring a self-triggered event (e.g. a service call that changes the
-- very entity the arrow listens to).
dropFirst :: Mealy eff (Event a) (Event a)
dropFirst = go False
where
go seen = Mealy mempty $ \_ _ input ->
case input of
Event _ | not seen -> pure (Pair Tick (go True))
_ -> pure (Pair input (go seen))
duration :: forall eff a. Mealy eff a NominalDiffTime 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 where
go :: Request -> Maybe (UTCTime, NominalDiffTime) -> a -> Maybe (UTCTime, NominalDiffTime) delta :: (SerializeUTCTime, SerializeUTCTime) -> NominalDiffTime
go req Nothing _ = Just (requestTime req, requestTime req `diffUTCTime` requestTime req) delta (SerializeUTCTime start, SerializeUTCTime end) = end `diffUTCTime` start
go req (Just (startTime, _)) _ = Just (startTime, requestTime req `diffUTCTime` startTime) 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 -- | Rollup, hold back bursty messages
-- --
-- Consider a case where you have a bursty set of data. You care to get an immediate response, -- 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 -- but don't want to spam the output.
rollup rollup
:: Int -- ^ How many items to pass through before burst protection :: (Serialize a, Eq a)
-> Int -- ` How many seconds to collect the bursty data => 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]) -> Mealy eff (Event a) (Event [a])
rollup limit seconds = mapAccumRequest go (Left Tick) (either id (\(_, _, _, ev) -> ev)) rollup limit seconds =
mapAccumRequest
go
(Left Tick)
(either id (\(_, _, _, ev) -> ev))
where where
e a = Endo ([a] ++) go
go :: Request -> Either (Event [a]) (UTCTime, Int, Endo [a], Event [a]) -> Event a -> Either (Event [a]) (UTCTime, Int, Endo [a], Event [a]) :: Request
go _ (Left _) Tick = Left Tick -> Either (Event [a]) (SerializeUTCTime, Int, Seq a, Event [a])
go req (Left _) (Event a) = Right (addUTCTime (fromIntegral seconds) (requestTime req), 1, mempty, Event [a]) -> Event a
go req (Right (end, n, acc, _)) Tick -> Either (Event [a]) (SerializeUTCTime, Int, Seq a, Event [a])
| requestTime req >= end = Left (Event $ appEndo acc [])
| otherwise = Right (end, n, acc, Tick) go _ (Left _) Tick =
go req (Right (end, n, acc, _)) (Event a) Left Tick
| requestTime req >= end = Left (Event $ appEndo acc [a])
| n < limit = Right (end, n+1, acc, Event [a]) go req (Left _) (Event a) =
| otherwise = Right (end, n+1, acc <> e a, Tick) 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 window into the events
sliding :: Int -> Mealy eff (Event a) [a] sliding :: (Serialize a, Eq a) => Int -> Mealy eff (Event a) [a]
sliding size = mapAccum go [] id sliding size = mapAccum go [] id
where where
go :: [a] -> Event a -> [a] go :: [a] -> Event a -> [a]
go acc Tick = acc go acc Tick = acc
go acc (Event a) = let xs = acc ++ [a] in drop (max 0 (length xs - size)) xs go acc (Event a) = let xs = acc ++ [a] in drop (max 0 (length xs - size)) xs
fixed :: Int -> Mealy eff (Event a) [a]
fixed seconds = mapAccumRequest go Nothing (maybe [] ((`appEndo` []) . snd))
where
e a = Endo ([a] ++)
go :: Request -> Maybe (UTCTime, Endo [a]) -> Event a -> Maybe (UTCTime, Endo [a])
go req Nothing Tick = Just (addUTCTime (fromIntegral seconds) (requestTime req), mempty)
go req Nothing (Event a) = Just (addUTCTime (fromIntegral seconds) (requestTime req), e a)
go req (Just (end, acc)) ev =
case ev of
Tick | requestTime req >= end -> Just (addUTCTime (fromIntegral seconds) end, mempty)
| otherwise -> Just (end, acc)
Event a | requestTime req >= end -> Just (addUTCTime (fromIntegral seconds) end, e a)
| otherwise -> Just (end, acc <> e a)
currentTime :: Mealy eff a LocalTime currentTime :: Mealy eff a LocalTime
currentTime = Mealy mempty $ \_ Request{requestTime, requestTimeZone} _ -> currentTime = Mealy mempty $ \_nt -> Fun $ \Request{requestTime, requestTimeZone} _ ->
pure (Pair (utcToLocalTime requestTimeZone requestTime) currentTime) utcToLocalTime requestTimeZone requestTime
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 :: Mealy eff a () -> Mealy eff (Event a) ()
onEvent f = events >>> (arr (const ()) ||| f) onEvent f = events >>> (arr (const ()) ||| f)
+11 -7
View File
@@ -27,7 +27,7 @@ module HomeAssistant.Controller
, Light(..) , Light(..)
) where ) where
import AFRP (Mealy (..), Pair (..), eff, Event(..), events, filterA, (>>|), toEvent, Request) import AFRP (Mealy (..), eff, Event(..), events, filterA, (>>|), toEvent, Request)
import Control.Arrow (Arrow(..), returnA) import Control.Arrow (Arrow(..), returnA)
import Control.Category ((>>>)) import Control.Category ((>>>))
import Data.Aeson (Value, object, (.=)) import Data.Aeson (Value, object, (.=))
@@ -37,6 +37,8 @@ import Control.Lens (has, only, (^?), to)
import Data.Aeson.Lens (key, _String, _Integral) import Data.Aeson.Lens (key, _String, _Integral)
import qualified Data.Text.Lens as TL import qualified Data.Text.Lens as TL
import Data.Bool (bool) import Data.Bool (bool)
import Data.Serialize (Serialize)
import GHC.Generics (Generic)
data Target = EntityId !T.Text | AreaId !T.Text data Target = EntityId !T.Text | AreaId !T.Text
deriving (Show,Eq,Ord) deriving (Show,Eq,Ord)
@@ -65,11 +67,11 @@ debug = proc x -> do
returnA -< x returnA -< x
traceEvent :: Show a => HASS (Event a) (Event a) traceEvent :: Show a => HASS (Event a) (Event a)
traceEvent = m traceEvent = proc ev -> do
where case ev of
m = Mealy mempty $ \nt req -> \case Event a -> eff Trace -< a
Event a -> nt (Trace req a) >>= \() -> pure (Pair (Event a) m) Tick -> returnA -< ()
Tick -> pure (Pair Tick m) returnA -< ev
traceValue :: Show a => HASS a a traceValue :: Show a => HASS a a
traceValue = proc x -> do traceValue = proc x -> do
@@ -77,7 +79,9 @@ traceValue = proc x -> do
returnA -< x returnA -< x
data DoorState = Open | Closed data DoorState = Open | Closed
deriving (Show, Eq) deriving (Show, Eq, Generic)
instance Serialize DoorState
data Presence = Occupied | Unoccupied data Presence = Occupied | Unoccupied
deriving (Show, Eq) deriving (Show, Eq)
+5 -1
View File
@@ -13,6 +13,8 @@ import Control.Category ((>>>))
import Data.Aeson (Value) import Data.Aeson (Value)
import HomeAssistant.Controller (entityRead, traceEvent, HASS) import HomeAssistant.Controller (entityRead, traceEvent, HASS)
import Control.Arrow (Arrow(..)) import Control.Arrow (Arrow(..))
import GHC.Generics (Generic)
import Data.Serialize (Serialize)
ruuviTemperatures :: Mealy eff (Event Value) Double ruuviTemperatures :: Mealy eff (Event Value) Double
ruuviTemperatures = entityRead @Double "sensor.ruuvitag_b168_temperature" >>> hold 0 ruuviTemperatures = entityRead @Double "sensor.ruuvitag_b168_temperature" >>> hold 0
@@ -21,7 +23,9 @@ ruuviPressures :: Mealy eff (Event Value) Double
ruuviPressures = entityRead "sensor.ruuvitag_b168_pressure" >>> hold 0 ruuviPressures = entityRead "sensor.ruuvitag_b168_pressure" >>> hold 0
data Ruuvi = Ruuvi { ruuviTemperature :: Double, ruuviPressure :: Double } data Ruuvi = Ruuvi { ruuviTemperature :: Double, ruuviPressure :: Double }
deriving (Show, Eq) deriving (Show, Eq, Generic)
instance Serialize Ruuvi
ruuvi :: Mealy eff (Event Value) (Event Ruuvi) ruuvi :: Mealy eff (Event Value) (Event Ruuvi)
ruuvi = (Ruuvi <$> ruuviTemperatures <*> ruuviPressures) >>> changes ruuvi = (Ruuvi <$> ruuviTemperatures <*> ruuviPressures) >>> changes
+17 -13
View File
@@ -13,7 +13,7 @@ module HomeAssistant.Runtime
, runController , runController
) where ) where
import AFRP (Event (..), Mealy (..), Pair (..), Request (..)) import AFRP (Event (..), Mealy (..), Request (..), Auto, stepAutoSerializing)
import Control.Concurrent.Async (async, waitAny) import Control.Concurrent.Async (async, waitAny)
import Control.Concurrent.STM (atomically, dupTChan, readTChan) import Control.Concurrent.STM (atomically, dupTChan, readTChan)
import Data.Aeson (Value) import Data.Aeson (Value)
@@ -31,18 +31,19 @@ import Data.UUID (UUID, toText)
import qualified Data.UUID.V4 as UUID.V4 import qualified Data.UUID.V4 as UUID.V4
import Katip (runKatipT, logF, sl, Severity (..), ls, Namespace (Namespace), runKatipContextT) import Katip (runKatipT, logF, sl, Severity (..), ls, Namespace (Namespace), runKatipContextT)
import Control.Monad.IO.Class (liftIO, MonadIO) import Control.Monad.IO.Class (liftIO, MonadIO)
import Control.Monad.Fix (MonadFix)
import HomeAssistant.Controller.Ruuvi (ruuviController) import HomeAssistant.Controller.Ruuvi (ruuviController)
import HomeAssistant.Controller.Children (schoolLightController) import HomeAssistant.Controller.Children (schoolLightController)
import Data.Maybe (fromMaybe) import Data.Maybe (fromMaybe)
import qualified System.Metrics import qualified System.Metrics
import qualified HomeAssistant.Runtime.Metrics import qualified HomeAssistant.Runtime.Metrics
import System.FilePath ((</>))
step :: (MonadFix m, MonadIO m) => (forall x. eff x -> m x) -> UUID -> Mealy eff a b -> a -> m (Pair b (Mealy eff a b)) step :: (MonadIO m) => FilePath -> UUID -> Auto m a b -> a -> m (b, Auto m a b)
step nt trace (Mealy _ f) a = do step path trace st a = do
now <- liftIO getCurrentTime now <- liftIO getCurrentTime
tz <- liftIO getCurrentTimeZone tz <- liftIO getCurrentTimeZone
f nt (Request now tz trace) a let req = Request now tz trace
stepAutoSerializing path st req a
data Controller = forall b. Controller T.Text (HASS (Event Value) b) Bool data Controller = forall b. Controller T.Text (HASS (Event Value) b) Bool
@@ -59,17 +60,19 @@ controllers =
-- | Steps the machine for every inbound message; service calls go to the -- | 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 -- bus. A restart re-dups the inbound channel and starts from the machine's
-- initial state; messages broadcast during the restart window are lost. -- initial state; messages broadcast during the restart window are lost.
runController :: Bus -> Controller -> IO Void runController :: FilePath -> Bus -> Controller -> IO Void
runController bus (Controller name machine _enabled) = do runController rootDir bus (Controller name machine _enabled) = do
inbound <- atomically (dupTChan (busInbound bus)) inbound <- atomically (dupTChan (busInbound bus))
go inbound machine let ns = Namespace [name]
let worker = runMealy machine (runKatipContextT (busLogEnv bus) () ns . channelHassEval bus)
let path = rootDir </> T.unpack name
go path inbound worker
where where
go inbound f = do go path inbound f = do
msg <- atomically (readTChan inbound) msg <- atomically (readTChan inbound)
uuid <- UUID.V4.nextRandom uuid <- UUID.V4.nextRandom
let ns = Namespace [name] (_, next) <- step path uuid f msg
Pair _ f' <- step (runKatipContextT (busLogEnv bus) () ns . channelHassEval bus) uuid f msg go path inbound next
go inbound f'
defaultMain :: IO () defaultMain :: IO ()
defaultMain = withSocketsDo $ do defaultMain = withSocketsDo $ do
@@ -77,6 +80,7 @@ defaultMain = withSocketsDo $ do
withBus severity $ \bus -> do withBus severity $ \bus -> do
token <- getEnv "HA_TOKEN" token <- getEnv "HA_TOKEN"
host <- getEnv "HA_HOST" host <- getEnv "HA_HOST"
rootPath <- fromMaybe "/tmp/" <$> lookupEnv "HA_LIB_DIR"
store <- System.Metrics.newStore store <- System.Metrics.newStore
System.Metrics.registerGcMetrics store System.Metrics.registerGcMetrics store
rrdPath <- fromMaybe "hass-controller.rrd" <$> lookupEnv "HA_RRD_PATH" rrdPath <- fromMaybe "hass-controller.rrd" <$> lookupEnv "HA_RRD_PATH"
@@ -86,7 +90,7 @@ defaultMain = withSocketsDo $ do
workers = workers =
[ ("reader", readerAction host 8123 token ents bus) [ ("reader", readerAction host 8123 token ents bus)
, ("writer", writerAction 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)] ++ [("metrics", HomeAssistant.Runtime.Metrics.metricsAction store rrdPath rrdtool)]
as <- mapM (\(name, act) -> async (supervised name defaultBackoff act)) workers as <- mapM (\(name, act) -> async (supervised name defaultBackoff act)) workers
(_, v) <- waitAny as (_, v) <- waitAny as
+5 -3
View File
@@ -17,9 +17,11 @@ import Test.Hspec.Hedgehog (forAll, forAllWith, hedgehog, (===))
-- sides on generated inputs. -- sides on generated inputs.
runPure :: Mealy Identity a b -> [a] -> [b] runPure :: Mealy Identity a b -> [a] -> [b]
runPure _ [] = [] runPure m = go (runMealy m id)
runPure m (a : as) = case runIdentity (runMealy m id fakeRequest a) of where
Pair b m' -> b : runPure m' as 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. -- | 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 :: Gen (Mealy Identity a b) -> PropertyT IO (Mealy Identity a b)
+41 -89
View File
@@ -9,9 +9,10 @@ import Data.Foldable (for_)
import AFRP import AFRP
import Data.Functor.Identity (Identity (..)) import Data.Functor.Identity (Identity (..))
import Data.List (sort) import Data.List (sort)
import Data.Serialize (get, put, runGet, runPut)
import qualified Data.Set as S import qualified Data.Set as S
import qualified Data.Text as T import qualified Data.Text as T
import Data.Time (NominalDiffTime, UTCTime (..), utc) import Data.Time (Day (..), NominalDiffTime, UTCTime (..), picosecondsToDiffTime, utc)
import Data.UUID (nil) import Data.UUID (nil)
import Hedgehog import Hedgehog
import qualified Hedgehog.Gen as Gen import qualified Hedgehog.Gen as Gen
@@ -25,17 +26,24 @@ fakeRequest = Request (sec 0) utc nil
sec :: Integer -> UTCTime sec :: Integer -> UTCTime
sec n = UTCTime (toEnum 0) (fromIntegral n) sec n = UTCTime (toEnum 0) (fromIntegral n)
untime :: SerializeUTCTime -> UTCTime
untime (SerializeUTCTime t) = t
runPure :: Mealy Identity a b -> [a] -> [b] runPure :: Mealy Identity a b -> [a] -> [b]
runPure _ [] = [] runPure m = go (runMealy m id)
runPure m (a : as) = case runIdentity (AFRP.runMealy m id fakeRequest a) of where
Pair b m' -> b : runPure m' as 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). -- | Run a Mealy with a per-step wall clock (seconds since the day-0 epoch).
runTimed :: Mealy Identity a b -> [(Integer, a)] -> [b] runTimed :: Mealy Identity a b -> [(Integer, a)] -> [b]
runTimed _ [] = [] runTimed m = go (runMealy m id)
runTimed m ((s, a) : as) = where
case runIdentity (AFRP.runMealy m id (Request (sec s) utc nil) a) of go _ [] = []
Pair b m' -> b : runTimed m' as 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). -- | A minimal State monad for observing effectful arrows (e.g. whenA gating).
newtype St a = St { unSt :: Int -> (a, Int) } newtype St a = St { unSt :: Int -> (a, Int) }
@@ -54,12 +62,11 @@ instance MonadFix St where
mfix f = St $ \s -> let (a, s') = unSt (f a) s in (a, s') mfix f = St $ \s -> let (a, s') = unSt (f a) s in (a, s')
runStEff :: Mealy St a b -> Int -> [a] -> ([b], Int) runStEff :: Mealy St a b -> Int -> [a] -> ([b], Int)
runStEff m s0 as = go m s0 as runStEff m = go (runMealy m id)
where where
go _ s [] = ([], s) go _ s [] = ([], s)
go m' s (a : rest) = go w s (a : rest) = case unSt (stepAuto w fakeRequest a) s of
case unSt (AFRP.runMealy m' id fakeRequest a) s of ((b, w'), s') -> let (bs, s'') = go w' s' rest in (b : bs, s'')
(Pair b m'', s') -> let (bs, s'') = go m'' s' rest in (b : bs, s'')
spec :: Spec spec :: Spec
spec = describe "AFRP" $ do spec = describe "AFRP" $ do
@@ -72,7 +79,6 @@ spec = describe "AFRP" $ do
lMergeSpec lMergeSpec
changesSpec changesSpec
edgeSpec edgeSpec
dropFirstSpec
filterASpec filterASpec
slidingSpec slidingSpec
mapAccumSpec mapAccumSpec
@@ -81,9 +87,8 @@ spec = describe "AFRP" $ do
delayEventSpec delayEventSpec
debounceSpec debounceSpec
rollupSpec rollupSpec
fixedSpec serializeSpec
effSpec effSpec
switchSpec
mapAccumRequestSpec mapAccumRequestSpec
preMapAccumRequestSpec preMapAccumRequestSpec
whenASpec whenASpec
@@ -211,19 +216,6 @@ edgeSpec = describe "edge" $ do
else o' === Tick else o' === Tick
_ -> failure _ -> failure
dropFirstSpec :: Spec
dropFirstSpec = describe "dropFirst" $ do
it "drops the first Event, passes the rest" $
runPure dropFirst [Tick, Event 1, Event 2, Event 3]
`shouldBe` [Tick, Tick, Event 2, Event 3 :: Event Int]
it "passes Tick through untouched before first Event" $
runPure dropFirst [Tick, Tick, Tick :: Event Int]
`shouldBe` [Tick, Tick, Tick :: Event Int]
it "drops only the first Event, Ticks before it are inert" $
runPure dropFirst [Tick, Tick, Event 'a', Tick, Event 'b']
`shouldBe` [Tick, Tick, Tick, Tick, Event 'b' :: Event Char]
filterASpec :: Spec filterASpec :: Spec
filterASpec = describe "filterA" $ do filterASpec = describe "filterA" $ do
@@ -407,36 +399,6 @@ rollupSpec = describe "rollup" $ do
emitted = concat [xs | Event xs <- out] emitted = concat [xs | Event xs <- out]
sort emitted === sort [x | Event x <- evs] sort emitted === sort [x | Event x <- evs]
fixedSpec :: Spec
fixedSpec = describe "fixed" $ do
it "accumulates events within a window and rolls over on expiry" $
runTimed (fixed 10)
[ (0, Event 'a'), (1, Event 'b'), (2, Tick)
, (12, Event 'c'), (13, Tick)
]
`shouldBe` [ ['a'], ['a', 'b'], ['a', 'b']
, ['c'], ['c']
]
it "starts a window even on a leading Tick" $
runTimed (fixed 10)
[ (0, Tick), (1, Event 'a')
, (12, Tick), (13, Tick)
, (25, Event 'b')
]
`shouldBe` [ [], ['a'], [], [], ['b'] ]
it "output is always the current window's accumulated list" $
hedgehog $ do
w <- forAll $ Gen.int (Range.constant 1 10)
evs <- forAll $ Gen.list (Range.linear 0 30) eventGen
let out = runTimed (fixed w) (zip [0 ..] evs)
expected =
[ [ x | Event x <- take (i - lo + 1) (drop lo evs) ]
| i <- [0 .. length evs - 1]
, let lo = (i `div` w) * w
]
out === expected
eventGen :: Gen (Event Char) eventGen :: Gen (Event Char)
eventGen = Gen.frequency eventGen = Gen.frequency
@@ -444,6 +406,20 @@ eventGen = Gen.frequency
, (1, pure Tick) , (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)
serializeSpec :: Spec
serializeSpec = describe "SerializeUTCTime" $ do
it "get (put x) == pure x" $
hedgehog $ do
x <- forAll timeGen
tripping x (runPut . put) (runGet get)
effSpec :: Spec effSpec :: Spec
effSpec = describe "eff" $ do effSpec = describe "eff" $ do
it "lifts a pure effect function into a stateless Mealy" $ it "lifts a pure effect function into a stateless Mealy" $
@@ -456,41 +432,17 @@ effSpec = describe "eff" $ do
let out = runPure (eff (\_ x -> Identity (x * 2))) xs let out = runPure (eff (\_ x -> Identity (x * 2))) xs
out === map (* 2) xs out === map (* 2) xs
switchSpec :: Spec
switchSpec = describe "switch" $ do
it "switches to the continuation at the first Event" $
runPure (switch (arr (\x -> (x, if x >= (3 :: Int) then Event () else Tick)))
(const (arr (const 99))))
[1, 2, 3, 4, 5]
`shouldBe` [1, 2, 99, 99, 99]
it "never switches if no Event is emitted" $
runPure (switch (arr (\x -> (x, Tick :: Event ())))
(const (arr (const 99))))
[1, 2, 3]
`shouldBe` [1, 2, 3 :: Int]
it "prefix outputs come from the first arrow, suffix from the continuation" $
hedgehog $ do
threshold <- forAll $ Gen.int (Range.linear (-20) 20)
xs <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear (-20) 20))
let firstArr = arr (\x -> (x, if x >= threshold then Event () else Tick))
out = runPure (switch firstArr (const (arr (const 99)))) xs
(pre, _post) = break (>= threshold) xs
take (length pre) out === pre
drop (length pre) out === replicate (length xs - length pre) 99
mapAccumRequestSpec :: Spec mapAccumRequestSpec :: Spec
mapAccumRequestSpec = describe "mapAccumRequest" $ do mapAccumRequestSpec = describe "mapAccumRequest" $ do
it "accumulates request times, post-state extraction" $ it "accumulates request times, post-state extraction" $
runTimed (mapAccumRequest (\req s _ -> s ++ [requestTime req]) [] id) runTimed (mapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
[(0, 'a'), (5, 'b'), (10, 'c')] [(0, 'a'), (5, 'b'), (10, 'c')]
`shouldBe` [ [sec 0], [sec 0, sec 5], [sec 0, sec 5, sec 10] ] `shouldBe` [ [sec 0], [sec 0, sec 5], [sec 0, sec 5, sec 10] ]
it "output i is every request time seen so far" $ it "output i is every request time seen so far" $
hedgehog $ do hedgehog $ do
secs' <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear 0 100)) secs' <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear 0 100))
let out = runTimed (mapAccumRequest (\req s _ -> s ++ [requestTime req]) [] id) let out = runTimed (mapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
[(fromIntegral s, ()) | s <- secs'] [(fromIntegral s, ()) | s <- secs']
expected = [ map (sec . fromIntegral) (take (i + 1) secs') | i <- [0 .. length secs' - 1] ] expected = [ map (sec . fromIntegral) (take (i + 1) secs') | i <- [0 .. length secs' - 1] ]
out === expected out === expected
@@ -498,14 +450,14 @@ mapAccumRequestSpec = describe "mapAccumRequest" $ do
preMapAccumRequestSpec :: Spec preMapAccumRequestSpec :: Spec
preMapAccumRequestSpec = describe "preMapAccumRequest" $ do preMapAccumRequestSpec = describe "preMapAccumRequest" $ do
it "accumulates request times, pre-state extraction" $ it "accumulates request times, pre-state extraction" $
runTimed (preMapAccumRequest (\req s _ -> s ++ [requestTime req]) [] id) runTimed (preMapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
[(0, 'a'), (5, 'b'), (10, 'c')] [(0, 'a'), (5, 'b'), (10, 'c')]
`shouldBe` [ [], [sec 0], [sec 0, sec 5] ] `shouldBe` [ [], [sec 0], [sec 0, sec 5] ]
it "output i is every request time before the current step" $ it "output i is every request time before the current step" $
hedgehog $ do hedgehog $ do
secs' <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear 0 100)) secs' <- forAll $ Gen.list (Range.linear 0 30) (Gen.int (Range.linear 0 100))
let out = runTimed (preMapAccumRequest (\req s _ -> s ++ [requestTime req]) [] id) let out = runTimed (preMapAccumRequest (\req s _ -> s ++ [SerializeUTCTime (requestTime req)]) [] (map untime))
[(fromIntegral s, ()) | s <- secs'] [(fromIntegral s, ()) | s <- secs']
expected = [ map (sec . fromIntegral) (take i secs') | i <- [0 .. length secs' - 1] ] expected = [ map (sec . fromIntegral) (take i secs') | i <- [0 .. length secs' - 1] ]
out === expected out === expected
@@ -569,11 +521,11 @@ sampleSpec = describe "sample" $ do
-- | A stateless arrow carrying a fixed entity set, for testing propagation. -- | A stateless arrow carrying a fixed entity set, for testing propagation.
subscribed :: S.Set T.Text -> Mealy Identity Int Int subscribed :: S.Set T.Text -> Mealy Identity Int Int
subscribed ents = Mealy ents $ \_ _ a -> pure (Pair a (subscribed ents)) subscribed ents = Mealy ents $ \_ -> Fun $ \_ a -> a
-- | Same as 'subscribed' but yields a function, for testing '<*>'. -- | Same as 'subscribed' but yields a function, for testing '<*>'.
subscribedF :: S.Set T.Text -> Mealy Identity Int (Int -> Int) subscribedF :: S.Set T.Text -> Mealy Identity Int (Int -> Int)
subscribedF ents = Mealy ents $ \_ _ a -> pure (Pair (a +) (subscribedF ents)) subscribedF ents = Mealy ents $ \_ -> Fun $ \_ a -> (a +)
entitiesSpec :: Spec entitiesSpec :: Spec
entitiesSpec = describe "entities" $ do entitiesSpec = describe "entities" $ do
@@ -590,7 +542,7 @@ entitiesSpec = describe "entities" $ do
entities (hold 'a') `shouldBe` S.empty entities (hold 'a') `shouldBe` S.empty
entities (changes @Int) `shouldBe` S.empty entities (changes @Int) `shouldBe` S.empty
entities edge `shouldBe` S.empty entities edge `shouldBe` S.empty
entities (sliding (3 :: Int)) `shouldBe` S.empty entities (sliding (3 :: Int) :: Mealy Identity (Event Int) [Int]) `shouldBe` S.empty
it "Category (.) unions entity sets" $ it "Category (.) unions entity sets" $
entities (subscribed (S.singleton "a") >>> subscribed (S.singleton "b")) entities (subscribed (S.singleton "a") >>> subscribed (S.singleton "b"))
+7 -5
View File
@@ -16,7 +16,7 @@ import Data.Aeson (Value, object, (.=))
import qualified Data.Text as T import qualified Data.Text as T
import Data.Time (UTCTime (..), utc) import Data.Time (UTCTime (..), utc)
import Data.UUID (nil) import Data.UUID (nil)
import AFRP (Mealy (..), Pair (..), Request (..)) import AFRP (Mealy (..), Request (..), stepAuto)
import HomeAssistant.Controller (HASSEff (..), Service) import HomeAssistant.Controller (HASSEff (..), Service)
fakeRequest :: Request fakeRequest :: Request
@@ -49,10 +49,12 @@ interp (Trace _ _) = pure ()
-- | Run a HASS arrow over a list of inputs, collecting per-step emitted services. -- | Run a HASS arrow over a list of inputs, collecting per-step emitted services.
runHASS :: Mealy HASSEff a b -> [a] -> [(b, [Service])] runHASS :: Mealy HASSEff a b -> [a] -> [(b, [Service])]
runHASS _ [] = [] runHASS m = go (runMealy m interp)
runHASS m (a : as) = where
case runAcc (runMealy m interp fakeRequest a) [] of go _ [] = []
(Pair b m', svcs) -> (b, svcs) : runHASS m' as go w (a : as) =
case runAcc (stepAuto w fakeRequest a) [] of
((b, w'), svcs) -> (b, svcs) : go w' as
services :: [(b, [Service])] -> [[Service]] services :: [(b, [Service])] -> [[Service]]
services = map snd services = map snd