packages feed

prodapi-core (empty) → 0.1.0.0

raw patch · 6 files changed

+436/−0 lines, 6 filesdep +aesondep +asyncdep +base

Dependencies added: aeson, async, base, bytestring, contravariant, text, time, unbounded-delays, unix-time

Files

+ CHANGELOG.md view
@@ -0,0 +1,5 @@+# Revision history for prodapi++## 0.1.0.0 -- YYYY-mm-dd++* First version. Released on an unsuspecting world.
+ LICENSE view
@@ -0,0 +1,30 @@+Copyright Lucas DiCioccio (c) 2021++All rights reserved.++Redistribution and use in source and binary forms, with or without+modification, are permitted provided that the following conditions are met:++    * Redistributions of source code must retain the above copyright+      notice, this list of conditions and the following disclaimer.++    * Redistributions in binary form must reproduce the above+      copyright notice, this list of conditions and the following+      disclaimer in the documentation and/or other materials provided+      with the distribution.++    * Neither the name of Author name here nor the names of other+      contributors may be used to endorse or promote products derived+      from this software without specific prior written permission.++THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS+"AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT+LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR+A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT+OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL,+SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT+LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE,+DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY+THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT+(INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE+OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
+ prodapi-core.cabal view
@@ -0,0 +1,36 @@+cabal-version:       >=1.10+name:                prodapi-core+version:             0.1.0.0+synopsis:            Core utilities for building Haskell services (tracers, background tasks, steppers).+description:         Core infrastructure components for building production Haskell services.+                     Includes contravariant tracers, background task management, and state-machine steppers.+                     These components have minimal dependencies (no Servant, no Prometheus).+license:             BSD3+license-file:        LICENSE+author:              Lucas DiCioccio+maintainer:          lucas@dicioccio.fr+category:            System+build-type:          Simple+extra-source-files:  CHANGELOG.md++library+  exposed-modules:+      Prod.Tracer+      Prod.Background+      Prod.Stepper+      Paths_prodapi_core+  hs-source-dirs:+      src+  default-extensions: OverloadedStrings DataKinds TypeApplications TypeOperators+  build-depends:+    base >= 4.7 && <5,+    aeson >= 2.2.1 && < 2.3,+    bytestring >= 0.12.1 && < 0.13,+    text >= 2.1.1 && < 2.2,+    contravariant >= 1.5.5 && < 1.6,+    time >= 1.12.2 && < 1.13,+    async >= 2.2.5 && < 2.3,+    unix-time >= 0.4.16 && < 0.5,+    unbounded-delays >= 0.1.1.1 && < 0.2+  default-language:    Haskell2010+
+ src/Prod/Background.hs view
@@ -0,0 +1,107 @@+{-# LANGUAGE DeriveFunctor #-}+{-# LANGUAGE ExistentialQuantification #-}+{-# LANGUAGE NumericUnderscores #-}++module Prod.Background (+    BackgroundVal,+    MicroSeconds,+    background,+    backgroundLoop,+    kill,+    link,+    readBackgroundVal,+    Track (..),+)+where++import Control.Concurrent (threadDelay)+import Control.Concurrent.Async (Async, async, cancel)+import qualified Control.Concurrent.Async as Async+import Control.Monad (forever)+import Control.Monad.IO.Class (MonadIO, liftIO)+import Data.IORef (IORef, atomicModifyIORef', newIORef, readIORef)+import Prod.Tracer (Tracer (..))++import GHC.Stack (CallStack, HasCallStack, callStack)++data Track r+    = Init r+    | RunStart+    | RunDone r r+    | Kill CallStack+    deriving (Show, Functor)++-- | A value that is coupled to an async in charge of updating the value.+data BackgroundVal a+    = forall r.+      BackgroundVal+    { transform :: r -> a+    -- ^ a transformation to apply to the background val, allows to turn the IORef into a functor+    , currentValue :: IORef r+    -- ^ a mutable reference for the latest value+    , backgroundTask :: Async ()+    -- ^ a background task responsible for updating the value, implementations should guarantee that once the Aync () is cancelled, currentValue is never updated+    , tracer :: Tracer IO (Track r)+    }++instance Functor BackgroundVal where+    fmap g (BackgroundVal f ioref task tracer) =+        BackgroundVal (g . f) ioref task tracer++-- | Starts a background task continuously updating a value.+background ::+    Tracer IO (Track a) ->+    -- | initial state+    b ->+    -- | initial value+    a ->+    -- | state-influenced task+    (b -> IO (a, b)) ->+    IO (BackgroundVal a)+background tracer initState defaultValue task = do+    trace (Init defaultValue)+    ref <- newIORef defaultValue+    BackgroundVal id ref <$> async (loop ref initState) <*> pure tracer+  where+    trace = runTracer tracer+    loop ref st0 = do+        trace (RunStart)+        (newVal, st1) <- task st0+        oldVal <- seq newVal $ atomicModifyIORef' ref (\old -> (newVal, old))+        trace (RunDone oldVal newVal)+        seq st1 $ loop ref st1++-- | Fantom type for annotating Int.+type MicroSeconds n = n++{- | Starts a background task continuously updating a value at a periodic interval.+This is implemented by interspersing a threadDelay before the task and calling background and hiding the 'state-passing' arguments.+-}+backgroundLoop ::+    Tracer IO (Track a) ->+    -- | initial value+    a ->+    -- | periodic task+    IO a ->+    -- | wait period between two executions+    MicroSeconds Int ->+    IO (BackgroundVal a)+backgroundLoop tracer defaultValue task usecs = do+    background tracer () defaultValue (\_ -> threadDelay usecs >> fmap adapt task)+  where+    adapt x = (x, ())++-- | Kills the watchdog by killing the underlying async.+readBackgroundVal :: (MonadIO m) => BackgroundVal a -> m a+readBackgroundVal (BackgroundVal f ioref _ _) =+    fmap f $ liftIO $ readIORef ioref++-- | Kills the watchdog by killing the underlying async.+kill :: (HasCallStack, MonadIO m) => BackgroundVal a -> m ()+kill bkg@(BackgroundVal _ _ _ tracer) = liftIO $ do+    runTracer tracer (Kill callStack)+    liftIO $ cancel . backgroundTask $ bkg++link :: BackgroundVal a -> BackgroundVal b -> IO ()+link b1 b2 = Async.link2 (backgroundTask b1) (backgroundTask b2)+
+ src/Prod/Stepper.hs view
@@ -0,0 +1,145 @@+{-# LANGUAGE DeriveFunctor #-}+{-# LANGUAGE DeriveGeneric #-}++-- | A simplistic stepper for delayable singly-threaded state machines.+module Prod.Stepper where++import qualified Control.Concurrent.Thread.Delay as ThreadDelay+import Control.Monad (when)+import Data.Aeson (FromJSON, ToJSON)+import Data.Coerce (coerce)+import Data.Int (Int64)+import Data.Time.Clock (UTCTime, diffUTCTime, getCurrentTime, nominalDiffTimeToSeconds)+import Data.UnixTime+import Foreign.C.Types (CTime (..))+import GHC.Generics (Generic)++import Prod.Tracer (Tracer (..), runTracer)++-------------------------------------------------------------------------------+data DelaySpec+    = DelaySafetyAmount+    | DelayUntil UTCTime+    | DelayUntilEpochInteger Int64+    deriving (Show, Eq, Ord, Generic)++instance ToJSON DelaySpec+instance FromJSON DelaySpec++type SafetySeconds = Int64++waitUntil :: SafetySeconds -> DelaySpec -> IO ()+waitUntil _ (DelayUntilEpochInteger time) = do+    now <- getUnixTime+    let diff = UnixTime (coerce time) 0 `diffUnixTime` now+    let delayV = udtSeconds diff + 10+    when (delayV > 0) $+        ThreadDelay.delay (toInteger (coerce $ delayV * 1000000 :: Int64))+waitUntil _ (DelayUntil time) = do+    now <- getCurrentTime+    let diff = time `diffUTCTime` now+    let delayV = round (nominalDiffTimeToSeconds diff + 10) :: Int64+    when (delayV > 0) $+        ThreadDelay.delay (toInteger (coerce $ delayV * 1000000 :: Int64))+waitUntil amount DelaySafetyAmount =+    when (amount > 0) $+        ThreadDelay.delay (toInteger (coerce $ amount * 1000000 :: Int64))++isExpired :: DelaySpec -> IO Bool+isExpired (DelayUntilEpochInteger time) = do+    now <- getUnixTime+    let diff = UnixTime (coerce time) 0 `diffUnixTime` now+    let delayV = udtSeconds diff + 10+    pure (delayV <= 0)+isExpired (DelayUntil time) = do+    now <- getCurrentTime+    let diff = time `diffUTCTime` now+    let delayV = round (nominalDiffTimeToSeconds diff + 10) :: Int64+    pure (delayV <= 0)+isExpired DelaySafetyAmount =+    pure False++-------------------------------------------------------------------------------+data Trace input output+    = Start input+    | Suspend DelaySpec input+    | Finish input+    | Delaying DelaySpec input output+    | Inlining input output+    deriving (Show, Eq, Ord, Generic)++instance (ToJSON i, ToJSON o) => ToJSON (Trace i o)+instance (FromJSON i, FromJSON o) => FromJSON (Trace i o)++-------------------------------------------------------------------------------+type BaseStepIO a b = (b -> IO ()) -> a -> IO ()++mapInput :: (c -> a) -> BaseStepIO a b -> BaseStepIO c b+mapInput f g = \consumeA cVal -> g consumeA (f cVal)++mapOutput :: (b -> z) -> BaseStepIO a b -> BaseStepIO a z+mapOutput f g = \consumeZ bVal -> g (consumeZ . f) bVal++-------------------------------------------------------------------------------+data Delayable a = Inline a | Delay DelaySpec a+    deriving (Show, Functor)++-- | Turn a pair of next function into a handler of Delayable.+nextStep :: (DelaySpec -> t1 -> t2) -> (t1 -> t2) -> Delayable t1 -> t2+nextStep _ inlineF (Inline val) = inlineF val+nextStep delayF _ (Delay d val) = delayF d val++-------------------------------------------------------------------------------+type StepIO a b = BaseStepIO (Delayable a) (Delayable b)++-- special case for a StepIO that will start immediately (removing the need to handle the Delayable branch)+type StartStepIO a b = BaseStepIO a (Delayable b)++-------------------------------------------------------------------------------+data ExecFunctions arg = ExecFunctions+    { inline :: arg -> IO ()+    , delay :: DelaySpec -> arg -> IO ()+    }++defineExecution ::+    -- | tracer to save the state+    Tracer IO (Trace z1 z2) ->+    -- | function to turn an input into a state we hope to resume+    (a -> z1) ->+    (b -> z2) ->+    -- | handler-defining function with a bag of delay/inline functions and a start value+    (ExecFunctions b -> a -> IO ()) ->+    StepIO a b+defineExecution tracer f g runStep = \next input -> do+    -- callback given to downstream consumer: so that we can delay+    let delayF val delaySpec nextObj = do+            let mappedOutput = g nextObj+            runTracer tracer (Delaying delaySpec val mappedOutput)+            next (Delay delaySpec nextObj)++    -- callback given to downstream consumer: so that we can inline+    let inlineF val nextObj = do+            let mappedOutput = g nextObj+            runTracer tracer (Inlining val mappedOutput)+            next (Inline nextObj)++    -- hanling of input+    case input of+        Inline val -> do+            let mappedInput = f val+            runTracer tracer (Start mappedInput)+            runStep (ExecFunctions (inlineF mappedInput) (delayF mappedInput)) val+            runTracer tracer (Finish mappedInput)+        Delay delaySpec val -> do+            let mappedInput = f val+            runTracer tracer (Suspend delaySpec mappedInput)++onSuspend :: (DelaySpec -> a -> IO ()) -> Tracer IO (Trace a b)+onSuspend action =+    Tracer go+  where+    go (Suspend delaySpec obj) = do+        action delaySpec obj+    go _ = do+        pure ()+
+ src/Prod/Tracer.hs view
@@ -0,0 +1,113 @@+-- https://www.youtube.com/watch?v=qzOQOmmkKEM&feature=emb_logo++module Prod.Tracer (+    Tracer (..),+    silent,+    traceIf,+    traceBoth,++    -- * common utilities+    tracePrint,+    traceHPrint,+    traceHPut,+    encodeJSON,+    pulls,++    -- * re-exports+    Contravariant (..),+    Divisible (..),+    Decidable (..),+) where++import Control.Monad ((>=>))+import Control.Monad.IO.Class (MonadIO, liftIO)+import Data.Aeson (ToJSON, encode)+import Data.ByteString.Lazy (ByteString, hPut)+import Data.Functor.Contravariant+import Data.Functor.Contravariant.Divisible+import System.IO (Handle, hPrint)++newtype Tracer m a = Tracer {runTracer :: (a -> m ())}++instance Contravariant (Tracer m) where+    contramap f (Tracer g) = Tracer (g . f)++instance (Applicative m) => Divisible (Tracer m) where+    conquer = silent+    divide = traceSplit++instance (Applicative m) => Decidable (Tracer m) where+    lose _ = silent+    choose = tracePick++-- | Disable Tracing.+{-# INLINE silent #-}+silent :: (Applicative m) => Tracer m a+silent = Tracer (const $ pure ())++{- | Splits a tracer into two chunks that are run sequentially.++This name can be confusing but it has to be thought backwards for Contravariant logging:+We compose a target tracer from two tracers but we split the content of the trace.++Note that the split function may actually duplicate inputs (that's how traceBoth works).+-}+{-# INLINEABLE traceSplit #-}+traceSplit :: (Applicative m) => (c -> (a, b)) -> Tracer m a -> Tracer m b -> Tracer m c+traceSplit split (Tracer f1) (Tracer f2) = Tracer (go . split)+  where+    go (b, c) = f1 b *> f2 c++{- | If you are given two tracers and want to pass both.+Composition occurs in sequence.+-}+{-# INLINEABLE traceBoth #-}+traceBoth :: (Applicative m) => Tracer m a -> Tracer m a -> Tracer m a+traceBoth t1 t2 = traceSplit (\x -> (x, x)) t1 t2++{- | Picks a tracer based on the emitted object.+Example logic that can be built is traceIf that silent messages.+-}+{-# INLINEABLE tracePick #-}+tracePick :: (Applicative m) => (c -> Either a b) -> Tracer m a -> Tracer m b -> Tracer m c+tracePick split (Tracer f1) (Tracer f2) = Tracer $ \a ->+    let e = split a+     in either f1 f2 e++-- | Filter by dynamically testing values.+{-# INLINEABLE traceIf #-}+traceIf :: (Applicative m) => (a -> Bool) -> Tracer m a -> Tracer m a+traceIf predicate t = tracePick (\x -> if predicate x then Left () else Right x) silent t++-- | A tracer that prints emitted events.+tracePrint :: (MonadIO m, Show a) => Tracer m a+tracePrint = Tracer (liftIO . print)++-- | A tracer that prints emitted to some handle.+traceHPrint :: (MonadIO m, Show a) => Handle -> Tracer m a+traceHPrint handle = Tracer (liftIO . hPrint handle)++-- | A tracer that puts some ByteString to some handle.+traceHPut :: (MonadIO m) => Handle -> Tracer m ByteString+traceHPut handle = Tracer (liftIO . hPut handle)++-- | A conversion encoding values to JSON.+{-# INLINE encodeJSON #-}+encodeJSON :: (ToJSON a) => Tracer m ByteString -> Tracer m a+encodeJSON = contramap encode++{- | Pulls a value to complete a trace when a trace occurs.++This function allows to combines pushed values with pulled values.  Hence,+performing some scheduling between behaviours.+Typical usage would be to annotate a trace with a background value, or perform+data augmentation in a pipelines of traces.++Note that if you rely on this function you need to pay attention of the+blocking effect of 'pulls': the traced value c is not forwarded until a+value b is available.+-}+{-# INLINE pulls #-}+pulls :: (Monad m) => (c -> m b) -> Tracer m b -> Tracer m c+pulls act (Tracer f1) = Tracer $ act >=> f1+