From b20320c779e4c5333797cc75ab6790c0ff910617 Mon Sep 17 00:00:00 2001 From: Mats Rauhala Date: Sat, 12 Sep 2026 20:09:42 +0300 Subject: [PATCH] Explicit state --- home-assistant-controller.cabal | 2 + src/AFRP.hs | 407 ++++++++++++++------------ src/HomeAssistant/Controller.hs | 18 +- src/HomeAssistant/Controller/Ruuvi.hs | 6 +- src/HomeAssistant/Runtime.hs | 30 +- test/AFRPSpec.hs | 45 --- test/Support.hs | 4 +- 7 files changed, 259 insertions(+), 253 deletions(-) diff --git a/home-assistant-controller.cabal b/home-assistant-controller.cabal index 497c0ff..5651e9e 100644 --- a/home-assistant-controller.cabal +++ b/home-assistant-controller.cabal @@ -97,6 +97,8 @@ library , process , directory , cereal + , containers + , filepath -- Directories containing source files. hs-source-dirs: src diff --git a/src/AFRP.hs b/src/AFRP.hs index 4553b7c..d60434c 100644 --- a/src/AFRP.hs +++ b/src/AFRP.hs @@ -3,12 +3,13 @@ module AFRP ( Mealy(..) + , Auto(..) , eff , withEntities , Event(..) , hold , events - , switch + -- , switch , preMapAccum , preMapAccumRequest , mapAccum @@ -23,7 +24,6 @@ module AFRP , Pair(..) , Request(..) , edge - , dropFirst , duration , tag , isEvent @@ -31,28 +31,32 @@ module AFRP , sample , rollup , sliding - , fixed , debounce , currentTime , onEvent , save , load + , stepAuto + , stepAutoSerializing ) where import Control.Category (Category(..), (>>>)) import Prelude hiding ((.), id) 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 Data.Either (fromLeft) import Data.Bool (bool) -import Data.Monoid (Endo(..)) import Data.UUID (UUID) import qualified Data.Set as S import qualified Data.Text as T -import Data.Serialize (Get, Putter, runPut, Serialize (put), runGet) +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) } @@ -99,38 +103,38 @@ data Request = Request data Auto m a b - = Fun (a -> b) -- Stateless variant, needed at least for 'id' - | forall s. Stateful !(Codec s) !(State s) !(State s -> a -> m (b, State s)) -- State is explicitly part of it + = 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 $ f . x - Stateful codec s x -> Stateful codec s $ \s' a -> do - (a',s'') <- x s' a + 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 (const a) + pure a = Fun (\_req -> const a) fa <*> fb = case (fa,fb) of - (Fun af, Fun bf) -> Fun (af <*> bf) + (Fun af, Fun bf) -> Fun $ \req -> (af req <*> bf req) (Stateful codec s af, Fun bf) -> Stateful codec s - (\s' x -> do - let a = bf x - (h, s'') <- af s' x + (\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' x -> do - (a, s'') <- bf s' x - let h = af x + (\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' x -> do - (a, as') <- bf (snd <$> s') x - (h, bs') <- af (fst <$> s') x + (\s' req x -> do + (a, as') <- bf (snd <$> s') req x + (h, bs') <- af (fst <$> s') req x pure (h a, mergeState bs' as') ) @@ -140,58 +144,62 @@ instance (Monad m, Semigroup b) => Semigroup (Auto m a b) where case (fa,fb) of (Fun af, Fun bf) -> Fun (af <> bf) (Stateful codec s af, Fun bf) -> Stateful codec s - (\s' a -> do - (ab, s'') <- af s' a - let bb = bf a + (\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' a -> do - let ab = af a - (bb, s'') <- bf s' a + (\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 a -> do - (ab, as'') <- af (fst <$> s) a - (bb, bs'') <- bf (snd <$> s) a + (\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 + id = Fun $ \_ -> id af . ag = case (af, ag) of - (Fun f, Fun g) -> Fun (f . g) - (Stateful codec s f, Fun g) -> Stateful codec s (\s' -> f s' . g) - (Fun f, Stateful codec s g) -> Stateful codec s (\s' -> fmap (first f) . g s') + (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 a -> do - (b, s') <- g (snd <$> s) a - (c, s'') <- f (fst <$> s) b + 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 = Fun + arr f = Fun $ const f first = \case - Fun f -> Fun $ first f - Stateful codec s f -> Stateful codec s $ \s' (b,d) -> do - (c, s'') <- f s' b + 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 $ + Fun f -> Fun $ \req -> \case - Left b -> Left $ f b + Left b -> Left $ f req b Right d -> Right d - Stateful codec s f -> Stateful codec s $ \s' -> \case + Stateful codec s f -> Stateful codec s $ \s' req -> \case Right d -> pure (Right d, s') Left b -> do - (c, s'') <- f s' b + (c, s'') <- f s' req b pure (Left c, s'') serialize :: Auto m a b -> B.ByteString @@ -240,27 +248,24 @@ load path a = handle defaultOnMissingFile (flip deserialize a <$> B.readFile pat -- trigger subscriptions. data Mealy eff a b = Mealy { entities :: S.Set T.Text - , runMealy :: forall m. Monad 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 - Mealy ast af <> Mealy bst bf = Mealy (ast <> bst) $ \nt r a -> do - Pair x af' <- af nt r a - Pair x' bf' <- bf nt r a - pure (Pair (x <> x') (af' <> bf')) + Mealy ast af <> Mealy bst bf = Mealy (ast <> bst) $ \nt -> af nt <> bf nt instance Monoid b => Monoid (Mealy eff a b) where - mempty = m - where - m = Mealy mempty $ \_ _ _ -> pure (Pair mempty m) + mempty = Mealy mempty $ \_nt -> mempty + -- where + -- m = Mealy mempty $ \_ _ _ -> pure (Pair mempty m) eff :: (Request -> a -> eff b) -> Mealy eff a b -eff f = m - where - m = Mealy mempty $ \nt req x -> - nt (f req x) >>= \b -> pure (Pair b m) +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 @@ -269,53 +274,32 @@ 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 mempty (\_ _ x -> pure (Pair x id)) - (Mealy ast f) . (Mealy bst g) = Mealy (ast <> bst) $ \nt t a -> do - Pair b g' <- g nt t a - Pair c f' <- f nt t b - pure (Pair 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 - where - 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')) + arr f = Mealy mempty $ \_nt -> arr f + first (Mealy st f) = Mealy st $ \nt -> first (f nt) instance ArrowChoice (Mealy eff) where - left (Mealy st f) = lm - 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) + left (Mealy st f) = Mealy st $ \nt -> left (f nt) --- 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 - fmap f (Mealy st g) = Mealy st $ \nt t a -> do - Pair b g' <- g nt t a - pure (Pair (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 mempty $ \_ _ _ -> pure (Pair b (pure b)) - Mealy ast f <*> Mealy bst x = Mealy (ast <> bst) $ \nt t a -> do - b' <- x nt t a - let Pair x' xNext = b' - Pair f' fNext <- f nt t a - pure (Pair (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, Eq, Functor, Foldable, Traversable) + deriving (Show, Eq, Functor, Foldable, Traversable, Generic) + +instance Serialize a => Serialize (Event a) instance Semigroup (Event a) where (<>) = lMerge @@ -323,10 +307,18 @@ instance Semigroup (Event a) where instance Monoid (Event a) where mempty = Tick -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')) +hold :: (Serialize a, Eq a) => a -> Mealy m (Event a) a +hold def = preMapAccum step def id + where + 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 = arr $ \case @@ -340,51 +332,77 @@ 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 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 +-- 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 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 -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) (State x False) (\s _req a -> pure $ step s a) where - go b = Mealy mempty $ \_ t a -> - let next = f t b a - in pure (Pair (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 - 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 -mapAccumRequest 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) (State x False) (\s _req a -> pure $ step s a) where - go b = Mealy mempty $ \_ t a -> - let next = f t b a - in pure (Pair (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')) + +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 { 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 + +-- TODO: Needs '\x -> pure x == get (put x)' test +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 = mapAccumRequest step initial output where @@ -396,17 +414,17 @@ 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 -debounce :: NominalDiffTime -> Mealy eff (Event a) (Event a) +debounce :: (Serialize a, Eq a) => NominalDiffTime -> Mealy eff (Event a) (Event a) debounce delay = mapAccumRequest step initial output where @@ -415,13 +433,13 @@ debounce delay = let now = requestTime req held = case input of Tick -> pending st - Event x -> Just (delay `addUTCTime` now, x) + Event x -> Just (SerializeUTCTime $ delay `addUTCTime` now, x) in case held of - Just (due, x) + Just (SerializeUTCTime due, x) | due <= now -> DelayState Nothing (Event x) _ -> 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) where go :: Eq a => Maybe (a, Event a) -> a -> Maybe (a, Event a) @@ -454,86 +472,103 @@ lMerge (Event a) _ = Event a lMerge Tick (Event a) = Event a edge :: Mealy eff Bool (Event ()) -edge = go False - where - go True = Mealy mempty $ \_ _ -> \case - True -> pure (Pair Tick (go True)) - False -> pure (Pair Tick (go False)) - go False = Mealy mempty $ \_ _ -> \case - True -> pure (Pair (Event ()) (go True)) - False -> pure (Pair Tick (go False)) +edge = + mapAccum + (\(_, current) new -> (current, new)) + (False, False) + (\(old, current) -> + if not old && current + then Event () + 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 = 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 +-- 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 - :: Int -- ^ How many items to pass through before burst protection - -> Int -- ` How many seconds to collect the bursty data + :: (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)) +rollup limit seconds = + mapAccumRequest + go + (Left Tick) + (either id (\(_, _, _, ev) -> ev)) where - e a = Endo ([a] ++) - go :: Request -> Either (Event [a]) (UTCTime, Int, Endo [a], Event [a]) -> Event a -> Either (Event [a]) (UTCTime, Int, Endo [a], Event [a]) - go _ (Left _) Tick = Left Tick - go req (Left _) (Event a) = Right (addUTCTime (fromIntegral seconds) (requestTime req), 1, mempty, Event [a]) - go req (Right (end, n, acc, _)) Tick - | requestTime req >= end = Left (Event $ appEndo acc []) - | otherwise = Right (end, n, acc, Tick) - go req (Right (end, n, acc, _)) (Event a) - | requestTime req >= end = Left (Event $ appEndo acc [a]) - | n < limit = Right (end, n+1, acc, Event [a]) - | otherwise = Right (end, n+1, acc <> e a, Tick) + 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 :: Int -> Mealy eff (Event a) [a] +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 -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 mempty $ \_ Request{requestTime, requestTimeZone} _ -> - pure (Pair (utcToLocalTime requestTimeZone requestTime) currentTime) +currentTime = Mealy mempty $ \_nt -> Fun $ \Request{requestTime, requestTimeZone} _ -> + 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 f = events >>> (arr (const ()) ||| f) diff --git a/src/HomeAssistant/Controller.hs b/src/HomeAssistant/Controller.hs index 8e3acd3..997b328 100644 --- a/src/HomeAssistant/Controller.hs +++ b/src/HomeAssistant/Controller.hs @@ -27,7 +27,7 @@ module HomeAssistant.Controller , Light(..) ) 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.Category ((>>>)) import Data.Aeson (Value, object, (.=)) @@ -37,6 +37,8 @@ import Control.Lens (has, only, (^?), to) 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,Ord) @@ -65,11 +67,11 @@ debug = proc x -> do returnA -< x traceEvent :: Show a => HASS (Event a) (Event a) -traceEvent = m - where - m = Mealy mempty $ \nt req -> \case - Event a -> nt (Trace req a) >>= \() -> pure (Pair (Event a) m) - Tick -> pure (Pair Tick m) +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 @@ -77,7 +79,9 @@ traceValue = proc x -> do returnA -< x data DoorState = Open | Closed - deriving (Show, Eq) + deriving (Show, Eq, Generic) + +instance Serialize DoorState data Presence = Occupied | Unoccupied deriving (Show, Eq) diff --git a/src/HomeAssistant/Controller/Ruuvi.hs b/src/HomeAssistant/Controller/Ruuvi.hs index 30e5bbc..52c8a14 100644 --- a/src/HomeAssistant/Controller/Ruuvi.hs +++ b/src/HomeAssistant/Controller/Ruuvi.hs @@ -13,6 +13,8 @@ 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 @@ -21,7 +23,9 @@ ruuviPressures :: Mealy eff (Event Value) Double ruuviPressures = entityRead "sensor.ruuvitag_b168_pressure" >>> hold 0 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 = (Ruuvi <$> ruuviTemperatures <*> ruuviPressures) >>> changes diff --git a/src/HomeAssistant/Runtime.hs b/src/HomeAssistant/Runtime.hs index 05610d8..c12c893 100644 --- a/src/HomeAssistant/Runtime.hs +++ b/src/HomeAssistant/Runtime.hs @@ -13,7 +13,7 @@ module HomeAssistant.Runtime , runController ) where -import AFRP (Event (..), Mealy (..), Pair (..), Request (..)) +import AFRP (Event (..), Mealy (..), Request (..), Auto, stepAutoSerializing) import Control.Concurrent.Async (async, waitAny) import Control.Concurrent.STM (atomically, dupTChan, readTChan) import Data.Aeson (Value) @@ -31,18 +31,19 @@ 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 (()) -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 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 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 @@ -59,17 +60,19 @@ controllers = -- | 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 worker = runMealy machine (runKatipContextT (busLogEnv bus) () ns . channelHassEval bus) + let path = rootDir T.unpack name + 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] - Pair _ f' <- step (runKatipContextT (busLogEnv bus) () ns . channelHassEval bus) uuid f msg - go inbound f' + (_, next) <- step path uuid f msg + go path inbound next defaultMain :: IO () defaultMain = withSocketsDo $ do @@ -77,6 +80,7 @@ defaultMain = withSocketsDo $ do withBus severity $ \bus -> do token <- getEnv "HA_TOKEN" host <- getEnv "HA_HOST" + rootPath <- fromMaybe "/tmp/" <$> lookupEnv "HA_LIB_DIR" store <- System.Metrics.newStore System.Metrics.registerGcMetrics store rrdPath <- fromMaybe "hass-controller.rrd" <$> lookupEnv "HA_RRD_PATH" @@ -86,7 +90,7 @@ defaultMain = withSocketsDo $ do 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 diff --git a/test/AFRPSpec.hs b/test/AFRPSpec.hs index 0e8f278..c5f9ee5 100644 --- a/test/AFRPSpec.hs +++ b/test/AFRPSpec.hs @@ -72,7 +72,6 @@ spec = describe "AFRP" $ do lMergeSpec changesSpec edgeSpec - dropFirstSpec filterASpec slidingSpec mapAccumSpec @@ -81,7 +80,6 @@ spec = describe "AFRP" $ do delayEventSpec debounceSpec rollupSpec - fixedSpec effSpec switchSpec mapAccumRequestSpec @@ -211,19 +209,6 @@ edgeSpec = describe "edge" $ do else o' === Tick _ -> 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 = describe "filterA" $ do @@ -407,36 +392,6 @@ rollupSpec = describe "rollup" $ do emitted = concat [xs | Event xs <- out] 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.frequency diff --git a/test/Support.hs b/test/Support.hs index 6b5bca3..1f34a25 100644 --- a/test/Support.hs +++ b/test/Support.hs @@ -18,6 +18,7 @@ import Data.Time (UTCTime (..), utc) import Data.UUID (nil) import AFRP (Mealy (..), Pair (..), Request (..)) import HomeAssistant.Controller (HASSEff (..), Service) +import AFRP (stepAuto) fakeRequest :: Request fakeRequest = Request (sec 0) utc nil @@ -51,7 +52,8 @@ interp (Trace _ _) = pure () runHASS :: Mealy HASSEff a b -> [a] -> [(b, [Service])] runHASS _ [] = [] runHASS m (a : as) = - case runAcc (runMealy m interp fakeRequest a) [] of + let w = runMealy m interp + in case runAcc (stepAuto w fakeRequest a) [] of (Pair b m', svcs) -> (b, svcs) : runHASS m' as services :: [(b, [Service])] -> [[Service]]