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 +5/−0
- LICENSE +30/−0
- prodapi-core.cabal +36/−0
- src/Prod/Background.hs +107/−0
- src/Prod/Stepper.hs +145/−0
- src/Prod/Tracer.hs +113/−0
+ 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+