dynamic-pipeline (empty) → 0.1.0.0
raw patch · 13 files changed
+1813/−0 lines, 13 filesdep +HListdep +QuickCheckdep +asyncsetup-changed
Dependencies added: HList, QuickCheck, async, base, bytestring, compdata, dynamic-pipeline, hspec, lens, optparse-applicative, relude, trifecta, unagi-chan
Files
- LICENSE +30/−0
- README.md +45/−0
- Setup.hs +2/−0
- dynamic-pipeline.cabal +226/−0
- examples/Graph/ConnComp.hs +69/−0
- examples/Graph/ConnectedComp.hs +80/−0
- examples/Main.hs +50/−0
- examples/Misc/RepeatedDP.hs +45/−0
- src/DynamicPipeline.hs +233/−0
- src/DynamicPipeline/Channel.hs +135/−0
- src/DynamicPipeline/Flow.hs +336/−0
- src/DynamicPipeline/Stage.hs +561/−0
- test/Spec.hs +1/−0
+ LICENSE view
@@ -0,0 +1,30 @@+Copyright Juan Pablo Royo Sales (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 Juan Pablo Royo Sales 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.
+ README.md view
@@ -0,0 +1,45 @@+# DynamicPipeline+++[](https://hackage.haskell.org/package/dynamic-pipeline)+++**DynamicPipeline** is a *Type Safe* Dynamic and Parallel Streaming Library, which is an implementation of **Dynamic Pipeline Paradigm (DPP)**+proposed in this paper [DPP](https://biblioteca.sistedes.es/articulo/the-dynamic-pipeline-paradigm/).++The aim of this Library is to provide all the **Type level** constructs to guide the user in building a *DPP* flow to solve any algorithm that fits on this computational model. + +This implementation has been developed using *Type Level Programming* techniques like `Type families`, `Defunctionalization`, `Existential Types` and `Dynamic Record Tagged Types` among others.++Using all this techniques, we provide a *High Level and Type Safe* **DPP** Library to build a Data Flow Algorithm avoiding as much as possible boilerplate code, but maintaining safety and robustness.++## Prerequisites++This project has been built and developed with the following versions:++- [Stack](https://docs.haskellstack.org/en/stable/README/) `2.1.3` or higher+- Stack LTS `17.7`+- [GHC](https://www.haskell.org/ghc/) `8.10.3` ++## Cloning the project ++```bash+git clone https://github.com/jproyo/dynamic-pipeline+```++## Issues ++https://github.com/jproyo/dynamic-pipeline/issues++## Pull Request++https://github.com/jproyo/dynamic-pipeline/pulls++## Licensing++`dynamic-pipeline` is an open source project available under a liberal [BSD-3-Clause license](./LICENSE)+++++
+ Setup.hs view
@@ -0,0 +1,2 @@+import Distribution.Simple+main = defaultMain
+ dynamic-pipeline.cabal view
@@ -0,0 +1,226 @@+cabal-version: 1.12++-- This file has been generated from package.yaml by hpack version 0.34.4.+--+-- see: https://github.com/sol/hpack+--+-- hash: d407255a4c8f1348472a5c3ffd16fa33771e9f71937861df2783d3045db1b031++name: dynamic-pipeline+version: 0.1.0.0+synopsis: Library Type Safe implementation of Dynamic Pipeline Paradigm (DPP).+description: @dynamic-pipeline@ is a __/Type Safe/__ Dynamic and Parallel Streaming Library, which is an implementation of __Dynamic Pipeline Paradigm (DPP)__ + proposed in this paper [DPP](https://biblioteca.sistedes.es/articulo/the-dynamic-pipeline-paradigm/).+ .+ The aim of this tool is to provide all the __Type level__ constructs to guide the user in building a /DPP/ flow to solve any algorithm that fits on + this computational model. + .+ The idea behind the implementation is similar to other __Type level__ programming libraries like [@servant@](https://hackage.haskell.org/package/servant)+ . + Some of the feature of this library includes:+ .+ * /Small Number of Dependencies/: Dependencies have been kept to minimum in order to reduce compilation time.+ .+ * /Type Level Encoding/: As it was explained before, there is an intensive use of Type Level techniques in the Library in order to guide the user in + defining the Flow Model and fill the gaps with the combinators provided and the Compiler indications. On the other hand some /Type Custom Errors/ have been+ encoded in the Library in order to provide meaningful messages about the Data Flow Grammar Defined in /DPP/.+ .+ * /Stream Parallelization/: __DPP__ model follows a /Parallelization Pipeline/ approach with the use of intermediate Channels. All the burden have been hidden + for the user in order to reduce complexity and focus on the algorithm problem implementations.+ .+ * /Safe Scoped Execution/: Use of Existential to define the Scope execution of /DPP/ reducing even more the complexity to the user.+ .+ Import "DynamicPipeline" for using the library.+ .+ * /Examples/: <src/examples examples directory> in the package+category: Data,Data Flow,Pipes,Streaming,Parallel,Type-safe+homepage: https://github.com/jproyo/dynamic-pipeline#readme+bug-reports: https://github.com/jproyo/dynamic-pipeline/issues+author: Juan Pablo Royo Sales+maintainer: juanpablo.royo@gmail.com+copyright: Copyright (C) 2021 Juan Pablo Royo Sales+license: BSD3+license-file: LICENSE+build-type: Simple+extra-source-files:+ README.md++source-repository head+ type: git+ location: https://github.com/jproyo/dynamic-pipeline++library+ exposed-modules:+ DynamicPipeline+ other-modules:+ DynamicPipeline.Channel+ DynamicPipeline.Flow+ DynamicPipeline.Stage+ hs-source-dirs:+ src+ default-extensions:+ ConstraintKinds+ DataKinds+ DefaultSignatures+ DeriveAnyClass+ DeriveFoldable+ DeriveFunctor+ DeriveGeneric+ DeriveLift+ DeriveTraversable+ DerivingStrategies+ DerivingVia+ EmptyCase+ FlexibleContexts+ FlexibleInstances+ FunctionalDependencies+ GADTs+ GeneralizedNewtypeDeriving+ LambdaCase+ MultiParamTypeClasses+ MultiWayIf+ NoImplicitPrelude+ OverloadedStrings+ PackageImports+ PartialTypeSignatures+ PolyKinds+ Rank2Types+ RankNTypes+ RecordWildCards+ ScopedTypeVariables+ StandaloneDeriving+ TemplateHaskell+ TupleSections+ TypeApplications+ TypeOperators+ TypeFamilies+ ghc-options: -Wall -fno-warn-partial-type-signatures -fconstraint-solver-iterations=0 -fspecialise-aggressively -fexpose-all-unfoldings -flate-specialise -Wincomplete-record-updates -Wincomplete-uni-patterns -Wredundant-constraints -Wno-unused-top-binds -Wno-unused-foralls -fno-warn-deprecations -fno-warn-orphans+ build-depends:+ HList >=0.5.0.0+ , async >=2.1.0 && <=2.2.3+ , base >=4.7 && <5+ , bytestring >=0.10.9.0 && <=0.11.1.0+ , compdata >=0.10 && <=0.12.1+ , lens >=4.19 && <=5.0.1+ , relude >=0.7.0.0 && <=1.0.0.1+ , unagi-chan >=0.4.1.0+ default-language: Haskell2010++executable examples+ main-is: Main.hs+ other-modules:+ Graph.ConnComp+ Graph.ConnectedComp+ Misc.RepeatedDP+ hs-source-dirs:+ examples+ default-extensions:+ ConstraintKinds+ DataKinds+ DefaultSignatures+ DeriveAnyClass+ DeriveFoldable+ DeriveFunctor+ DeriveGeneric+ DeriveLift+ DeriveTraversable+ DerivingStrategies+ DerivingVia+ EmptyCase+ FlexibleContexts+ FlexibleInstances+ FunctionalDependencies+ GADTs+ GeneralizedNewtypeDeriving+ LambdaCase+ MultiParamTypeClasses+ MultiWayIf+ NoImplicitPrelude+ OverloadedStrings+ PackageImports+ PartialTypeSignatures+ PolyKinds+ Rank2Types+ RankNTypes+ RecordWildCards+ ScopedTypeVariables+ StandaloneDeriving+ TemplateHaskell+ TupleSections+ TypeApplications+ TypeOperators+ TypeFamilies+ ghc-options: -Wall -fno-warn-partial-type-signatures -fconstraint-solver-iterations=0 -fspecialise-aggressively -fexpose-all-unfoldings -flate-specialise -Wincomplete-record-updates -Wincomplete-uni-patterns -Wredundant-constraints -Wno-unused-top-binds -Wno-unused-foralls -fno-warn-deprecations -fno-warn-orphans -threaded -O3 -rtsopts -with-rtsopts=-N+ build-depends:+ HList+ , async+ , base >=4.7 && <5+ , bytestring >=0.10.9.0 && <=0.11.1.0+ , compdata >=0.10 && <=0.12.1+ , dynamic-pipeline+ , lens >=4.19 && <=5.0.1+ , optparse-applicative+ , relude >=0.7.0.0 && <=1.0.0.1+ , trifecta+ , unagi-chan >=0.4.1.0+ default-language: Haskell2010++test-suite conn-comp-tests+ type: exitcode-stdio-1.0+ main-is: Spec.hs+ other-modules:+ Paths_dynamic_pipeline+ hs-source-dirs:+ test+ default-extensions:+ ConstraintKinds+ DataKinds+ DefaultSignatures+ DeriveAnyClass+ DeriveFoldable+ DeriveFunctor+ DeriveGeneric+ DeriveLift+ DeriveTraversable+ DerivingStrategies+ DerivingVia+ EmptyCase+ FlexibleContexts+ FlexibleInstances+ FunctionalDependencies+ GADTs+ GeneralizedNewtypeDeriving+ LambdaCase+ MultiParamTypeClasses+ MultiWayIf+ NoImplicitPrelude+ OverloadedStrings+ PackageImports+ PartialTypeSignatures+ PolyKinds+ Rank2Types+ RankNTypes+ RecordWildCards+ ScopedTypeVariables+ StandaloneDeriving+ TemplateHaskell+ TupleSections+ TypeApplications+ TypeOperators+ TypeFamilies+ ghc-options: -Wall -fno-warn-partial-type-signatures -fconstraint-solver-iterations=0 -fspecialise-aggressively -fexpose-all-unfoldings -flate-specialise -Wincomplete-record-updates -Wincomplete-uni-patterns -Wredundant-constraints -Wno-unused-top-binds -Wno-unused-foralls -fno-warn-deprecations -fno-warn-orphans -threaded -O2 -rtsopts -with-rtsopts=-N -Wno-unused-local-binds -Wno-unused-matches+ build-tool-depends:+ hspec-discover:hspec-discover+ build-depends:+ HList >=0.5.0.0+ , QuickCheck+ , async >=2.1.0 && <=2.2.3+ , base >=4.7 && <5+ , bytestring >=0.10.9.0 && <=0.11.1.0+ , compdata >=0.10 && <=0.12.1+ , dynamic-pipeline+ , hspec ==2.*+ , lens >=4.19 && <=5.0.1+ , relude >=1.0.0.0+ , unagi-chan >=0.4.1.0+ default-language: Haskell2010
+ examples/Graph/ConnComp.hs view
@@ -0,0 +1,69 @@+-- |+-- Module : Graph.ConnComp+-- Copyright : (c) 2021 Juan Pablo Royo Sales+--+-- License : BSD3+-- Maintainer : juanpablo.royo@gmail.com+-- Stability : experimental+-- Portability : GHC++module Graph.ConnComp+ ( Edge(..)+ , ConnectedComponents+ , toEdge+ , toConnectedComp+ , isIncident+ , addToConnectedComp+ , includedIncident+ , Graph.ConnComp.null+ , intersect+ , member+ ) where++import qualified Data.IntSet as S+import Relude+import Text.Trifecta+import Text.Trifecta.Parser as P++newtype Edge = Edge (Int, Int)+ deriving (Show, Eq, Ord)++nonEdge :: Edge+nonEdge = Edge (-1, -1)++newtype ConnectedComponents = ConnectedComponents IntSet+ deriving newtype (Monoid, Semigroup, Show, Eq)++toEdge :: String -> Edge+toEdge = foldResult (const nonEdge) identity . toEdge'++toEdge' :: String -> Text.Trifecta.Result Edge+toEdge' = P.parseString parseEdge mempty++parseInt :: Parser Int+parseInt = fromInteger <$> integer ++parseEdge :: Parser Edge+parseEdge = fmap Edge . (,) <$> (whiteSpace *> parseInt <* whiteSpace) <*> parseInt++toConnectedComp :: Edge -> ConnectedComponents+toConnectedComp (Edge (a, b)) = ConnectedComponents (a `S.insert` S.singleton b)++member :: Int -> ConnectedComponents -> Bool +member x (ConnectedComponents cc) = x `S.member`cc++addToConnectedComp :: Edge -> ConnectedComponents -> ConnectedComponents+addToConnectedComp (Edge (a, b)) (ConnectedComponents set) = ConnectedComponents $ a `S.insert` (b `S.insert` set)++isIncident :: Edge -> Edge -> Bool+isIncident (Edge (a, b)) (Edge (c, d)) = a == c || b == c || a == d || b == d++includedIncident :: Edge -> ConnectedComponents -> Bool+includedIncident (Edge (a, b)) (ConnectedComponents set) = S.member a set || S.member b set++null :: ConnectedComponents -> Bool+null (ConnectedComponents s) = S.null s++intersect :: ConnectedComponents -> ConnectedComponents -> Bool+intersect (ConnectedComponents s1) (ConnectedComponents s2) = not $ S.null $ S.intersection s1 s2 +
+ examples/Graph/ConnectedComp.hs view
@@ -0,0 +1,80 @@+-- |+-- Module : Graph.ConnectedComp+-- Copyright : (c) 2021 Juan Pablo Royo Sales+--+-- License : BSD3+-- Maintainer : juanpablo.royo@gmail.com+-- Stability : experimental+-- Portability : GHC+module Graph.ConnectedComp where++import DynamicPipeline+import Graph.ConnComp+import Relude++-- brittany-disable-next-binding+type DPConnComp = Source (Channel (Edge :<+> ConnectedComponents :<+> Eof))+ :>> Generator (Channel (Edge :<+> ConnectedComponents :<+> Eof))+ :>> Sink++source' :: FilePath+ -> Stage+ (WriteChannel Edge -> WriteChannel ConnectedComponents -> DP st ())+source' filePath = withSource @DPConnComp+ $ \edgeOut _ -> unfoldFile filePath edgeOut (toEdge . decodeUtf8)++sink' :: Stage (ReadChannel Edge -> ReadChannel ConnectedComponents -> DP st ())+sink' = withSink @DPConnComp $ \_ cc -> withDP $ foldM cc print++generator' :: GeneratorStage DPConnComp ConnectedComponents Edge st+generator' =+ let gen = withGenerator @DPConnComp genAction+ in mkGenerator gen filterTemplate++filterTemplate :: Filter DPConnComp ConnectedComponents Edge st+filterTemplate = actor actor1 |>> actor actor2++actor1 :: Edge+ -> ReadChannel Edge+ -> ReadChannel ConnectedComponents+ -> WriteChannel Edge+ -> WriteChannel ConnectedComponents+ -> StateT ConnectedComponents (DP st) ()+actor1 _ readEdge _ writeEdge _ = + foldM readEdge $ \e -> get >>= doActor e+ where+ doActor v conn+ | toConnectedComp v `intersect` conn = modify (toConnectedComp v <>)+ | otherwise = push v writeEdge++actor2 :: Edge+ -> ReadChannel Edge+ -> ReadChannel ConnectedComponents+ -> WriteChannel Edge+ -> WriteChannel ConnectedComponents+ -> StateT ConnectedComponents (DP st) ()+actor2 _ _ readCC _ writeCC = do + foldM' readCC pushMemory $ \e -> get >>= doActor e++ where+ pushMemory = get >>= flip push writeCC++ doActor cc conn+ | cc `intersect` conn = modify (cc <>)+ | otherwise = push cc writeCC+++genAction :: Filter DPConnComp ConnectedComponents Edge st+ -> ReadChannel Edge+ -> ReadChannel ConnectedComponents+ -> WriteChannel Edge+ -> WriteChannel ConnectedComponents+ -> DP st ()+genAction filter' readEdge readCC _ writeCC = do+ let unfoldFilter = mkUnfoldFilterForAll filter' toConnectedComp readEdge (readCC .*. HNil) + results <- unfoldF unfoldFilter+ foldM (hHead results) (`push` writeCC)+++program :: FilePath -> IO ()+program file = runDP $ mkDP @DPConnComp (source' file) generator' sink'
+ examples/Main.hs view
@@ -0,0 +1,50 @@+module Main where++import Graph.ConnectedComp as CC+import Misc.RepeatedDP as Repeated+import Options.Applicative as Opt+import Relude+++data ProgramOptions = RepeatedElements+ | ConnectedComponents FilePath+ deriving(Show)+++programDesc :: ParserInfo ProgramOptions+programDesc = info (everyProgram <**> helper) (fullDesc <> header "Examples on dynamic-pipeline library")+ where+ everyProgram =+ subparser (command "repeated-elements" subCmdRepeated <> command "connected-components" subCmdConnectedComp)++subCmdRepeated :: ParserInfo ProgramOptions+subCmdRepeated = info+ (pure RepeatedElements)+ (fullDesc <> header "repeated-elements - Dynamic Pipeline Examples" <> progDesc+ "Given a list of 2000 repeated Integers filter and output unique 1000 integers"+ )+++subCmdConnectedComp :: ParserInfo ProgramOptions+subCmdConnectedComp = info+ (helper <*> toplevelOpt)+ ( fullDesc+ <> header "connected-components - Dynamic Pipeline Examples"+ <> progDesc+ "Given a file path passed by parameters to this command argument with a list of edges, calculate and output the Sets of Connected Components"+ )+ where toplevelOpt = ConnectedComponents <$> fileOption+++fileOption :: Parser FilePath+fileOption =+ Opt.strOption (Opt.long "filepath" <> Opt.short 'f' <> Opt.metavar "FilePath" <> Opt.help "File with Set of Edges")++main :: IO ()+main = execParser programDesc >>= main'++main' :: ProgramOptions -> IO ()+main' = \case+ RepeatedElements -> Repeated.program+ ConnectedComponents file -> CC.program file+
+ examples/Misc/RepeatedDP.hs view
@@ -0,0 +1,45 @@+-- |+-- Module : Misc.RepeatedDP+-- Copyright : (c) 2021 Juan Pablo Royo Sales+--+-- License : BSD3+-- Maintainer : juanpablo.royo@gmail.com+-- Stability : experimental+-- Portability : GHC+module Misc.RepeatedDP where++import DynamicPipeline+import Relude++type DPExample = Source (Channel (Int :<+> Eof)) :>> Generator (Channel (Int :<+> Eof)) :>> Sink++source' :: Stage (WriteChannel Int -> DP s ())+source' = withSource @DPExample $ \cout -> unfoldT ([1 .. 1000] <> [1 .. 1000]) cout identity++generator' :: GeneratorStage DPExample (Maybe Int) Int s+generator' =+ let gen = withGenerator @DPExample genAction+ in mkGenerator gen filterTemp++genAction :: Filter DPExample (Maybe Int) Int s + -> ReadChannel Int+ -> WriteChannel Int+ -> DP s ()+genAction filter' cin cout = + let unfoldFilter = mkUnfoldFilterForAll' (`push` cout) filter' Just cin HNil + in void $ unfoldF unfoldFilter++filterTemp :: Filter DPExample (Maybe Int) Int s +filterTemp = mkFilter actorRepeted++actorRepeted :: Int+ -> ReadChannel Int+ -> WriteChannel Int+ -> StateT (Maybe Int) (DP s) ()+actorRepeted i rc wc = foldM rc $ \e -> if e /= i then push e wc else pure ()++sink' :: Stage (ReadChannel Int -> DP s ())+sink' = withSink @DPExample $ flip foldM print++program :: IO ()+program = runDP $ mkDP @DPExample source' generator' sink'
+ src/DynamicPipeline.hs view
@@ -0,0 +1,233 @@+-- |+-- Module : DynamicPipeline+-- Copyright : (c) 2021 Juan Pablo Royo Sales+--+-- License : BSD3+-- Maintainer : juanpablo.royo@gmail.com+-- Stability : experimental+-- Portability : GHC+--+-- __DynamicPipeline__ is a __/Type Safe/__ Dynamic and Parallel Streaming Library, which is an implementation of __Dynamic Pipeline Paradigm (DPP)__ +-- proposed in this paper [DPP](https://biblioteca.sistedes.es/articulo/the-dynamic-pipeline-paradigm/).+-- The aim of this Library is to provide all the __Type level__ constructs to guide the user in building a /DPP/ flow to solve any algorithm that fits on +-- this computational model. +-- +-- This implementation has been developed using /Type Level Programming/ techniques like @Type families@, @Defunctionalization@, @Existential Types@ and +-- @Dynamic Record Tagged Types@ among others.+-- Using all this techniques, we provide a /High Level and Type Safe/ DynamicPipeline Library to build a Data Flow Algorithm avoiding as much as possible +-- boilerplate code, but maintaining safety and robustness.+-- +-- Example of Filtering Repeated elements of a Stream+-- +-- @+-- import "DynamicPipeline"+--+-- type DPExample = 'Source' ('Channel' (Int ':<+>' 'Eof')) ':>>' 'Generator' ('Channel' (Int ':<+>' 'Eof')) ':>>' 'Sink'+-- +-- source' :: 'Stage' ('WriteChannel' Int -> 'DP' s ())+-- source' = 'withSource' @DPExample $ \cout -> 'unfoldT' ([1 .. 1000] <> [1 .. 1000]) cout identity+-- +-- generator' :: 'GeneratorStage' DPExample (Maybe Int) Int s+-- generator' =+-- let gen = 'withGenerator' @DPExample genAction+-- in 'mkGenerator' gen filterTemp+-- +-- genAction :: 'Filter' DPExample (Maybe Int) Int s +-- -> 'ReadChannel' Int+-- -> 'WriteChannel' Int+-- -> 'DP' s ()+-- genAction filter\' cin cout = +-- let unfoldFilter = 'mkUnfoldFilterForAll'' (\`'push'` cout) filter' Just cin 'HNil' +-- in void $ 'unfoldF' unfoldFilter+-- +-- filterTemp :: 'Filter' DPExample (Maybe Int) Int s +-- filterTemp = 'mkFilter' actorRepeted+-- +-- actorRepeted :: Int+-- -> 'ReadChannel' Int+-- -> 'WriteChannel' Int+-- -> StateT (Maybe Int) ('DP' s) ()+-- actorRepeted i rc wc = do+-- liftIO $ 'foldM' rc $ \e -> if e /= i then 'push' e wc else pure ()+-- +-- sink\' :: 'Stage' ('ReadChannel' Int -> 'DP' s ())+-- sink\' = 'withSink' @DPExample $ flip 'foldM' print+-- +-- program :: IO ()+-- program = 'runDP' $ 'mkDP' @DPExample source\' generator\' sink\'+-- @+--+module DynamicPipeline + ( -- * DP Flow Grammar #grammar#+ -- $grammar+ + -- * Building 'DynamicPipeline' #dp#+ -- $dp++ -- ** Generator and Filter #genfilter#+ -- $generator+ + -- * Types Flow definition+ Eof,+ Sink,+ Generator,+ Source,+ Channel,+ type (:>>)(..),+ type (:<+>)(..), + -- * Smart Constructors + DynamicPipeline,+ Filter,+ Actor,+ GeneratorStage,+ Stage,+ ValidDP,+ IsDP,+ DP, + UnFoldFilter, + withDP, + mkGenerator,+ mkFilter,+ single,+ actor,+ (|>>>),+ (|>>),+ withSource,+ withGenerator,+ withSink,+ mkDP,+ runDP,+ unfoldF,+ mkUnfoldFilter,+ mkUnfoldFilter',+ mkUnfoldFilterForAll,+ mkUnfoldFilterForAll',+ (.*.), HList(HNil), hHead, + -- * Channels+ ReadChannel,+ WriteChannel,+ foldM,+ foldM',+ push,+ pull,+ unfoldM,+ unfoldFile,+ unfoldT+ )+ where++import Data.HList ((.*.), HList(HNil), hHead)+import DynamicPipeline.Flow+import DynamicPipeline.Channel+import DynamicPipeline.Stage++-- $grammar+-- The following is the Regular Grammar allowed to build a /DPP/ Flow definition:+-- +-- @+-- __DP__ = 'Source' __CHANS__ ':>>' 'Generator' __CHANS__ ':>>' 'Sink'+-- __CHANS__ = 'Channel' __CH__+-- __CH__ = 'Eof' | 'Type' ':<+>' __CH__+-- @+--+-- Example: +-- +-- @ 'Source' ('Channel' (Int ':<+>' Int)) ':>>' 'Generator' ('Channel' (Int ':<+>' Int)) ':>>' 'Sink' @+--+--+-- $dp+-- 'DynamicPipeline' Data type is the point where all the information is contained in order the library can run our /DP/ Algorithm.+-- +-- This Data type contains three fundamental pieces: 'Source', 'Generator' and 'Sink'. But all these are dynamic based on the +-- defined Flow.+-- One of the fundamental feature of this Library is to provide several combinators that deduce from the Definition Flow, what are the +-- Function Signatures the user must fulfill according to his definition.+-- +-- All these combinators work in the same manner which based on the flow definition present to the user at compile time what is the function that must +-- be provided.+-- Lets see an example based on the "Misc.RepeatedDP", which basically filter out repeated elements in a stream.+--+-- >>> import Relude+-- >>> import DynamicPipeline+-- >>> type DPEx = Source (Channel (Int :<+> Eof)) :>> Generator (Channel (Int :<+> Eof)) :>> Sink+-- >>> :t withSource @DPEx+-- withSource @DPEx+-- :: forall k (st :: k).+-- (WriteChannel Int -> DP st ())+-- -> Stage (WriteChannel Int -> DP st ())+--+-- In @type DPEx = ..@ we are defining a Flow which contains a 'Source' that is going to have an 'Int' Channel that is going to feed the 'Generator'. +-- Therefore the 'Source' should write on that channel and because of that we are being asked to provide a Function that @WriteChannel Int -> DP st ()@.+-- Remember that our Monadic context is always 'DP'.+-- +-- Having that we can provide that function and have all the pieces together for 'Source'.+--+-- >>> let source' = withSource @DPEx $ \wc -> unfoldT ([1..10] <> [1..10] <> [1..10] <> [1..10]) wc identity+-- >>> :t source'+-- source' :: forall k (st :: k). Stage (WriteChannel Int -> DP st ())+--+-- So we are done. we provide that function.+-- Now we can do the same for 'Sink' which is the other opposite part of the Stream because 'Generator' is a little different as we can see in the documentation.+--+-- >>> let sink' = withSink @DPEx $ \rc -> foldM rc $ putStr . show+-- >>> :t sink'+-- sink' :: forall k (st :: k). Stage (ReadChannel Int -> DP st ())+--+-- Done with 'Sink'.+--+-- $generator+-- Now we reach to the last piece which needs more work to be done because it is the core of /DPP/ which dynamically adds Parallel computations between the 'Generator' Stage+-- and previous 'Filter's and 'Source'.+--+-- Fortunately we have the same combinator 'withGenerator' but it is not so straightforward what to put there. So, lets go step by step.+--+-- >>> :t withGenerator @DPEx+-- withGenerator @DPEx+-- :: forall k filter (st :: k).+-- (filter -> ReadChannel Int -> WriteChannel Int -> DP st ())+-- -> Stage+-- (filter -> ReadChannel Int -> WriteChannel Int -> DP st ())+--+-- At the first Glance it is asking for some similar function that is going to return our desired 'Stage' but there is some type parameter which is +-- not so obvious __@filter@__.+-- Fortunately we have combinators for that as well that can save us a lot of time and effort.+--+-- /Note: We could have done a Generator with an Empty 'Filter' but we are not taking advantage of DPP in building a Pipeline Parallelization Computational Algorithm/+--+-- In the case of 'Filter' we have several combinators at our disposal.+-- +-- * Use 'mkFilter' if your /DPP/ contains 1 actor per Filter+--+-- * Use '|>>' and '|>>>' if your /DPP/ contains more than 1 actor+--+-- In our example we are going to use 1 actor only that is going to discard repeated elements+-- >>> :t mkFilter @DPEx actor1+-- Variable not in scope:+-- actor1+-- :: filterParam+-- -> ReadChannel Int+-- -> WriteChannel Int+-- -> StateT filterState (DP st) ()+--+-- First lets fill in the gaps.+--+-- >>> let filter' = mkFilter @DPEx (\i rc wc -> foldM rc $ \e -> if e /= i then push e wc else pure ())+-- >>> :t filter'+-- filter' :: forall k filterState (st :: k). Filter DPEx filterState Int st+--+-- Basically we are checking if the element that we are reding from the Channel (Remember that we can have multiple 'Filter' on front writing to us),+-- is equal to the First Element that was read by the 'Generator' and on which this 'Filter' was instantiated with (a.k.a. @filterParam@). +-- If the element is not equal we 'push' it to the next 'Filter' or 'Generator', otherwise we discarded.+--+--+-- >>> let gen' = mkGenerator (withGenerator @DPEx $ \f r w -> let unf = mkUnfoldFilterForAll' (`push` w) f Just r HNil in void $ unfoldF unf) filter'+-- >>> :t gen'+-- gen' :: forall k (st :: k). GeneratorStage DPEx (Maybe Int) Int st+--+-- Now we have everything in place we only need to call 'runDP' and 'mkDP' +--+-- >>> runDP $ mkDP @DPEx source' gen' sink'+-- 12345678910+-- ++
+ src/DynamicPipeline/Channel.hs view
@@ -0,0 +1,135 @@+{-# LANGUAGE AllowAmbiguousTypes #-}+{-# LANGUAGE UndecidableInstances #-}+-- |+-- Module : DynamicPipeline.Channel+-- Copyright : (c) 2021 Juan Pablo Royo Sales+--+-- License : BSD3+-- Maintainer : juanpablo.royo@gmail.com+-- Stability : experimental+-- Portability : GHC+--+module DynamicPipeline.Channel + ( ReadChannel,+ WriteChannel,+ DynamicPipeline.Channel.foldM,+ foldM',+ push,+ pull,+ unfoldM,+ unfoldFile,+ unfoldT, + newChannel,+ end+ ) where++import qualified Control.Concurrent as CC+import Control.Concurrent.Chan.Unagi.NoBlocking+import Control.Lens hiding ((<|))+import Data.ByteString as B+import Data.Comp.Algebra (CoalgM, anaM)+import Data.Foldable as F+import Data.HList+import Relude as R+++-- | 'WriteChannel' can only write values into some Channel Queue+-- +-- [@a@]: Type that this Channel can write+newtype WriteChannel a = WriteChannel { unWrite :: InChan (Maybe a) }++-- | 'ReadChannel' can only read values of a previously written Channel. It is connected to a 'WriteChannel' but hidden for the user +-- +-- [@a@]: Type that this Channel can read+newtype ReadChannel a = ReadChannel { unRead :: OutChan (Maybe a) }++-- | 'foldM' is a /Catamorphism/ for consuming a 'ReadChannel' and do some Monadic @m@ computation with each element+{-# INLINE foldM #-}+foldM :: MonadIO m + => ReadChannel a -- ^'ReadChannel'+ -> (a -> m ()) -- ^Computation to do with read element+ -> m ()+foldM = flip foldM' (pure ())++-- | Idem 'foldM' but allows pass a monadic computation to perform at the end of the Channel+{-# INLINE foldM' #-}+foldM' :: MonadIO m + => ReadChannel a -- ^'ReadChannel'+ -> m () -- ^Computation to do at the end of the channel+ -> (a -> m ()) -- ^Computation to do with read element+ -> m ()+foldM' = loop'+ where+ loop' c onNothing io = + maybe onNothing (\e -> io e >> loop' c onNothing io) =<< liftIO (pull c)++-- | Push element @a@ into 'WriteChannel'+{-# INLINE push #-}+push :: MonadIO m => a -> WriteChannel a -> m ()+push a c = liftIO $ writeChan (unWrite c) (Just a)++-- | Pull element @Maybe a@ from 'ReadChannel'+{-# INLINE pull #-}+pull :: MonadIO m => ReadChannel a -> m (Maybe a)+pull = liftIO . readChan (CC.threadDelay 100) . unRead++-- | Coalgebra with Monadic computation to Feed some 'WriteChannel'+--+-- [@m@]: Monadic computation wrapping Coalgebra+--+-- [@a@]: Element get from some Source and to be write in some Channel+--+{-# WARNING SourceFeedCoalgM "INTERNAL USE" #-}+data SourceFeedCoalgM m a = Done -- ^ Termination 'Term' of the Coalgebra+ | Computation -- ^ Continuation 'Term' of the Coalgebra+ { _cSeed :: m a -- ^ Computation that Seeds the /Anamorphism/ + , _cStop :: m Bool -- ^ Stop signal+ , _cOnElem :: a -> m () -- ^ Computation on Read Element+ }++-- | 'SourceFeedCoalgM' in terms of 'CoalgM' +{-# WARNING sourceFeedCoalgM "INTERNAL USE" #-}+sourceFeedCoalgM :: MonadIO m => CoalgM m Maybe (SourceFeedCoalgM m a)+sourceFeedCoalgM Done = return Nothing+sourceFeedCoalgM c@Computation{..} = ifM _cStop + (return $ Just Done) + ( _cSeed >>= _cOnElem >> return (Just c) )++-- | unfold from a Monadic seed @m a@ to a 'WriteChannel'+{-# INLINE unfoldM #-}+unfoldM :: forall m a b. MonadIO m + => m a -- ^Monadic Seed + -> (a -> b) -- ^Map input from seed to something to be written in Channel+ -> m Bool -- ^When stop unfolding+ -> WriteChannel b -- ^'WriteChannel' to write input seed elements+ -> m ()+unfoldM seed fn stopIfM writeChannel =+ let onElem = flip push writeChannel . fn+ in anaM sourceFeedCoalgM (Computation seed stopIfM onElem) >> pure ()++-- | Using 'unfoldM', unfold from file+{-# INLINE unfoldFile #-}+unfoldFile :: MonadIO m + => FilePath -- ^Seed 'FilePath' to read from+ -> WriteChannel b -- ^'WriteChannel' to write File contents+ -> (ByteString -> b) -- ^Transform 'ByteString' read from File to something meaningful for your App+ -> m ()+unfoldFile file writeChannel fn = liftIO $+ R.withFile file ReadMode $ \h ->+ unfoldM (B.hGetLine h) fn (R.hIsEOF h) writeChannel++-- | Idem 'unfoldM' but for 'Foldable', for example a List @[a]@. Useful for testing purpose+{-# INLINE unfoldT #-}+unfoldT :: (MonadIO m, Foldable t) => t a -> WriteChannel b -> (a -> b) -> m ()+unfoldT ts writeChannel fn = forM_ ts (flip push writeChannel . fn)++{-# WARNING newChannel "INTERNAL USE" #-}+{-# NOINLINE newChannel #-}+newChannel :: forall a. IO (WriteChannel a, ReadChannel a)+newChannel = bimap WriteChannel ReadChannel <$> newChan++{-# WARNING end "INTERNAL USE" #-}+{-# INLINE end #-}+end :: WriteChannel a -> IO ()+end = flip writeChan Nothing . unWrite+
+ src/DynamicPipeline/Flow.hs view
@@ -0,0 +1,336 @@+{-# LANGUAGE AllowAmbiguousTypes #-}+{-# LANGUAGE UndecidableInstances #-}+-- |+-- Module : DynamicPipeline.Flow+-- Copyright : (c) 2021 Juan Pablo Royo Sales+--+-- License : BSD3+-- Maintainer : juanpablo.royo@gmail.com+-- Stability : experimental+-- Portability : GHC+--+module DynamicPipeline.Flow + ( Eof,+ Sink,+ Generator,+ Source,+ Channel,+ type (:>>)(..),+ type (:<+>)(..),+ ChanIn,+ ChanOut,+ ChanOutIn,+ ChansFilter,+ ChanWriteSource,+ ChanReadWriteGen,+ ChanReadOut,+ MkCh(..),+ MkChans(..),+ ExpandGenToCh, + ExpandSinkToCh,+ ExpandSourceToCh,+ ExpandFilterToCh,+ AllChans,+ FilterChans,+ inGenOut,+ makeChans,+ getFilterChannels+ ) where++import Control.Lens hiding ((<|))+import Data.Foldable as F+import Data.HList+import Data.HList.Labelable+import DynamicPipeline.Channel+import Relude as R++-- | 'Source' contains the 'Source' Stage its Channels definitions in the DP definition Flow.+-- +-- @ a ~ 'Channel' @+data Source (a :: Type)++-- | 'Generator' contains the 'Generator' Stage its Channels definitions in the DP definition Flow.+-- +-- @ a ~ 'Channel' @+data Generator (a :: Type)++-- | 'Sink' contains the 'Sink' Stage end of Flow of DP definition.+data Sink++-- |'Eof' is the __End of Channel__ mark in the DP Definition Flow+data Eof++-- |'Channel' is the Container Type of /Open Union Type/ which is going to be defined with ':<+>'.+--+-- @ a ~ (Type ':<+>' Type ':<+>' ... ':<+>' Eof) @+data Channel (a :: Type)++-- | This is the Type level function of the /Open Union Type/ for Channels. +-- +-- Channels forms an /Open Union Type/ in each stage because according to __DPP__ we can have multiple /In/ and /Out/ Channels +-- in a Single Stage. +--+-- 'Eof' should be the last Channel of the /Open Union Type/ to indicate termination of the Grammar.+--+-- @ chann1 ~ Type @+--+-- @ chann2 ~ Type @+data chann1 :<+> chann2 = chann1 :<+> chann2+ deriving (Typeable, Eq, Show, Functor, Traversable, Foldable, Bounded)+infixr 5 :<+>++-- | This is the Type level function of the /Open Union Type/ for Stages. +-- +-- This should have the form:+--+-- @ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink' @+data a :>> b = a :>> b+ deriving (Typeable, Eq, Show, Functor, Traversable, Foldable, Bounded)+infixr 5 :>>++-- Internal Data Types for expanding function based on Channel definitions+{-# WARNING ChanIn "INTERNAL USE" #-}+data ChanIn (a :: Type)+{-# WARNING ChanOut "INTERNAL USE" #-}+data ChanOut (a :: Type)+{-# WARNING ChanOutIn "INTERNAL USE" #-}+data ChanOutIn (a :: Type) (b :: Type)+{-# WARNING ChansFilter "INTERNAL USE" #-}+data ChansFilter (a :: Type)+{-# WARNING ChanWriteSource "INTERNAL USE" #-}+data ChanWriteSource (a :: Type)+{-# WARNING ChanReadWriteGen "INTERNAL USE" #-}+data ChanReadWriteGen (a :: Type)+{-# WARNING ChanReadOut "INTERNAL USE" #-}+data ChanReadOut (a :: Type)++-- Type encoding for Building Chans. Only for internal use in the Associated Type Family and combinators of MkCh and MkChans+-- For accessing Dynamic Indexed Records of Channels+{-# WARNING inLabel "INTERNAL USE" #-}+inLabel :: Label "Source"+inLabel = Label++{-# WARNING genLabel "INTERNAL USE" #-}+genLabel :: Label "generator"+genLabel = Label++{-# WARNING outLabel "INTERNAL USE" #-}+outLabel :: Label "Sink"+outLabel = Label++{-# WARNING inChLabel "INTERNAL USE" #-}+inChLabel :: Label "in-ch"+inChLabel = Label++{-# WARNING outChLabel "INTERNAL USE" #-}+outChLabel :: Label "out-ch"+outChLabel = Label++-- Associated Type Family: Building Source and Sink Channels+{-# WARNING MkCh "INTERNAL USE" #-}+class MkCh (a :: Type) where+ type HChI a :: [Type]+ type HChO a :: [Type]+ mkCh :: Proxy a -> IO (HList (HChI a), HList (HChO a))++instance MkCh more => MkCh (a :<+> more) where+ type HChI (a :<+> more) = WriteChannel a ': HChI more+ type HChO (a :<+> more) = ReadChannel a ': HChO more+ mkCh _ = do+ (i, o) <- newChannel @a+ (il, ol) <- mkCh (Proxy @more)+ return (i .*. il, o .*. ol)++instance MkCh Eof where+ type HChI Eof = '[]+ type HChO Eof = '[]+ mkCh _ = return (HNil, HNil)++-- Type Family Defunctionalization to Expand Source, Generator and Sinks to its own HList Channel types.+{-# WARNING ExpandToHList "INTERNAL USE" #-}+type family ExpandToHList (a :: Type) (param :: Type) :: [Type]+type instance ExpandToHList (ChanWriteSource ( Source (Channel inToGen)+ :>> Generator (Channel genToOut)+ :>> Sink )+ ) _ = HChI inToGen++type instance ExpandToHList (ChanReadWriteGen ( Source (Channel inToGen)+ :>> Generator (Channel genToOut)+ :>> Sink)+ ) filter = filter ': HAppendListR (HChO inToGen) (HChI genToOut)++type instance ExpandToHList (ChanReadOut ( Source (Channel inToGen)+ :>> Generator (Channel genToOut)+ :>> Sink )+ ) filter = HChO genToOut++{-# WARNING ExpandSourceToCh "INTERNAL USE" #-}+type ExpandSourceToCh a = ExpandToHList (ChanWriteSource a) Void+{-# WARNING ExpandGenToCh "INTERNAL USE" #-}+type ExpandGenToCh a filter = ExpandToHList (ChanReadWriteGen a) filter+{-# WARNING ExpandFilterToCh "INTERNAL USE" #-}+type ExpandFilterToCh a param = ExpandGenToCh a param+{-# WARNING ExpandSinkToCh "INTERNAL USE" #-}+type ExpandSinkToCh a = ExpandToHList (ChanReadOut a) Void++-- Class for building Channels base on a DP Definition on `a` Type+{-# WARNING MkChans "INTERNAL USE" #-}+class MkChans (a :: Type) where+ type HChan a :: Type+ mkChans :: Proxy a -> IO (HChan a)++-- Instance for Building Channels for all the Chain Source :>> Generator :>> Sink+instance ( MkCh inToGen+ , MkCh genToOut)+ => MkChans (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink) where++ type HChan (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink)+ = Record '[ Tagged "Source" (Record '[ Tagged "in-ch" (HList (HChI inToGen))+ , Tagged "out-ch" (HList (HChO inToGen))+ ]+ )+ , Tagged "generator" (Record '[ Tagged "in-ch" (HList (HChI genToOut))+ , Tagged "out-ch" (HList (HChO genToOut))+ ]+ )+ , Tagged "Sink" (Record '[ Tagged "in-ch" (HList (HChI genToOut))])+ ]++ mkChans _ = do+ (ii, io) <- mkCh (Proxy @inToGen)+ (gi, go) <- mkCh (Proxy @genToOut)+ (oi, _) <- mkCh (Proxy @genToOut)+ return $ (inLabel .=. (inChLabel .=. ii .*. outChLabel .=. io .*. emptyRecord))+ .*.+ (genLabel .=. (inChLabel .=. gi .*. outChLabel .=. go .*. emptyRecord))+ .*.+ (outLabel .=. (inChLabel .=. oi .*. emptyRecord))+ .*.+ emptyRecord++-- Instance for Building Only Channels for Filters on each Generator action+instance MkCh inToGen+ => MkChans (ChansFilter (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink)) where++ type HChan (ChansFilter (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink))+ = Record '[ Tagged "in-ch" (HList (HChO inToGen))+ , Tagged "out-ch" (HList (HChI inToGen))+ ]++ mkChans _ = do+ (writes', reads') <- mkCh (Proxy @inToGen)+ return $ mkRecord (inChLabel .=. reads' .*. outChLabel .=. writes' .*. HNil)+++{-# WARNING makeChans "INTERNAL USE" #-}+{-# INLINE makeChans #-}+makeChans :: forall (a :: Type). MkChans a => IO (HChan a)+makeChans = mkChans (Proxy @a)++-- Ugly Dynamic Indexed Record Viewer to generate specific list of channels+{-# WARNING sourceChans "INTERNAL USE" #-}+{-# INLINE sourceChans #-}+sourceChans :: ( LabeledOpticF (LabelableTy r1) (Const t1)+ , LabeledOpticP (LabelableTy r1) (->)+ , LabeledOpticTo (LabelableTy r1) "in-ch" (->)+ , LabeledOpticF (LabelableTy r2) (Const t1)+ , LabeledOpticP (LabelableTy r2) (->)+ , LabeledOpticTo (LabelableTy r2) "Source" (->)+ , Labelable "in-ch" r1 s t2 t1 t1+ , Labelable "Source" r2 t3 t3 (r1 s) (r1 t2))+ => r2 t3 -> t1+sourceChans = let inl = hLens' inLabel+ inch = hLens' inChLabel+ in view (inl . inch)++{-# WARNING generatorChans "INTERNAL USE" #-}+{-# INLINE generatorChans #-}+generatorChans :: ( LabeledOpticF (LabelableTy r1) (Const (HList l1))+ , LabeledOpticP (LabelableTy r1) (->)+ , LabeledOpticTo (LabelableTy r1) "out-ch" (->)+ , LabeledOpticF (LabelableTy r2) (Const (HList l2))+ , LabeledOpticP (LabelableTy r2) (->)+ , LabeledOpticTo (LabelableTy r2) "in-ch" (->)+ , LabeledOpticF (LabelableTy r3) (Const (HList l2))+ , LabeledOpticP (LabelableTy r3) (->)+ , LabeledOpticTo (LabelableTy r3) "generator" (->)+ , LabeledOpticF (LabelableTy r3) (Const (HList l1))+ , LabeledOpticTo (LabelableTy r3) "Source" (->)+ , HAppendList l1 l2+ , Labelable "generator" r3 t1 t1 (r2 s1) (r2 t2)+ , Labelable "in-ch" r2 s1 t2 (HList l2) (HList l2)+ , Labelable "Source" r3 t1 t1 (r1 s2) (r1 t3)+ , Labelable "out-ch" r1 s2 t3 (HList l1) (HList l1))+ => r3 t1 -> HList (HAppendListR l1 l2)+generatorChans ch = let inl = hLens' inLabel+ genl = hLens' genLabel+ inch = hLens' inChLabel+ outch = hLens' outChLabel+ outsIn = view (inl . outch) ch+ insGen = view (genl . inch) ch+ in outsIn `hAppendList` insGen++{-# WARNING sinkChans "INTERNAL USE" #-}+{-# INLINE sinkChans #-}+sinkChans :: ( LabeledOpticF (LabelableTy r1) (Const t1)+ , LabeledOpticP (LabelableTy r1) (->)+ , LabeledOpticTo (LabelableTy r1) "out-ch" (->)+ , LabeledOpticF (LabelableTy r2) (Const t1)+ , LabeledOpticP (LabelableTy r2) (->)+ , LabeledOpticTo (LabelableTy r2) "generator" (->)+ , Labelable "generator" r2 t2 t2 (r1 s) (r1 t3)+ , Labelable "out-ch" r1 s t3 t1 t1)+ => r2 t2 -> t1+sinkChans = let genl = hLens' genLabel+ outch = hLens' outChLabel+ in view (genl . outch)++{-# WARNING AllChans "INTERNAL USE" #-}+type AllChans r2 r3 l1 r4 l2 t2 s t1 s2 t5 l3 l4 = (LabeledOpticTo (LabelableTy r2) "in-ch" (->),+ LabeledOpticF (LabelableTy r3) (Const (HList l3)),+ LabeledOpticP (LabelableTy r3) (->),+ LabeledOpticTo (LabelableTy r3) "Source" (->),+ LabeledOpticF (LabelableTy r2) (Const (HList l1)),+ LabeledOpticP (LabelableTy r2) (->),+ LabeledOpticTo (LabelableTy r2) "out-ch" (->),+ LabeledOpticTo (LabelableTy r4) "in-ch" (->),+ LabeledOpticF (LabelableTy r3) (Const (HList l2)),+ LabeledOpticTo (LabelableTy r3) "generator" (->),+ LabeledOpticF (LabelableTy r3) (Const (HList l1)),+ LabeledOpticF (LabelableTy r2) (Const (HList l3)),+ LabeledOpticF (LabelableTy r4) (Const (HList l4)),+ LabeledOpticP (LabelableTy r4) (->),+ LabeledOpticTo (LabelableTy r4) "out-ch" (->),+ LabeledOpticF (LabelableTy r3) (Const (HList l4)),+ LabeledOpticF (LabelableTy r4) (Const (HList l2)),+ HAppendList l1 l2, Labelable "generator" r3 t2 t2 (r4 s) (r4 t1),+ Labelable "in-ch" r2 s2 t5 (HList l3) (HList l3),+ Labelable "in-ch" r4 s t1 (HList l2) (HList l2),+ Labelable "Source" r3 t2 t2 (r2 s2) (r2 t5),+ Labelable "out-ch" r2 s2 t5 (HList l1) (HList l1),+ Labelable "out-ch" r4 s t1 (HList l4) (HList l4))++{-# WARNING inGenOut "INTERNAL USE" #-}+{-# INLINE inGenOut #-}+inGenOut :: AllChans r2 r3 l1 r4 l2 t2 s t1 s2 t5 l3 l4 => r3 t2 -> (HList l3, HList (HAppendListR l1 l2), HList l4)+inGenOut ch = (sourceChans ch, generatorChans ch, sinkChans ch)+++{-# WARNING FilterChans "INTERNAL USE" #-}+type FilterChans r b t a = (LabeledOpticF (LabelableTy r) (Const b),+ LabeledOpticTo (LabelableTy r) "out-ch" (->),+ Labelable "out-ch" r t t b b,+ LabeledOpticF (LabelableTy r) (Const a),+ LabeledOpticP (LabelableTy r) (->),+ LabeledOpticTo (LabelableTy r) "in-ch" (->),+ Labelable "in-ch" r t t a a)++{-# WARNING getFilterChannels "INTERNAL USE" #-}+{-# INLINE getFilterChannels #-}+getFilterChannels :: FilterChans r b t a => r t -> (a, b)+getFilterChannels ch =+ let inch = hLens' inChLabel+ outch = hLens' outChLabel+ reads' = ch^.inch+ writes' = ch^.outch+ in (reads', writes')
+ src/DynamicPipeline/Stage.hs view
@@ -0,0 +1,561 @@+{-# LANGUAGE AllowAmbiguousTypes #-}+{-# LANGUAGE UndecidableInstances #-}+-- |+-- Module : DynamicPipeline.Channel+-- Copyright : (c) 2021 Juan Pablo Royo Sales+--+-- License : BSD3+-- Maintainer : juanpablo.royo@gmail.com+-- Stability : experimental+-- Portability : GHC+--+module DynamicPipeline.Stage+ ( DynamicPipeline,+ Filter,+ Actor,+ GeneratorStage,+ Stage,+ ValidDP,+ IsDP,+ DP,+ UnFoldFilter,+ withDP,+ mkGenerator,+ mkFilter,+ single,+ actor,+ (|>>>),+ (|>>),+ withSource,+ withGenerator,+ withSink,+ mkDP,+ runDP,+ unfoldF,+ mkUnfoldFilter,+ mkUnfoldFilter',+ mkUnfoldFilterForAll,+ mkUnfoldFilterForAll'+ ) where++import Control.Concurrent.Async+import Control.Lens hiding ((<|))+import Data.HList+import Data.List.NonEmpty+import DynamicPipeline.Channel+import DynamicPipeline.Flow+import GHC.TypeLits+import Relude as R+++-- | FCF - Type Level Defunctionalization: Boolean 'And' Type Level Function+type family And (a :: Bool) (b :: Bool) :: Bool where+ And 'True 'True = 'True+ And a b = 'False++-- | FCF - Type Level Defunctionalization+-- 'IsDP' Validates if /DP/ Flow Type Level Definition is Correct according to the Grammar+--+-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+type family IsDP (dpDefinition :: k) :: Bool where+ IsDP (Source (Channel inToGen)+ :>> Generator (Channel genToOut)+ :>> Sink)+ = And (IsDP (Source (Channel inToGen))) (IsDP (Generator (Channel genToOut)))+ IsDP (Source (Channel (a :<+> more))) = IsDP (Source (Channel more))+ IsDP (Source (Channel Eof)) = 'True+ IsDP (Generator (Channel (a :<+> more))) = IsDP (Generator (Channel more))+ IsDP (Generator (Channel a)) = 'True+ IsDP x = 'False+++-- | FCF - Type Level Defunctionalization+-- 'ValidDP' Check if 'IsDP' is True+--+-- [@a@]: @IsDP dpDefinition ~ 'True@+--+-- Throw a 'TypeError' if Grammar is not correct+type family ValidDP (a :: Bool) :: Constraint where+ ValidDP 'True = ()+ ValidDP 'False = TypeError+ ( 'Text "Invalid Semantic for Building DP Program"+ ':$$: 'Text "Language Grammar:"+ ':$$: 'Text "DP = Source CHANS :>> Generator CHANS :>> Sink"+ ':$$: 'Text "CHANS = Channel CH"+ ':$$: 'Text "CH = Type | Type :<+> CH"+ ':$$: 'Text "Example: 'Source (Channel (Int :<+> Int)) :>> Generator (Channel (Int :<+> Int)) :>> Sink'"+ )++-- Inductive Type Family for Expanding and building Source, Generator, Filter and Sink Functions Signatures+type family WithSource (dpDefinition :: Type) (monadicAction :: Type -> Type) :: Type where+ WithSource (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink) monadicAction+ = WithSource (ChanIn inToGen) monadicAction+ WithSource (ChanIn (dpDefinition :<+> more)) monadicAction = WriteChannel dpDefinition -> WithSource (ChanIn more) monadicAction+ WithSource (ChanIn Eof) monadicAction = monadicAction ()+ WithSource (ChanOutIn (dpDefinition :<+> more) ins) monadicAction = ReadChannel dpDefinition -> WithSource (ChanOutIn more ins) monadicAction+ WithSource (ChanOutIn Eof ins) monadicAction = WithSource (ChanIn ins) monadicAction+ WithSource dpDefinition _ = TypeError+ ( 'Text "Invalid Semantic for Source Stage"+ ':$$: 'Text "in the DP Definition '"+ ':<>: 'ShowType dpDefinition+ ':<>: 'Text "'"+ ':$$: 'Text "Language Grammar:"+ ':$$: 'Text "DP = Source CHANS :>> Generator CHANS :>> Sink"+ ':$$: 'Text "CHANS = Channel CH"+ ':$$: 'Text "CH = Type | Type :<+> CH"+ ':$$: 'Text "Example: 'Source (Channel (Int :<+> Int)) :>> Generator (Channel (Int :<+> Int)) :>> Sink'"+ )++type family WithGenerator (a :: Type) (filter :: Type) (monadicAction :: Type -> Type) :: Type where+ WithGenerator (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink) filter monadicAction+ = filter -> WithGenerator (ChanOutIn inToGen genToOut) filter monadicAction+ WithGenerator (ChanIn (a :<+> more)) filter monadicAction = WriteChannel a -> WithGenerator (ChanIn more) filter monadicAction+ WithGenerator (ChanIn Eof) filter monadicAction = monadicAction ()+ WithGenerator (ChanOutIn (a :<+> more) ins) filter monadicAction = ReadChannel a -> WithGenerator (ChanOutIn more ins) filter monadicAction+ WithGenerator (ChanOutIn Eof ins) filter monadicAction = WithGenerator (ChanIn ins) filter monadicAction+ WithGenerator dpDefinition _ _ = TypeError+ ( 'Text "Invalid Semantic for Generator Stage"+ ':$$: 'Text "in the DP Definition '"+ ':<>: 'ShowType dpDefinition+ ':<>: 'Text "'"+ ':$$: 'Text "Language Grammar:"+ ':$$: 'Text "DP = Source CHANS :>> Generator CHANS :>> Sink"+ ':$$: 'Text "CHANS = Channel CH"+ ':$$: 'Text "CH = Type | Type :<+> CH"+ ':$$: 'Text "Example: 'Source (Channel (Int :<+> Int)) :>> Generator (Channel (Int :<+> Int)) :>> Sink'"+ )++type family WithFilter (dpDefinition :: Type) (param :: Type) (monadicAction :: Type -> Type) :: Type where+ WithFilter (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink) param monadicAction+ = param -> WithFilter (ChanOutIn inToGen genToOut) param monadicAction+ WithFilter (ChanIn (dpDefinition :<+> more)) param monadicAction = WriteChannel dpDefinition -> WithFilter (ChanIn more) param monadicAction+ WithFilter (ChanIn Eof) param monadicAction = monadicAction ()+ WithFilter (ChanOutIn (dpDefinition :<+> more) ins) param monadicAction = ReadChannel dpDefinition -> WithFilter (ChanOutIn more ins) param monadicAction+ WithFilter (ChanOutIn Eof ins) param m = WithFilter (ChanIn ins) param m+ WithFilter dpDefinition _ _ = TypeError+ ( 'Text "Invalid Semantic Semantic for Generator Stage"+ ':$$: 'Text "in the DP Definition '"+ ':<>: 'ShowType dpDefinition+ ':<>: 'Text "'"+ ':$$: 'Text "Language Grammar:"+ ':$$: 'Text "DP = Source CHANS :>> Generator CHANS :>> Sink"+ ':$$: 'Text "CHANS = Channel CH"+ ':$$: 'Text "CH = Type | Type :<+> CH"+ ':$$: 'Text "Example: 'Source (Channel (Int :<+> Int)) :>> Generator (Channel (Int :<+> Int)) :>> Sink'"+ )++type family WithSink (dpDefinition :: Type) (monadicAction :: Type -> Type) :: Type where+ WithSink (Source (Channel inToGen) :>> Generator (Channel genToOut) :>> Sink) monadicAction+ = WithSink (ChanOut genToOut) monadicAction+ WithSink (ChanOut (dpDefinition :<+> more)) monadicAction = ReadChannel dpDefinition -> WithSink (ChanOut more) monadicAction+ WithSink (ChanOut Eof) monadicAction = monadicAction ()+ WithSink dpDefinition _ = TypeError+ ( 'Text "Invalid Semantic for Sink Stage"+ ':$$: 'Text "in the DP Definition '"+ ':<>: 'ShowType dpDefinition+ ':<>: 'Text "'"+ ':$$: 'Text "Language Grammar:"+ ':$$: 'Text "DP = Source CHANS :>> Generator CHANS :>> Sink"+ ':$$: 'Text "CHANS = Channel CH"+ ':$$: 'Text "CH = Type | Type :<+> CH"+ ':$$: 'Text "Example: 'Source (Channel (Int :<+> Int)) :>> Generator (Channel (Int :<+> Int)) :>> Sink'"+ )+++-- | 'DP' is the only Monadic Action allowed to run a /DP/ Defined Flow.+-- It is restricted on Scope by its Existential Type @st@ in order to not escape out from this Monadic Context.+--+-- [@st@]: Existential Type to Ensure context of Monadic 'DP'+--+-- [@a@]: Any Type that carries the Monadic Context 'DP'+--+newtype DP st a = DP+ { runStage :: IO a+ } deriving newtype (Functor, Applicative, Monad, MonadIO)++-- | Smart Constructor of 'DP' from 'IO' action+withDP :: IO a -> DP s a+withDP = DP++-- Defunctionalization+data Stage a where+ Stage :: Proxy a -> a -> Stage a++{-# INLINE mkStage #-}+mkStage :: forall a. Proxy a -> a -> Stage a+mkStage = Stage @a++{-# INLINE mkStage' #-}+mkStage' :: forall a. a -> Stage a+mkStage' = Stage (Proxy @a)++class EvalC l t | l -> t where+ run :: l -> t++instance forall a b. (a ~ b) => EvalC (Stage a) b where+ run (Stage _ f) = f++{-# INLINE runStageWith #-}+runStageWith :: forall (n :: HNat) f (xs :: [*]) s.+ (HCurry' n f xs (DP s ()), ArityFwd f n, ArityRev f n, CloseList xs)+ => Stage f -> HList xs -> DP s (Async ())+runStageWith fn cIns = withDP $ async (runStage (hUncurry (run fn) cIns) >> closeList cIns)++{-# INLINE runStageWith' #-}+runStageWith' :: forall (n :: HNat) f (xs :: [*]) (ss :: [*]) s.+ (HCurry' n f xs (DP s ()), ArityFwd f n, ArityRev f n, CloseList ss)+ => Stage f -> HList xs -> HList ss -> DP s (Async ())+runStageWith' fn cIns cClose = withDP (async (runStage (hUncurry (run fn) cIns) >> closeList cClose))++-- | 'DynamicPipeline' data type which contains all the three Stages definitions that have been generated by other combinators like 'withSource',+-- 'withGenerator' and 'withSink'.+--+-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@filterState@]: State of the 'StateT' 'Monad' that is the local State of the Filter execution+--+-- [@filterParam@]: Type of the First Parameter that is pass to the Filter when it is created by the Generator /Anamorphism/. Generator can change the type received from the Reader Channels.+--+-- [@st@]: Existential Scope of 'DP' 'Monad'.+data DynamicPipeline dpDefinition filterState filterParam st =+ DynamicPipeline+ { source :: Stage (WithSource dpDefinition (DP st))+ , generator :: GeneratorStage dpDefinition filterState filterParam st+ , sink :: Stage (WithSink dpDefinition (DP st))+ }++-- | 'GeneartorStage' is a special 'Stage' data type according to /DPP/ Definition which contains a 'Filter' template definition,+-- in orther to know how to spawn a new 'Filter' if it is needed, and the 'Stage' of the Generator to allow the user to perform some computation+-- in that case.+--+-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@filterState@]: State of the 'StateT' 'Monad' that is the local State of the Filter execution+--+-- [@filterParam@]: Type of the First Parameter that is pass to the Filter when it is created by the Generator /Anamorphism/. Generator can change the type received from the Reader Channels.+--+-- [@st@]: Existential Scope of 'DP' 'Monad'.+data GeneratorStage dpDefinition filterState filterParam st = GeneratorStage+ { _gsGenerator :: Stage (WithGenerator dpDefinition (Filter dpDefinition filterState filterParam st) (DP st))+ , _gsFilterTemplate :: Filter dpDefinition filterState filterParam st+ }++-- | 'Filter' Is the template definition of the 'Filter' that may be spawned when reading elements on the Stream.+--+-- * 'Filter' is a 'NonEmpty' List of 'Actor's. +-- +-- * Each 'Actor' is executed sequentially on the that List when an Element arrive to that 'Filter' instance.+-- +-- * All the 'Filter' execution (a.k.a. @forM_ actors runStage@) executes in a 'StateT' 'Monad' to share an internal state among 'Actor's.+--+-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@filterState@]: State of the 'StateT' 'Monad' that is the local State of the Filter execution+--+-- [@filterParam@]: Type of the First Parameter that is pass to the Filter when it is created by the Generator /Anamorphism/. Generator can change the type received from the Reader Channels.+--+-- [@st@]: Existential Scope of 'DP' 'Monad'.+newtype Filter dpDefinition filterState filterParam st =+ Filter { unFilter :: NonEmpty (Actor dpDefinition filterState filterParam (StateT filterState (DP st))) }+ deriving Generic++instance Wrapped (Filter s' s a param)++-- | 'Actor' Is a particular 'Stage' computation inside a 'Filter'.+--+-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@filterState@]: State of the 'StateT' 'Monad' that is the local State of the Filter execution+--+-- [@filterParam@]: Type of the First Parameter that is pass to the Filter when it is created by the Generator /Anamorphism/. Generator can change the type received from the Reader Channels.+--+-- [@monadicAction@]: 'Monad' Wrapped in 'StateT'.+newtype Actor dpDefinition filterState filterParam monadicAction =+ Actor { unActor :: MonadState filterState monadicAction => Stage (WithFilter dpDefinition filterParam monadicAction) }+++-- | Smart Constructor of 'GeneratorStage'.+{-# INLINE mkGenerator #-}+mkGenerator :: Stage (WithGenerator dpDefinition (Filter dpDefinition filterState filterParam st) (DP st)) -- ^Generator 'Stage'+ -> Filter dpDefinition filterState filterParam st -- ^'Filter' template+ -> GeneratorStage dpDefinition filterState filterParam st+mkGenerator = GeneratorStage++-- | Smart Constructor of 'Filter'.+{-# INLINE mkFilter #-}+mkFilter :: forall dpDefinition filterState filterParam st. + WithFilter dpDefinition filterParam (StateT filterState (DP st)) -- ^Associated type family to Generate Function Signature+ -> Filter dpDefinition filterState filterParam st+mkFilter = Filter . single++-- | Smart Constructor of Single 'Actor' Wrapped in 'NonEmpty' List. +{-# INLINE single #-}+single :: forall dpDefinition filterState filterParam st. + WithFilter dpDefinition filterParam (StateT filterState (DP st)) -- ^Associated type family to Generate Function Signature+ -> NonEmpty (Actor dpDefinition filterState filterParam (StateT filterState (DP st)))+single = one . actor++-- | Smart Constructor of 'Actor'.+{-# INLINE actor #-}+actor :: forall dpDefinition filterState filterParam st.+ WithFilter dpDefinition filterParam (StateT filterState (DP st)) -- ^Associated type family to Generate Function Signature+ -> Actor dpDefinition filterState filterParam (StateT filterState (DP st))+actor = Actor . mkStage' @(WithFilter dpDefinition filterParam (StateT filterState (DP st)))++-- | Combinator to build 'Filter' in a /DSL/ approach.+-- Add a new 'Actor' to an already existing 'Filter'.+{-# INLINE (|>>>) #-}+(|>>>) :: forall dpDefinition filterState filterParam st. + Actor dpDefinition filterState filterParam (StateT filterState (DP st)) -- ^New 'Actor' to put on front+ -> Filter dpDefinition filterState filterParam st -- ^Existing 'Filter'+ -> Filter dpDefinition filterState filterParam st+(|>>>) a f = f & _Wrapped' %~ (a <|)+infixr 5 |>>>++-- | Combinator to build 'Filter' in a /DSL/ approach .+-- Given 2 'Actor's build a 'Filter'.+{-# INLINE (|>>) #-}+(|>>) :: forall dpDefinition filterState filterParam st. + Actor dpDefinition filterState filterParam (StateT filterState (DP st)) -- ^'Actor' 1+ -> Actor dpDefinition filterState filterParam (StateT filterState (DP st)) -- ^'Actor' 2+ -> Filter dpDefinition filterState filterParam st+(|>>) a1 a2 = Filter (a1 <|one a2)+infixr 5 |>>++{-# INLINE runActor #-}+runActor :: ( MonadState filterState monadicAction+ , HCurry' n (WithFilter dpDefinition filterParam monadicAction) xs r+ , ArityFwd (WithFilter dpDefinition filterParam monadicAction) n+ , ArityRev (WithFilter dpDefinition filterParam monadicAction) n+ ) => Actor dpDefinition filterState filterParam monadicAction -> HList xs -> r+runActor = hUncurry . run . unActor++{-# INLINE runFilter #-}+runFilter :: ( CloseList ss+ , HCurry' n (WithFilter dpDefinition filterParam (StateT filterState (DP st))) xs (StateT filterState2 (DP st) ())+ , ArityFwd (WithFilter dpDefinition filterParam (StateT filterState (DP st))) n+ , ArityRev (WithFilter dpDefinition filterParam (StateT filterState (DP st))) n+ ) => Filter dpDefinition filterState filterParam st -> filterState2 -> HList xs -> HList ss -> DP st (Async ())+runFilter f s clist cClose = DP $ async $ do+ void . runStage . flip evalStateT s . mapM_ (`runActor` clist) . unFilter $ f+ closeList cClose++-- | Combinator for Building a 'Source' Stage. It uses an Associated Type Class to deduce the Function Signature required to the user+-- taken from /DP/ Type Level Flow Definition +-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@st@]: Existential Scope of 'DP' 'Monad'.+{-# INLINE withSource #-}+withSource :: forall (dpDefinition :: Type) st. + WithSource dpDefinition (DP st) -- ^Associated type family to Generate Function Signature+ -> Stage (WithSource dpDefinition (DP st))+withSource = mkStage' @(WithSource dpDefinition (DP st))++-- | Combinator for Building a 'Generator' Stage. It uses an Associated Type Class to deduce the Function Signature required to the user+-- taken from /DP/ Type Level Flow Definition +-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@filter@]: 'Filter' template type+--+-- [@st@]: Existential Scope of 'DP' 'Monad'.+{-# INLINE withGenerator #-}+withGenerator :: forall (dpDefinition :: Type) (filter :: Type) st. + WithGenerator dpDefinition filter (DP st) -- ^Associated type family to Generate Function Signature + -> Stage (WithGenerator dpDefinition filter (DP st))+withGenerator = mkStage' @(WithGenerator dpDefinition filter (DP st))++-- | Combinator for Building a 'Sink' Stage. It uses an Associated Type Class to deduce the Function Signature required to the user+-- taken from /DP/ Type Level Flow Definition +-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@st@]: Existential Scope of 'DP' 'Monad'.+{-# INLINE withSink #-}+withSink :: forall (dpDefinition :: Type) st. + WithSink dpDefinition (DP st) -- ^Associated type family to Generate Function Signature + -> Stage (WithSink dpDefinition (DP st))+withSink = mkStage' @(WithSink dpDefinition (DP st))++{-# INLINE mkDP' #-}+mkDP' :: forall dpDefinition filterState filterParam st.+ Stage (WithSource dpDefinition (DP st))+ -> GeneratorStage dpDefinition filterState filterParam st+ -> Stage (WithSink dpDefinition (DP st))+ -> DynamicPipeline dpDefinition filterState filterParam st+mkDP' = DynamicPipeline @dpDefinition++-- Hiding DP Constraint for running DP+type DPConstraint dpDefinition filterState st filterParam filter iparams gparams oparams r2 r3 l1 r4 l2 t2 s1 t1 s2 t5 l3 l4 =+ ( MkChans dpDefinition+ , HChan dpDefinition ~ r3 t2+ , Filter dpDefinition filterState filterParam st ~ filter+ , CloseList l3+ , CloseList l4+ , CloseList (HAppendListR l1 l2)+ , iparams ~ WithSource dpDefinition (DP st)+ , gparams ~ WithGenerator dpDefinition filter (DP st)+ , oparams ~ WithSink dpDefinition (DP st)+ , ArityRev iparams (HLength (ExpandSourceToCh dpDefinition))+ , ArityFwd iparams (HLength (ExpandSourceToCh dpDefinition))+ , HCurry' (HLength (ExpandSourceToCh dpDefinition)) iparams l3 (DP st ())+ , ArityRev gparams (HLength (ExpandGenToCh dpDefinition filter))+ , ArityFwd gparams (HLength (ExpandGenToCh dpDefinition filter))+ , HCurry' (HLength (ExpandGenToCh dpDefinition filter)) gparams (filter ': HAppendListR l1 l2) (DP st ())+ , ArityRev oparams (HLength (ExpandSinkToCh dpDefinition))+ , ArityFwd oparams (HLength (ExpandSinkToCh dpDefinition))+ , HCurry' (HLength (ExpandSinkToCh dpDefinition)) oparams l4 (DP st ())+ , AllChans r2 r3 l1 r4 l2 t2 s1 t1 s2 t5 l3 l4)++{-# INLINE buildDPProg #-}+buildDPProg :: forall dpDefinition filterState st filterParam filter iparams gparams oparams r2 r3 l1 r4 l2 t2 s1 t1 s2 t5 l3 l4.+ DPConstraint dpDefinition filterState st filterParam filter iparams gparams oparams r2 r3 l1 r4 l2 t2 s1 t1 s2 t5 l3 l4+ => DynamicPipeline dpDefinition filterState filterParam st -> DP st ()+buildDPProg DynamicPipeline{..} = do+ (cIns, cGen, cOut) <- inGenOut <$> withDP (makeChans @dpDefinition)+ let genWithFilter = _gsFilterTemplate generator .*. cGen+ runStageWith source cIns+ >> runStageWith' @(HLength (ExpandGenToCh dpDefinition filter)) @gparams (_gsGenerator generator) genWithFilter cGen+ >> runStageWith sink cOut >>= DP . wait++-- | Smart constructor for 'DynamicPipeline' Definition+{-# INLINE mkDP #-}+mkDP :: forall dpDefinition filterState st filterParam filter iparams gparams oparams r2 r3 l1 r4 l2 t2 s1 t1 s2 t5 l3 l4.+ DPConstraint dpDefinition filterState st filterParam filter iparams gparams oparams r2 r3 l1 r4 l2 t2 s1 t1 s2 t5 l3 l4+ => Stage (WithSource dpDefinition (DP st)) -- ^ 'Source' Stage generated by 'withSource' combinator+ -> GeneratorStage dpDefinition filterState filterParam st -- ^ 'Generator' Stage generated by 'withGenerator' combinator+ -> Stage (WithSink dpDefinition (DP st)) -- ^ 'Sink' Stage generated by 'withSink' combinator+ -> DP st ()+mkDP inS gS oS = buildDPProg (mkDP' inS gS oS)++-- | Run 'DP' 'Monad' to final 'IO' result+{-# INLINE runDP #-}+runDP :: (forall st. DP st a) -> IO a+runDP = runStage++-- Closable Automatic Write Channels+data NotClose (a :: Type)++class CloseList xs where+ closeList :: HList xs -> IO ()++instance (IsClosable x, CloseList xs) => CloseList (x ': xs) where+ closeList (HCons x xs) = close x >> closeList xs++instance CloseList '[] where+ closeList _ = pure ()++class IsClosable f where+ close :: f -> IO ()++instance IsClosable (WriteChannel a) where+ close = end++instance IsClosable (ReadChannel a) where+ close = const $ pure ()++-- | 'SpawnFilterConstraint' Constraint type alias+type SpawnFilterConstraint dpDefinition readElem st filterState filterParam l r t l1 b0 l2 l3 b2 b3 l4 =+ ( MkChans (ChansFilter dpDefinition)+ , FilterChans r (HList l3) t (HList (ReadChannel readElem : l1))+ , l1 ~ l+ , CloseList (ReadChannel readElem ': l4)+ , HAppendList l l3+ , l4 ~ HAppendListR l l3+ , l2 ~ (readElem ': ReadChannel readElem ': l4)+ , HChan (ChansFilter dpDefinition) ~ r t+ , WithFilter dpDefinition filterParam (StateT filterState (DP st)) ~ (b2 -> ReadChannel b2 -> b3)+ , HLength (ExpandFilterToCh dpDefinition filterParam) ~ HLength l2+ , HCurry' (HLength l2) (WithFilter dpDefinition filterParam (StateT filterState (DP st))) l2 (StateT filterState (DP st) ())+ , ArityFwd (WithFilter dpDefinition filterParam (StateT filterState (DP st))) (HLength (ExpandFilterToCh dpDefinition filterParam))+ , ArityRev b3 (HLength l4)+ )++-- | 'UnFoldFilter' is a wrapper Data Type that contains all the information needed to spawn 'Filter' instances according to /DPP/.+-- The user will have the capability to select how those filters are going to be spawned, for example on each read element, how to setup+-- initial states of 'StateT' Monad on 'Actor' computations in filters, among others.+--+-- [@dpDefinition ~ 'Source' ('Channel' ..) ':>>' 'Generator' ('Channel' ..) ':>>' 'Sink'@]: /DP/ Type level Flow Definition+--+-- [@readElem@]: Type of the element that is being read from the Selected Channel in the 'Generator' Stage+--+-- [@st@]: Existential Scope of 'DP' 'Monad'.+--+-- [@filterState@]: State of the 'StateT' 'Monad' that is the local State of the Filter execution+--+-- [@filterParam@]: Type of the First Parameter that is pass to the Filter when it is created by the Generator /Anamorphism/. Generator can change the type received from the Reader Channels.+--+data UnFoldFilter dpDefinition readElem st filterState filterParam l = + UnFoldFilter + { _ufSpawnIf :: readElem -> Bool -- ^Given a new Element determine if we need to interpose a new Filter or not+ , _ufOnElem :: readElem -> DP st () -- ^For each element that the Filter is consuming allow to do something outside the filter with that element. For example trace or debug+ , _ufFilter :: Filter dpDefinition filterState filterParam st -- ^'Filter' Template+ , _ufInitState :: readElem -> filterState -- ^Given the First element in this Filter Instance how to Initiate Internal 'Filter' 'StateT' (Memory)+ , _ufReadChannel :: ReadChannel readElem -- ^Main 'ReadChannel' to feed filter+ , _ufRsChannels :: HList l -- ^'HList' with the rest of the ReadChannels if There are needed or 'HNil' if it only contians 1 read channel+ }+ +-- | Smart Constructor for 'UnFoldFilter'+mkUnfoldFilter :: (readElem -> Bool) -- ^Given a new Element determine if we need to interpose a new Filter or not+ -> (readElem -> DP st ()) -- ^For each element that the Filter is consuming allow to do something outside the filter with that element. For example trace or debug+ -> Filter dpDefinition filterState filterParam st -- ^'Filter' Template+ -> (readElem -> filterState) -- ^Given the First element in this Filter Instance how to Initiate Internal 'Filter' 'StateT' (Memory)+ -> ReadChannel readElem -- ^Main 'ReadChannel' to feed filter+ -> HList l -- ^'HList' with the rest of the ReadChannels if There are needed or 'HNil' if it only contians 1 read channel+ -> UnFoldFilter dpDefinition readElem st filterState filterParam l+mkUnfoldFilter = UnFoldFilter++-- | Smart Constructor for 'UnFoldFilter' which bypass to do something externally on each read element+mkUnfoldFilter' :: (readElem -> Bool) -- ^ + -> Filter dpDefinition filterState filterParam st -- ^ + -> (readElem -> filterState) -- ^ + -> ReadChannel readElem -- ^ + -> HList l -- ^ + -> UnFoldFilter dpDefinition readElem st filterState filterParam l+mkUnfoldFilter' spawnIf = mkUnfoldFilter spawnIf (const $ pure ())++-- | Smart Constructor for 'UnFoldFilter' That creates a 'Filter' for each element on the Read Channel and interpose on Front of 'Generator' Stage+-- and Last 'Filter'+--+-- @ Source ---> Filter1 ---> Filter2 ... ---> FilterN ---> Generator ---> Sink @+--+mkUnfoldFilterForAll :: Filter dpDefinition filterState filterParam st -- ^ + -> (readElem -> filterState) -- ^ + -> ReadChannel readElem -- ^ + -> HList l -- ^ + -> UnFoldFilter dpDefinition readElem st filterState filterParam l+mkUnfoldFilterForAll = mkUnfoldFilter' (const True)++-- | Idem for 'mkUnfoldFilterForAll' but do something on each Element externally+--+mkUnfoldFilterForAll' :: (readElem -> DP st ()) -- ^ + -> Filter dpDefinition filterState filterParam st -- ^ + -> (readElem -> filterState) -- ^ + -> ReadChannel readElem -- ^ + -> HList l -- ^ + -> UnFoldFilter dpDefinition readElem st filterState filterParam l+mkUnfoldFilterForAll' = mkUnfoldFilter (const True)++-- | Run 'UnFoldFilter'+{-# INLINE unfoldF #-}+unfoldF :: forall dpDefinition readElem st filterState filterParam l r t l1 b0 l2 l3 b2 b3 l4.+ SpawnFilterConstraint dpDefinition readElem st filterState filterParam l r t l1 b0 l2 l3 b2 b3 l4+ => UnFoldFilter dpDefinition readElem st filterState filterParam l -- ^ 'UnFoldFilter'+ -> DP st (HList l) -- ^Return the list of 'ReadChannel's with the results to be read for the 'Generator' at the end. You can use this to pass the results to 'Sink'+unfoldF = loopSpawn++ where+ loopSpawn uf@UnFoldFilter{..} =+ maybe (pure _ufRsChannels) (loopSpawn <=< doOnElem uf) =<< DP (pull _ufReadChannel)++ doOnElem uf@UnFoldFilter{..} elem' = do+ _ufOnElem elem'+ if _ufSpawnIf elem'+ then do+ (reads', writes' :: HList l3) <- getFilterChannels <$> DP (makeChans @(ChansFilter dpDefinition))+ let hlist = elem' .*. _ufReadChannel .*. (_ufRsChannels `hAppendList` writes')+ void $ runFilter _ufFilter (_ufInitState elem') hlist (_ufReadChannel .*. (_ufRsChannels `hAppendList` writes'))+ return $ uf { _ufReadChannel = hHead reads', _ufRsChannels = hTail reads' }+ else return uf+++
+ test/Spec.hs view
@@ -0,0 +1,1 @@+{-# OPTIONS_GHC -F -pgmF hspec-discover #-}