packages feed

dph-prim-par (empty) → 0.5.1.1

raw patch · 24 files changed

+2826/−0 lines, 24 filesdep +basedep +dph-basedep +dph-prim-interfacesetup-changed

Dependencies added: base, dph-base, dph-prim-interface, dph-prim-seq, old-time, random, vector

Files

+ Data/Array/Parallel/Unlifted.hs view
@@ -0,0 +1,266 @@+{-# LANGUAGE PackageImports, CPP #-}++-- | Primitive parallel combinators that work on flat, unlifted arrays.+--   Some of them don't actually have parallel implementations, so we bail out+--   to the regular sequential ones.+--+--   This set of combinators is used when the program is comiled with @-fdph-par@.+--   When compiling with @-fdph-seq@, the ones in the @dph-prim-seq@ package are used+--   instead. The @dph-prim-seq package@ exports the same names, but all combinators+--   are implemented sequentially.+--+--   The API is defined in @DPH_Header.h@ and @DPH_Interface.h@ to ensure that both+--   @dph-prim-par@ and @dph-prim-seq@ really do export the same symbols.++#include "DPH_Header.h"++import Data.Array.Parallel.Unlifted.Parallel+import Data.Array.Parallel.Base.TracePrim+import Data.Array.Parallel.Unlifted.Distributed ( DT )+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector         as Seq+import qualified Data.Array.Parallel.Unlifted.Sequential.Segmented      as Seq+import Data.Array.Parallel.Unlifted.Sequential.Vector (Unbox,Vector)+import Prelude (($!))++#include "DPH_Interface.h"++class (Unbox a, DT a) => Elt a++type Array      = Vector+type Segd       = UPSegd+type Sel2       = UPSel2+type SelRep2    = UPSelRep2+++-- Constant time operations ---------------------------------------------------+--   We don't bother tracing these.++length                  = Seq.length+empty                   = Seq.empty+zip                     = Seq.zip+unzip                   = Seq.unzip+fsts                    = Seq.fsts+snds                    = Seq.snds+(!:)                    = (Seq.!)++elementsSel2_0          = elementsUPSel2_0+elementsSel2_1          = elementsUPSel2_1+repSel2                 = repUPSel2++mkSelRep2               = mkUPSelRep2+indicesSelRep2          = indicesUPSelRep2+elementsSelRep2_0       = elementsUPSelRep2_0+elementsSelRep2_1       = elementsUPSelRep2_1++mkSegd                  = mkUPSegd+lengthSegd              = lengthUPSegd+lengthsSegd             = lengthsUPSegd+indicesSegd             = indicesUPSegd+elementsSegd            = elementsUPSegd+++-------------------------------------------------------------------------------+-- These take least O(n) time in the length of the vector.+--   NOTE: That actual tracing is only enabled when +--         dph-base/D/A/P/Config.tracePrimEnabled is set to True,+--         otherwise tracePrim is a no-op.++replicate n val +        =  tracePrim (TraceReplicate n)+        $! replicateUP n val+++repeat n _ arr+        =  tracePrim (TraceRepeat n (Seq.length arr))+        $! repeatUP n arr+++extract arr i n+        =  tracePrim (TraceExtract (Seq.length arr) i n)+        $! Seq.extract arr i n+++drop n arr+        =  tracePrim (TraceDrop n (Seq.length arr))+        $! dropUP n arr+++permute arrSrc arrIxs+        =  tracePrim (TracePermute (Seq.length arrSrc))+        $! Seq.permute arrSrc arrIxs+++bpermuteDft len f arrIxs+        =  tracePrim (TraceBPermuteDft len)+        $! Seq.bpermuteDft len f arrIxs+++bpermute arrSrc arrIxs+        =  tracePrim (TraceBPermute (Seq.length arrSrc))+        $! bpermuteUP arrSrc arrIxs+++mbpermute f arrSrc streamIxs+        =  tracePrim (TraceMBPermute (Seq.length arrSrc))+        $! Seq.mbpermute f arrSrc streamIxs+++update arrSrc arrNew+        =  tracePrim (TraceUpdate (Seq.length arrSrc) (Seq.length arrNew))+        $! updateUP arrSrc arrNew+++(+:+) arr1 arr2+        =  tracePrim (TraceAppend (Seq.length arr1 + Seq.length arr2))+        $! (Seq.++) arr1 arr2+++interleave arr1 arr2+        =  tracePrim (TraceInterleave (Seq.length arr1 + Seq.length arr2))+        $! interleaveUP arr1 arr2+++pack arrSrc arrFlag+        =  tracePrim (TracePack (Seq.length arrSrc))+        $! packUP arrSrc arrFlag+++combine arrSel arr1 arr2+        =  tracePrim (TraceCombine (Seq.length arrSel))+        $! combineUP arrSel arr1 arr2+++combine2 arrTag sel arr1 arr2+        =  tracePrim (TraceCombine2 (Seq.length arrTag))+        $! combine2UP arrTag sel arr1 arr2+++map f arr+        =  tracePrim (TraceMap (Seq.length arr))+        $! mapUP f arr+++filter f src+ = let  dst     = filterUP f src+   in   tracePrim (TraceFilter (Seq.length src) (Seq.length dst)) dst+++zipWith f arr1 arr2+        =  tracePrim (TraceZipWith (Seq.length arr1) (Seq.length arr2))+        $! zipWithUP f arr1 arr2+++fold f x arr+        =  tracePrim (TraceFold (Seq.length arr))+        $! foldUP f x arr++        +fold1 f arr+        =  tracePrim (TraceFold1 (Seq.length arr))+        $! Seq.fold1 f arr+++and arr =  tracePrim (TraceAnd (Seq.length arr))+        $! andUP arr+++sum arr =  tracePrim (TraceSum (Seq.length arr))+        $! sumUP arr++                        +scan f x arr+        =  tracePrim (TraceScan (Seq.length arr))+        $! scanUP f x arr++        +indexed arr+        =  tracePrim (TraceIndexed (Seq.length arr))+        $! indexedUP arr+++enumFromTo from to+ = let  arr     = enumFromToUP from to+   in   tracePrim (TraceEnumFromTo (Seq.length arr)) arr++        +enumFromThenTo from thn to+ = let  arr     = enumFromThenToUP from thn to+   in   tracePrim (TraceEnumFromThenTo (Seq.length arr)) arr++   +enumFromStepLen from step len+ = let  arr     = enumFromStepLenUP from step len+   in   tracePrim (TraceEnumFromStepLen (Seq.length arr)) arr+++enumFromStepLenEach n starts steps lens+ = let  arr     = enumFromStepLenEachUP n starts steps lens+   in   tracePrim (TraceEnumFromStepLenEach (Seq.length arr)) arr+++mkSel2 tag is n0 n1 rep+        =  tracePrim (TraceMkSel2 (Seq.length is))+        $! mkUPSel2 tag is n0 n1 rep+++tagsSel2 sel+ = let  tags    = tagsUPSel2 sel+   in   tracePrim (TraceTagsSel2 (Seq.length tags)) tags+++indicesSel2 sel      + = let  arr     = indicesUPSel2 sel+   in   tracePrim (TraceIndicesSel2 (Seq.length arr)) arr+++replicate_s segd arr+        =  tracePrim (TraceReplicate_s (Seq.length arr))+        $! replicateSUP segd arr+++replicate_rs n arr+        =  tracePrim (TraceReplicate_rs n (Seq.length arr))+        $! replicateRSUP n arr+++append_s segd xd xs yd ys+ = let  arr     = appendSUP segd xd xs yd ys+   in   tracePrim (TraceAppend_s (Seq.length arr)) arr++        +fold_s f x segd arr+        =  tracePrim (TraceFold_s (Seq.length arr))+        $! foldSUP f x segd arr++        +fold1_s f segd arr+        =  tracePrim (TraceFold1_s (Seq.length arr))+        $! fold1SUP f segd arr+++fold_r f z segSize arr+        =  tracePrim (TraceFold_r (Seq.length arr))+        $! Seq.foldlRU f z segSize arr+++sum_r x arr+        =  tracePrim (TraceSum_r (Seq.length arr))+        $! sumRUP x arr+++indices_s segd+ = let  arr     = indicesSUP segd+   in   tracePrim (TraceIndices_s (Seq.length arr)) arr+++-- Random arrays ------------------------------------------+randoms                 = Seq.random+randomRs                = Seq.randomR+++-- IO -----------------------------------------------------+class Seq.UIO a => IOElt a+hPut                    = Seq.hPut+hGet                    = Seq.hGet+toList                  = Seq.toList+fromList                = Seq.fromList
+ Data/Array/Parallel/Unlifted/Distributed.hs view
@@ -0,0 +1,47 @@+-- | Distributed types and operations.+module Data.Array.Parallel.Unlifted.Distributed (+  -- * Gang operations+  Gang, forkGang, gangSize,++  -- * Gang hacks+  theGang,++  -- * Distributed types and classes+  DT(..),++  -- * Higher-order combinators+  mapD, zipWithD, foldD, scanD,++  -- * Equality+  eqD, neqD,++  -- * Distributed scalars+  scalarD,+  andD, orD, sumD,++  -- * Distributed pairs+  zipD, unzipD, fstD, sndD,++  -- * Distributed arrays+  lengthD, splitLenD, splitLenIdxD,+  splitD, splitAsD, joinLengthD, joinD, splitJoinD, joinDM,+  splitSegdD, splitSegdD', splitSD,+  lengthUSegdD, lengthsUSegdD, indicesUSegdD, elementsUSegdD,+  Distribution, balanced, unbalanced,++  -- * Permutations+  permuteD, bpermuteD, atomicUpdateD,++  -- * Debugging+  fromD, toD, debugD+) where++import Data.Array.Parallel.Unlifted.Distributed.Gang (+  Gang, forkGang, gangSize)+import Data.Array.Parallel.Unlifted.Distributed.TheGang+import Data.Array.Parallel.Unlifted.Distributed.Types+import Data.Array.Parallel.Unlifted.Distributed.Combinators+import Data.Array.Parallel.Unlifted.Distributed.Scalars+import Data.Array.Parallel.Unlifted.Distributed.Arrays+import Data.Array.Parallel.Unlifted.Distributed.Basics+
+ Data/Array/Parallel/Unlifted/Distributed/Arrays.hs view
@@ -0,0 +1,351 @@+{-# LANGUAGE EmptyDataDecls, ScopedTypeVariables #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Operations on distributed arrays.+module Data.Array.Parallel.Unlifted.Distributed.Arrays (+  lengthD, splitLenD, splitLenIdxD,+  splitAsD, splitD, joinLengthD, joinD, splitJoinD, joinDM,+  splitSegdD, splitSegdD', splitSD,++  permuteD, bpermuteD, atomicUpdateD,++  Distribution, balanced, unbalanced+) where+import Data.Array.Parallel.Base ( ST, runST)+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Sequential.Segmented+import Data.Array.Parallel.Unlifted.Distributed.Gang (+  Gang, gangSize, seqGang)+import Data.Array.Parallel.Unlifted.Distributed.DistST (+  DistST, stToDistST, myIndex )+import Data.Array.Parallel.Unlifted.Distributed.Types (+  DT, Dist, mkDPrim, indexD, lengthD, newD, writeMD, zipD, unzipD, fstD, sndD,+  elementsUSegdD,+  checkGangD)+import Data.Array.Parallel.Unlifted.Distributed.Basics+import Data.Array.Parallel.Unlifted.Distributed.Combinators+import Data.Array.Parallel.Unlifted.Distributed.Scalars (+  sumD)++import Data.Bits ( shiftR )+import Control.Monad ( when )++import GHC.Base ( quotInt, remInt )++here s = "Data.Array.Parallel.Unlifted.Distributed.Arrays." Prelude.++ s+++-- Distribution ---------------------------------------------------------------+-- | This is a phantom parameter used to record whether a distributed value+--   is balanced evenly among the threads. It's used to signal this property+--   between RULES, but the actual value is never used.+data Distribution++balanced :: Distribution+{-# NOINLINE balanced #-}+balanced = error $ here "balanced: touched"++unbalanced :: Distribution+{-# NOINLINE unbalanced #-}+unbalanced = error $ here "unbalanced: touched"+++-- Splitting and Joining array lengths ----------------------------------------+-- | Distribute an array length over a 'Gang'.+--   Each thread holds the number of elements it's reponsible for.+splitLenD :: Gang -> Int -> Dist Int+{-# INLINE splitLenD #-}+splitLenD g n = generateD_cheap g len+  where+    !p = gangSize g+    !l = n `quotInt` p+    !m = n `remInt` p++    {-# INLINE [0] len #-}+    len i | i < m     = l+1+          | otherwise = l+++-- | Distribute an array length over a 'Gang'.+--   Each thread holds the number of elements it's responsible for, +--   and the index of the start of its chunk.+splitLenIdxD :: Gang -> Int -> Dist (Int,Int)+{-# INLINE splitLenIdxD #-}+splitLenIdxD g n = generateD_cheap g len_idx+  where+    !p = gangSize g+    !l = n `quotInt` p+    !m = n `remInt` p++    {-# INLINE [0] len_idx #-}+    len_idx i | i < m     = (l+1, i*(l+1))+              | otherwise = (l,   i*l + m)+++-- | Get the overall length of a distributed array.+--   We ask each thread for its chunk length, and sum them all up.+joinLengthD :: Unbox a => Gang -> Dist (Vector a) -> Int+{-# INLINE joinLengthD #-}+joinLengthD g = sumD g . lengthD+                                               ++-- Splitting and Joining arrays -----------------------------------------------+-- | Distribute an array over a 'Gang' such that each threads gets the given+--   number of elements.+splitAsD :: Unbox a => Gang -> Dist Int -> Vector a -> Dist (Vector a)+{-# INLINE_DIST splitAsD #-}+splitAsD g dlen !arr = zipWithD (seqGang g) (Seq.slice arr) is dlen+  where+    is = fst $ scanD g (+) 0 dlen+++-- | Distribute an array over a 'Gang'.+--+--   NOTE: This is defined in terms of splitD_impl to avoid introducing loops+--         through RULES. Without it, splitJoinD would be a loop breaker.+splitD :: Unbox a => Gang -> Distribution -> Vector a -> Dist (Vector a)+{-# INLINE_DIST splitD #-}+splitD g _ arr = splitD_impl g arr++splitD_impl :: Unbox a => Gang -> Vector a -> Dist (Vector a)+{-# INLINE_DIST splitD_impl #-}+splitD_impl g !arr = generateD_cheap g (\i -> Seq.slice arr (idx i) (len i))+  where+    n = Seq.length arr+    !p = gangSize g+    !l = n `quotInt` p+    !m = n `remInt` p++    {-# INLINE [0] idx #-}+    idx i | i < m     = (l+1)*i+          | otherwise = l*i + m++    {-# INLINE [0] len #-}+    len i | i < m     = l+1+          | otherwise = l+++-- | Join a distributed array.+--+--   NOTE: This is defined in terms of joinD_impl to avoid introducing loops+--         through RULES. Without it, splitJoinD would be a loop breaker.+joinD :: Unbox a => Gang -> Distribution -> Dist (Vector a) -> Vector a+{-# INLINE CONLIKE [1] joinD #-}+joinD g _ darr  = joinD_impl g darr++joinD_impl :: forall a. Unbox a => Gang -> Dist (Vector a) -> Vector a+{-# INLINE_DIST joinD_impl #-}+joinD_impl g !darr = checkGangD (here "joinD") g darr $+                     Seq.new n (\ma -> zipWithDST_ g (copy ma) di darr)+  where+    (!di,!n) = scanD g (+) 0 $ lengthD darr+    copy :: forall s. MVector s a -> Int -> Vector a -> DistST s ()+    copy ma i arr = stToDistST (Seq.copy (mslice i (Seq.length arr) ma) arr)+++-- | Split a vector over a gang, run a distributed computation, then+--   join the pieces together again.+splitJoinD+        :: (Unbox a, Unbox b)+        => Gang+        -> (Dist (Vector a) -> Dist (Vector b))+        -> Vector a+        -> Vector b+{-# INLINE_DIST splitJoinD #-}+splitJoinD g f !xs = joinD_impl g (f (splitD_impl g xs))++++-- | Join a distributed array, yielding a mutable global array+joinDM :: Unbox a => Gang -> Dist (Vector a) -> ST s (MVector s a)+{-# INLINE joinDM #-}+joinDM g darr = checkGangD (here "joinDM") g darr $+                do+                  marr <- Seq.newM n+                  zipWithDST_ g (copy marr) di darr+                  return marr+  where+    (!di,!n) = scanD g (+) 0 $ lengthD darr+    --+    copy ma i arr = stToDistST (Seq.copy (mslice i (Seq.length arr) ma) arr)+++{-# RULES++"splitD[unbalanced]/joinD" forall g b da.+  splitD g unbalanced (joinD g b da) = da++"splitD[balanced]/joinD" forall g da.+  splitD g balanced (joinD g balanced da) = da++"splitD/splitJoinD" forall g b f xs.+  splitD g b (splitJoinD g f xs) = f (splitD g b xs)++"splitJoinD/joinD" forall g b f da.+  splitJoinD g f (joinD g b da) = joinD g b (f da)++"splitJoinD/splitJoinD" forall g f1 f2 xs.+  splitJoinD g f1 (splitJoinD g f2 xs) = splitJoinD g (f1 . f2) xs++  #-}++{-# RULES++"Seq.zip/joinD[1]" forall g xs ys.+  Seq.zip (joinD g balanced xs) ys+    = joinD g balanced (zipWithD g Seq.zip xs (splitD g balanced ys))++"Seq.zip/joinD[2]" forall g xs ys.+  Seq.zip xs (joinD g balanced ys)+    = joinD g balanced (zipWithD g Seq.zip (splitD g balanced xs) ys)++"Seq.zip/splitJoinD" forall gang f g xs ys.+  Seq.zip (splitJoinD gang (imapD gang f) xs) (splitJoinD gang (imapD gang g) ys)+    = splitJoinD gang (imapD gang (\i zs -> let (as,bs) = Seq.unzip zs+                                            in Seq.zip (f i as) (g i bs)))+                      (Seq.zip xs ys)++  #-}+++-- Permutation ----------------------------------------------------------------+-- | Permute for distributed arrays.+permuteD :: forall a. Unbox a => Gang -> Dist (Vector a) -> Dist (Vector Int) -> Vector a+{-# INLINE_DIST permuteD #-}+permuteD g darr dis = Seq.new n (\ma -> zipWithDST_ g (permute ma) darr dis)+  where+    n = joinLengthD g darr+    permute :: forall s. MVector s a -> Vector a -> Vector Int -> DistST s ()+    permute ma arr is = stToDistST (Seq.mpermute ma arr is)+++-- NOTE: The bang is necessary because the array must be fully evaluated+-- before we pass it to the parallel computation.+bpermuteD :: Unbox a => Gang -> Vector a -> Dist (Vector Int) -> Dist (Vector a)+{-# INLINE bpermuteD #-}+bpermuteD g !as ds = mapD g (Seq.bpermute as) ds+++-- Update ---------------------------------------------------------------------+-- NB: This does not (and cannot) try to prevent two threads from writing to+-- the same position. We probably want to consider this an (unchecked) user+-- error.+atomicUpdateD :: forall a. Unbox a+             => Gang -> Dist (Vector a) -> Dist (Vector (Int,a)) -> Vector a+{-# INLINE atomicUpdateD #-}+atomicUpdateD g darr upd = runST (+  do+    marr <- joinDM g darr+    mapDST_ g (update marr) upd+    Seq.unsafeFreeze marr+  )+  where+    update :: forall s. MVector s a -> Vector (Int,a) -> DistST s ()+    update marr arr = stToDistST (Seq.mupdate marr arr)+++--- Splitting and Joining segment descriptors ---------------------------------+splitSegdD :: Gang -> USegd -> Dist USegd+{-# NOINLINE splitSegdD #-}+splitSegdD g !segd = mapD g lengthsToUSegd+                   $ splitAsD g d lens+  where+    !d = snd+       . mapAccumLD g chunk 0+       . splitLenD g+       $ elementsUSegd segd++    n    = lengthUSegd segd+    lens = lengthsUSegd segd++    chunk !i !k = let !j = go i k+                  in (j,j-i)++    go !i !k | i >= n    = i+             | m == 0    = go (i+1) k+             | k <= 0    = i+             | otherwise = go (i+1) (k-m)+      where+        m = lens ! i+++search :: Int -> Vector Int -> Int+search !x ys = go 0 (Seq.length ys)+  where+    go i n | n <= 0        = i+           | (ys!mid) < x = go (mid+1) (n-half-1)+           | otherwise     = go i half+      where+        half = n `shiftR` 1+        mid  = i + half+++chunk :: USegd -> Int -> Int -> Bool -> (# Vector Int, Int, Int #)+chunk !segd !di !dn is_last+  = (# lens', k-left_len, left_off #)+  where+    !lens' = runST (do+                      mlens' <- Seq.newM n'+                      when (left /= 0) $ Seq.write mlens' 0 left+                      Seq.copy (Seq.mdrop left_len mlens')+                               (Seq.slice lens k (k'-k))+                      when (right /= 0) $ Seq.write mlens' (n' - 1) right+                      Seq.unsafeFreeze mlens')++    lens = lengthsUSegd segd+    idxs = indicesUSegd segd+    n    = Seq.length lens++    k  = search di idxs+    k' | is_last   = n+       | otherwise = search (di+dn) idxs++    left  | k == n    = dn+          | otherwise = min ((idxs!k) - di) dn++    right | k' == k   = 0+          | otherwise = di + dn - (idxs ! (k'-1))++    left_len | left == 0   = 0+             | otherwise   = 1++    left_off | left == 0   = 0+             | otherwise   = di - idxs ! (k-1)++    n' = left_len + (k'-k)+++splitSegdD' :: Gang -> USegd -> Dist ((USegd,Int),Int)+{-# INLINE splitSegdD' #-}+splitSegdD' g !segd = imapD g mk+                         (splitLenIdxD g+                         (elementsUSegd segd))+  where+    !p = gangSize g++    mk i (dn,di) = case chunk segd di dn (i == p-1) of+                     (# lens, l, o #) -> ((lengthsToUSegd lens,l),o)+++joinSegD :: Gang -> Dist USegd -> USegd+{-# INLINE_DIST joinSegD #-}+joinSegD g = lengthsToUSegd+           . joinD g unbalanced+           . mapD g lengthsUSegd+++splitSD :: Unbox a => Gang -> Dist USegd -> Vector a -> Dist (Vector a)+{-# INLINE_DIST splitSD #-}+splitSD g dsegd xs = splitAsD g (elementsUSegdD dsegd) xs++{-# RULES++"splitSD/splitJoinD" forall g d f xs.+  splitSD g d (splitJoinD g f xs) = f (splitSD g d xs)++"splitSD/Seq.zip" forall g d xs ys.+  splitSD g d (Seq.zip xs ys) = zipWithD g Seq.zip (splitSD g d xs)+                                             (splitSD g d ys)++  #-}
+ Data/Array/Parallel/Unlifted/Distributed/Basics.hs view
@@ -0,0 +1,44 @@+-- | Basic operations on distributed types.+module Data.Array.Parallel.Unlifted.Distributed.Basics (+  eqD, neqD, toD, fromD+) where+import Data.Array.Parallel.Unlifted.Distributed.Gang (+  Gang, gangSize)+import Data.Array.Parallel.Unlifted.Distributed.Types (+  DT, Dist, indexD, newD, writeMD,+  checkGangD)+import Data.Array.Parallel.Unlifted.Distributed.Combinators (+  zipWithD)+import Data.Array.Parallel.Unlifted.Distributed.Scalars (+  andD, orD)+import Control.Monad ( zipWithM_ )+++here s = "Data.Array.Parallel.Unlifted.Distributed.Basics." ++ s+++-- | Test whether to distributed values are equal. +--   This requires a 'Gang' and hence can't be defined in terms of 'Eq'.+eqD :: (Eq a, DT a) => Gang -> Dist a -> Dist a -> Bool+eqD g dx dy = andD g (zipWithD g (==) dx dy)+++-- | Test whether to distributed values are not equal.+--   This requires a 'Gang' and hence can't be defined in terms of 'Eq'.+neqD :: (Eq a, DT a) => Gang -> Dist a -> Dist a -> Bool+neqD g dx dy = orD g (zipWithD g (/=) dx dy)++++-- | Generate a distributed value from the first @p@ elements of a list.+--   /NOTE:/ For debugging only, don't use in production code.+toD :: DT a => Gang -> [a] -> Dist a+toD g xs = newD g (\md -> zipWithM_ (writeMD md) [0 .. gangSize g - 1] xs)+++-- | Yield all elements of a distributed value.+--   /NOTE:/ For debugging only, don't use in production code.+fromD :: DT a => Gang -> Dist a -> [a]+fromD g dt = checkGangD (here "fromDT") g dt $+             map (indexD dt) [0 .. gangSize g - 1]+
+ Data/Array/Parallel/Unlifted/Distributed/Combinators.hs view
@@ -0,0 +1,204 @@+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Standard combinators for distributed types.+module Data.Array.Parallel.Unlifted.Distributed.Combinators (+  generateD, generateD_cheap,+  imapD, mapD, zipD, unzipD, fstD, sndD, zipWithD, izipWithD,+  foldD, scanD, mapAccumLD,++  -- * Monadic combinators+  mapDST_, mapDST, zipWithDST_, zipWithDST+) where+import Data.Array.Parallel.Base ( ST, runST)+import Data.Array.Parallel.Unlifted.Distributed.Gang (+  Gang, gangSize)+import Data.Array.Parallel.Unlifted.Distributed.Types (+  DT, Dist, MDist, indexD, zipD, unzipD, fstD, sndD, deepSeqD,+  newMD, writeMD, unsafeFreezeMD,+  checkGangD, measureD, debugD)+import Data.Array.Parallel.Unlifted.Distributed.DistST+import Debug.Trace+++here s = "Data.Array.Parallel.Unlifted.Distributed.Combinators." ++ s+++-- | Create a distributed value, given a function that makes the value in each thread.+generateD :: DT a => Gang -> (Int -> a) -> Dist a+{-# NOINLINE generateD #-}+generateD g f +        = runDistST g (myIndex >>= return . f)+++-- | Create a distributed value, but run it sequentially (I think?)+generateD_cheap :: DT a => Gang -> (Int -> a) -> Dist a+{-# NOINLINE generateD_cheap #-}+generateD_cheap g f +        = runDistST_seq g (myIndex >>= return . f)+++-- Mapping --------------------------------------------------------------------+-- | Map a function across all elements of a distributed value.+--   The worker function also gets the current thread index.+--   As opposed to `imapD'` this version also deepSeqs each element before+--   passing it to the function.+imapD :: (DT a, DT b) => Gang -> (Int -> a -> b) -> Dist a -> Dist b+{-# INLINE [0] imapD #-}+imapD g f d = imapD' g (\i x -> x `deepSeqD` f i x) d+++-- | Map a function across all elements of a distributed value.+--   The worker function also gets the current thread index.+imapD' :: (DT a, DT b) => Gang -> (Int -> a -> b) -> Dist a -> Dist b+{-# NOINLINE imapD' #-}+imapD' g f !d = checkGangD (here "imapD") g d+                (runDistST g (do+                                i <- myIndex+                                x <- myD d+                                return (f i x)))+++-- | Map a function over a distributed value.+mapD :: (DT a, DT b) => Gang -> (a -> b) -> Dist a -> Dist b+{-# INLINE mapD #-}+mapD g = imapD g . const++{-# RULES++"imapD/generateD" forall gang f g.+  imapD gang f (generateD gang g) = generateD gang (\i -> f i (g i))++"imapD/generateD_cheap" forall gang f g.+  imapD gang f (generateD_cheap gang g) = generateD gang (\i -> f i (g i))++"imapD/imapD" forall gang f g d.+  imapD gang f (imapD gang g d) = imapD gang (\i x -> f i (g i x)) d++  #-}+++-- Zipping --------------------------------------------------------------------+-- | Combine two distributed values with the given function.+zipWithD :: (DT a, DT b, DT c)+         => Gang -> (a -> b -> c) -> Dist a -> Dist b -> Dist c+{-# INLINE zipWithD #-}+zipWithD g f dx dy = mapD g (uncurry f) (zipD dx dy)+++-- | Combine two distributed values with the given function.+--   The worker function also gets the index of the current thread.+izipWithD :: (DT a, DT b, DT c)+          => Gang -> (Int -> a -> b -> c) -> Dist a -> Dist b -> Dist c+{-# INLINE izipWithD #-}+izipWithD g f dx dy = imapD g (\i -> uncurry (f i)) (zipD dx dy)++{-# RULES+"zipD/imapD[1]" forall gang f xs ys.+  zipD (imapD gang f xs) ys+    = imapD gang (\i (x,y) -> (f i x,y)) (zipD xs ys)++"zipD/imapD[2]" forall gang f xs ys.+  zipD xs (imapD gang f ys)+    = imapD gang (\i (x,y) -> (x, f i y)) (zipD xs ys)++"zipD/generateD[1]" forall gang f xs.+  zipD (generateD gang f) xs+    = imapD gang (\i x -> (f i, x)) xs++"zipD/generateD[2]" forall gang f xs.+  zipD xs (generateD gang f)+    = imapD gang (\i x -> (x, f i)) xs++  #-}+++-- Folding --------------------------------------------------------------------+-- | Fold a distributed value.+foldD :: DT a => Gang -> (a -> a -> a) -> Dist a -> a+{-# NOINLINE foldD #-}+foldD g f !d = checkGangD ("here foldD") g d $+              fold 1 (d `indexD` 0)+  where+    !n = gangSize g+    --+    fold i x | i == n    = x+             | otherwise = fold (i+1) (f x $ d `indexD` i)+++-- | Prefix sum of a distributed value.+scanD :: forall a. DT a => Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)+{-# NOINLINE scanD #-}+scanD g f z !d = checkGangD (here "scanD") g d $+                 runST (do+                   md <- newMD g+                   s  <- scan md 0 z+                   d' <- unsafeFreezeMD md+                   return (d',s))+  where+    !n = gangSize g+    scan :: forall s. MDist a s -> Int -> a -> ST s a+    scan md i !x | i == n    = return x+                 | otherwise = do+                                 writeMD md i x+                                 scan md (i+1) (f x $ d `indexD` i)++-- | Combination of map and fold.+mapAccumLD :: forall a b acc. (DT a, DT b)+           => Gang -> (acc -> a -> (acc,b))+                   -> acc -> Dist a -> (acc,Dist b)+{-# INLINE_DIST mapAccumLD #-}+mapAccumLD g f acc !d = checkGangD (here "mapAccumLD") g d $+                        runST (do+                          md   <- newMD g+                          acc' <- go md 0 acc+                          d'   <- unsafeFreezeMD md+                          return (acc',d'))+  where+    !n = gangSize g+    go :: MDist b s -> Int -> acc -> ST s acc+    go md i acc | i == n    = return acc+                | otherwise = case f acc (d `indexD` i) of+                                (acc',b) -> do+                                              writeMD md i b+                                              go md (i+1) acc'+                                ++-- Versions that work on DistST -----------------------------------------------+-- NOTE: The following combinators must be strict in the Dists because if they+-- are not, the Dist might be evaluated (in parallel) when it is requested in+-- the current computation which, again, is parallel. This would break our+-- model andlead to a deadlock. Hence the bangs.++mapDST_ :: DT a => Gang -> (a -> DistST s ()) -> Dist a -> ST s ()+{-# INLINE mapDST_ #-}+mapDST_ g p d = mapDST_' g (\x -> x `deepSeqD` p x) d+++mapDST_' :: DT a => Gang -> (a -> DistST s ()) -> Dist a -> ST s ()+mapDST_' g p !d = checkGangD (here "mapDST_") g d $+                  distST_ g (myD d >>= p)+++mapDST :: (DT a, DT b) => Gang -> (a -> DistST s b) -> Dist a -> ST s (Dist b)+{-# INLINE mapDST #-}+mapDST g p !d = mapDST' g (\x -> x `deepSeqD` p x) d++mapDST' :: (DT a, DT b) => Gang -> (a -> DistST s b) -> Dist a -> ST s (Dist b)+mapDST' g p !d = checkGangD (here "mapDST_") g d $+                 distST g (myD d >>= p)+++zipWithDST_ :: (DT a, DT b)+            => Gang -> (a -> b -> DistST s ()) -> Dist a -> Dist b -> ST s ()+{-# INLINE zipWithDST_ #-}+zipWithDST_ g p !dx !dy = mapDST_ g (uncurry p) (zipD dx dy)+++zipWithDST :: (DT a, DT b, DT c)+           => Gang+           -> (a -> b -> DistST s c) -> Dist a -> Dist b -> ST s (Dist c)+{-# INLINE zipWithDST #-}+zipWithDST g p !dx !dy = mapDST g (uncurry p) (zipD dx dy)+
+ Data/Array/Parallel/Unlifted/Distributed/DistST.hs view
@@ -0,0 +1,112 @@+{-# LANGUAGE ScopedTypeVariables #-}+-- | Distributed ST computations.+--+--  Computations of type 'DistST' are data-parallel computations which+--  are run on each thread of a gang. At the moment, they can only access the+--  element of a (possibly mutable) distributed value owned by the current+--  thread.+--+-- /TODO:/ Add facilities for implementing parallel scans etc.+module Data.Array.Parallel.Unlifted.Distributed.DistST (+  DistST, stToDistST, distST_, distST, runDistST, runDistST_seq, traceDistST,+  myIndex, myD, readMyMD, writeMyMD+) where+import Data.Array.Parallel.Base (+  ST, runST)+import Data.Array.Parallel.Unlifted.Distributed.Gang+import Data.Array.Parallel.Unlifted.Distributed.Types (+  DT(..), Dist, MDist)++import Control.Monad (liftM)+++-- | Data-parallel computations.+--   When applied to a thread gang, the computation implicitly knows the index+--   of the thread it's working on. Alternatively, if we know the thread index+--   then we can make a regular ST computation.+newtype DistST s a = DistST { unDistST :: Int -> ST s a }++instance Monad (DistST s) where+  {-# INLINE return #-}+  return         = DistST . const . return ++  {-# INLINE (>>=) #-}+  DistST p >>= f = DistST $ \i -> do+                                    x <- p i+                                    unDistST (f x) i+++-- | Yields the index of the current thread within its gang.+myIndex :: DistST s Int+{-# INLINE myIndex #-}+myIndex = DistST return+++-- | Lifts an 'ST' computation into the 'DistST' monad.+--   The lifted computation should be data parallel.+stToDistST :: ST s a -> DistST s a+{-# INLINE stToDistST #-}+stToDistST p = DistST $ \i -> p+++-- | Yields the 'Dist' element owned by the current thread.+myD :: DT a => Dist a -> DistST s a+{-# NOINLINE myD #-}+myD dt = liftM (indexD dt) myIndex+++-- | Yields the 'MDist' element owned by the current thread.+readMyMD :: DT a => MDist a s -> DistST s a+{-# NOINLINE readMyMD #-}+readMyMD mdt + = do	i <- myIndex+	stToDistST $ readMD mdt i+++-- | Writes the 'MDist' element owned by the current thread.+writeMyMD :: DT a => MDist a s -> a -> DistST s ()+{-# NOINLINE writeMyMD #-}+writeMyMD mdt x + = do	i <- myIndex+	stToDistST $ writeMD mdt i x+++-- | Execute a data-parallel computation on a 'Gang'.+--   The same DistST comutation runs on each thread.+distST_ :: Gang -> DistST s () -> ST s ()+{-# INLINE distST_ #-}+distST_ g = gangST g . unDistST+++-- | Execute a data-parallel computation, yielding the distributed result.+distST :: DT a => Gang -> DistST s a -> ST s (Dist a)+{-# INLINE distST #-}+distST g p + = do	md <- newMD g+        distST_ g $ writeMyMD md =<< p+        unsafeFreezeMD md+++-- | Run a data-parallel computation, yielding the distributed result.+runDistST :: DT a => Gang -> (forall s. DistST s a) -> Dist a+{-# NOINLINE runDistST #-}+runDistST g p = runST (distST g p)++runDistST_seq :: forall a. DT a => Gang -> (forall s. DistST s a) -> Dist a+{-# NOINLINE runDistST_seq #-}+runDistST_seq g p = runST (+  do+     md <- newMD g+     go md 0+     unsafeFreezeMD md)                           +  where+    !n = gangSize g+    go :: forall s. MDist a s -> Int -> ST s ()+    go md i | i < n     = do+                            writeMD md i =<< unDistST p i+                            go md (i+1)+            | otherwise = return ()++traceDistST :: String -> DistST s ()+traceDistST s = DistST $ \n -> traceGangST ("Worker " ++ show n ++ ": " ++ s)+
+ Data/Array/Parallel/Unlifted/Distributed/Gang.hs view
@@ -0,0 +1,236 @@+{-# LANGUAGE CPP #-}++-- | Gang primitives.+--+-- /TODO:/+--+-- * Implement busy waiting.+--+-- * Benchmark.+--+-- * Generalise thread indices?++#define SEQ_IF_GANG_BUSY 1+#define TRACE_GANG 0++module Data.Array.Parallel.Unlifted.Distributed.Gang (+  Gang, seqGang, forkGang, gangSize, gangIO, gangST, traceGang, traceGangST +) where++import GHC.IO+import GHC.ST+import GHC.Conc                  ( forkOn )+import GHC.Exts                  ( traceEvent )++import Control.Concurrent.MVar+import Control.Exception         ( assert )+import Control.Monad             ( zipWithM, zipWithM_ )++import System.Time ( ClockTime(..), getClockTime )++-- Requests and operations on them --------------------------------------------+-- | The 'Req' type encapsulates work requests for individual members of a gang. +data Req +	-- | Instruct the worker to run the given action then signal it's done+	--   by writing to the MVar.+	= ReqDo	       (Int -> IO ()) (MVar ())++	-- | Tell the worker that we're shutting the gang down. The worker should+	--   signal that it's received the equest down by writing to the MVar before+	--   returning to its caller (forkGang) 	+	| ReqShutdown  (MVar ())+++-- | Create a new request for the given action.+newReq :: (Int -> IO ()) -> IO Req+newReq p + = do	mv	<- newEmptyMVar+	return	$ ReqDo p mv+++-- | Block until a thread request has been executed.+--   NOTE: only one thread can wait for the request.+waitReq :: Req -> IO ()+waitReq req+ = case req of+	ReqDo     fn varDone	-> takeMVar varDone+	ReqShutdown varDone	-> takeMVar varDone+++-- Thread gangs and operations on them ----------------------------------------+-- | A 'Gang' is a group of threads which execute arbitrary work requests.+--   To get the gang to do work, write Req-uest values+--   to its MVars+data Gang = Gang !Int           -- Number of 'Gang' threads+                 [MVar Req]     -- One 'MVar' per thread+                 (MVar Bool)    -- Indicates whether the 'Gang' is busy+++instance Show Gang where+  showsPrec p (Gang n _ _) +	= showString "<<"+        . showsPrec p n+        . showString " threads>>"+++-- | A sequential gang has no threads.+seqGang :: Gang -> Gang+seqGang (Gang n _ mv) = Gang n [] mv+++-- | The worker thread of a 'Gang'.+--   The threads blocks on the MVar waiting for a work request.+gangWorker :: Int -> MVar Req -> IO ()+gangWorker threadId varReq+ = do	traceGang $ "Worker " ++ show threadId ++ " waiting for request."+	req	<- takeMVar varReq+	+	case req of+	 ReqDo action varDone+	  -> do	traceGang $ "Worker " ++ show threadId ++ " begin"+		start 	<- getGangTime+		action threadId+		end 	<- getGangTime+		traceGang $ "Worker " ++ show threadId ++ " end (" ++ diffTime start end ++ ")"+		+		putMVar varDone ()+		gangWorker threadId varReq++	 ReqShutdown varDone+	  -> do	traceGang $ "Worker " ++ show threadId ++ " shutting down."+		putMVar varDone ()+++-- | Finaliser for worker threads.+--   We want to shutdown the corresponding thread when it's MVar becomes unreachable.+--     Without this Repa programs can complain about "Blocked indefinitely on an MVar"+--     because worker threads are still blocked on the request MVars when the program ends.+--     Whether the finalizer is called or not is very racey. It happens about 1 in 10 runs+--     when for the repa-edgedetect benchmark, and less often with the others.+-- +--   We're relying on the comment in System.Mem.Weak that says+--    "If there are no other threads to run, the runtime system will check for runnable+--     finalizers before declaring the system to be deadlocked."+-- +--   If we were creating and destroying the gang cleanly we wouldn't need this, but theGang +--     is created with a top-level unsafePerformIO. Hacks beget hacks beget hacks...+--+finaliseWorker :: MVar Req -> IO ()+finaliseWorker varReq+ = do	varDone <- newEmptyMVar+	putMVar varReq (ReqShutdown varDone) +	takeMVar varDone+	return ()+++-- | Fork a 'Gang' with the given number of threads (at least 1).+forkGang :: Int -> IO Gang+forkGang n+ = assert (n > 0) + $ do	+	-- Create the vars we'll use to issue work requests.+	mvs	<- sequence . replicate n $ newEmptyMVar+	+	-- Add finalisers so we can shut the workers down cleanly if they become unreachable.+	mapM_ (\var -> addMVarFinalizer var (finaliseWorker var)) mvs++	-- Create all the worker threads+	zipWithM_ forkOn [0..] +		$ zipWith gangWorker [0 .. n-1] mvs++	-- The gang is currently idle.+	busy	<- newMVar False+	+	return $ Gang n mvs busy+++-- | The number of threads in the 'Gang'.+gangSize :: Gang -> Int+gangSize (Gang n _ _) = n+++-- | Issue work requests for the 'Gang' and wait until they have been executed.+--   If the gang is already busy then just run the action in the+--   requesting thread. +--+--   TODO: We might want to print a configurable warning that this is happening.+--+gangIO	:: Gang+	-> (Int -> IO ())+	-> IO ()++gangIO (Gang n [] busy)  p = mapM_ p [0 .. n-1]+#if SEQ_IF_GANG_BUSY+gangIO (Gang n mvs busy) p + = do	traceGang   "gangIO: issuing work requests (SEQ_IF_GANG_BUSY)"+	b <- swapMVar busy True++	traceGang $ "gangIO: gang is currently " ++ (if b then "busy" else "idle")+	if b+	 then mapM_ p [0 .. n-1]+	 else do+		parIO n mvs p+		swapMVar busy False+		return ()+#else+gangIO (Gang n mvs busy) p = parIO n mvs p+#endif++-- | Issue some requests to the worker threads and wait for them to complete.+parIO 	:: Int			-- ^ Number of threads in the gang.+	-> [MVar Req]		-- ^ Request vars for worker threads.+	-> (Int -> IO ())	-- ^ Action to run in all the workers, it's given the ix of+				--   the particular worker thread it's running on.+	-> IO ()++parIO n mvs p + = do	traceGang "parIO: begin"++	start 	<- getGangTime+	reqs	<- sequence . replicate n $ newReq p++	traceGang "parIO: issuing requests"+	zipWithM putMVar mvs reqs++	traceGang "parIO: waiting for requests to complete"+	mapM_ waitReq reqs+	end 	<- getGangTime++	traceGang $ "parIO: end " ++ diffTime start end+++-- | Same as 'gangIO' but in the 'ST' monad.+gangST :: Gang -> (Int -> ST s ()) -> ST s ()+gangST g p = unsafeIOToST . gangIO g $ unsafeSTToIO . p+++-- Tracing -------------------------------------------------------------------+#if TRACE_GANG+getGangTime :: IO Integer+getGangTime+ = do	TOD sec pico <- getClockTime+	return (pico + sec * 1000000000000)++diffTime :: Integer -> Integer -> String+diffTime x y = show (y-x)++traceGang :: String -> IO ()+traceGang s+ = do	t <- getGangTime+	traceEvent $ show t ++ " @ " ++ s++#else+getGangTime :: IO ()+getGangTime = return ()++diffTime :: () -> () -> String+diffTime _ _ = ""++traceGang :: String -> IO ()+traceGang _ = return ()++#endif++traceGangST :: String -> ST s ()+traceGangST s = unsafeIOToST (traceGang s)+
+ Data/Array/Parallel/Unlifted/Distributed/Scalars.hs view
@@ -0,0 +1,37 @@+-- | Distributed scalars.+--   With a distributed value like (Dist Int), each thread has its own integer, +--   which may or may not have the same values as the ones on other threads.+module Data.Array.Parallel.Unlifted.Distributed.Scalars (+  unitD, scalarD,+  orD, andD, sumD+) where++import Data.Array.Parallel.Unlifted.Distributed.Gang (+  Gang)+import Data.Array.Parallel.Unlifted.Distributed.Types (+  DT, Dist, unitD)+import Data.Array.Parallel.Unlifted.Distributed.Combinators (+  mapD, foldD)+++-- | Distribute a scalar.+--   Each thread gets its own copy of the same value.+scalarD :: DT a => Gang -> a -> Dist a+scalarD g x = mapD g (const x) (unitD g)+++-- | OR together all instances of a distributed 'Bool'.+orD :: Gang -> Dist Bool -> Bool+orD g = foldD g (||)+++-- | AND together all instances of a distributed 'Bool'.+andD :: Gang -> Dist Bool -> Bool+andD g = foldD g (&&)+++-- | Sum all instances of a distributed number.+sumD :: (Num a, DT a) => Gang -> Dist a -> a+sumD g = foldD g (+)++
+ Data/Array/Parallel/Unlifted/Distributed/TheGang.hs view
@@ -0,0 +1,20 @@++-- DPH programs always used a single, shared gang of threads.+-- The gang exists at top level, and is initialised unsafely.+-- +-- The Vectoriser guarantees that the gang is only used by a single+-- computation at a time.+--+module Data.Array.Parallel.Unlifted.Distributed.TheGang (+  theGang+) where++import Data.Array.Parallel.Unlifted.Distributed.Gang ++import System.IO.Unsafe (unsafePerformIO)+import GHC.Conc (numCapabilities)++theGang :: Gang+{-# NOINLINE theGang #-}+theGang = unsafePerformIO (forkGang numCapabilities)+
+ Data/Array/Parallel/Unlifted/Distributed/Types.hs view
@@ -0,0 +1,495 @@+{-# OPTIONS -fno-warn-incomplete-patterns #-}+{-# LANGUAGE CPP #-}++#include "fusion-phases.h"++-- | Distributed types.+module Data.Array.Parallel.Unlifted.Distributed.Types (+  -- * Distributed types+  DT, Dist, MDist, DPrim(..),++  -- * Operations on immutable distributed types+  indexD, unitD, zipD, unzipD, fstD, sndD, lengthD,+  newD,+  -- zipSD, unzipSD, fstSD, sndSD,+  deepSeqD,++  lengthUSegdD, lengthsUSegdD, indicesUSegdD, elementsUSegdD,++  -- * Operations on mutable distributed types+  newMD, readMD, writeMD, unsafeFreezeMD,++  -- * Assertions+  checkGangD, checkGangMD,++  -- * Debugging functions+  sizeD, sizeMD, measureD, debugD+) where++import Data.Array.Parallel.Unlifted.Distributed.Gang (+  Gang, gangSize )+import Data.Array.Parallel.Unlifted.Sequential.Vector ( Unbox, Vector )+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as V+import Data.Array.Parallel.Unlifted.Sequential.Segmented+import Data.Array.Parallel.Base++import qualified Data.Vector.Unboxed as V+import qualified Data.Vector.Unboxed.Mutable as MV+import qualified Data.Vector as BV+import qualified Data.Vector.Mutable as MBV++import Data.Word     (Word8)+import Control.Monad (liftM, liftM2, liftM3)++import Data.List ( intercalate )++infixl 9 `indexD`++here s = "Data.Array.Parallel.Unlifted.Distributed.Types." ++ s+++-- Distributed Types ----------------------------------------------------------+-- | Class of distributable types. Instances of 'DT' can be+--   distributed across all workers of a 'Gang'. +--   All such types must be hyperstrict as we do not want to pass thunks+--   into distributed computations.+class DT a where+  data Dist  a+  data MDist a :: * -> *++  -- | Extract a single element of an immutable distributed value.+  indexD         :: Dist a -> Int -> a++  -- | Create an unitialised distributed value for the given 'Gang'.+  --   The gang is used (only) to know how many elements are needed+  --   in the distributed value.+  newMD          :: Gang                  -> ST s (MDist a s)++  -- | Extract an element from a mutable distributed value.+  readMD         :: MDist a s -> Int      -> ST s a++  -- | Write an element of a mutable distributed value.+  writeMD        :: MDist a s -> Int -> a -> ST s ()++  -- | Unsafely freeze a mutable distributed value.+  unsafeFreezeMD :: MDist a s             -> ST s (Dist a)++  deepSeqD       :: a -> b -> b+  deepSeqD = seq+++  -- Debugging ------------------------+  -- | Number of elements in the distributed value.+  --   For debugging only, as we shouldn't depend on the size of the gang.+  sizeD :: Dist a -> Int++  -- | Number of elements in the mutable distributed value.+  --   For debugging only, as we shouldn't care about the actual number.+  sizeMD :: MDist a s -> Int++  -- | Show a distributed value.+  --   For debugging only.+  measureD :: a -> String+  measureD _ = "None"++-- Show instance (for debugging only)+instance (Show a, DT a) => Show (Dist a) where+  show d = show (Prelude.map (indexD d) [0 .. sizeD d - 1])++++-- Checking -------------------------------------------------------------------+-- | Check that the sizes of the 'Gang' and of the distributed value match.+checkGangD :: DT a => String -> Gang -> Dist a -> b -> b+checkGangD loc g d v = checkEq loc "Wrong gang" (gangSize g) (sizeD d) v+++-- | Check that the sizes of the 'Gang' and of the mutable distributed value match.+checkGangMD :: DT a => String -> Gang -> MDist a s -> b -> b+checkGangMD loc g d v = checkEq loc "Wrong gang" (gangSize g) (sizeMD d) v+++-- Operations -----------------------------------------------------------------+-- | Given a computation that can write its result to a mutable distributed value, +--   run the computation to generate an immutable distributed value.+newD :: DT a => Gang -> (forall s . MDist a s -> ST s ()) -> Dist a+newD g init =+  runST (do+           mdt <- newMD g+           init mdt+           unsafeFreezeMD mdt)++-- | Show all members of a distributed value.+debugD :: DT a => Dist a -> String+debugD d = "["+         ++ intercalate "," [measureD (indexD d i) | i <- [0 .. sizeD d-1]]+         ++ "]"+++-- DPrim ----------------------------------------------------------------------+-- | For distributed primitive values, we can just store all the members in+--   a vector. The vector has the same length as the number of threads in the gang.+--+class Unbox e => DPrim e where++  -- | Make an immutable distributed value.+  mkDPrim :: V.Vector e -> Dist  e++  -- | Unpack an immutable distributed value back into a vector.+  unDPrim :: Dist  e -> V.Vector e++  -- | Make a mutable distributed value.+  mkMDPrim :: MV.STVector s e -> MDist  e s++  -- | Unpack a mutable distributed value back into a vector.+  unMDPrim :: MDist  e s -> MV.STVector s e+++-- | Get the member corresponding to a thread index.+primIndexD :: DPrim a => Dist a -> Int -> a+{-# INLINE primIndexD #-}+primIndexD = (V.!) . unDPrim+++-- | Create a new distributed value, having as many members as threads+--   in the given 'Gang'.+primNewMD :: DPrim a => Gang -> ST s (MDist a s)+{-# INLINE primNewMD #-}+primNewMD = liftM mkMDPrim . MV.new . gangSize+++-- | Read the member of a distributed value corresponding to the given thread index.+primReadMD :: DPrim a => MDist a s -> Int -> ST s a+{-# INLINE primReadMD #-}+primReadMD = MV.read . unMDPrim+++-- | Write the member of a distributed value corresponding to the given thread index.+primWriteMD :: DPrim a => MDist a s -> Int -> a -> ST s ()+{-# INLINE primWriteMD #-}+primWriteMD = MV.write . unMDPrim+++-- | Freeze a mutable distributed value to an immutable one.+--   You promise not to update the mutable one any further.+primUnsafeFreezeMD :: DPrim a => MDist a s -> ST s (Dist a)+{-# INLINE primUnsafeFreezeMD #-}+primUnsafeFreezeMD = liftM mkDPrim . V.unsafeFreeze . unMDPrim+++-- | Get the size of a distributed value, that is, the number of threads+--   in the gang that it was created for.+primSizeD :: DPrim a => Dist a -> Int+{-# INLINE primSizeD #-}+primSizeD = V.length . unDPrim+++-- | Get the size of a distributed mutable value, that is, the number of threads+--   in the gang it was created for.+primSizeMD :: DPrim a => MDist a s -> Int+{-# INLINE primSizeMD #-}+primSizeMD = MV.length . unMDPrim+++-- Unit -----------------------------------------------------------------------+instance DT () where+  data Dist ()    = DUnit  !Int+  data MDist () s = MDUnit !Int++  indexD  (DUnit n) i       = check (here "indexD[()]") n i $ ()+  newMD                     = return . MDUnit . gangSize+  readMD   (MDUnit n) i     = check (here "readMD[()]")  n i $+                               return ()+  writeMD  (MDUnit n) i ()  = check (here "writeMD[()]") n i $+                               return ()+  unsafeFreezeMD (MDUnit n) = return $ DUnit n++  sizeD  = error "dph-prim-par:sizeD[()] undefined"+  sizeMD = error "dph-prim-par:sizeMD[()] undefined"++-- | Yield a distributed unit.+unitD :: Gang -> Dist ()+{-# INLINE_DIST unitD #-}+unitD = DUnit . gangSize+++-- Bool -----------------------------------------------------------------------+instance DPrim Bool where+  mkDPrim           = DBool+  unDPrim (DBool a) = a++  mkMDPrim            = MDBool+  unMDPrim (MDBool a) = a++instance DT Bool where+  data Dist  Bool   = DBool  !(V.Vector    Bool)+  data MDist Bool s = MDBool !(MV.STVector s Bool)++  indexD         = primIndexD+  newMD          = primNewMD+  readMD         = primReadMD+  writeMD        = primWriteMD+  unsafeFreezeMD = primUnsafeFreezeMD+  sizeD          = primSizeD+  sizeMD         = primSizeMD+++-- Char -----------------------------------------------------------------------+instance DPrim Char where+  mkDPrim           = DChar+  unDPrim (DChar a) = a++  mkMDPrim            = MDChar+  unMDPrim (MDChar a) = a++instance DT Char where+  data Dist  Char   = DChar  !(V.Vector    Char)+  data MDist Char s = MDChar !(MV.STVector s Char)++  indexD         = primIndexD+  newMD          = primNewMD+  readMD         = primReadMD+  writeMD        = primWriteMD+  unsafeFreezeMD = primUnsafeFreezeMD+  sizeD          = primSizeD+  sizeMD         = primSizeMD+++-- Int ------------------------------------------------------------------------+instance DPrim Int where+  mkDPrim          = DInt+  unDPrim (DInt a) = a++  mkMDPrim            = MDInt+  unMDPrim (MDInt a) = a++instance DT Int where+  data Dist  Int   = DInt  !(V.Vector    Int)+  data MDist Int s = MDInt !(MV.STVector s Int)++  indexD         = primIndexD+  newMD          = primNewMD+  readMD         = primReadMD+  writeMD        = primWriteMD+  unsafeFreezeMD = primUnsafeFreezeMD+  sizeD          = primSizeD+  sizeMD         = primSizeMD++  measureD n = "Int " ++ show n+++-- Word8 ----------------------------------------------------------------------+instance DPrim Word8 where+  mkDPrim            = DWord8+  unDPrim (DWord8 a) = a++  mkMDPrim             = MDWord8+  unMDPrim (MDWord8 a) = a++instance DT Word8 where+  data Dist  Word8   = DWord8  !(V.Vector    Word8)+  data MDist Word8 s = MDWord8 !(MV.STVector s Word8)++  indexD         = primIndexD+  newMD          = primNewMD+  readMD         = primReadMD+  writeMD        = primWriteMD+  unsafeFreezeMD = primUnsafeFreezeMD+  sizeD          = primSizeD+  sizeMD         = primSizeMD+++-- Float ----------------------------------------------------------------------+instance DPrim Float where+  mkDPrim            = DFloat+  unDPrim (DFloat a) = a++  mkMDPrim             = MDFloat+  unMDPrim (MDFloat a) = a++instance DT Float where+  data Dist  Float   = DFloat  !(V.Vector    Float)+  data MDist Float s = MDFloat !(MV.STVector s Float)++  indexD         = primIndexD+  newMD          = primNewMD+  readMD         = primReadMD+  writeMD        = primWriteMD+  unsafeFreezeMD = primUnsafeFreezeMD+  sizeD          = primSizeD+  sizeMD         = primSizeMD+++-- Double ---------------------------------------------------------------------+instance DPrim Double where+  mkDPrim             = DDouble+  unDPrim (DDouble a) = a++  mkMDPrim              = MDDouble+  unMDPrim (MDDouble a) = a++instance DT Double where+  data Dist  Double   = DDouble  !(V.Vector    Double)+  data MDist Double s = MDDouble !(MV.STVector s Double)++  indexD         = primIndexD+  newMD          = primNewMD+  readMD         = primReadMD+  writeMD        = primWriteMD+  unsafeFreezeMD = primUnsafeFreezeMD+  sizeD          = primSizeD+  sizeMD         = primSizeMD+++-- Pairs ----------------------------------------------------------------------+instance (DT a, DT b) => DT (a,b) where+  data Dist  (a,b)   = DProd  !(Dist a)    !(Dist b)+  data MDist (a,b) s = MDProd !(MDist a s) !(MDist b s)++  indexD d i               = (fstD d `indexD` i,sndD d `indexD` i)+  newMD g                  = liftM2 MDProd (newMD g) (newMD g)+  readMD  (MDProd xs ys) i = liftM2 (,) (readMD xs i) (readMD ys i)+  writeMD (MDProd xs ys) i (x,y)+                            = writeMD xs i x >> writeMD ys i y+  unsafeFreezeMD (MDProd xs ys)+                            = liftM2 DProd (unsafeFreezeMD xs)+                                           (unsafeFreezeMD ys)++  {-# INLINE deepSeqD #-}+  deepSeqD (x,y) z = deepSeqD x (deepSeqD y z)++  sizeD  (DProd  x _) = sizeD  x+  sizeMD (MDProd x _) = sizeMD x++  measureD (x,y) = "Pair " ++ "(" ++ measureD x ++ ") (" ++  measureD y ++ ")"+++-- | Pairing of distributed values.+-- /The two values must belong to the same/ 'Gang'.+zipD :: (DT a, DT b) => Dist a -> Dist b -> Dist (a,b)+{-# INLINE [0] zipD #-}+zipD !x !y = checkEq (here "zipDT") "Size mismatch" (sizeD x) (sizeD y) $+             DProd x y++-- | Unpairing of distributed values.+unzipD :: (DT a, DT b) => Dist (a,b) -> (Dist a, Dist b)+{-# INLINE_DIST unzipD #-}+unzipD (DProd dx dy) = (dx,dy)++-- | Extract the first elements of a distributed pair.+fstD :: (DT a, DT b) => Dist (a,b) -> Dist a+{-# INLINE_DIST fstD #-}+fstD = fst . unzipD++-- | Extract the second elements of a distributed pair.+sndD :: (DT a, DT b) => Dist (a,b) -> Dist b+{-# INLINE_DIST sndD #-}+sndD = snd . unzipD+++-- Maybe ----------------------------------------------------------------------+instance DT a => DT (Maybe a) where+  data Dist  (Maybe a)   = DMaybe  !(Dist  Bool)   !(Dist  a)+  data MDist (Maybe a) s = MDMaybe !(MDist Bool s) !(MDist a s)++  indexD (DMaybe bs as) i+    | bs `indexD` i       = Just $ as `indexD` i+    | otherwise           = Nothing+  newMD g = liftM2 MDMaybe (newMD g) (newMD g)+  readMD (MDMaybe bs as) i =+    do+      b <- readMD bs i+      if b then liftM Just $ readMD as i+           else return Nothing+  writeMD (MDMaybe bs as) i Nothing  = writeMD bs i False+  writeMD (MDMaybe bs as) i (Just x) = writeMD bs i True+                                     >> writeMD as i x+  unsafeFreezeMD (MDMaybe bs as) = liftM2 DMaybe (unsafeFreezeMD bs)+                                                 (unsafeFreezeMD as)++  {-# INLINE deepSeqD #-}+  deepSeqD Nothing  z = z+  deepSeqD (Just x) z = deepSeqD x z++  sizeD  (DMaybe  b _) = sizeD  b+  sizeMD (MDMaybe b _) = sizeMD b++  measureD Nothing = "Nothing"+  measureD (Just x) = "Just (" ++ measureD x ++ ")"+++-- Vector ---------------------------------------------------------------------+instance Unbox a => DT (V.Vector a) where+  data Dist  (Vector a)   = DVector  !(Dist  Int)   !(BV.Vector      (Vector a))+  data MDist (Vector a) s = MDVector !(MDist Int s) !(MBV.STVector s (Vector a))++  indexD (DVector _ a) i = a BV.! i+  newMD g = liftM2 MDVector (newMD g) (MBV.replicate (gangSize g)+                                         (error "MDist (Vector a) - uninitalised"))+  readMD (MDVector _ marr) = MBV.read marr+  writeMD (MDVector mlen marr) i a =+    do+      writeMD mlen i (V.length a)+      MBV.write marr i $! a+  unsafeFreezeMD (MDVector len a) = liftM2 DVector (unsafeFreezeMD len)+                                               (BV.unsafeFreeze a)+  sizeD  (DVector  _ a) = BV.length  a+  sizeMD (MDVector _ a) = MBV.length a++  measureD xs = "Vector " ++ show (V.length xs)+++-- | Yield the distributed length of a distributed array.+lengthD :: Unbox a => Dist (Vector a) -> Dist Int+lengthD (DVector l _) = l+++-- USegd ----------------------------------------------------------------------+instance DT USegd where+  data Dist  USegd   = DUSegd  !(Dist (Vector Int))+                               !(Dist (Vector Int))+                               !(Dist Int)+  data MDist USegd s = MDUSegd !(MDist (Vector Int) s)+                               !(MDist (Vector Int) s)+                               !(MDist Int        s)++  indexD (DUSegd lens idxs eles) i+          = mkUSegd (indexD lens i) (indexD idxs i) (indexD eles i)+  newMD g = liftM3 MDUSegd (newMD g) (newMD g) (newMD g)+  readMD (MDUSegd lens idxs eles) i+          = liftM3 mkUSegd (readMD lens i) (readMD idxs i) (readMD eles i)+  writeMD (MDUSegd lens idxs eles) i segd+          = do+              writeMD lens i (lengthsUSegd  segd)+              writeMD idxs i (indicesUSegd  segd)+              writeMD eles i (elementsUSegd segd)+  unsafeFreezeMD (MDUSegd lens idxs eles)+          = liftM3 DUSegd (unsafeFreezeMD lens)+                          (unsafeFreezeMD idxs)+                          (unsafeFreezeMD eles)++  deepSeqD segd z = deepSeqD (lengthsUSegd  segd)+                  $ deepSeqD (indicesUSegd  segd)+                  $ deepSeqD (elementsUSegd segd) z++  sizeD  (DUSegd  _ _ eles) = sizeD eles+  sizeMD (MDUSegd _ _ eles) = sizeMD eles++  measureD segd = "Segd " ++ show (lengthUSegd segd) ++ " " ++ show (elementsUSegd segd)++lengthUSegdD :: Dist USegd -> Dist Int+{-# INLINE_DIST lengthUSegdD #-}+lengthUSegdD (DUSegd lens _ _) = lengthD lens++lengthsUSegdD :: Dist USegd -> Dist (Vector Int)+{-# INLINE_DIST lengthsUSegdD #-}+lengthsUSegdD (DUSegd lens _ _ ) = lens++indicesUSegdD :: Dist USegd -> Dist (Vector Int)+{-# INLINE_DIST indicesUSegdD #-}+indicesUSegdD (DUSegd _ idxs _) = idxs++elementsUSegdD :: Dist USegd -> Dist Int+{-# INLINE_DIST elementsUSegdD #-}+elementsUSegdD (DUSegd _ _ dns) = dns+
+ Data/Array/Parallel/Unlifted/Parallel.hs view
@@ -0,0 +1,39 @@+-- | Parallel operations on unlifted arrays+module Data.Array.Parallel.Unlifted.Parallel (+  UPSegd, UPSel2, UPSelRep2,++  bpermuteUP, updateUP,++  enumFromToUP, enumFromThenToUP, enumFromStepLenUP, enumFromStepLenEachUP,++  mapUP, filterUP, packUP, combineUP, combine2UP,+  zipWithUP, foldUP, scanUP,++  andUP, sumUP,++  tagsUPSel2, indicesUPSel2, elementsUPSel2_0, elementsUPSel2_1,+  selUPSel2, repUPSel2, mkUPSel2,+  mkUPSelRep2, indicesUPSelRep2, elementsUPSelRep2_0, elementsUPSelRep2_1,++  lengthUPSegd, lengthsUPSegd, indicesUPSegd, elementsUPSegd,+  segdUPSegd, distUPSegd,+  lengthsToUPSegd, mkUPSegd,+ +  replicateSUP, replicateRSUP, appendSUP, indicesSUP,+  foldSUP, foldRUP, fold1SUP, sumSUP, sumRUP,++  indexedUP, replicateUP, repeatUP, interleaveUP,++  dropUP+) where+import Data.Array.Parallel.Unlifted.Parallel.Permute+import Data.Array.Parallel.Unlifted.Parallel.Combinators+import Data.Array.Parallel.Unlifted.Parallel.Basics+import Data.Array.Parallel.Unlifted.Parallel.Sums+import Data.Array.Parallel.Unlifted.Parallel.Enum+import Data.Array.Parallel.Unlifted.Parallel.Segmented+import Data.Array.Parallel.Unlifted.Parallel.Subarrays+import Data.Array.Parallel.Unlifted.Parallel.UPSegd+import Data.Array.Parallel.Unlifted.Parallel.UPSel+import Data.Array.Parallel.Unlifted.Parallel.Text ()+
+ Data/Array/Parallel/Unlifted/Parallel/Basics.hs view
@@ -0,0 +1,67 @@+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Basic operations on parallel unlifted arrays.+module Data.Array.Parallel.Unlifted.Parallel.Basics (+  lengthUP, nullUP, indexedUP,+  replicateUP, repeatUP, interleaveUP+) where++import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Parallel.Combinators ( mapUP )+import Data.Array.Parallel.Unlifted.Parallel.Enum        ( enumFromToUP )+import Data.Array.Parallel.Unlifted.Parallel.Permute     ( bpermuteUP )++import GHC.Base ( remInt )++-- NOTE: some of the functions are exactly the same as the U version++-- | Test whether the given array is empty+nullUP :: Unbox e => Vector e -> Bool+nullUP  = (== 0) . Seq.length++-- | Yield an empty array+emptyUP :: Unbox e => Vector e+emptyUP = Seq.new 0 (const $ return ())++lengthUP :: Unbox e => Vector e -> Int+lengthUP = Seq.length+++-- | Yield an array where all elements contain the same value+replicateUP :: Unbox e => Int -> e -> Vector e+{-# INLINE_UP replicateUP #-}+replicateUP n !e = joinD theGang balanced+                 . mapD theGang (\n ->Seq.replicate n e)+                 $ splitLenD theGang n+++-- | Repeat an array the given number of times.+repeatUP :: Unbox e => Int -> Vector e -> Vector e+{-# INLINE_UP repeatUP #-}+repeatUP n es = seq m+              . bpermuteUP es+              . mapUP (\i -> i `remInt` m)+              $ enumFromToUP 0 (m*n-1)+  where+    m = Seq.length es++-- | Interleave elements of two arrays+interleaveUP :: Unbox e => Vector e -> Vector e -> Vector e+{-# INLINE_UP interleaveUP #-}+interleaveUP xs ys = joinD theGang unbalanced+                     (zipWithD theGang Seq.interleave+                       (splitD theGang balanced xs)+                       (splitD theGang balanced ys))+   +-- | Associate each element of the array with its index+indexedUP :: (DT e, Unbox e) => Vector e -> Vector (Int,e)+{-# INLINE_U indexedUP #-}+indexedUP = splitJoinD theGang indexedFn +  where+    sizes  arr   = fst $ scanD theGang (+) 0 $ lengthD arr+    indexedFn = \arr -> (zipWithD theGang (\o -> Seq.map (\(x,y) -> (x + o, y))) (sizes arr) $ +                        mapD theGang Seq.indexed arr)++
+ Data/Array/Parallel/Unlifted/Parallel/Combinators.hs view
@@ -0,0 +1,164 @@+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Parallel combinators for unlifted arrays+module Data.Array.Parallel.Unlifted.Parallel.Combinators (+  mapUP, filterUP, packUP, combineUP, combine2UP,+  zipWithUP, foldUP, fold1UP, foldl1UP, scanUP+) where++import Data.Array.Parallel.Base+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Parallel.UPSel+++-- | Apply a worker to all elements of a vector.+mapUP :: (Unbox a, Unbox b) => (a -> b) -> Vector a -> Vector b+{-# INLINE mapUP #-}+mapUP f xs +        = splitJoinD theGang (mapD theGang (Seq.map f)) xs+++-- | Keep elements that match the given predicate.+filterUP :: Unbox a => (a -> Bool) -> Vector a -> Vector a+{-# INLINE filterUP #-}+filterUP f+        = joinD  theGang unbalanced+        . mapD   theGang (Seq.filter f)+        . splitD theGang unbalanced+++-- | Take elements of an array where a flag value is true, and pack them into+--   the result. +-- +--   * The souce and flag arrays must have the same length, but this is not checked.+--+packUP :: Unbox e => Vector e -> Vector Bool -> Vector e+{-# INLINE_UP packUP #-}+packUP xs flags +        = Seq.fsts . filterUP snd $ Seq.zip xs flags+++-- | Combine two vectors based on a selector. +--   If the selector is true then take the element from the first vector, +--   otherwise take it from the second.+--+--   * The data vectors must have enough elements to satisfy the flag vector, +--     but this is not checked.+--  +combineUP :: Unbox a => Vector Bool -> Vector a -> Vector a -> Vector a+{-# INLINE combineUP #-}+combineUP flags xs ys +        = combine2UP tags (mkUPSelRep2 tags) xs ys+        where tags = Seq.map (fromBool . not) flags+++-- | Combine two vectors based on a selector. +--+--   * The data vectors must have enough elements to satisfy the selector,+--     but this is not checked.+--+--   TODO: What is the difference between the Tag and the UPSelRep2?+--+combine2UP :: Unbox a => Vector Tag -> UPSelRep2 -> Vector a -> Vector a -> Vector a+{-# INLINE_UP combine2UP #-}+combine2UP tags rep !xs !ys +        = joinD    theGang balanced+        $ zipWithD theGang go rep+        $ splitD   theGang balanced tags+        where   go ((i,j), (m,n)) ts +                 = Seq.combine2ByTag ts +                        (Seq.slice xs i m)+                        (Seq.slice ys j n)+    ++-- | Combine two vectors into a third.+zipWithUP :: (Unbox a, Unbox b, Unbox c) +          => (a -> b -> c) -> Vector a -> Vector b -> Vector c+{-# INLINE zipWithUP #-}+zipWithUP f xs ys+        = splitJoinD theGang +                (mapD theGang (Seq.map (uncurry f))) +                (Seq.zip xs ys)+++-- | Undirected fold.+--   Note that this function has more constraints on its parameters than the+--   standard fold function from the Haskell Prelude.+--+--   * The worker function must be associative.+--   * The provided starting element must be neutral with respect to the worker.+--     For example 0 is neutral wrt (+) and 1 is neutral wrt (*).+--+--   We need these constraints so that we can partition the fold across +--   several threads. Each thread folds a chunk of the input vector, +--   then we fold together all the results in the main thread.+--+foldUP  :: (Unbox a, DT a) => (a -> a -> a) -> a -> Vector a -> a+{-# INLINE foldUP #-}+foldUP f !z xs+        = foldD theGang f+                (mapD   theGang (Seq.fold f z)+                (splitD theGang unbalanced xs))+++-- | Left fold over an array. +--+--   * If the vector is empty then this returns the provided neural element.+--   * The worker function must be associative.+--   * The provided starting element must be neutral with respect to the worker,+--     see `foldUP` for discussion.+--+foldlUP :: (DT a, Unbox a) => (a -> a -> a) -> a -> Vector a -> a+{-# INLINE_UP foldlUP #-}+foldlUP f z arr +  | Seq.null arr = z+  | otherwise    = foldl1UP f arr+++-- | Alias for `foldl1UP`+fold1UP :: (DT a, Unbox a) => (a -> a -> a) -> Vector a -> a+{-# INLINE fold1UP #-}+fold1UP = foldl1UP+++-- | Left fold over an array, using the first element of the vector as the+--   neural element.+--+--   * If the vector contains no elements then you'll get a bounds-check error.+--   * The worker function must be associative.+--   * The provided starting element must be neutral with respect to the worker,+--     see `foldUP` for discussion.+--+--   TODO: The two type class constraints are in a different order. Does that matter?+--+foldl1UP :: (DT a, Unbox a) => (a -> a -> a) -> Vector a -> a+{-# INLINE_U foldl1UP #-}+foldl1UP f arr +        = (maybe z (f z)+        . foldD  theGang combine+        . mapD   theGang (Seq.foldl1Maybe f)+        . splitD theGang unbalanced) arr+        where+                z = arr ! 0+                combine (Just x) (Just y) = Just (f x y)+                combine (Just x) Nothing  = Just x+                combine Nothing  (Just y) = Just y+                combine Nothing  Nothing  = Nothing+++-- | Prefix scan. Similar to fold, but produce an array of the intermediate states.+--+--   * The worker function must be associative.+--   * The provided starting element must be neutral with respect to the worker,+--     see `foldUP` for discussion.+--+scanUP :: (DT a, Unbox a) => (a -> a -> a) -> a -> Vector a -> Vector a+{-# INLINE_UP scanUP #-}+scanUP f z +        = splitJoinD theGang go+        where   go xs = let (ds,zs) = unzipD $ mapD theGang (Seq.scanRes f z) xs+                            zs'     = fst (scanD theGang f z zs)+                        in  zipWithD theGang (Seq.map . f) zs' ds+
+ Data/Array/Parallel/Unlifted/Parallel/Enum.hs view
@@ -0,0 +1,65 @@+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Enum-related parallel operations on unlifted arrays+module Data.Array.Parallel.Unlifted.Parallel.Enum (+  enumFromToUP, enumFromThenToUP, enumFromStepLenUP, enumFromStepLenEachUP    +) where+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Distributed (+  mapD, scanD, zipD, splitLenIdxD, joinD, splitD, balanced, unbalanced,+  theGang)+import Data.Array.Parallel.Unlifted.Parallel.Combinators (+  mapUP)+import GHC.Base ( divInt )+++delay_inline :: a -> a+{-# INLINE [0] delay_inline #-}+delay_inline x = x+++enumFromToUP :: (Unbox a, Enum a) => a -> a -> Vector a+{-# INLINE enumFromToUP #-}+enumFromToUP start end = mapUP toEnum (enumFromStepLenUP start' 1 len)+  where+    start' = fromEnum start+    end'   = fromEnum end+    len    = delay_inline max (end' - start' + 1) 0+++enumFromThenToUP :: (Unbox a, Enum a) => a -> a -> a -> Vector a+{-# INLINE enumFromThenToUP #-}+enumFromThenToUP start next end =+  mapUP toEnum (enumFromStepLenUP start' delta len)+  where+    start' = fromEnum start+    next'  = fromEnum next+    end'   = fromEnum end+    delta  = next' - start'+    -- distance between start' and end' expressed in deltas+    dist   = (end' - start' + delta) `divInt` delta+    len    = max dist 0+++enumFromStepLenUP :: Int -> Int -> Int -> Vector Int+{-# INLINE enumFromStepLenUP #-}+enumFromStepLenUP start delta len =+  joinD theGang balanced+  (mapD theGang gen+  (splitLenIdxD theGang len))+  where+    gen (n,i) = Seq.enumFromStepLen (i * delta + start) delta n+++enumFromStepLenEachUP :: Int -> Vector Int -> Vector Int -> Vector Int -> Vector Int+{-# INLINE enumFromStepLenEachUP #-}+enumFromStepLenEachUP n starts steps lens+  = joinD theGang unbalanced+  $ mapD theGang enum+  $ splitD theGang unbalanced (Seq.zip (Seq.zip starts steps) lens)+  where+    enum ps = let (qs, llens) = Seq.unzip ps+                  (lstarts, lsteps) = Seq.unzip qs+              in Seq.enumFromStepLenEach (Seq.sum llens) lstarts lsteps llens+
+ Data/Array/Parallel/Unlifted/Parallel/Permute.hs view
@@ -0,0 +1,42 @@+{-# LANGUAGE ScopedTypeVariables #-}++-- | Parallel permutations for unlifted arrays+module Data.Array.Parallel.Unlifted.Parallel.Permute (+  bpermuteUP, updateUP+) where+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Distributed+++bpermuteUP :: Unbox a => Vector a -> Vector Int -> Vector a+{-# INLINE bpermuteUP #-}+bpermuteUP as is = splitJoinD theGang (bpermuteD theGang as) is++{-+  We can't support this for arbitrary types. The problem is:+  what happens if the second array maps multiple elements to the same position?+  I don't know what the semantics is supposed to be in Nesl, the spec don't+  seem to say anything. Note that it is not sufficient to say, e.g., that it+  is unspecified which value gets written; if we have an array of pairs,+  for instance, we might well get the first and second components from+  different values.++  We could require that the second array maps at most one element to each index.+  However, this is not what is wanted most of the time, at least not in the+  algorithms I've seen.++  So we only do the update in parallel if writing an element into the array is+  atomic. Otherwise, we do a sequential update.+-}++updateUP :: forall a. Unbox a => Vector a -> Vector (Int,a) -> Vector a+{-# INLINE updateUP #-}+updateUP as us+  {- hasAtomicWriteMU (undefined :: a) +  = atomicUpdateD theGang (splitD theGang unbalanced as)+                          (splitD theGang unbalanced us)+  -}++  | otherwise+  = Seq.update as us+
+ Data/Array/Parallel/Unlifted/Parallel/Segmented.hs view
@@ -0,0 +1,213 @@+{-# LANGUAGE CPP #-}++#include "fusion-phases.h"++-- | Parallel combinators for segmented unboxed arrays+module Data.Array.Parallel.Unlifted.Parallel.Segmented (+  replicateSUP, replicateRSUP, appendSUP, indicesSUP,+  foldSUP, foldRUP, fold1SUP, sumSUP, sumRUP+) where+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Sequential.Segmented+import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Parallel.Combinators (+  mapUP, zipWithUP, packUP, combineUP)+import Data.Array.Parallel.Unlifted.Parallel.Sums (+  sumUP )+import Data.Array.Parallel.Unlifted.Parallel.Basics (+  replicateUP, repeatUP)+import Data.Array.Parallel.Unlifted.Parallel.Enum+import Data.Array.Parallel.Unlifted.Parallel.Permute ( bpermuteUP )+import Data.Array.Parallel.Unlifted.Parallel.UPSegd+import qualified Data.Vector.Fusion.Stream as S+import Data.Vector.Fusion.Stream.Monadic ( Stream(..), Step(..) )+import Data.Vector.Fusion.Stream.Size    ( Size(..) )+import Control.Monad.ST ( ST, runST )+++-- replicate ------------------------------------------------------------------++-- | Segmented replication, using a segment descriptor.+replicateSUP :: Unbox a => UPSegd -> Vector a -> Vector a+{-# INLINE_UP replicateSUP #-}+replicateSUP segd !xs +  = joinD theGang balanced+  . mapD theGang rep+  $ distUPSegd segd+  where+    rep ((dsegd,di),_)+      = replicateSU dsegd (Seq.slice xs di (lengthUSegd dsegd))+++-- | Segmented replication.+--   Each element in the vector is replicated the given number of times.+--   +--   @replicateRSUP 2 [1, 2, 3, 4, 5] = [1, 1, 2, 2, 3, 3, 4, 4, 5, 5]@+--+--   TODO: make this efficient+-- +replicateRSUP :: Unbox a => Int -> Vector a -> Vector a+{-# INLINE_UP replicateRSUP #-}+replicateRSUP n xs+        = replicateSUP (lengthsToUPSegd (replicateUP (Seq.length xs) n)) xs+++-- append ---------------------------------------------------------------------+-- | Segmented append.+appendSUP+        :: Unbox a+        => UPSegd       -- ^ segment descriptor of result array+        -> UPSegd       -- ^ segment descriptor of first array+        -> Vector a     -- ^ data of first array+        -> UPSegd       -- ^ segment descriptor of second array+        -> Vector a     -- ^ data of first array+        -> Vector a++{-# INLINE_UP appendSUP #-}+appendSUP segd !xd !xs !yd !ys+  = joinD theGang balanced+  . mapD  theGang append+  $ distUPSegd segd+  where append ((segd,seg_off),el_off)+         = Seq.unstream+         $ appendSegS (segdUPSegd xd) xs+                      (segdUPSegd yd) ys+                      (elementsUSegd segd)+                      seg_off el_off+++appendSegS+        :: Unbox a      +        => USegd        -- ^ segment descriptor of first array+        -> Vector a     -- ^ data of first array+        -> USegd        -- ^ segment descriptor of second array+        -> Vector a     -- ^ data of second array+        -> Int          -- +        -> Int+        -> Int+        -> S.Stream a++{-# INLINE_STREAM appendSegS #-}+appendSegS !xd !xs !yd !ys !n seg_off el_off+  = Stream next state (Exact n)+  where+    !xlens = lengthsUSegd xd+    !ylens = lengthsUSegd yd++    state+      | n == 0 = Nothing+      | el_off < xlens ! seg_off+      = let i = (indicesUSegd xd ! seg_off) + el_off+            j = indicesUSegd yd ! seg_off+            k = (lengthsUSegd xd ! seg_off) - el_off+        in  Just (False, seg_off, i, j, k, n)++      | otherwise+      = let -- NOTE: *not* indicesUSegd xd ! (seg_off+1) since seg_off+1+            -- might be out of bounds+            i       = (indicesUSegd xd ! seg_off) + (lengthsUSegd xd ! seg_off)+            el_off' = el_off - lengthsUSegd xd ! seg_off+            j       = (indicesUSegd yd ! seg_off) + el_off'+            k       = (lengthsUSegd yd ! seg_off) - el_off'+        in  Just (True, seg_off, i, j, k, n)++    {-# INLINE next #-}+    next Nothing = return Done++    next (Just (False, seg, i, j, k, n))+      | n == 0    = return Done+      | k == 0    = return $ Skip (Just (True, seg, i, j, ylens ! seg, n))+      | otherwise = return $ Yield (xs!i) (Just (False, seg, i+1, j, k-1, n-1))++    next (Just (True, seg, i, j, k, n))+      | n == 0    = return Done+      | k == 0+      = let !seg' = seg+1+        in  return $ Skip (Just (False, seg', i, j, xlens ! seg', n))++      | otherwise = return $ Yield (ys!j) (Just (True, seg, i, j+1, k-1, n-1))+++-- fold -----------------------------------------------------------------------+fixupFold :: Unbox a => (a -> a -> a) -> MVector s a+          -> Dist (Int,Vector a) -> ST s ()+{-# NOINLINE fixupFold #-}+fixupFold f !mrs !dcarry = go 1+  where+    !p = gangSize theGang++    go i | i >= p = return ()+         | Seq.null c = go (i+1)+         | otherwise   = do+                           x <- Seq.read mrs k+                           Seq.write mrs k (f x (c ! 0))+                           go (i+1)+      where+        (k,c) = indexD dcarry i+++folds :: Unbox a => (a -> a -> a)+              -> (USegd -> Vector a -> Vector a) -> UPSegd -> Vector a -> Vector a+{-# INLINE folds #-}+folds f g segd xs = dcarry `seq` drs `seq` runST (+  do+    mrs <- joinDM theGang drs+    fixupFold f mrs dcarry+    Seq.unsafeFreeze mrs)+  where+    (dcarry,drs)+          = unzipD+          $ mapD theGang partial+          $ zipD (distUPSegd segd)+                 (splitD theGang balanced xs)++    partial (((segd,k),off), as)+      = let rs = g segd as+            {-# INLINE [0] n #-}+            n | off == 0  = 0+              | otherwise = 1+        in+        ((k, Seq.take n rs), Seq.drop n rs)+++foldSUP :: Unbox a => (a -> a -> a) -> a -> UPSegd -> Vector a -> Vector a+{-# INLINE foldSUP #-}+foldSUP f !z = folds f (foldlSU f z)+++fold1SUP :: Unbox a => (a -> a -> a) -> UPSegd -> Vector a -> Vector a+{-# INLINE fold1SUP #-}+fold1SUP f = folds f (fold1SU f)+++sumSUP :: (Num e, Unbox e) => UPSegd -> Vector e -> Vector e+{-# INLINE sumSUP #-}+sumSUP = foldSUP (+) 0+++sumRUP :: (Num e, Unbox e) => Int -> Vector e -> Vector e+{-# INLINE sumRUP #-}+sumRUP = foldRUP (+) 0+++foldRUP :: (Unbox a, Unbox b) => (b -> a -> b) -> b -> Int -> Vector a -> Vector b+{-# INLINE foldRUP #-}+foldRUP f z !segSize xs = +   joinD theGang unbalanced+    (mapD theGang +      (foldlRU f z segSize)+      (splitAsD theGang (mapD theGang (*segSize) dlen) xs))+  where+    noOfSegs = Seq.length xs `div` segSize+    dlen = splitLenD theGang noOfSegs+++-- indices --------------------------------------------------------------------+indicesSUP :: UPSegd -> Vector Int+{-# INLINE_UP indicesSUP #-}+indicesSUP = joinD theGang balanced+           . mapD theGang indices+           . distUPSegd+  where+    indices ((segd,k),off) = indicesSU' off segd+
+ Data/Array/Parallel/Unlifted/Parallel/Subarrays.hs view
@@ -0,0 +1,15 @@+{-# LANGUAGE CPP #-}++#include "fusion-phases.h"++-- | Subarrays of flat unlifted arrays.+module Data.Array.Parallel.Unlifted.Parallel.Subarrays (+  dropUP+) where+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Distributed+++dropUP :: Unbox e => Int -> Vector e -> Vector e+dropUP n xs = Seq.slice xs (min (max 0 n) (Seq.length xs)) (min (Seq.length xs) (Seq.length xs - n)) +{-# INLINE_U dropUP #-}
+ Data/Array/Parallel/Unlifted/Parallel/Sums.hs view
@@ -0,0 +1,68 @@+-- | Sum-like parallel combinators for unlifted arrays+module Data.Array.Parallel.Unlifted.Parallel.Sums (+  andUP, orUP, sumUP+) where++import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Parallel.Combinators (+  foldUP, foldl1UP, fold1UP, mapUP)+import Data.Array.Parallel.Unlifted.Parallel.Basics ( +  indexedUP)++-- | Compute the logical AND of all the elements in a array.+andUP :: Vector Bool -> Bool+{-# INLINE andUP #-}+andUP = foldUP (&&) True+++-- | Compute the logical OR of all the elements in a array.+orUP :: Vector Bool -> Bool+{-# INLINE orUP #-}+orUP = foldUP (||) False++-- | Check whether all the elements in a array meet the given predicate.+allUP :: Unbox e => (e -> Bool) -> Vector e -> Bool+{-# INLINE allUP #-}+allUP p = andUP . mapUP p++-- | Check whether any of the elements in a array meet the given predicate.+anyUP :: Unbox e => (e -> Bool) -> Vector e -> Bool+{-# INLINE anyUP #-}+anyUP p =  orUP . mapUP p+++-- | Compute the sum all the elements of a array.+sumUP :: (Unbox a, DT a, Num a) => Vector a -> a+{-# INLINE sumUP #-}+sumUP = foldUP (+) 0+++-- | Compute the product of all the elements of an array.+productUP :: (DT e, Num e, Unbox e) => Vector e -> e+{-# INLINE productUP #-}+productUP = foldUP (*) 1+++-- | Determine the maximum element in an array.+maximumUP :: (DT e, Ord e, Unbox e) => Vector e -> e+{-# INLINE maximumUP #-}+maximumUP = fold1UP max+++--  |Determine the maximum element in an array under the given ordering+maximumByUP :: (DT e, Unbox e) => (e -> e -> Ordering) -> Vector e -> e+{-# INLINE maximumByUP #-}+maximumByUP = fold1UP . maxBy+  where+    maxBy compare x y = case x `compare` y of+                          LT -> y+                          _  -> x++-- | Determine the index of the maximum element in an array under the given+--   ordering+maximumIndexByUP :: (DT e, Unbox e) => (e -> e -> Ordering) -> Vector e -> Int+{-# INLINE maximumIndexByUP #-}+maximumIndexByUP cmp = fst . maximumByUP cmp' . indexedUP+  where+    cmp' (_,x) (_,y) = cmp x y
+ Data/Array/Parallel/Unlifted/Parallel/Text.hs view
@@ -0,0 +1,11 @@+-- | Read\/Show instances for segmented unlifted arrays.+module Data.Array.Parallel.Unlifted.Parallel.Text ()+where++import Data.Array.Parallel.Base (+  Read(..), showsApp)+import Data.Array.Parallel.Unlifted.Parallel.UPSegd (+  UPSegd, lengthsUPSegd )++instance Show UPSegd where+  showsPrec k = showsApp k "toUPSegd" . lengthsUPSegd
+ Data/Array/Parallel/Unlifted/Parallel/UPSegd.hs view
@@ -0,0 +1,68 @@+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Parallel segment descriptors.+module Data.Array.Parallel.Unlifted.Parallel.UPSegd (++  -- * Types+  UPSegd,++  -- * Operations on segment descriptors+  lengthUPSegd, lengthsUPSegd, indicesUPSegd, elementsUPSegd,+  segdUPSegd, distUPSegd,+  lengthsToUPSegd, mkUPSegd+) where++import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Sequential.Segmented.USegd+import Data.Array.Parallel.Unlifted.Distributed++data UPSegd = UPSegd { upsegd_usegd :: !USegd+                     , upsegd_dsegd :: Dist ((USegd,Int),Int)+                     }+++lengthUPSegd :: UPSegd -> Int+{-# INLINE lengthUPSegd #-}+lengthUPSegd = lengthUSegd . upsegd_usegd+++lengthsUPSegd :: UPSegd -> Vector Int+{-# INLINE lengthsUPSegd #-}+lengthsUPSegd = lengthsUSegd . upsegd_usegd+++indicesUPSegd :: UPSegd -> Vector Int+{-# INLINE indicesUPSegd #-}+indicesUPSegd = indicesUSegd . upsegd_usegd+++elementsUPSegd :: UPSegd -> Int+{-# INLINE elementsUPSegd #-}+elementsUPSegd = elementsUSegd . upsegd_usegd+++segdUPSegd :: UPSegd -> USegd+{-# INLINE segdUPSegd #-}+segdUPSegd = upsegd_usegd+++distUPSegd :: UPSegd -> Dist ((USegd,Int),Int)+{-# INLINE distUPSegd #-}+distUPSegd = upsegd_dsegd+++lengthsToUPSegd :: Vector Int -> UPSegd+{-# INLINE lengthsToUPSegd #-}+lengthsToUPSegd = toUPSegd . lengthsToUSegd+++mkUPSegd :: Vector Int -> Vector Int -> Int -> UPSegd+{-# INLINE mkUPSegd #-}+mkUPSegd lens idxs n = toUPSegd (mkUSegd lens idxs n)+++toUPSegd :: USegd -> UPSegd+{-# INLINE toUPSegd #-}+toUPSegd segd = UPSegd segd (splitSegdD' theGang segd)+
+ Data/Array/Parallel/Unlifted/Parallel/UPSel.hs view
@@ -0,0 +1,170 @@+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Parallel selectors.+module Data.Array.Parallel.Unlifted.Parallel.UPSel (+  -- * Types+  UPSel2, UPSelRep2,++  -- * Operations on segment descriptors+  tagsUPSel2, indicesUPSel2, elementsUPSel2_0, elementsUPSel2_1,+  selUPSel2, repUPSel2, mkUPSel2,+  mkUPSelRep2, indicesUPSelRep2, elementsUPSelRep2_0, elementsUPSelRep2_1,+) where+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Data.Array.Parallel.Unlifted.Sequential.USel+import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Base (Tag, tagToInt)+++-- | Contains a selector `USel2`, as well as an `USelRep2` which says how +--   to distribute this selector across the PEs. +--+--   See @dph-prim-seq:Data.Array.Parallel.Unlifted.Sequential.Segmented.USel@+--   for more discussion of what selectors are for.+--+data UPSel2 +        = UPSel2 +        { upsel2_usel :: USel2+        , upsel2_rep  :: UPSelRep2 }+++-- | A `UPSelRep2` describes how to distribute the two data vectors+--   corresponding to a `UPSel2` across several PEs.+--+--   Suppose we want to perform the following combine operation:+--+-- @+--    combine [0,0,1,1,0,1,0,0,1] [A0,A1,A2,A3,A4] [B0,B1,B2,B3] +--     = [A0,A1,B0,B1,A2,B2,A3,A4,B3]+-- @+--+--   The first array is the tags array, that says which of the data arrays to+--   get each successive element from. As `combine` is difficult to compute+--   in parallel, if we are going to perform several combines with the same+--   tag array, we can precompute a selector that tells us where to get each+--   element. The selector contains the original tags, as well as the source+--   index telling us where to get each element for the result array.+-- +-- @+--    [0,0,1,1,0,1,0,0,1]      -- tags    (which data vector to take the elem from)+--    [0,1,0,1,2,2,3,4,3]      -- indices (where in the vector to take the elem from)+-- @+--+--  Suppose we want to distribute the combine operation across 3 PEs. It's+--  easy to split the selector like so:+--+-- @       +--     PE0                PE1               PE2+--    [0,0,1]            [1,0,1]           [0,0,1]   -- tags+--    [0,1,0]            [1,2,2]           [3,4,3]   -- indices+-- @+--+--  We now need to split the two data arrays. Each PE needs slices of the data+--  arrays that correspond to the parts of the selector that were given to it.+--  For the current example we get:+--+-- @+--    PE0                PE1               PE2+--    [A0,A1]            [A2]              [A3,A4]+--    [B0]               [B1,B2]           [B3]+-- @+--+--  The `UPSelRep2` contains the starting index and length of each of of these+--  slices:+--+-- @+--         PE0                PE1               PE2+--    ((0, 0), (2, 1))   ((2, 1), (1, 2))  ((3, 3), (2, 1))+--    indices   lens      indices  lens    indices  lens+-- @+--+type UPSelRep2+        = Dist ((Int,Int), (Int,Int))+++-- | O(1). Get the tags of a selector.+tagsUPSel2 :: UPSel2 -> Vector Tag+{-# INLINE tagsUPSel2 #-}+tagsUPSel2 = tagsUSel2 .  upsel2_usel+++-- | O(1). Get the indices of a selector.+indicesUPSel2 :: UPSel2 -> Vector Int+{-# INLINE indicesUPSel2 #-}+indicesUPSel2 = indicesUSel2 . upsel2_usel+++-- | O(1). TODO: What is this for?+elementsUPSel2_0 :: UPSel2 -> Int+{-# INLINE elementsUPSel2_0 #-}+elementsUPSel2_0 = elementsUSel2_0 . upsel2_usel+++-- | O(1). TODO: What is this for?+elementsUPSel2_1 :: UPSel2 -> Int+{-# INLINE elementsUPSel2_1 #-}+elementsUPSel2_1 = elementsUSel2_1 . upsel2_usel+++-- | O(1). TODO: What is this for?+selUPSel2 :: UPSel2 -> USel2+{-# INLINE selUPSel2 #-}+selUPSel2 = upsel2_usel+++-- | O(1). TODO: What is this for?+repUPSel2 :: UPSel2 -> UPSelRep2+{-# INLINE repUPSel2 #-}+repUPSel2 = upsel2_rep+++-- Representation selectors ---------------------------------------------------++-- | Computes a `UPSelRep2` from an array of tags. This is used when parallelising+--   a `combine` operation. See the docs for `UPSelRep2` for details.+mkUPSelRep2 :: Vector Tag -> UPSelRep2+{-# INLINE mkUPSelRep2 #-}+mkUPSelRep2 tags = zipD idxs lens+  where+    lens = mapD   theGang count+         $ splitD theGang balanced tags++    idxs = fst+         $ scanD theGang add (0,0) lens++    count bs = let ones = Seq.sum (Seq.map tagToInt bs)+               in (Seq.length bs - ones,ones)++    add (x1,y1) (x2,y2) = (x1+x2, y1+y2)+++indicesUPSelRep2 :: Vector Tag -> UPSelRep2 -> Vector Int+{-# INLINE indicesUPSelRep2 #-}+indicesUPSelRep2 tags rep = joinD theGang balanced+                          $ zipWithD theGang indices+                                             (splitD theGang balanced tags)+                                              rep+  where+    indices tags ((i,j), (m,n))+      = Seq.combine2ByTag tags (Seq.enumFromStepLen i 1 m)+                               (Seq.enumFromStepLen j 1 n)+++-- | O(n).+elementsUPSelRep2_0 :: Vector Tag -> UPSelRep2 -> Int+{-# INLINE elementsUPSelRep2_0 #-}+elementsUPSelRep2_0 _ = sumD theGang . fstD . sndD+++-- | O(n).+elementsUPSelRep2_1 :: Vector Tag -> UPSelRep2 -> Int+{-# INLINE elementsUPSelRep2_1 #-}+elementsUPSelRep2_1 _ = sumD theGang . sndD . sndD+++-- | O(1). Construct a selector. Wrapper for `UPSel2`.+mkUPSel2 :: Vector Tag -> Vector Int -> Int -> Int -> UPSelRep2 -> UPSel2+{-# INLINE mkUPSel2 #-}+mkUPSel2 tags is n0 n1 rep = UPSel2 (mkUSel2 tags is n0 n1) rep+
+ LICENSE view
@@ -0,0 +1,37 @@+Copyright (c) 2001-2011, The DPH Team+All rights reserved.++The DPH Team is:+  Manuel M T Chakravarty+  Gabriele Keller+  Roman Leshchinskiy+  Ben Lippmeier+  George Roldugin++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 name of the University nor the names of its contributors may be+used to endorse or promote products derived from this software without+specific prior written permission. ++THIS SOFTWARE IS PROVIDED BY THE UNIVERSITY COURT OF THE UNIVERSITY OF+GLASGOW AND THE 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+UNIVERSITY COURT OF THE UNIVERSITY OF GLASGOW OR THE 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.+
+ Setup.hs view
@@ -0,0 +1,3 @@+import Distribution.Simple+main = defaultMain+
+ dph-prim-par.cabal view
@@ -0,0 +1,52 @@+Name:           dph-prim-par+Version:        0.5.1.1+License:        BSD3+License-File:   LICENSE+Author:         The DPH Team+Maintainer:     Ben Lippmeier <benl@cse.unsw.edu.au>+Homepage:       http://www.haskell.org/haskellwiki/GHC/Data_Parallel_Haskell+Category:       Data Structures+Synopsis:       Parallel Primitives for Data-Parallel Haskell.++Cabal-Version:  >= 1.6+Build-Type:     Simple++Library+  Exposed-Modules:+        Data.Array.Parallel.Unlifted.Distributed+        Data.Array.Parallel.Unlifted.Parallel+        Data.Array.Parallel.Unlifted+  Other-Modules:+        Data.Array.Parallel.Unlifted.Distributed.Gang+        Data.Array.Parallel.Unlifted.Distributed.TheGang+        Data.Array.Parallel.Unlifted.Distributed.DistST+        Data.Array.Parallel.Unlifted.Distributed.Types+        Data.Array.Parallel.Unlifted.Distributed.Combinators+        Data.Array.Parallel.Unlifted.Distributed.Scalars+        Data.Array.Parallel.Unlifted.Distributed.Arrays+        Data.Array.Parallel.Unlifted.Distributed.Basics+        Data.Array.Parallel.Unlifted.Parallel.Combinators+        Data.Array.Parallel.Unlifted.Parallel.Sums+        Data.Array.Parallel.Unlifted.Parallel.Basics+        Data.Array.Parallel.Unlifted.Parallel.Permute+        Data.Array.Parallel.Unlifted.Parallel.Enum+        Data.Array.Parallel.Unlifted.Parallel.Segmented+        Data.Array.Parallel.Unlifted.Parallel.Subarrays+        Data.Array.Parallel.Unlifted.Parallel.UPSegd+        Data.Array.Parallel.Unlifted.Parallel.UPSel+        Data.Array.Parallel.Unlifted.Parallel.Text++  Exposed: False++  Extensions: TypeFamilies, GADTs, RankNTypes,+              BangPatterns, MagicHash, UnboxedTuples, TypeOperators+  GHC-Options: -Odph -funbox-strict-fields -fcpr-off++  Build-Depends:  +        base     == 4.4.*,+        random   == 1.0.*,+        vector   == 0.7.*,+        old-time == 1.0.*,+        dph-base == 0.5.*,+        dph-prim-interface == 0.5.*,+        dph-prim-seq       == 0.5.*