dph-prim-par 0.6.1.1 → 0.7.0.1
raw patch · 56 files changed
+3404/−2390 lines, 56 filesdep ~basedep ~dph-basedep ~dph-prim-interfacePVP ok
version bump matches the API change (PVP)
Dependency ranges changed: base, dph-base, dph-prim-interface, dph-prim-seq, vector
API changes (from Hackage documentation)
- Data.Array.Parallel.Unlifted.Distributed: glueSegdD :: Gang -> Dist ((USegd, Int), Int) -> Dist USegd
- Data.Array.Parallel.Unlifted.Distributed.Gang: data Gang
- Data.Array.Parallel.Unlifted.Distributed.Gang: forkGang :: Int -> IO Gang
- Data.Array.Parallel.Unlifted.Distributed.Gang: gangIO :: Gang -> (Int -> IO ()) -> IO ()
- Data.Array.Parallel.Unlifted.Distributed.Gang: gangST :: Gang -> (Int -> ST s ()) -> ST s ()
- Data.Array.Parallel.Unlifted.Distributed.Gang: gangSize :: Gang -> Int
- Data.Array.Parallel.Unlifted.Distributed.Gang: instance Show Gang
- Data.Array.Parallel.Unlifted.Distributed.Gang: seqGang :: Gang -> Gang
- Data.Array.Parallel.Unlifted.Distributed.Gang: traceGang :: String -> IO ()
- Data.Array.Parallel.Unlifted.Distributed.Gang: traceGangST :: String -> ST s ()
- Data.Array.Parallel.Unlifted.Distributed.TheGang: theGang :: Gang
+ Data.Array.Parallel.Unlifted: append_vs :: (Elt a, Elts a) => Segd -> VSegd -> Arrays a -> VSegd -> Arrays a -> Array a
+ Data.Array.Parallel.Unlifted: zipWith5 :: (Elt a, Elt b, Elt c, Elt d, Elt e, Elt f) => (a -> b -> c -> d -> e -> f) -> Array a -> Array b -> Array c -> Array d -> Array e -> Array f
+ Data.Array.Parallel.Unlifted: zipWith6 :: (Elt a, Elt b, Elt c, Elt d, Elt e, Elt f, Elt g) => (a -> b -> c -> d -> e -> f -> g) -> Array a -> Array b -> Array c -> Array d -> Array e -> Array f -> Array g
+ Data.Array.Parallel.Unlifted: zipWith7 :: (Elt a, Elt b, Elt c, Elt d, Elt e, Elt f, Elt g, Elt h) => (a -> b -> c -> d -> e -> f -> g -> h) -> Array a -> Array b -> Array c -> Array d -> Array e -> Array f -> Array g -> Array h
+ Data.Array.Parallel.Unlifted: zipWith8 :: (Elt a, Elt b, Elt c, Elt d, Elt e, Elt f, Elt g, Elt h, Elt i) => (a -> b -> c -> d -> e -> f -> g -> h -> i) -> Array a -> Array b -> Array c -> Array d -> Array e -> Array f -> Array g -> Array h -> Array i
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: atomicUpdateD :: Unbox a => Gang -> Dist (Vector a) -> Dist (Vector (Int, a)) -> Vector a
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: balanced :: Distribution
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: bpermuteD :: Unbox a => Gang -> Vector a -> Dist (Vector Int) -> Dist (Vector a)
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: carryD :: (Unbox a, DT a) => Gang -> (a -> a -> a) -> a -> Dist Bool -> Dist (Vector a) -> (Dist (Vector a), a)
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: data Distribution
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: joinD :: Unbox a => Gang -> Distribution -> Dist (Vector a) -> Vector a
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: joinDM :: Unbox a => Gang -> Dist (Vector a) -> ST s (MVector s a)
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: joinLengthD :: Unbox a => Gang -> Dist (Vector a) -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: lengthD :: Unbox a => Dist (Vector a) -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: permuteD :: Unbox a => Gang -> Dist (Vector a) -> Dist (Vector Int) -> Vector a
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: splitAsD :: Unbox a => Gang -> Dist Int -> Vector a -> Dist (Vector a)
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: splitD :: Unbox a => Gang -> Distribution -> Vector a -> Dist (Vector a)
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: splitJoinD :: (Unbox a, Unbox b) => Gang -> (Dist (Vector a) -> Dist (Vector b)) -> Vector a -> Vector b
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: splitLenD :: Gang -> Int -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: splitLenIdxD :: Gang -> Int -> Dist (Int, Int)
+ Data.Array.Parallel.Unlifted.Distributed.Arrays: unbalanced :: Distribution
+ Data.Array.Parallel.Unlifted.Distributed.Basics: eqD :: (Eq a, DT a) => Gang -> Dist a -> Dist a -> Bool
+ Data.Array.Parallel.Unlifted.Distributed.Basics: fromD :: DT a => Gang -> Dist a -> [a]
+ Data.Array.Parallel.Unlifted.Distributed.Basics: neqD :: (Eq a, DT a) => Gang -> Dist a -> Dist a -> Bool
+ Data.Array.Parallel.Unlifted.Distributed.Basics: toD :: DT a => Gang -> [a] -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WBpermute :: What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WFMapGen :: What -> What -> What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WFMapMap :: What -> What -> What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WFZipMap :: What -> What -> What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WJoinCopy :: Int -> What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WLength :: What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WLengthIdx :: What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WScalar :: What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WSlice :: What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: WZip :: What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: What :: String -> What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: data What
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: foldD :: DT a => What -> Gang -> (a -> a -> a) -> Dist a -> a
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: fstD :: (DT a, DT b) => Dist (a, b) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: imapD :: (DT a, DT b) => What -> Gang -> (Int -> a -> b) -> Dist a -> Dist b
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: izipWithD :: (DT a, DT b, DT c) => What -> Gang -> (Int -> a -> b -> c) -> Dist a -> Dist b -> Dist c
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: mapAccumLD :: (DT a, DT b) => Gang -> (acc -> a -> (acc, b)) -> acc -> Dist a -> (acc, Dist b)
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: mapD :: (DT a, DT b) => What -> Gang -> (a -> b) -> Dist a -> Dist b
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: scanD :: DT a => What -> Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: sndD :: (DT a, DT b) => Dist (a, b) -> Dist b
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: unzipD :: (DT a, DT b) => Dist (a, b) -> (Dist a, Dist b)
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: zipD :: (DT a, DT b) => Dist a -> Dist b -> Dist (a, b)
+ Data.Array.Parallel.Unlifted.Distributed.Combinators: zipWithD :: (DT a, DT b, DT c) => What -> Gang -> (a -> b -> c) -> Dist a -> Dist b -> Dist c
+ Data.Array.Parallel.Unlifted.Distributed.Data.Bool: andD :: Gang -> Dist Bool -> Bool
+ Data.Array.Parallel.Unlifted.Distributed.Data.Bool: instance DPrim Bool
+ Data.Array.Parallel.Unlifted.Distributed.Data.Bool: instance DT Bool
+ Data.Array.Parallel.Unlifted.Distributed.Data.Bool: orD :: Gang -> Dist Bool -> Bool
+ Data.Array.Parallel.Unlifted.Distributed.Data.Maybe: instance DT a => DT (Maybe a)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Ordering: instance DPrim Ordering
+ Data.Array.Parallel.Unlifted.Distributed.Data.Ordering: instance DT Ordering
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: class DT a where data family Dist a data family MDist a :: * -> * deepSeqD = seq measureD _ = "None"
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: deepSeqD :: DT a => a -> b -> b
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: indexD :: DT a => String -> Dist a -> Int -> a
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: instance PprPhysical (Dist Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: measureD :: DT a => a -> String
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: newMD :: DT a => Gang -> ST s (MDist a s)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: readMD :: DT a => MDist a s -> Int -> ST s a
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: scalarD :: DT a => Gang -> a -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: sizeD :: DT a => Dist a -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: sizeMD :: DT a => MDist a s -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: sumD :: (Num a, DT a) => Gang -> Dist a -> a
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: unsafeFreezeMD :: DT a => MDist a s -> ST s (Dist a)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar: writeMD :: DT a => MDist a s -> Int -> a -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DPrim Char
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DPrim Double
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DPrim Float
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DPrim Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DPrim Integer
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DPrim Word8
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DT Char
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DT Double
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DT Float
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DT Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DT Integer
+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base: instance DT Word8
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: fstD :: (DT a, DT b) => Dist (a, b) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: instance (DT a, DT b) => DT (a, b)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: instance (DT a, DT b, DT c) => DT (a, b, c)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: instance (PprPhysical (Dist a), PprPhysical (Dist b)) => PprPhysical (Dist (a, b))
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: sndD :: (DT a, DT b) => Dist (a, b) -> Dist b
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: unzip3D :: (DT a, DT b, DT c) => Dist (a, b, c) -> (Dist a, Dist b, Dist c)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: unzipD :: (DT a, DT b) => Dist (a, b) -> (Dist a, Dist b)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: zip3D :: (DT a, DT b, DT c) => Dist a -> Dist b -> Dist c -> Dist (a, b, c)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple: zipD :: (DT a, DT b) => Dist a -> Dist b -> Dist (a, b)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: lengthD :: Dist USSegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: splitSSegdOnElemsD :: Gang -> USSegd -> Dist ((USSegd, Int), Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: takeElementsD :: Dist USSegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: takeIndicesD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: takeLengthsD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: takeSourcesD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: takeStartsD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd: takeUSegdD :: Dist USSegd -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base: lengthD :: Dist USSegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base: takeElementsD :: Dist USSegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base: takeIndicesD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base: takeLengthsD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base: takeSourcesD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base: takeStartsD :: Dist USSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base: takeUSegdD :: Dist USSegd -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.DT: instance DT USSegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.DT: instance PprPhysical (Dist USSegd)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Split: splitSSegdOnElemsD :: Gang -> USSegd -> Dist ((USSegd, Int), Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: glueSegdD :: Gang -> Dist ((USegd, Int), Int) -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: joinSegdD :: Gang -> Dist USegd -> USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: lengthD :: Dist USegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: mkDUSegdD :: Dist (Vector Int) -> Dist (Vector Int) -> Dist Int -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: splitSD :: Unbox a => Gang -> Dist USegd -> Vector a -> Dist (Vector a)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: splitSegdOnElemsD :: Gang -> USegd -> Dist ((USegd, Int), Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: splitSegdOnSegsD :: Gang -> USegd -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: takeElementsD :: Dist USegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: takeIndicesD :: Dist USegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd: takeLengthsD :: Dist USegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base: lengthD :: Dist USegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base: mkDUSegdD :: Dist (Vector Int) -> Dist (Vector Int) -> Dist Int -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base: takeElementsD :: Dist USegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base: takeIndicesD :: Dist USegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base: takeLengthsD :: Dist USegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.DT: instance DT USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.DT: instance PprPhysical (Dist USegd)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split: glueSegdD :: Gang -> Dist ((USegd, Int), Int) -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split: joinSegdD :: Gang -> Dist USegd -> USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split: splitSD :: Unbox a => Gang -> Dist USegd -> Vector a -> Dist (Vector a)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split: splitSegdOnElemsD :: Gang -> USegd -> Dist ((USegd, Int), Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split: splitSegdOnSegsD :: Gang -> USegd -> Dist USegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: instance DT UVSegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: instance PprPhysical (Dist UVSegd)
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: lengthD :: Dist UVSegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: takeElementsD :: Dist UVSegd -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: takeIndicesD :: Dist UVSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: takeLengthsD :: Dist UVSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: takeSourcesD :: Dist UVSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: takeStartsD :: Dist UVSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: takeUSSegdD :: Dist UVSegd -> Dist USSegd
+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd: takeVSegidsD :: Dist UVSegd -> Dist (Vector Int)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Unit: instance DT ()
+ Data.Array.Parallel.Unlifted.Distributed.Data.Unit: unitD :: Gang -> Dist ()
+ Data.Array.Parallel.Unlifted.Distributed.Data.Vector: instance (Unbox a, Show a) => PprPhysical (Dist (Vector a))
+ Data.Array.Parallel.Unlifted.Distributed.Data.Vector: instance Unbox a => DT (Vector a)
+ Data.Array.Parallel.Unlifted.Distributed.Data.Vector: lengthD :: Unbox a => Dist (Vector a) -> Dist Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: checkGangD :: DT a => String -> Gang -> Dist a -> b -> b
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: class DT a where data family Dist a data family MDist a :: * -> * deepSeqD = seq measureD _ = "None"
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: data Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: debugD :: DT a => Dist a -> String
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: deepSeqD :: DT a => a -> b -> b
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: foldD :: DT a => What -> Gang -> (a -> a -> a) -> Dist a -> a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: forkGang :: Int -> IO Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: gangSize :: Gang -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: generateD :: DT a => What -> Gang -> (Int -> a) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: generateD_cheap :: DT a => What -> Gang -> (Int -> a) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: imapD' :: (DT a, DT b) => What -> Gang -> (Int -> a -> b) -> Dist a -> Dist b
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: indexD :: DT a => String -> Dist a -> Int -> a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: measureD :: DT a => a -> String
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: newD :: DT a => Gang -> (forall s. MDist a s -> ST s ()) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: newMD :: DT a => Gang -> ST s (MDist a s)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: readMD :: DT a => MDist a s -> Int -> ST s a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: scanD :: DT a => What -> Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: seqGang :: Gang -> Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: sizeD :: DT a => Dist a -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: sizeMD :: DT a => MDist a s -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: theGang :: Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: unsafeFreezeMD :: DT a => MDist a s -> ST s (Dist a)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive: writeMD :: DT a => MDist a s -> Int -> a -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: class Unbox e => DPrim e
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: mkDPrim :: DPrim e => Vector e -> Dist e
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: mkMDPrim :: DPrim e => STVector s e -> MDist e s
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: primIndexD :: DPrim a => String -> Dist a -> Int -> a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: primNewMD :: DPrim a => Gang -> ST s (MDist a s)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: primReadMD :: DPrim a => MDist a s -> Int -> ST s a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: primSizeD :: DPrim a => Dist a -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: primSizeMD :: DPrim a => MDist a s -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: primUnsafeFreezeMD :: DPrim a => MDist a s -> ST s (Dist a)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: primWriteMD :: DPrim a => MDist a s -> Int -> a -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: unDPrim :: DPrim e => Dist e -> Vector e
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim: unMDPrim :: DPrim e => MDist e s -> STVector s e
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: checkGangD :: DT a => String -> Gang -> Dist a -> b -> b
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: checkGangMD :: DT a => String -> Gang -> MDist a s -> b -> b
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: class DT a where data family Dist a data family MDist a :: * -> * deepSeqD = seq measureD _ = "None"
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: debugD :: DT a => Dist a -> String
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: deepSeqD :: DT a => a -> b -> b
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: indexD :: DT a => String -> Dist a -> Int -> a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: instance (Show a, DT a) => Show (Dist a)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: measureD :: DT a => a -> String
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: newD :: DT a => Gang -> (forall s. MDist a s -> ST s ()) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: newMD :: DT a => Gang -> ST s (MDist a s)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: readMD :: DT a => MDist a s -> Int -> ST s a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: sizeD :: DT a => Dist a -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: sizeMD :: DT a => MDist a s -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: unsafeFreezeMD :: DT a => MDist a s -> ST s (Dist a)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT: writeMD :: DT a => MDist a s -> Int -> a -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: data DistST s a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: distST :: DT a => Comp -> Gang -> DistST s a -> ST s (Dist a)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: distST_ :: Comp -> Gang -> DistST s () -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: instance Monad (DistST s)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: mapDST :: (DT a, DT b) => What -> Gang -> (a -> DistST s b) -> Dist a -> ST s (Dist b)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: mapDST_ :: DT a => What -> Gang -> (a -> DistST s ()) -> Dist a -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: myD :: DT a => Dist a -> DistST s a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: myIndex :: DistST s Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: readMyMD :: DT a => MDist a s -> DistST s a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: runDistST :: DT a => Comp -> Gang -> (forall s. DistST s a) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: runDistST_seq :: DT a => Gang -> (forall s. DistST s a) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: stToDistST :: ST s a -> DistST s a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: writeMyMD :: DT a => MDist a s -> a -> DistST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: zipWithDST :: (DT a, DT b, DT c) => What -> Gang -> (a -> b -> DistST s c) -> Dist a -> Dist b -> ST s (Dist c)
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST: zipWithDST_ :: (DT a, DT b) => What -> Gang -> (a -> b -> DistST s ()) -> Dist a -> Dist b -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: WorkCopy :: Int -> Workload
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: WorkUnknown :: Workload
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: data Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: data Workload
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: forkGang :: Int -> IO Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: gangIO :: Gang -> String -> Workload -> (Int -> IO ()) -> IO ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: gangST :: Gang -> String -> Workload -> (Int -> ST s ()) -> ST s ()
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: gangSize :: Gang -> Int
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: instance Eq Workload
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: instance Show Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: instance Show Workload
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang: seqGang :: Gang -> Gang
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators: foldD :: DT a => What -> Gang -> (a -> a -> a) -> Dist a -> a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators: generateD :: DT a => What -> Gang -> (Int -> a) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators: generateD_cheap :: DT a => What -> Gang -> (Int -> a) -> Dist a
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators: imapD' :: (DT a, DT b) => What -> Gang -> (Int -> a -> b) -> Dist a -> Dist b
+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators: scanD :: DT a => What -> Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)
+ Data.Array.Parallel.Unlifted.Distributed.What: CDist :: What -> Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: CFold :: What -> Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: CGen :: Bool -> What -> Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: CMap :: What -> Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: CScan :: What -> Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: WBpermute :: What
+ Data.Array.Parallel.Unlifted.Distributed.What: WFMapGen :: What -> What -> What
+ Data.Array.Parallel.Unlifted.Distributed.What: WFMapMap :: What -> What -> What
+ Data.Array.Parallel.Unlifted.Distributed.What: WFZipMap :: What -> What -> What
+ Data.Array.Parallel.Unlifted.Distributed.What: WJoinCopy :: Int -> What
+ Data.Array.Parallel.Unlifted.Distributed.What: WLength :: What
+ Data.Array.Parallel.Unlifted.Distributed.What: WLengthIdx :: What
+ Data.Array.Parallel.Unlifted.Distributed.What: WScalar :: What
+ Data.Array.Parallel.Unlifted.Distributed.What: WSlice :: What
+ Data.Array.Parallel.Unlifted.Distributed.What: WZip :: What
+ Data.Array.Parallel.Unlifted.Distributed.What: What :: String -> What
+ Data.Array.Parallel.Unlifted.Distributed.What: data Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: data What
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Eq Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Eq What
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Ord Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Ord What
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Read Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Read What
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Show Comp
+ Data.Array.Parallel.Unlifted.Distributed.What: instance Show What
+ Data.Array.Parallel.Unlifted.Parallel: appendSUPV :: (Unboxes a, Unbox a) => UPSegd -> UPVSegd -> Vectors a -> UPVSegd -> Vectors a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel: extractsFromVectorsUPVSegdP :: (Unbox a, Unboxes a) => UPVSegd -> Vectors a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel: indexsFromVectorsUPVSegdP :: (Unbox a, Unboxes a) => Vectors a -> UPVSegd -> Vector (Int, Int) -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Basics: emptyUP :: Unbox e => Vector e
+ Data.Array.Parallel.Unlifted.Parallel.Basics: indexedUP :: (DT e, Unbox e) => Vector e -> Vector (Int, e)
+ Data.Array.Parallel.Unlifted.Parallel.Basics: interleaveUP :: Unbox e => Vector e -> Vector e -> Vector e
+ Data.Array.Parallel.Unlifted.Parallel.Basics: lengthUP :: Unbox e => Vector e -> Int
+ Data.Array.Parallel.Unlifted.Parallel.Basics: nullUP :: Unbox e => Vector e -> Bool
+ Data.Array.Parallel.Unlifted.Parallel.Basics: repeatUP :: Unbox e => Int -> Vector e -> Vector e
+ Data.Array.Parallel.Unlifted.Parallel.Basics: replicateUP :: Unbox e => Int -> e -> Vector e
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: combine2UP :: Unbox a => Vector Tag -> UPSelRep2 -> Vector a -> Vector a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: combineUP :: Unbox a => Vector Bool -> Vector a -> Vector a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: filterUP :: Unbox a => (a -> Bool) -> Vector a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: fold1UP :: (DT a, Unbox a) => (a -> a -> a) -> Vector a -> a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: foldUP :: (Unbox a, DT a) => (a -> a -> a) -> a -> Vector a -> a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: foldl1UP :: (DT a, Unbox a) => (a -> a -> a) -> Vector a -> a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: foldlUP :: (DT a, Unbox a) => (a -> a -> a) -> a -> Vector a -> a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: mapUP :: (Unbox a, Unbox b) => (a -> b) -> Vector a -> Vector b
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: packUP :: Unbox e => Vector e -> Vector Bool -> Vector e
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: scanUP :: (DT a, Unbox a) => (a -> a -> a) -> a -> Vector a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Combinators: zipWithUP :: (Unbox a, Unbox b, Unbox c) => (a -> b -> c) -> Vector a -> Vector b -> Vector c
+ Data.Array.Parallel.Unlifted.Parallel.Enum: enumFromStepLenEachUP :: Int -> Vector Int -> Vector Int -> Vector Int -> Vector Int
+ Data.Array.Parallel.Unlifted.Parallel.Enum: enumFromStepLenUP :: Int -> Int -> Int -> Vector Int
+ Data.Array.Parallel.Unlifted.Parallel.Enum: enumFromThenToUP :: (Unbox a, Enum a) => a -> a -> a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Enum: enumFromToUP :: (Unbox a, Enum a) => a -> a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: extractsFromNestedUPSSegd :: Unbox a => UPSSegd -> Vector (Vector a) -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: extractsFromVectorsUPSSegd :: (Unbox a, Unboxes a) => UPSSegd -> Vectors a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: extractsFromVectorsUPSSegdSegmap :: (Unbox a, Unboxes a) => UPSSegd -> Vectors a -> Vector Int -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: extractsFromVectorsUPVSegd :: (Unbox a, Unboxes a) => UPVSegd -> Vectors a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: extractsFromVectorsUPVSegdP :: (Unbox a, Unboxes a) => UPVSegd -> Vectors a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: indexsFromVector :: Unbox a => Vector a -> Vector Int -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: indexsFromVectorsUPVSegd :: (Unbox a, Unboxes a) => Vectors a -> UPVSegd -> Vector (Int, Int) -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Extracts: indexsFromVectorsUPVSegdP :: (Unbox a, Unboxes a) => Vectors a -> UPVSegd -> Vector (Int, Int) -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Permute: bpermuteUP :: Unbox a => Vector a -> Vector Int -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Permute: updateUP :: Unbox a => Vector a -> Vector (Int, a) -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Segmented: appendSUP :: Unbox a => UPSegd -> UPSegd -> Vector a -> UPSegd -> Vector a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Segmented: appendSUPV :: (Unboxes a, Unbox a) => UPSegd -> UPVSegd -> Vectors a -> UPVSegd -> Vectors a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Segmented: appendSUP_old :: Unbox a => UPSegd -> UPSegd -> Vector a -> UPSegd -> Vector a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Segmented: foldRUP :: (Unbox a, Unbox b) => (b -> a -> b) -> b -> Int -> Vector a -> Vector b
+ Data.Array.Parallel.Unlifted.Parallel.Segmented: replicateRSUP :: Unbox a => Int -> Vector a -> Vector a
+ Data.Array.Parallel.Unlifted.Parallel.Segmented: sumRUP :: (Num e, Unbox e) => Int -> Vector e -> Vector e
+ Data.Array.Parallel.Unlifted.Parallel.Subarrays: dropUP :: Unbox e => Int -> Vector e -> Vector e
+ Data.Array.Parallel.Unlifted.Parallel.Sums: allUP :: Unbox e => (e -> Bool) -> Vector e -> Bool
+ Data.Array.Parallel.Unlifted.Parallel.Sums: andUP :: Vector Bool -> Bool
+ Data.Array.Parallel.Unlifted.Parallel.Sums: anyUP :: Unbox e => (e -> Bool) -> Vector e -> Bool
+ Data.Array.Parallel.Unlifted.Parallel.Sums: maximumByUP :: (DT e, Unbox e) => (e -> e -> Ordering) -> Vector e -> e
+ Data.Array.Parallel.Unlifted.Parallel.Sums: maximumIndexByUP :: (DT e, Unbox e) => (e -> e -> Ordering) -> Vector e -> Int
+ Data.Array.Parallel.Unlifted.Parallel.Sums: maximumUP :: (DT e, Ord e, Unbox e) => Vector e -> e
+ Data.Array.Parallel.Unlifted.Parallel.Sums: orUP :: Vector Bool -> Bool
+ Data.Array.Parallel.Unlifted.Parallel.Sums: productUP :: (DT e, Num e, Unbox e) => Vector e -> e
+ Data.Array.Parallel.Unlifted.Parallel.Sums: sumUP :: (Unbox a, DT a, Num a) => Vector a -> a
+ Data.Array.Parallel.Unlifted.Parallel.Text: instance Show UPSegd
+ Data.Array.Parallel.Unlifted.Parallel.UPVSegd: takeDistributed :: UPVSegd -> Dist ((USegd, Int), Int)
- Data.Array.Parallel.Unlifted.Distributed: foldD :: DT a => Gang -> (a -> a -> a) -> Dist a -> a
+ Data.Array.Parallel.Unlifted.Distributed: foldD :: DT a => What -> Gang -> (a -> a -> a) -> Dist a -> a
- Data.Array.Parallel.Unlifted.Distributed: mapD :: (DT a, DT b) => Gang -> (a -> b) -> Dist a -> Dist b
+ Data.Array.Parallel.Unlifted.Distributed: mapD :: (DT a, DT b) => What -> Gang -> (a -> b) -> Dist a -> Dist b
- Data.Array.Parallel.Unlifted.Distributed: scanD :: DT a => Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)
+ Data.Array.Parallel.Unlifted.Distributed: scanD :: DT a => What -> Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)
- Data.Array.Parallel.Unlifted.Distributed: zipWithD :: (DT a, DT b, DT c) => Gang -> (a -> b -> c) -> Dist a -> Dist b -> Dist c
+ Data.Array.Parallel.Unlifted.Distributed: zipWithD :: (DT a, DT b, DT c) => What -> Gang -> (a -> b -> c) -> Dist a -> Dist b -> Dist c
Files
- Data/Array/Parallel/Unlifted.hs +7/−4
- Data/Array/Parallel/Unlifted/Distributed.hs +4/−6
- Data/Array/Parallel/Unlifted/Distributed/Arrays.hs +103/−71
- Data/Array/Parallel/Unlifted/Distributed/Basics.hs +7/−6
- Data/Array/Parallel/Unlifted/Distributed/Combinators.hs +87/−153
- Data/Array/Parallel/Unlifted/Distributed/Data/Bool.hs +50/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/Maybe.hs +63/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/Ordering.hs +30/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/Scalar.hs +38/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/Scalar/Base.hs +129/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/Tuple.hs +173/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USSegd.hs +18/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USSegd/Base.hs +71/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USSegd/DT.hs +88/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USSegd/Split.hs +250/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USegd.hs +20/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USegd/Base.hs +56/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USegd/DT.hs +84/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/USegd/Split.hs +361/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/UVSegd.hs +147/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/Unit.hs +56/−0
- Data/Array/Parallel/Unlifted/Distributed/Data/Vector.hs +85/−0
- Data/Array/Parallel/Unlifted/Distributed/DistST.hs +0/−119
- Data/Array/Parallel/Unlifted/Distributed/Gang.hs +0/−243
- Data/Array/Parallel/Unlifted/Distributed/Primitive.hs +47/−0
- Data/Array/Parallel/Unlifted/Distributed/Primitive/DPrim.hs +86/−0
- Data/Array/Parallel/Unlifted/Distributed/Primitive/DT.hs +108/−0
- Data/Array/Parallel/Unlifted/Distributed/Primitive/DistST.hs +205/−0
- Data/Array/Parallel/Unlifted/Distributed/Primitive/Gang.hs +296/−0
- Data/Array/Parallel/Unlifted/Distributed/Primitive/Operators.hs +127/−0
- Data/Array/Parallel/Unlifted/Distributed/Scalars.hs +0/−37
- Data/Array/Parallel/Unlifted/Distributed/TheGang.hs +0/−22
- Data/Array/Parallel/Unlifted/Distributed/Types.hs +0/−21
- Data/Array/Parallel/Unlifted/Distributed/Types/Base.hs +0/−108
- Data/Array/Parallel/Unlifted/Distributed/Types/Maybe.hs +0/−45
- Data/Array/Parallel/Unlifted/Distributed/Types/Prim.hs +0/−250
- Data/Array/Parallel/Unlifted/Distributed/Types/Tuple.hs +0/−145
- Data/Array/Parallel/Unlifted/Distributed/Types/USSegd.hs +0/−130
- Data/Array/Parallel/Unlifted/Distributed/Types/USegd.hs +0/−113
- Data/Array/Parallel/Unlifted/Distributed/Types/UVSegd.hs +0/−128
- Data/Array/Parallel/Unlifted/Distributed/Types/Unit.hs +0/−45
- Data/Array/Parallel/Unlifted/Distributed/Types/Vector.hs +0/−61
- Data/Array/Parallel/Unlifted/Distributed/USSegd.hs +0/−245
- Data/Array/Parallel/Unlifted/Distributed/USegd.hs +0/−349
- Data/Array/Parallel/Unlifted/Distributed/What.hs +34/−0
- Data/Array/Parallel/Unlifted/Parallel.hs +5/−0
- Data/Array/Parallel/Unlifted/Parallel/Basics.hs +12/−7
- Data/Array/Parallel/Unlifted/Parallel/Combinators.hs +20/−13
- Data/Array/Parallel/Unlifted/Parallel/Enum.hs +3/−2
- Data/Array/Parallel/Unlifted/Parallel/Extracts.hs +84/−6
- Data/Array/Parallel/Unlifted/Parallel/Segmented.hs +338/−8
- Data/Array/Parallel/Unlifted/Parallel/UPSSegd.hs +4/−2
- Data/Array/Parallel/Unlifted/Parallel/UPSegd.hs +6/−5
- Data/Array/Parallel/Unlifted/Parallel/UPSel.hs +4/−3
- Data/Array/Parallel/Unlifted/Parallel/UPVSegd.hs +57/−12
- dph-prim-par.cabal +41/−31
Data/Array/Parallel/Unlifted.hs view
@@ -18,10 +18,9 @@ -- @dph-prim-par@ and @dph-prim-seq@ packages. -- #include "DPH_Header.h"-import Data.Array.Parallel.Unlifted.Parallel import Data.Array.Parallel.Base.TracePrim+import Data.Array.Parallel.Unlifted.Parallel import Data.Array.Parallel.Unlifted.Distributed ( DT )- import Data.Array.Parallel.Unlifted.Sequential.Vector (Unbox, Vector) import Data.Array.Parallel.Unlifted.Vectors (Unboxes, Vectors) import Data.Array.Parallel.Unlifted.Parallel.UPSel@@ -62,6 +61,10 @@ = let arr = appendSUP segd xd xs yd ys in tracePrim (TraceAppend_s (Seq.length arr)) arr +append_vs segd xd xs yd ys+ = let arr = appendSUPV segd xd xs yd ys+ in tracePrim (TraceAppend_vs (Seq.length arr)) arr+ replicate n val = tracePrim (TraceReplicate n) $! replicateUP n val@@ -108,7 +111,7 @@ length = Seq.length index = Seq.index indexs = indexsFromVector-indexs_avs = indexsFromVectorsUPVSegd+indexs_avs = indexsFromVectorsUPVSegdP extract arr i n = tracePrim (TraceExtract (Seq.length arr) i n)@@ -116,7 +119,7 @@ extracts_nss = extractsFromNestedUPSSegd extracts_ass = extractsFromVectorsUPSSegd-extracts_avs = extractsFromVectorsUPVSegd+extracts_avs = extractsFromVectorsUPVSegdP drop n arr = tracePrim (TraceDrop n (Seq.length arr))
Data/Array/Parallel/Unlifted/Distributed.hs view
@@ -40,7 +40,6 @@ , joinD , splitJoinD , joinDM- , glueSegdD , carryD , Distribution@@ -56,12 +55,11 @@ -- * Debugging , fromD, toD, debugD) where-import Data.Array.Parallel.Unlifted.Distributed.TheGang import Data.Array.Parallel.Unlifted.Distributed.Combinators-import Data.Array.Parallel.Unlifted.Distributed.Scalars+import Data.Array.Parallel.Unlifted.Distributed.Data.Bool+import Data.Array.Parallel.Unlifted.Distributed.Data.Scalar+import Data.Array.Parallel.Unlifted.Distributed.Data.Ordering () import Data.Array.Parallel.Unlifted.Distributed.Arrays-import Data.Array.Parallel.Unlifted.Distributed.USegd import Data.Array.Parallel.Unlifted.Distributed.Basics-import Data.Array.Parallel.Unlifted.Distributed.Types-import Data.Array.Parallel.Unlifted.Distributed.Gang (Gang, forkGang, gangSize)+import Data.Array.Parallel.Unlifted.Distributed.Primitive
Data/Array/Parallel/Unlifted/Distributed/Arrays.hs view
@@ -23,13 +23,13 @@ -- * Carry , carryD) where-import Data.Array.Parallel.Base (ST, runST)-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Unlifted.Distributed.DistST-import Data.Array.Parallel.Unlifted.Distributed.Types+import Data.Array.Parallel.Unlifted.Distributed.Data.Scalar+import Data.Array.Parallel.Unlifted.Distributed.Data.Vector+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST+import Data.Array.Parallel.Unlifted.Distributed.Primitive import Data.Array.Parallel.Unlifted.Distributed.Combinators-import Data.Array.Parallel.Unlifted.Distributed.Scalars import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector, MVector, Unbox)+import Data.Array.Parallel.Base (ST, runST) import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as Seq import GHC.Base ( quotInt, remInt ) import Control.Monad@@ -65,16 +65,18 @@ -- = [128,128,128,127]@ -- splitLenD :: Gang -> Int -> Dist Int-splitLenD g n = generateD_cheap g len- where- !p = gangSize g+splitLenD gang n + = generateD_cheap WLength gang len+ where+ !p = gangSize gang !l = n `quotInt` p !m = n `remInt` p {-# INLINE [0] len #-} len i | i < m = l+1 | otherwise = l-{-# INLINE splitLenD #-}+{-# NOINLINE splitLenD #-}+-- NOINLINE because it's cheap and doesn't need to fuse with anything. -- | O(threads).@@ -86,16 +88,18 @@ -- = [(128,0),(128,128),(128,256),(127,384)]@ -- splitLenIdxD :: Gang -> Int -> Dist (Int, Int)-splitLenIdxD g n = generateD_cheap g len_idx- where- !p = gangSize g+splitLenIdxD gang n + = generateD_cheap WLengthIdx gang len_idx+ where+ !p = gangSize gang !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)-{-# INLINE splitLenIdxD #-}+{-# NOINLINE splitLenIdxD #-}+-- NOINLINE because it's cheap and doesn't need to fuse with anything. -- | O(threads).@@ -104,9 +108,12 @@ -- and summing them up. joinLengthD :: Unbox a => Gang -> Dist (Vector a) -> Int joinLengthD g = sumD g . lengthD-{-# INLINE joinLengthD #-}- +{-# NOINLINE joinLengthD #-}+-- NOINLINE because it's cheap and doesn't need to fuse with anything.+-- No operations are performed on the elements, so we don't need+-- to specialise for the element type. + -- Splitting and Joining arrays ----------------------------------------------- -- | Distribute an array over a 'Gang' such that each threads gets the given -- number of elements.@@ -114,11 +121,14 @@ -- @splitAsD theGangN4 (splitLenD theGangN4 10) [1 2 3 4 5 6 7 8 9 0] -- = [[1 2 3] [4 5 6] [7 8] [9 0]]@ -- -splitAsD :: Unbox a => Gang -> Dist Int -> Vector a -> Dist (Vector a)-splitAsD g dlen !arr - = zipWithD (seqGang g) (Seq.slice "splitAsD" arr) is dlen+splitAsD + :: Unbox a + => Gang -> Dist Int -> Vector a -> Dist (Vector a)++splitAsD gang dlen !arr + = zipWithD WSlice (seqGang gang) (Seq.slice "splitAsD" arr) is dlen where- is = fst $ scanD g (+) 0 dlen+ is = fst $ scanD (What "splitAsD") gang (+) 0 dlen {-# INLINE_DIST splitAsD #-} @@ -128,29 +138,46 @@ -- through RULES. Without it, splitJoinD would be a loop breaker. -- splitD :: Unbox a => Gang -> Distribution -> Vector a -> Dist (Vector a)-splitD g _ arr = splitD_impl g arr+splitD g _ arr + = splitD_impl g arr {-# INLINE_DIST splitD #-} splitD_impl :: Unbox a => Gang -> Vector a -> Dist (Vector a) splitD_impl g !arr - = generateD_cheap g (\i -> Seq.slice "splitD_impl" arr (idx i) (len i))- where- n = Seq.length arr- !p = gangSize g- !l = n `quotInt` p- !m = n `remInt` p+ = generateD_cheap WSlice g + (\i -> Seq.slice "splitD_impl" arr (idx i) (len i)) - {-# INLINE [0] idx #-}- idx i | i < m = (l+1)*i- | otherwise = l*i + m+ where n = Seq.length arr+ !p = gangSize g+ !l = n `quotInt` p+ !m = n `remInt` p - {-# INLINE [0] len #-}- len i | i < m = l+1- | otherwise = l+ {-# 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 {-# INLINE_DIST splitD_impl #-} +-- SplitJoin ------------------------------------------------------------------+-- | 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+splitJoinD g f !xs + = joinD_impl g (f (splitD_impl g xs))+{-# INLINE_DIST splitJoinD #-}+++-- Join ----------------------------------------------------------------------- -- | Join a distributed array. -- Join sums up the array lengths of each chunk, allocates a new result array, -- and copies each chunk into the result.@@ -164,42 +191,37 @@ joinD_impl :: forall a. Unbox a => Gang -> Dist (Vector a) -> Vector a-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 (Seq.mslice i (Seq.length arr) ma) arr)-{-# INLINE_DIST joinD_impl #-}-+joinD_impl gang !darr + = let -- Determine where each thread's local chunk should go+ -- in the result vector, and count the total number of elements.+ (!di,!n) = scanD (What "joinD_impl/count") gang (+) 0 + $ lengthD darr --- | 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-splitJoinD g f !xs - = joinD_impl g (f (splitD_impl g xs))-{-# INLINE_DIST splitJoinD #-}+ copy :: forall s. MVector s a -> Int -> Vector a -> DistST s ()+ copy ma i arr + = stToDistST (Seq.copy (Seq.mslice i (Seq.length arr) ma) arr)+ {-# INLINE copy #-} + in Seq.new n $ \ma + -> zipWithDST_ + (WJoinCopy n) + gang (copy ma) di darr+{-# INLINE_DIST joinD_impl #-} -- | Join a distributed array, yielding a mutable global array joinDM :: Unbox a => Gang -> Dist (Vector a) -> ST s (MVector s a)-joinDM g darr - = checkGangD (here "joinDM") g darr +joinDM gang darr + = checkGangD (here "joinDM") gang darr $ do marr <- Seq.newM n- zipWithDST_ g (copy marr) di darr+ zipWithDST_ (WJoinCopy n) gang (copy marr) di darr return marr where- (!di,!n) = scanD g (+) 0 $ lengthD darr+ (!di,!n) = scanD (What "joinDM/count") gang (+) 0 + $ lengthD darr - copy ma i arr = stToDistST (Seq.copy (Seq.mslice i (Seq.length arr) ma) arr)+ copy ma i arr + = stToDistST (Seq.copy (Seq.mslice i (Seq.length arr) ma) arr) {-# INLINE joinDM #-} @@ -226,17 +248,21 @@ "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))+ = joinD g balanced (zipWithD WZip 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)+ = joinD g balanced (zipWithD WZip 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)+"Seq.zip/splitJoinD" + forall what1 what2 gang f g xs ys+ . Seq.zip (splitJoinD gang (imapD what1 gang f) xs) + (splitJoinD gang (imapD what2 gang g) ys)+ = splitJoinD gang + (imapD (WFZipMap what1 what2)+ gang (\i zs -> let (as,bs) = Seq.unzip zs+ in Seq.zip (f i as) (g i bs)))+ (Seq.zip xs ys) #-} @@ -247,7 +273,7 @@ :: forall a. Unbox a => Gang -> Dist (Vector a) -> Dist (Vector Int) -> Vector a permuteD g darr dis - = Seq.new n (\ma -> zipWithDST_ g (permute ma) darr dis)+ = Seq.new n (\ma -> zipWithDST_ (What "permuteD") g (permute ma) darr dis) where n = joinLengthD g darr @@ -258,11 +284,17 @@ -- 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)-bpermuteD g !as ds = mapD g (Seq.bpermute as) ds-{-# INLINE bpermuteD #-}+bpermuteD :: Unbox a + => Gang + -> Vector a + -> Dist (Vector Int) + -> Dist (Vector a) +bpermuteD gang !as ds + = mapD WBpermute gang (Seq.bpermute as) ds+{-# INLINE_DIST bpermuteD #-} + -- 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@@ -272,12 +304,12 @@ atomicUpdateD g darr upd = runST $ do marr <- joinDM g darr- mapDST_ g (update marr) upd+ mapDST_ (What "atomicUpdateD") 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)-{-# INLINE atomicUpdateD #-}+{-# INLINE_DIST atomicUpdateD #-} -- Carry ----------------------------------------------------------------------
Data/Array/Parallel/Unlifted/Distributed/Basics.hs view
@@ -4,10 +4,10 @@ module Data.Array.Parallel.Unlifted.Distributed.Basics (eqD, neqD, toD, fromD) where-import Data.Array.Parallel.Unlifted.Distributed.Gang -import Data.Array.Parallel.Unlifted.Distributed.Types+import Data.Array.Parallel.Unlifted.Distributed.Primitive import Data.Array.Parallel.Unlifted.Distributed.Combinators -import Data.Array.Parallel.Unlifted.Distributed.Scalars+import Data.Array.Parallel.Unlifted.Distributed.Data.Bool+import Data.Array.Parallel.Unlifted.Distributed.Data.Scalar import Control.Monad ( zipWithM_ ) here :: String -> String@@ -18,21 +18,22 @@ -- 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)+ = andD g (zipWithD (What "eq") 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)+ = orD g (zipWithD (What "neq") g (/=) dx dy) -- | Generate a distributed value from the first @p@ elements of a list. -- -- * 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)+toD g xs+ = newD g (\md -> zipWithM_ (writeMD md) [0 .. gangSize g - 1] xs) -- | Yield all elements of a distributed value.
Data/Array/Parallel/Unlifted/Distributed/Combinators.hs view
@@ -5,68 +5,36 @@ -- | Standard combinators for distributed types. module Data.Array.Parallel.Unlifted.Distributed.Combinators - ( generateD, generateD_cheap+ ( W.What (..) , imapD, mapD , zipD, unzipD , fstD, sndD , zipWithD, izipWithD , foldD , scanD- , mapAccumLD-- -- * Monadic combinators- , mapDST_, mapDST, zipWithDST_, zipWithDST)+ , mapAccumLD) where import Data.Array.Parallel.Base ( ST, runST)-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Unlifted.Distributed.Types-import Data.Array.Parallel.Unlifted.Distributed.DistST+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Unlifted.Distributed.Data.Tuple+import Data.Array.Parallel.Unlifted.Distributed.Data.Maybe ()+import qualified Data.Array.Parallel.Unlifted.Distributed.What as W here s = "Data.Array.Parallel.Unlifted.Distributed.Combinators." ++ s --- | Create a distributed value, given a function to create the instance--- for each thread.-generateD :: DT a => Gang -> (Int -> a) -> Dist a-generateD g f - = runDistST g (myIndex >>= return . f)-{-# NOINLINE generateD #-} ---- | Create a distributed value, but do it sequentially.--- --- This function is used when we want to operate on a distributed value, but--- there isn't much data involved. For example, if we want to distribute --- a single integer to each thread, then there's no need to fire up the --- gang for this.--- -generateD_cheap :: DT a => Gang -> (Int -> a) -> Dist a-generateD_cheap g f - = runDistST_seq g (myIndex >>= return . f)-{-# NOINLINE generateD_cheap #-}-- -- 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-imapD g f d = imapD' g (\i x -> x `deepSeqD` f i x) d-{-# INLINE [0] imapD #-}----- | 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-imapD' g f !d - = checkGangD (here "imapD") g d- $ runDistST g - (do i <- myIndex- x <- myD d- return (f i x))-{-# NOINLINE imapD' #-}-+--+-- Fusing maps+-- ~~~~~~~~~~~+-- The staging here is important. +-- Our rewrite rules only operate on the imapD form, so fusion between the worker+-- functions of consecutive maps takes place before phase [0]. +--+-- At phase [0] we then inline imapD which introduces the call to imapD' which+-- uses the gang to evaluate its (now fused) worker.+-- -- | Map a function to every instance of a distributed value. --@@ -75,21 +43,53 @@ -- -- @mapD theGang (V.map (+ 1)) :: Dist (Vector Int) -> Dist (Vector Int)@ ---mapD :: (DT a, DT b) => Gang -> (a -> b) -> Dist a -> Dist b-mapD g = imapD g . const+mapD :: (DT a, DT b) + => W.What -- ^ What is the worker function doing.+ -> Gang + -> (a -> b) + -> Dist a + -> Dist b++mapD wFn gang+ = imapD wFn gang . const {-# INLINE mapD #-}+-- INLINE because this is just a convenience wrapper for imapD.+-- None of our rewrite rules are particular to mapD. +-- | 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) + => W.What -- ^ What is the worker function doing.+ -> Gang + -> (Int -> a -> b) + -> Dist a -> Dist b+imapD wFn gang f d + = imapD' wFn gang (\i x -> x `deepSeqD` f i x) d+{-# INLINE [0] imapD #-}+-- INLINE [0] because we want to wait until phase [0] before introducing+-- the call to imapD'. Our rewrite rules operate directly on the imapD+-- formp, so once imapD is inlined no more fusion can take place.++ {-# RULES -"imapD/generateD" forall gang f g.- imapD gang f (generateD gang g) = generateD gang (\i -> f i (g i))+"imapD/generateD" + forall wMap wGen gang f g+ . imapD wMap gang f (generateD wGen gang g) + = generateD (W.WFMapGen wMap wGen) 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/generateD_cheap" + forall wMap wGen gang f g+ . imapD wMap gang f (generateD_cheap wGen gang g) + = generateD (W.WFMapGen wMap wGen) 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+"imapD/imapD" + forall wMap1 wMap2 gang f g d+ . imapD wMap1 gang f (imapD wMap2 gang g d) + = imapD (W.WFMapMap wMap1 wMap2) gang (\i x -> f i (g i x)) d #-} @@ -97,74 +97,54 @@ -- 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-zipWithD g f dx dy = mapD g (uncurry f) (zipD dx dy)+ => W.What -- ^ What is the worker function doing.+ -> Gang + -> (a -> b -> c) + -> Dist a -> Dist b -> Dist c++zipWithD what g f dx dy + = mapD what g (uncurry f) (zipD dx dy) {-# INLINE zipWithD #-} -- | 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-izipWithD g f dx dy = imapD g (\i -> uncurry (f i)) (zipD dx dy)+ => W.What -- ^ What is the worker function doing.+ -> Gang + -> (Int -> a -> b -> c) + -> Dist a -> Dist b -> Dist c++izipWithD what g f dx dy + = imapD what g (\i -> uncurry (f i)) (zipD dx dy) {-# INLINE izipWithD #-} {-# 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[1]" + forall gang f xs ys what+ . zipD (imapD what gang f xs) ys+ = imapD what 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/imapD[2]" + forall gang f xs ys what+ . zipD xs (imapD what gang f ys)+ = imapD what 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[1]" + forall gang f xs what+ . zipD (generateD what gang f) xs+ = imapD what 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+"zipD/generateD[2]" + forall gang f xs what+ . zipD xs (generateD what gang f)+ = imapD what gang (\i x -> (x, f i)) xs #-} --- Folding ----------------------------------------------------------------------- | Fold all the instances of a distributed value.-foldD :: DT a => Gang -> (a -> a -> a) -> Dist a -> a-foldD g f !d - = checkGangD ("here foldD") g d - $ fold 1 (indexD (here "foldD") d 0)- where- !n = gangSize g- --- fold i x | i == n = x- | otherwise = fold (i+1) (f x $ indexD (here "foldD") d i)-{-# NOINLINE foldD #-}----- | Prefix sum of the instances of a distributed value.-scanD :: forall a. DT a => Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)-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 $ indexD (here "scanD") d i)-{-# NOINLINE scanD #-}--+-- MapAccumL ------------------------------------------------------------------ -- | Combination of map and fold. mapAccumLD :: forall a b acc. (DT a, DT b)@@ -192,49 +172,3 @@ go md (i+1) acc'' {-# INLINE_DIST mapAccumLD #-} ---- 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 ()-mapDST_ g p d - = mapDST_' g (\x -> x `deepSeqD` p x) d-{-# INLINE mapDST_ #-}---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)-mapDST g p !d = mapDST' g (\x -> x `deepSeqD` p x) d-{-# INLINE mapDST #-}---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 ()-zipWithDST_ g p !dx !dy - = mapDST_ g (uncurry p) (zipD dx dy)-{-# INLINE zipWithDST_ #-}---zipWithDST - :: (DT a, DT b, DT c)- => Gang- -> (a -> b -> DistST s c) -> Dist a -> Dist b -> ST s (Dist c)-zipWithDST g p !dx !dy - = mapDST g (uncurry p) (zipD dx dy)-{-# INLINE zipWithDST #-}-
+ Data/Array/Parallel/Unlifted/Distributed/Data/Bool.hs view
@@ -0,0 +1,50 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++module Data.Array.Parallel.Unlifted.Distributed.Data.Bool+ ( orD+ , andD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base ()+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import qualified Data.Array.Parallel.Unlifted.Distributed.What as W+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as V+import qualified Data.Vector.Unboxed.Mutable as MV+import Prelude as P+++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+++-- | OR together all instances of a distributed 'Bool'.+orD :: Gang -> Dist Bool -> Bool+orD g = foldD (W.What "orD") g (||)+{-# INLINE_DIST orD #-}+++-- | AND together all instances of a distributed 'Bool'.+andD :: Gang -> Dist Bool -> Bool+andD g = foldD (W.What "andD") g (&&)+{-# INLINE_DIST andD #-}++
+ Data/Array/Parallel/Unlifted/Distributed/Data/Maybe.hs view
@@ -0,0 +1,63 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Maybes.+module Data.Array.Parallel.Unlifted.Distributed.Data.Maybe +where+import Data.Array.Parallel.Unlifted.Distributed.Data.Bool ()+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim ()+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Control.Monad+++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 str (DMaybe bs as) i+ | indexD (str ++ "/indexD[Maybe]") bs i+ = Just $ indexD (str ++ "/indexD[Maybe]" ++ str) as i++ | otherwise = Nothing+ {-# INLINE_DIST indexD #-}++ newMD g+ = liftM2 MDMaybe (newMD g) (newMD g)+ {-# INLINE_DIST newMD #-}++ readMD (MDMaybe bs as) i + = do b <- readMD bs i+ if b then liftM Just $ readMD as i+ else return Nothing+ {-# INLINE_DIST readMD #-}++ writeMD (MDMaybe bs _) i Nothing + = writeMD bs i False++ writeMD (MDMaybe bs as) i (Just x)+ = do writeMD bs i True+ writeMD as i x+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDMaybe bs as)+ = liftM2 DMaybe (unsafeFreezeMD bs)+ (unsafeFreezeMD as)+ {-# INLINE_DIST unsafeFreezeMD #-}++ deepSeqD Nothing z = z+ deepSeqD (Just x) z = deepSeqD x z+ {-# INLINE_DIST deepSeqD #-}++ sizeD (DMaybe b _) + = sizeD b+ {-# INLINE_DIST sizeD #-}++ sizeMD (MDMaybe b _) + = sizeMD b+ {-# INLINE_DIST sizeMD #-}++ measureD Nothing = "Nothing"+ measureD (Just x) = "Just (" ++ measureD x ++ ")"+ {-# NOINLINE measureD #-}+ -- NOINLINE because this is only used for debugging.
+ Data/Array/Parallel/Unlifted/Distributed/Data/Ordering.hs view
@@ -0,0 +1,30 @@++-- | Distribution of values of primitive types.+module Data.Array.Parallel.Unlifted.Distributed.Data.Ordering+where+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as V+import qualified Data.Vector.Unboxed.Mutable as MV+import Prelude as P+++instance DPrim Ordering where+ mkDPrim = DOrdering+ unDPrim (DOrdering a) = a++ mkMDPrim = MDOrdering+ unMDPrim (MDOrdering a) = a+++instance DT Ordering where+ data Dist Ordering = DOrdering !(V.Vector Ordering)+ data MDist Ordering s = MDOrdering !(MV.STVector s Ordering)++ indexD = primIndexD+ newMD = primNewMD+ readMD = primReadMD+ writeMD = primWriteMD+ unsafeFreezeMD = primUnsafeFreezeMD+ sizeD = primSizeD+ sizeMD = primSizeMD
+ Data/Array/Parallel/Unlifted/Distributed/Data/Scalar.hs view
@@ -0,0 +1,38 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of values of primitive types.+module Data.Array.Parallel.Unlifted.Distributed.Data.Scalar+ ( DT(..), Dist(..)+ , scalarD+ , sumD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base+import Data.Array.Parallel.Unlifted.Distributed.Data.Unit+import Data.Array.Parallel.Unlifted.Distributed.Combinators+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import qualified Data.Array.Parallel.Unlifted.Distributed.What as W+import Data.Array.Parallel.Pretty+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as V+import Prelude as P+++instance PprPhysical (Dist Int) where+ pprp (DInt xs)+ = text "DInt" <+> text (show $ V.toList xs)+ {-# NOINLINE pprp #-}+++-- | Distribute a scalar.+-- Each thread gets its own copy of the same value.+-- Example: scalarD theGangN4 10 = [10, 10, 10, 10] +scalarD :: DT a => Gang -> a -> Dist a+scalarD gang x + = mapD W.WScalar gang (const x) (unitD gang)+{-# INLINE_DIST scalarD #-}++-- | Sum all instances of a distributed number.+sumD :: (Num a, DT a) => Gang -> Dist a -> a+sumD g = foldD (W.What "sumD") g (+)+{-# INLINE_DIST sumD #-}
+ Data/Array/Parallel/Unlifted/Distributed/Data/Scalar/Base.hs view
@@ -0,0 +1,129 @@+{-# OPTIONS -fno-warn-orphans #-}+-- | `DPrim` and `DT` instances for scalar types. +module Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base+where+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Word+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as V+import qualified Data.Vector.Unboxed.Mutable as MV+import Prelude as P+++-- Integer -----------------------------------------------------------------------+-- FIXME: fake instances+instance DPrim Integer+instance DT Integer+++-- 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 " P.++ 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+
+ Data/Array/Parallel/Unlifted/Distributed/Data/Tuple.hs view
@@ -0,0 +1,173 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Tuples+module Data.Array.Parallel.Unlifted.Distributed.Data.Tuple + ( -- * Pairs+ zipD, unzipD, fstD, sndD+ + -- * Triples+ , zip3D, unzip3D)+where+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Data.Array.Parallel.Base+import Data.Array.Parallel.Pretty+import Control.Monad++here :: String -> String+here s = "Data.Array.Parallel.Unlifted.Distributed.Types.Tuple." ++ s+++-- 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 str d i+ = ( indexD (str ++ "/indexD[Tuple2]") (fstD d) i+ , indexD (str ++ "/indexD[Tuple2]") (sndD d) i)+ {-# INLINE_DIST indexD #-}++ newMD g+ = liftM2 MDProd (newMD g) (newMD g)+ {-# INLINE_DIST newMD #-}++ readMD (MDProd xs ys) i+ = liftM2 (,) (readMD xs i) (readMD ys i)+ {-# INLINE_DIST readMD #-}++ writeMD (MDProd xs ys) i (x,y)+ = do writeMD xs i x+ writeMD ys i y+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDProd xs ys)+ = liftM2 DProd (unsafeFreezeMD xs)+ (unsafeFreezeMD ys)+ {-# INLINE_DIST unsafeFreezeMD #-}++ deepSeqD (x, y) z + = deepSeqD x (deepSeqD y z)+ {-# INLINE deepSeqD #-}++ sizeD (DProd x _) + = sizeD x+ {-# INLINE_DIST sizeD #-}++ sizeMD (MDProd x _) + = sizeMD x+ {-# INLINE_DIST sizeMD #-}++ measureD (x, y) + = "Pair " ++ "(" ++ measureD x ++ ") (" ++ measureD y ++ ")"+ {-# NOINLINE measureD #-}+ -- NOINLINE beacuse this is only used during debugging.+++instance (PprPhysical (Dist a), PprPhysical (Dist b)) + => PprPhysical (Dist (a, b)) where+ pprp (DProd xs ys)+ = text "DProd"+ $$ (nest 8 $ vcat+ [ pprp xs+ , pprp ys ])+ {-# NOINLINE pprp #-}+ -- NOINLINE because this is only used during debugging.+++-- | 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)+zipD !x !y + = checkEq (here "zipDT") "Size mismatch" (sizeD x) (sizeD y) + $ DProd x y+{-# INLINE [0] zipD #-} -- TODO: why is this INLINE [0]???+++-- | Unpairing of distributed values.+unzipD :: (DT a, DT b) => Dist (a,b) -> (Dist a, Dist b)+unzipD (DProd dx dy) = (dx,dy)+{-# INLINE_DIST unzipD #-}+++-- | Extract the first elements of a distributed pair.+fstD :: (DT a, DT b) => Dist (a,b) -> Dist a+fstD = fst . unzipD+{-# INLINE_DIST fstD #-}+++-- | Extract the second elements of a distributed pair.+sndD :: (DT a, DT b) => Dist (a,b) -> Dist b+sndD = snd . unzipD+{-# INLINE_DIST sndD #-}+++-- Triples --------------------------------------------------------------------+instance (DT a, DT b, DT c) => DT (a,b,c) where+ data Dist (a,b,c) = DProd3 !(Dist a) !(Dist b) !(Dist c)+ data MDist (a,b,c) s = MDProd3 !(MDist a s) !(MDist b s) !(MDist c s)++ indexD str (DProd3 xs ys zs) i+ = ( indexD (here $ "indexD[Tuple3]/" ++ str) xs i+ , indexD (here $ "indexD[Tuple3]/" ++ str) ys i+ , indexD (here $ "indexD[Tuple3]/" ++ str) zs i)+ {-# INLINE_DIST indexD #-}++ newMD g+ = liftM3 MDProd3 (newMD g) (newMD g) (newMD g)+ {-# INLINE_DIST newMD #-}++ readMD (MDProd3 xs ys zs) i+ = liftM3 (,,) (readMD xs i) (readMD ys i) (readMD zs i)+ {-# INLINE_DIST readMD #-}++ writeMD (MDProd3 xs ys zs) i (x,y,z)+ = do writeMD xs i x+ writeMD ys i y+ writeMD zs i z+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDProd3 xs ys zs)+ = liftM3 DProd3 (unsafeFreezeMD xs)+ (unsafeFreezeMD ys)+ (unsafeFreezeMD zs)+ {-# INLINE_DIST unsafeFreezeMD #-}++ deepSeqD (x,y,z) k + = deepSeqD x (deepSeqD y (deepSeqD z k))+ {-# INLINE_DIST deepSeqD #-}++ sizeD (DProd3 x _ _) + = sizeD x+ {-# INLINE_DIST sizeD #-}++ sizeMD (MDProd3 x _ _) + = sizeMD x+ {-# INLINE_DIST sizeMD #-}++ measureD (x,y,z)+ = "Triple " + ++ "(" ++ measureD x ++ ") "+ ++ "(" ++ measureD y ++ ") "+ ++ "(" ++ measureD z ++ ")"+ {-# NOINLINE measureD #-}+ -- NOINLINE because this is only used for debugging.+++-- | Pairing of distributed values.+-- /The two values must belong to the same/ 'Gang'.+zip3D :: (DT a, DT b, DT c) => Dist a -> Dist b -> Dist c -> Dist (a,b,c)+zip3D !x !y !z+ = checkEq (here "zip3DT") "Size mismatch" (sizeD x) (sizeD y) + $ checkEq (here "zip3DT") "Size mismatch" (sizeD x) (sizeD z) + $ DProd3 x y z+{-# INLINE [0] zip3D #-} -- TODO: Why is this INLINE[0]??+++-- | Unpairing of distributed values.+unzip3D :: (DT a, DT b, DT c) => Dist (a,b,c) -> (Dist a, Dist b, Dist c)+unzip3D (DProd3 dx dy dz) + = (dx,dy,dz)+{-# INLINE_DIST unzip3D #-}+
+ Data/Array/Parallel/Unlifted/Distributed/Data/USSegd.hs view
@@ -0,0 +1,18 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USSegd + ( lengthD+ , takeLengthsD+ , takeIndicesD+ , takeElementsD+ , takeStartsD+ , takeSourcesD+ , takeUSegdD+ , splitSSegdOnElemsD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.DT ()+import Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base+import Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Split
+ Data/Array/Parallel/Unlifted/Distributed/Data/USSegd/Base.hs view
@@ -0,0 +1,71 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base+ ( lengthD+ , takeLengthsD+ , takeIndicesD+ , takeElementsD+ , takeStartsD+ , takeSourcesD+ , takeUSegdD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.DT+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd)+import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)+import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)+import Prelude as P+import qualified Data.Array.Parallel.Unlifted.Distributed.Data.USegd as DUSegd+import qualified Data.Array.Parallel.Unlifted.Distributed.Data.Vector as DV+++-- | O(1). Yield the overall number of segments.+lengthD :: Dist USSegd -> Dist Int+lengthD (DUSSegd starts _ _) + = DV.lengthD starts+{-# INLINE_DIST lengthD #-}+++-- | O(1). Yield the lengths of the individual segments.+takeLengthsD :: Dist USSegd -> Dist (Vector Int)+takeLengthsD (DUSSegd _ _ usegds)+ = DUSegd.takeLengthsD usegds+{-# INLINE_DIST takeLengthsD #-}+++-- | O(1). Yield the segment indices.+takeIndicesD :: Dist USSegd -> Dist (Vector Int)+takeIndicesD (DUSSegd _ _ usegds)+ = DUSegd.takeIndicesD usegds+{-# INLINE_DIST takeIndicesD #-}+++-- | O(1). Yield the number of data elements.+takeElementsD :: Dist USSegd -> Dist Int+takeElementsD (DUSSegd _ _ usegds)+ = DUSegd.takeElementsD usegds+{-# INLINE_DIST takeElementsD #-}+++-- | O(1). Yield the starting indices.+takeStartsD :: Dist USSegd -> Dist (Vector Int)+takeStartsD (DUSSegd starts _ _)+ = starts+{-# INLINE_DIST takeStartsD #-}+ ++-- | O(1). Yield the source ids+takeSourcesD :: Dist USSegd -> Dist (Vector Int)+takeSourcesD (DUSSegd _ sources _)+ = sources+{-# INLINE_DIST takeSourcesD #-}+++-- | O(1). Yield the USegd+takeUSegdD :: Dist USSegd -> Dist USegd+takeUSegdD (DUSSegd _ _ usegd)+ = usegd+{-# INLINE_DIST takeUSegdD #-}
+ Data/Array/Parallel/Unlifted/Distributed/Data/USSegd/DT.hs view
@@ -0,0 +1,88 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.DT+where+import Data.Array.Parallel.Unlifted.Distributed.Data.USegd ()+import Data.Array.Parallel.Unlifted.Distributed.Data.Vector ()+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd)+import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)+import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)+import Data.Array.Parallel.Pretty+import Control.Monad+import Prelude as P+import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd++instance DT USSegd where+ data Dist USSegd + = DUSSegd !(Dist (Vector Int)) -- segment starts+ !(Dist (Vector Int)) -- segment sources+ !(Dist USegd) -- distributed usegd++ data MDist USSegd s + = MDUSSegd !(MDist (Vector Int) s) -- segment starts+ !(MDist (Vector Int) s) -- segment sources+ !(MDist USegd s) -- distributed usegd++ indexD str (DUSSegd starts sources usegds) i+ = USSegd.mkUSSegd+ (indexD (str ++ "/indexD[USSegd]") starts i)+ (indexD (str ++ "/indexD[USSegd]") sources i)+ (indexD (str ++ "/indexD[USSegd]") usegds i)+ {-# INLINE_DIST indexD #-}++ newMD g+ = liftM3 MDUSSegd (newMD g) (newMD g) (newMD g)+ {-# INLINE_DIST newMD #-}++ readMD (MDUSSegd starts sources usegds) i+ = liftM3 USSegd.mkUSSegd (readMD starts i) (readMD sources i) (readMD usegds i)+ {-# INLINE_DIST readMD #-}++ writeMD (MDUSSegd starts sources usegds) i ussegd+ = do writeMD starts i (USSegd.takeStarts ussegd)+ writeMD sources i (USSegd.takeSources ussegd)+ writeMD usegds i (USSegd.takeUSegd ussegd)+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDUSSegd starts sources usegds)+ = liftM3 DUSSegd (unsafeFreezeMD starts)+ (unsafeFreezeMD sources)+ (unsafeFreezeMD usegds)+ {-# INLINE_DIST unsafeFreezeMD #-}++ deepSeqD ussegd z+ = deepSeqD (USSegd.takeStarts ussegd)+ $ deepSeqD (USSegd.takeSources ussegd)+ $ deepSeqD (USSegd.takeUSegd ussegd) z+ {-# INLINE_DIST deepSeqD #-}++ sizeD (DUSSegd _ _ usegd) + = sizeD usegd+ {-# INLINE_DIST sizeD #-}++ sizeMD (MDUSSegd _ _ usegd) + = sizeMD usegd+ {-# INLINE_DIST sizeMD #-}++ measureD ussegd + = "USSegd " P.++ show (USSegd.takeStarts ussegd)+ P.++ " " P.++ show (USSegd.takeSources ussegd)+ P.++ " " P.++ measureD (USSegd.takeUSegd ussegd)+ {-# NOINLINE measureD #-}+ -- NOINLINE because this is only used for debugging.+++instance PprPhysical (Dist USSegd) where+ pprp (DUSSegd starts sources usegds)+ = text "DUSSegd"+ $$ (nest 7 $ vcat+ [ text "starts: " <+> pprp starts+ , text "sources: " <+> pprp sources+ , text "usegds: " <+> pprp usegds])+ {-# NOINLINE pprp #-}+ -- NOINLINE because this is only used for debugging.+
+ Data/Array/Parallel/Unlifted/Distributed/Data/USSegd/Split.hs view
@@ -0,0 +1,250 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Operations on Distributed Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Split + (splitSSegdOnElemsD)+where+import Data.Array.Parallel.Unlifted.Distributed.Arrays+import Data.Array.Parallel.Unlifted.Distributed.Combinators+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd)+import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)+import Data.Array.Parallel.Base+import Data.Bits (shiftR)+import Control.Monad (when)+import Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.DT ()+import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd+import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Debug.Trace++here :: String -> String+here s = "Data.Array.Parallel.Unlifted.Distributed.USSegd." ++ s++-------------------------------------------------------------------------------+-- | Split a segment descriptor across the gang, element wise.+-- We try to put the same number of elements on each thread, which means+-- that segments are sometimes split across threads.+--+-- Each thread gets a slice of segment descriptor, the segid of the first +-- slice, and the offset of the first slice in its segment.+-- +-- Example:+-- In this picture each X represents 5 elements, and we have 5 segements in total.+--+-- @ segs: ----------------------- --- ------- --------------- -------------------+-- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|+-- | thread1 | thread2 | thread3 | thread4 |+-- segid: 0 0 3 4+-- offset: 0 45 0 5+--+-- pprp $ splitSegdOnElemsD theGang +-- $ lengthsToUSegd $ fromList [60, 10, 20, 40, 50 :: Int]+--+-- segd: DUSegd lengths: DVector lengths: [1,3,2,1]+-- chunks: [[45],[15,10,20],[40,5],[45]]+-- indices: DVector lengths: [1,3,2,1]+-- chunks: [[0], [0,15,25], [0,40],[0]]+-- elements: DInt [45,45,45,45]+--+-- segids: DInt [0,0,3,4] (segment id of first slice on thread)+-- offsets: DInt [0,45,0,5] (offset of that slice in its segment)+-- @+--+splitSSegdOnElemsD :: Gang -> USSegd -> Dist ((USSegd,Int),Int)+splitSSegdOnElemsD g !segd + = {-# SCC "splitSSegdOnElemsD" #-}+ traceEvent ("dph-prim-par: USSegd.splitSSegdOnElems")+ $ imapD (What "UPSSegd.splitSSegdOnElems/splitLenIx") g mk + (splitLenIdxD g (USegd.takeElements $ USSegd.takeUSegd segd))+ where + -- Number of threads in gang.+ !nThreads = gangSize g+++ -- Build a USSegd from just the lengths, starts and sources fields.+ -- The indices and elems fields of the contained USegd are + -- generated from the lengths.+ buildUSSegd :: Vector Int -> Vector Int -> Vector Int -> USSegd+ buildUSSegd lengths starts sources+ = USSegd.mkUSSegd starts sources+ $ USegd.fromLengths lengths++ -- Determine what elements go on a thread+ mk :: Int -- Thread index.+ -> (Int, Int) -- Number of elements on this thread,+ -- and starting offset into the flat array.+ -> ((USSegd, Int), Int) -- Segd for this thread, segid of first slice,+ -- and offset of first slice.++ mk i (nElems, ixStart) + = case chunk segd ixStart nElems (i == nThreads - 1) of+ (# lengths, starts, sources, l, o #) + -> ((buildUSSegd lengths starts sources, l), o)++{-# NOINLINE splitSSegdOnElemsD #-}+-- NOINLINE because it's complicated and won't fuse with anything.+-- This function has a large body of code and we don't want to blow up+-- the client modules by inlining it everywhere.+++-------------------------------------------------------------------------------+-- | Determine what elements go on a thread.+-- The 'chunk' refers to the a chunk of the flat array, and is defined+-- by a set of segment slices. +--+-- Example:+-- In this picture each X represents 5 elements, and we have 5 segements in total.+--+-- @ segs: ----------------------- --- ------- --------------- -------------------+-- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|+-- | thread1 | thread2 | thread3 | thread4 |+-- segid: 0 0 3 4+-- offset: 0 45 0 5+-- k: 0 1 3 5+-- k': 1 3 5 5+-- left: 0 15 0 45+-- right: 45 20 5 0+-- left_len: 0 1 0 1+-- left_off: 0 45 0 5+-- n': 1 3 2 1+-- @+chunk :: USSegd -- ^ Segment descriptor of entire array.+ -> Int -- ^ Starting offset into the flat array for the first+ -- slice on this thread.+ -> Int -- ^ Number of elements in this thread.+ -> Bool -- ^ Whether this is the last thread in the gang.+ -> (# Vector Int -- Lengths of segment slices, + , Vector Int -- Starting index of data in its vector+ , Vector Int -- Source id+ , Int -- segid of first slice+ , Int #) -- offset of first slice.++chunk !ussegd !nStart !nElems is_last+ = (# lengths', starts', sources', k-left_len, left_off #)+ where+ -- Lengths of all segments.+ -- eg: [60, 10, 20, 40, 50]+ lengths = USSegd.takeLengths ussegd++ -- Indices indices of all segments.+ -- eg: [0, 60, 70, 90, 130]+ indices = USSegd.takeIndices ussegd++ -- Starting indices for all segments.+ starts = USSegd.takeStarts ussegd++ -- Source ids for all segments.+ sources = USSegd.takeSources ussegd+ + -- Total number of segments defined by segment descriptor.+ -- eg: 5+ n = Seq.length lengths++ -- Segid of the first seg that starts after the left of this chunk.+ k = search nStart indices++ -- Segid of the first seg that starts after the right of this chunk.+ k' | is_last = n+ | otherwise = search (nStart + nElems) indices++ -- The length of the left-most slice of this chunk.+ left | k == n = nElems+ | otherwise = min ((Seq.index (here "chunk") indices k) - nStart) nElems++ -- The length of the right-most slice of this chunk.+ length_right + | k' == k = 0+ | otherwise = nStart + nElems - (Seq.index (here "chunk") indices (k'-1))++ -- Whether the first element in this chunk is an internal element of+ -- of a segment. Alternatively, indicates that the first element of + -- the chunk is not the first element of a segment. + left_len | left == 0 = 0+ | otherwise = 1++ -- If the first element of the chunk starts within a segment, + -- then gives the index within that segment, otherwise 0.+ left_off | left == 0 = 0+ | otherwise = nStart - (Seq.index (here "chunk") indices (k-1))++ -- How many segments this chunk straddles.+ n' = left_len + (k'-k)++ -- Create the lengths for this chunk by first copying out the lengths+ -- from the original segment descriptor. If the slices on the left+ -- and right cover partial segments, then we update the corresponding+ -- lengths.+ (!lengths', !starts', !sources')+ = runST (do+ -- Create a new array big enough to hold all the lengths for this chunk.+ mlengths' <- Seq.newM n'+ msources' <- Seq.newM n'+ mstarts' <- Seq.newM n'++ -- If the first element is inside a segment, + -- then update the length to be the length of the slice.+ when (left /= 0) + $ do Seq.write mlengths' 0 left+ Seq.write mstarts' 0 (Seq.index (here "chunk") starts (k - left_len) + left_off)+ Seq.write msources' 0 (Seq.index (here "chunk") sources (k - left_len))++ -- Copy out array lengths for this chunk.+ Seq.copy (Seq.mdrop left_len mlengths') (Seq.slice (here "chunk") lengths k (k'-k))+ Seq.copy (Seq.mdrop left_len mstarts') (Seq.slice (here "chunk") starts k (k'-k))+ Seq.copy (Seq.mdrop left_len msources') (Seq.slice (here "chunk") sources k (k'-k))++ -- If the last element is inside a segment, + -- then update the length to be the length of the slice.+ when (length_right /= 0)+ $ do Seq.write mlengths' (n' - 1) length_right++ clengths' <- Seq.unsafeFreeze mlengths'+ cstarts' <- Seq.unsafeFreeze mstarts'+ csources' <- Seq.unsafeFreeze msources'+ return (clengths', cstarts', csources'))++{- = trace + (render $ vcat+ [ text "CHUNK"+ , pprp segd+ , text "nStart: " <+> int nStart+ , text "nElems: " <+> int nElems+ , text "k: " <+> int k+ , text "k': " <+> int k'+ , text "left: " <+> int left+ , text "right: " <+> int right+ , text "left_len:" <+> int left_len+ , text "left_off:" <+> int left_off+ , text "n': " <+> int n'+ , text ""]) lens'+-}++{-# INLINE_DIST chunk #-}+-- INLINE_DIST even though it should be inlined into splitSSegdOnElemsD anyway+-- because that function contains the only use.+++-------------------------------------------------------------------------------+-- O(log n).+-- Given a monotonically increasing vector of `Int`s,+-- find the first element that is larger than the given value.+-- +-- eg search 75 [0, 60, 70, 90, 130] = 90+-- search 43 [0, 60, 70, 90, 130] = 60+--+search :: Int -> Vector Int -> Int+search !x ys = go 0 (Seq.length ys)+ where+ go i n | n <= 0 = i+ | Seq.index (here "search") ys mid < x+ = go (mid + 1) (n - half - 1)+ | otherwise = go i half+ where+ half = n `shiftR` 1+ mid = i + half+{-# INLINE_DIST search #-}+-- INLINE_DIST because we want it inlined into both uses in 'chunk' above.+
+ Data/Array/Parallel/Unlifted/Distributed/Data/USegd.hs view
@@ -0,0 +1,20 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USegd + ( mkDUSegdD+ , lengthD+ , takeLengthsD+ , takeIndicesD+ , takeElementsD+ , splitSegdOnSegsD+ , splitSegdOnElemsD+ , splitSD+ , joinSegdD+ , glueSegdD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.USegd.DT ()+import Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base+import Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split
+ Data/Array/Parallel/Unlifted/Distributed/Data/USegd/Base.hs view
@@ -0,0 +1,56 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base+ ( mkDUSegdD+ , lengthD+ , takeLengthsD+ , takeIndicesD+ , takeElementsD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.USegd.DT+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)+import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)+import qualified Data.Array.Parallel.Unlifted.Distributed.Data.Vector as DV+++-- | O(1). Construct a distributed segment descriptor+mkDUSegdD + :: Dist (Vector Int) -- ^ segment lengths+ -> Dist (Vector Int) -- ^ segment indices+ -> Dist Int -- ^ number of elements in each chunk+ -> Dist USegd++mkDUSegdD = DUSegd+{-# INLINE_DIST mkDUSegdD #-}+++-- | O(1). Yield the overall number of segments.+lengthD :: Dist USegd -> Dist Int+lengthD (DUSegd lens _ _) + = DV.lengthD lens+{-# INLINE_DIST lengthD #-}+++-- | O(1). Yield the lengths of the individual segments.+takeLengthsD :: Dist USegd -> Dist (Vector Int)+takeLengthsD (DUSegd lens _ _ )+ = lens+{-# INLINE_DIST takeLengthsD #-}+++-- | O(1). Yield the segment indices of a segment descriptor.+takeIndicesD :: Dist USegd -> Dist (Vector Int)+takeIndicesD (DUSegd _ idxs _)+ = idxs+{-# INLINE_DIST takeIndicesD #-}+++-- | O(1). Yield the number of data elements.+takeElementsD :: Dist USegd -> Dist Int+takeElementsD (DUSegd _ _ dns)+ = dns+{-# INLINE_DIST takeElementsD #-}
+ Data/Array/Parallel/Unlifted/Distributed/Data/USegd/DT.hs view
@@ -0,0 +1,84 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USegd.DT+where+import Data.Array.Parallel.Unlifted.Distributed.Data.Vector ()+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)+import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)+import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd+import Prelude as P+import Data.Array.Parallel.Pretty+import Control.Monad++instance DT USegd where+ data Dist USegd + = DUSegd !(Dist (Vector Int)) -- segment lengths+ !(Dist (Vector Int)) -- segment indices+ !(Dist Int) -- number of elements in this chunk++ data MDist USegd s + = MDUSegd !(MDist (Vector Int) s) -- segment lengths+ !(MDist (Vector Int) s) -- segment indices+ !(MDist Int s) -- number of elements in this chunk++ indexD str (DUSegd lens idxs eles) i+ = USegd.mkUSegd+ (indexD (str ++ "/indexD[USegd]") lens i)+ (indexD (str ++ "/indexD[USegd]") idxs i)+ (indexD (str ++ "/indexD[USegd]") eles i)+ {-# INLINE_DIST indexD #-}++ newMD g+ = liftM3 MDUSegd (newMD g) (newMD g) (newMD g)+ {-# INLINE_DIST newMD #-}++ readMD (MDUSegd lens idxs eles) i+ = liftM3 USegd.mkUSegd (readMD lens i) (readMD idxs i) (readMD eles i)+ {-# INLINE_DIST readMD #-}++ writeMD (MDUSegd lens idxs eles) i segd+ = do writeMD lens i (USegd.takeLengths segd)+ writeMD idxs i (USegd.takeIndices segd)+ writeMD eles i (USegd.takeElements segd)+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDUSegd lens idxs eles)+ = liftM3 DUSegd (unsafeFreezeMD lens)+ (unsafeFreezeMD idxs)+ (unsafeFreezeMD eles)+ {-# INLINE_DIST unsafeFreezeMD #-}++ deepSeqD segd z+ = deepSeqD (USegd.takeLengths segd)+ $ deepSeqD (USegd.takeIndices segd)+ $ deepSeqD (USegd.takeElements segd) z+ {-# INLINE_DIST deepSeqD #-}++ sizeD (DUSegd _ _ eles) + = sizeD eles+ {-# INLINE_DIST sizeD #-}++ sizeMD (MDUSegd _ _ eles) + = sizeMD eles+ {-# INLINE_DIST sizeMD #-}++ measureD segd + = "Segd " P.++ show (USegd.length segd)+ P.++ " " P.++ show (USegd.takeElements segd)+ {-# NOINLINE measureD #-}+ -- NOINLINE because this is only used for debugging.+++instance PprPhysical (Dist USegd) where+ pprp (DUSegd lens indices elements)+ = text "DUSegd"+ $$ (nest 7 $ vcat+ [ text "lengths: " <+> pprp lens+ , text "indices: " <+> pprp indices+ , text "elements:" <+> pprp elements])+ {-# NOINLINE pprp #-}+ -- NOINLINE because this is only used for debugging.
+ Data/Array/Parallel/Unlifted/Distributed/Data/USegd/Split.hs view
@@ -0,0 +1,361 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Operations on Distributed Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split+ ( splitSegdOnSegsD+ , splitSegdOnElemsD+ , splitSD+ , joinSegdD+ , glueSegdD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base+import Data.Array.Parallel.Unlifted.Distributed.Arrays+import Data.Array.Parallel.Unlifted.Distributed.Combinators+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)+import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector, Unbox)+import Data.Array.Parallel.Base+import Data.Bits (shiftR)+import Control.Monad (when)+import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import Debug.Trace++here :: String -> String+here s = "Data.Array.Parallel.Unlifted.Distributed.USegd." ++ s++-------------------------------------------------------------------------------+-- | Split a segment descriptor across the gang, segment wise.+-- Whole segments are placed on each thread, and we try to balance out+-- the segments so each thread has the same number of array elements.+--+-- We don't split segments across threads, as this would limit our ability+-- to perform intra-thread fusion of lifted operations. The down side+-- of this is that if we have few segments with an un-even size distribution+-- then large segments can cause the gang to become unbalanced.+--+-- In the following example the segment with size 100 dominates and+-- unbalances the gang. There is no reason to put any segments on the+-- the last thread because we need to wait for the first to finish anyway.+--+-- @ > pprp $ splitSegdOnSegsD theGang+-- $ lengthsToUSegd $ fromList [100, 10, 20, 40, 50 :: Int]+-- +-- DUSegd lengths: DVector lengths: [ 1, 3, 1, 0]+-- chunks: [[100],[10,20,40],[50],[]]+-- +-- indices: DVector lengths: [1,3,1,0]+-- chunks: [[0], [0,10,30], [0], []]+--+-- elements: DInt [100,70,50,0]+-- @+--+-- NOTE: This splitSegdOnSegsD function isn't currently used.+--+splitSegdOnSegsD :: Gang -> USegd -> Dist USegd+splitSegdOnSegsD g !segd + = traceEvent ("dph-prim-par: USegd.splitSegdOnSegsD")+ $ mapD (What "USegd.splitSegdOnSegds/fromLengths") g USegd.fromLengths+ $ splitAsD g d lens+ where+ !d = snd+ . mapAccumLD g chunks 0+ . splitLenD g+ $ USegd.takeElements segd++ n = USegd.length segd+ lens = USegd.takeLengths segd++ chunks !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 = Seq.index (here "splitSegdOnSegsD") lens i+{-# NOINLINE splitSegdOnSegsD #-}+-- NOINLINE because it won't fuse with anything.+++-------------------------------------------------------------------------------+-- | Split a segment descriptor across the gang, element wise.+-- We try to put the same number of elements on each thread, which means+-- that segments are sometimes split across threads.+--+-- Each thread gets a slice of segment descriptor, the segid of the first +-- slice, and the offset of the first slice in its segment.+-- +-- Example:+-- In this picture each X represents 5 elements, and we have 5 segements in total.+--+-- @ segs: ----------------------- --- ------- --------------- -------------------+-- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|+-- | thread1 | thread2 | thread3 | thread4 |+-- segid: 0 0 3 4+-- offset: 0 45 0 5+--+-- pprp $ splitSegdOnElemsD theGang4+-- $ lengthsToUSegd $ fromList [60, 10, 20, 40, 50 :: Int]+--+-- segd: DUSegd lengths: DVector lengths: [1,3,2,1]+-- chunks: [[45],[15,10,20],[40,5],[45]]+-- indices: DVector lengths: [1,3,2,1]+-- chunks: [[0], [0,15,25], [0,40],[0]]+-- elements: DInt [45,45,45,45]+--+-- segids: DInt [0,0,3,4] (segment id of first slice on thread)+-- offsets: DInt [0,45,0,5] (offset of that slice in its segment)+-- @+--+splitSegdOnElemsD :: Gang -> USegd -> Dist ((USegd,Int),Int)+splitSegdOnElemsD g !segd + = {-# SCC "splitSegdOnElemsD" #-} + traceEvent ("dph-prim-par: USegd.splitSegdOnElemsD")+ $ imapD (What "USegd.splitSegdOnElemsD/splitLenIdx") + g mk (splitLenIdxD g (USegd.takeElements segd))+ where + -- Number of threads in gang.+ !nThreads = gangSize g++ -- Determine what elements go on a thread+ mk :: Int -- Thread index.+ -> (Int, Int) -- Number of elements on this thread,+ -- and starting offset into the flat array.+ -> ((USegd, Int), Int) -- Segd for this thread, segid of first slice,+ -- and offset of first slice.++ mk i (nElems, ixStart) + = case getChunk segd ixStart nElems (i == nThreads - 1) of+ (# lens, l, o #) -> ((USegd.fromLengths lens, l), o)++{-# NOINLINE splitSegdOnElemsD #-}+-- NOINLINE because it won't fuse with anything.+++-------------------------------------------------------------------------------+-- | Determine what elements go on a thread.+-- The 'chunk' refers to the a chunk of the flat array, and is defined+-- by a set of segment slices. +--+-- Example:+-- In this picture each X represents 5 elements, and we have 5 segements in total.+--+-- @+-- segs: ----------------------- --- ------- --------------- -------------------+-- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|+-- | thread1 | thread2 | thread3 | thread4 |+-- segid: 0 0 3 4+-- offset: 0 45 0 5+-- k: 0 1 3 5+-- k': 1 3 5 5+-- left: 0 15 0 45+-- right: 45 20 5 0+-- left_len: 0 1 0 1+-- left_off: 0 45 0 5+-- n': 1 3 2 1+-- @+getChunk+ :: USegd -- ^ Segment descriptor of entire array.+ -> Int -- ^ Starting offset into the flat array for the first+ -- slice on this thread.+ -> Int -- ^ Number of elements in this thread.+ -> Bool -- ^ Whether this is the last thread in the gang.+ -> (# Vector Int -- Lengths of segment slices, + , Int -- segid of first slice,+ , Int #) -- offset of first slice.++getChunk !segd !nStart !nElems is_last+ = (# lens'', k-left_len, left_off #)+ where+ -- Lengths of all segments.+ -- eg: [60, 10, 20, 40, 50]+ !lens = USegd.takeLengths segd++ -- Indices indices of all segments.+ -- eg: [0, 60, 70, 90, 130]+ !idxs = USegd.takeIndices segd+ + -- Total number of segments defined by segment descriptor.+ -- eg: 5+ !n = Seq.length lens++ -- Segid of the first seg that starts after the left of this chunk.+ !k = search nStart idxs++ -- Segid of the first seg that starts after the right of this chunk.+ !k' | is_last = n+ | otherwise = search (nStart + nElems) idxs++ -- The length of the left-most slice of this chunk.+ !left | k == n = nElems+ | otherwise = min ((Seq.index (here "getChunk") idxs k) - nStart) nElems++ -- The length of the right-most slice of this chunk.+ !right | k' == k = 0+ | otherwise = nStart + nElems - (Seq.index (here "getChunk") idxs (k'-1))++ -- Whether the first element in this chunk is an internal element of+ -- of a segment. Alternatively, indicates that the first element of + -- the chunk is not the first element of a segment. + !left_len | left == 0 = 0+ | otherwise = 1++ -- If the first element of the chunk starts within a segment, + -- then gives the index within that segment, otherwise 0.+ !left_off | left == 0 = 0+ | otherwise = nStart - (Seq.index (here "getChunk") idxs (k-1))++ -- How many segments this chunk straddles.+ !n' = left_len + (k'-k)++ -- Create the lengths for this chunk by first copying out the lengths+ -- from the original segment descriptor. If the slices on the left+ -- and right cover partial segments, then we update the corresponding+ -- lengths.+ !lens' + = runST (do+ -- Create a new array big enough to hold all the lengths for this chunk.+ !mlens' <- Seq.newM n'++ -- If the first element is inside a segment, + -- then update the length to be the length of the slice.+ when (left /= 0) + $ Seq.write mlens' 0 left++ -- Copy out array lengths for this chunk.+ Seq.copy (Seq.mdrop left_len mlens')+ (Seq.slice "getChunk" lens k (k'-k))++ -- If the last element is inside a segment, + -- then update the length to be the length of the slice.+ when (right /= 0)+ $ Seq.write mlens' (n' - 1) right++ Seq.unsafeFreeze mlens')++ !lens'' = lens'+{- = trace + (render $ vcat+ [ text "CHUNK"+ , pprp segd+ , text "nStart: " <+> int nStart+ , text "nElems: " <+> int nElems+ , text "k: " <+> int k+ , text "k': " <+> int k'+ , text "left: " <+> int left+ , text "right: " <+> int right+ , text "left_len:" <+> int left_len+ , text "left_off:" <+> int left_off+ , text "n': " <+> int n'+ , text ""]) lens'+-}++{-# INLINE_DIST getChunk #-}+-- INLINE_DIST because we want this inlined into splitSegdOnElemsD+-- above, which is the only use.++-------------------------------------------------------------------------------+-- O(log n). Given a monotonically increasing vector of `Int`s,+-- find the first element that is larger than the given value.+-- +-- eg search 75 [0, 60, 70, 90, 130] = 90+-- search 43 [0, 60, 70, 90, 130] = 60+--+search :: Int -> Vector Int -> Int+search !x ys = go 0 (Seq.length ys)+ where+ go i n | n <= 0 = i++ | Seq.index (here "search") ys mid < x + = go (mid + 1) (n - half - 1)++ | otherwise = go i half+ where+ half = n `shiftR` 1+ mid = i + half+{-# INLINE_DIST search #-}+-- INLINE_DIST because we want this inlined into both uses in getChunk.+++-------------------------------------------------------------------------------+-- | time O(segs)+-- Join a distributed segment descriptor into a global one.+-- This simply joins the distributed lengths and indices fields, but does+-- not reconstruct the original segment descriptor as it was before splitting.+-- +-- @ > pprp $ joinSegdD theGang4 +-- $ fstD $ fstD $ splitSegdOnElemsD theGang+-- $ lengthsToUSegd $ fromList [60, 10, 20, 40, 50]+-- +-- USegd lengths: [45,15,10,20,40,5,45]+-- indices: [0,45,60,70,90,130,135]+-- elements: 180+-- @+-- +-- TODO: sequential runtime is O(segs) due to application of lengthsToUSegd+-- +joinSegdD :: Gang -> Dist USegd -> USegd+joinSegdD gang+ = USegd.fromLengths+ . joinD gang unbalanced+ . mapD (What "joinSegdD/takeLengths") gang USegd.takeLengths+{-# INLINE_DIST joinSegdD #-}+++-------------------------------------------------------------------------------+-- | Glue a distributed segment descriptor back into the original global one.+-- Prop: glueSegdD gang $ splitSegdOnElems gang usegd = usegd+--+-- NOTE: This is runs sequentially and should only be used for testing purposes.+--+glueSegdD :: Gang -> Dist ((USegd, Int), Int) -> Dist USegd+glueSegdD gang bundle+ = let !usegd = fstD $ fstD $ bundle+ !lengths = takeLengthsD usegd+ + !firstSegOffsets = sndD bundle++ -- | Whether the last segment in this chunk extends into the next chunk.+ segSplits :: Dist Bool+ !segSplits+ = generateD_cheap (What "glueSegdD/segd_offsegs") gang $ \ix + -> if ix >= sizeD lengths - 1+ then False+ else indexD (here "glueSegdD") firstSegOffsets (ix + 1) /= 0++ !lengths' = fst $ carryD gang (+) 0 segSplits lengths+ !dusegd' = mapD (What "glueSegdD/fromLenghts") gang + USegd.fromLengths lengths'++ in dusegd'+{-# INLINE_DIST glueSegdD #-}+++ -- TODO: Shift this into a+ -- separate Fusion.hs module++splitSD :: Unbox a => Gang -> Dist USegd -> Vector a -> Dist (Vector a)+splitSD g dsegd xs+ = splitAsD g (takeElementsD dsegd) xs+{-# INLINE_DIST splitSD #-}++{-# 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 WZip g Seq.zip + (splitSD g d xs)+ (splitSD g d ys)++ #-}
+ Data/Array/Parallel/Unlifted/Distributed/Data/UVSegd.hs view
@@ -0,0 +1,147 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Virtual Segment Descriptors+module Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd + ( lengthD+ , takeLengthsD+ , takeIndicesD+ , takeElementsD+ , takeStartsD+ , takeSourcesD+ , takeVSegidsD+ , takeUSSegdD)+where+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Data.Array.Parallel.Unlifted.Sequential.UVSegd (UVSegd)+import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd)+import Data.Array.Parallel.Unlifted.Sequential.Vector+import Data.Array.Parallel.Pretty+import Control.Monad+import Prelude as P+import qualified Data.Array.Parallel.Unlifted.Sequential.UVSegd as UVSegd+import qualified Data.Array.Parallel.Unlifted.Distributed.Data.USSegd as DUSegd+++-------------------------------------------------------------------------------+instance DT UVSegd where+ data Dist UVSegd + = DUVSegd !(Dist (Vector Int)) -- vsegids+ !(Dist USSegd) -- distributed ussegd++ data MDist UVSegd s + = MDUVSegd !(MDist (Vector Int) s) -- vsegids+ !(MDist USSegd s) -- distributed ussegd++ indexD str (DUVSegd vsegids ussegds) i+ = UVSegd.mkUVSegd+ (indexD (str P.++ "/indexD[UVSegd]") vsegids i)+ (indexD (str P.++ "/indexD[UVSegd]") ussegds i)+ {-# INLINE_DIST indexD #-}++ newMD g+ = liftM2 MDUVSegd (newMD g) (newMD g)+ {-# INLINE_DIST newMD #-}++ readMD (MDUVSegd vsegids ussegds) i+ = liftM2 UVSegd.mkUVSegd (readMD vsegids i) (readMD ussegds i)+ {-# INLINE_DIST readMD #-}++ writeMD (MDUVSegd vsegids ussegds) i uvsegd+ = do writeMD vsegids i (UVSegd.takeVSegids uvsegd)+ writeMD ussegds i (UVSegd.takeUSSegd uvsegd)+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDUVSegd vsegids ussegds)+ = liftM2 DUVSegd (unsafeFreezeMD vsegids)+ (unsafeFreezeMD ussegds)+ {-# INLINE_DIST unsafeFreezeMD #-}++ deepSeqD uvsegd z+ = deepSeqD (UVSegd.takeVSegids uvsegd)+ $ deepSeqD (UVSegd.takeUSSegd uvsegd) z+ {-# INLINE_DIST deepSeqD #-}++ sizeD (DUVSegd _ ussegd) + = sizeD ussegd+ {-# INLINE_DIST sizeD #-}++ sizeMD (MDUVSegd _ ussegd) + = sizeMD ussegd+ {-# INLINE_DIST sizeMD #-}++ measureD uvsegd + = "UVSegd " P.++ show (UVSegd.takeVSegids uvsegd)+ P.++ " " P.++ measureD (UVSegd.takeUSSegd uvsegd)+ {-# NOINLINE measureD #-}+ -- NOINLINE because this is only used during debugging.+++-------------------------------------------------------------------------------+instance PprPhysical (Dist UVSegd) where+ pprp (DUVSegd vsegids ussegds)+ = text "DUVSegd"+ $$ (nest 7 $ vcat+ [ text "vsegids: " <+> pprp vsegids+ , text "ussegds: " <+> pprp ussegds])+ {-# NOINLINE pprp #-}+ -- NOINLINE because this is only used during debugging.+++-------------------------------------------------------------------------------+-- | O(1). Yield the overall number of segments.+lengthD :: Dist UVSegd -> Dist Int+lengthD (DUVSegd _ ussegd) + = DUSegd.lengthD ussegd+{-# INLINE_DIST lengthD #-}+++-- | O(1). Yield the lengths of the individual segments.+takeLengthsD :: Dist UVSegd -> Dist (Vector Int)+takeLengthsD (DUVSegd _ ussegd)+ = DUSegd.takeLengthsD ussegd+{-# INLINE_DIST takeLengthsD #-}+++-- | O(1). Yield the segment indices.+takeIndicesD :: Dist UVSegd -> Dist (Vector Int)+takeIndicesD (DUVSegd _ ussegd)+ = DUSegd.takeIndicesD ussegd+{-# INLINE_DIST takeIndicesD #-}+++-- | O(1). Yield the number of data elements.+takeElementsD :: Dist UVSegd -> Dist Int+takeElementsD (DUVSegd _ ussegd)+ = DUSegd.takeElementsD ussegd+{-# INLINE_DIST takeElementsD #-}+++-- | O(1). Yield the starting indices.+takeStartsD :: Dist UVSegd -> Dist (Vector Int)+takeStartsD (DUVSegd _ ussegd)+ = DUSegd.takeStartsD ussegd+{-# INLINE_DIST takeStartsD #-}+ + +-- | O(1). Yield the source ids+takeSourcesD :: Dist UVSegd -> Dist (Vector Int)+takeSourcesD (DUVSegd _ ussegd)+ = DUSegd.takeSourcesD ussegd+{-# INLINE_DIST takeSourcesD #-}+++-- | O(1). Yield the vsegids+takeVSegidsD :: Dist UVSegd -> Dist (Vector Int)+takeVSegidsD (DUVSegd vsegids _)+ = vsegids+{-# INLINE_DIST takeVSegidsD #-}+++-- | O(1). Yield the USSegd+takeUSSegdD :: Dist UVSegd -> Dist USSegd+takeUSSegdD (DUVSegd _ ussegd)+ = ussegd+{-# INLINE_DIST takeUSSegdD #-}+
+ Data/Array/Parallel/Unlifted/Distributed/Data/Unit.hs view
@@ -0,0 +1,56 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of unit values.+module Data.Array.Parallel.Unlifted.Distributed.Data.Unit + (unitD)+where+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Base++here :: String -> String+here s = "Data.Array.Parallel.Unlifted.Distributed.Types.Unit." ++ s+++instance DT () where+ data Dist () = DUnit !Int+ data MDist () s = MDUnit !Int++ indexD str (DUnit n) i+ = check (str ++ "/indexD[Unit]") n i ()+ {-# INLINE_DIST indexD #-}++ newMD+ = return . MDUnit . gangSize+ {-# INLINE_DIST newMD #-}++ readMD (MDUnit n) i+ = check (here "readMD") n i+ $ return ()+ {-# INLINE_DIST readMD #-}++ writeMD (MDUnit n) i ()+ = check (here "writeMD") n i+ $ return ()+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDUnit n)+ = return $ DUnit n+ {-# INLINE_DIST unsafeFreezeMD #-}++ sizeD + = error $ here "sizeD undefined"+ {-# NOINLINE sizeD #-}+ -- NOINLINE because this is only used for debugging.++ sizeMD + = error $ here "sizeMD undefined"+ {-# NOINLINE sizeMD #-}+ -- NOINLINE because this is only used for debugging.+++-- | Yield a distributed unit.+unitD :: Gang -> Dist ()+unitD = DUnit . gangSize+{-# INLINE_DIST unitD #-}
+ Data/Array/Parallel/Unlifted/Distributed/Data/Vector.hs view
@@ -0,0 +1,85 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Distribution of Vectors.+module Data.Array.Parallel.Unlifted.Distributed.Data.Vector+ (lengthD)+where+import Data.Array.Parallel.Unlifted.Distributed.Data.Scalar+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Pretty+import Data.Array.Parallel.Unlifted.Sequential.Vector as V+import qualified Data.Vector as BV+import qualified Data.Vector.Mutable as MBV+import qualified Data.Array.Parallel.Base as B+import Prelude as P+import Control.Monad++here :: String -> String+here s = "Data.Array.Parallel.Unlifted.Distributed.Types.Vector." P.++ s++-------------------------------------------------------------------------------+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 str (DVector _ a) i+ = B.check (here ("indexD[Vector]/" P.++ str)) (BV.length a) i $ a BV.! i+ {-# INLINE_DIST indexD #-}++ newMD g+ = liftM2 MDVector+ (newMD g) + (MBV.replicate (gangSize g) (error "MDist (Vector a) - uninitalised"))+ {-# INLINE_DIST newMD #-}++ readMD (MDVector _ marr)+ = MBV.read marr+ {-# INLINE_DIST readMD #-}++ writeMD (MDVector mlen marr) i a + = do writeMD mlen i (V.length a)+ MBV.write marr i $! a+ {-# INLINE_DIST writeMD #-}++ unsafeFreezeMD (MDVector len a)+ = liftM2 DVector (unsafeFreezeMD len)+ (BV.unsafeFreeze a)+ {-# INLINE_DIST unsafeFreezeMD #-}++ sizeD (DVector _ a) + = BV.length a+ {-# INLINE_DIST sizeD #-}++ sizeMD (MDVector _ a) + = MBV.length a+ {-# INLINE_DIST sizeMD #-}++ measureD xs+ = "Vector " P.++ show (V.length xs)+ {-# NOINLINE measureD #-}+ -- NOINLINE because this is only used for debugging.+++instance (Unbox a, Show a) => PprPhysical (Dist (V.Vector a)) where+ pprp (DVector (DInt lengths) chunks)+ = text "DVector"+ $$ (nest 8 $ vcat+ [ text "lengths:" <+> (text $ show $ V.toList lengths)+ , text "chunks: " <+> (text $ show $ BV.toList $ BV.map V.toList chunks) ])+ {-# NOINLINE pprp #-}+ -- NOINLINE because this is only used for debugging.+++-------------------------------------------------------------------------------+-- | Yield the distributed length of a distributed array.+lengthD :: Unbox a => Dist (Vector a) -> Dist Int+lengthD (DVector l _) = l+{-# INLINE_DIST lengthD #-}+
− Data/Array/Parallel/Unlifted/Distributed/DistST.hs
@@ -1,119 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# 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-myIndex = DistST return-{-# INLINE myIndex #-}----- | Lifts an 'ST' computation into the 'DistST' monad.--- The lifted computation should be data parallel.-stToDistST :: ST s a -> DistST s a-stToDistST p = DistST $ \_ -> p-{-# INLINE stToDistST #-}----- | Yields the 'Dist' element owned by the current thread.-myD :: DT a => Dist a -> DistST s a-myD dt = liftM (indexD "myD" dt) myIndex-{-# NOINLINE myD #-}----- | Yields the 'MDist' element owned by the current thread.-readMyMD :: DT a => MDist a s -> DistST s a-readMyMD mdt - = do i <- myIndex- stToDistST $ readMD mdt i-{-# NOINLINE readMyMD #-}----- | Writes the 'MDist' element owned by the current thread.-writeMyMD :: DT a => MDist a s -> a -> DistST s ()-writeMyMD mdt x - = do i <- myIndex- stToDistST $ writeMD mdt i x-{-# NOINLINE writeMyMD #-}----- | Execute a data-parallel computation on a 'Gang'.--- The same DistST comutation runs on each thread.-distST_ :: Gang -> DistST s () -> ST s ()-distST_ g = gangST g . unDistST-{-# INLINE distST_ #-}----- | Execute a data-parallel computation, yielding the distributed result.-distST :: DT a => Gang -> DistST s a -> ST s (Dist a)-distST g p - = do md <- newMD g- distST_ g $ writeMyMD md =<< p- unsafeFreezeMD md-{-# INLINE distST #-}----- | Run a data-parallel computation, yielding the distributed result.-runDistST :: DT a => Gang -> (forall s. DistST s a) -> Dist a-runDistST g p = runST (distST g p)-{-# NOINLINE runDistST #-}---runDistST_seq :: forall a. DT a => Gang -> (forall s. DistST s a) -> Dist a-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 ()-{-# NOINLINE runDistST_seq #-}---traceDistST :: String -> DistST s ()-traceDistST s = DistST $ \n -> traceGangST ("Worker " ++ show n ++ ": " ++ s)-
− Data/Array/Parallel/Unlifted/Distributed/Gang.hs
@@ -1,243 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}---- If a work request is sent to the gang while another is already running--- then just run it sequentially instead of dying.-#define SEQ_IF_GANG_BUSY 1---- Trace all work requests sent to the gang.-#define TRACE_GANG 0---- | Gang primitives.-module Data.Array.Parallel.Unlifted.Distributed.Gang - ( Gang- , seqGang- , forkGang- , gangSize- , gangIO, gangST- , traceGang, traceGangST )-where-import GHC.IO-import GHC.ST-import Control.Concurrent (forkOn)-import Control.Concurrent.MVar-import Control.Exception (assert)-import Control.Monad--#if TRACE_GANG-import Debug.Trace (traceEventIO)-import System.Time ( ClockTime(..), getClockTime )-#endif ---- 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 _ 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.-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 the program can compilain about --- "Blocked indefinitely on an MVar" because worker threads are still--- blocked on the request MVars when the program ends. Whether this finalizer--- is called or not is very racey. It can happen 1 in 10 times, or less often.--- --- 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----- | O(1). Yield 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. -gangIO :: Gang- -> (Int -> IO ())- -> IO ()--gangIO (Gang n [] _) 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)---- | Emit a GHC event for debugging.-traceGang :: String -> IO ()-traceGang s- = do t <- getGangTime- traceEventIO $ show t ++ " @ " ++ s--#else-getGangTime :: IO ()-getGangTime = return ()--diffTime :: () -> () -> String-diffTime _ _ = ""---- | Emit a GHC event for debugging.-traceGang :: String -> IO ()-traceGang _ = return ()-#endif----- | Emit a GHC event for debugging, in the `ST` monad.-traceGangST :: String -> ST s ()-traceGangST s = unsafeIOToST (traceGang s)-
+ Data/Array/Parallel/Unlifted/Distributed/Primitive.hs view
@@ -0,0 +1,47 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Standard combinators for distributed types.+module Data.Array.Parallel.Unlifted.Distributed.Primitive+ ( -- * The Gang+ Gang+ , gangSize+ , seqGang+ , forkGang+ , theGang++ -- * Distributed Types+ , DT (..)+ , newD+ , debugD+ , checkGangD++ -- * Primitive Distributed Operators.+ , generateD+ , generateD_cheap+ , imapD'+ , foldD+ , scanD)+where+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang+import Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators +import Control.Concurrent (getNumCapabilities)+import System.IO.Unsafe (unsafePerformIO)+++-- | DPH programs use this single, shared gang of threads.+-- The gang exists at top level, and is initialised at program start.+-- +-- The vectoriser guarantees that the gang is only used by a single+-- computation at a time. This is true because the program produced+-- by the vector only uses flat parallelism, so parallel computations+-- don't invoke further parallel computations. If the vectorised program+-- tries to use nested parallelism then there is a bug in the vectoriser,+-- and the code will run sequentially.+--+theGang :: Gang+theGang = unsafePerformIO (getNumCapabilities >>= forkGang)+{-# NOINLINE theGang #-}+
+ Data/Array/Parallel/Unlifted/Distributed/Primitive/DPrim.hs view
@@ -0,0 +1,86 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}++-- | Distribution of values of primitive types.+module Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim + ( DPrim (..)+ , primIndexD+ , primNewMD+ , primReadMD+ , primWriteMD+ , primUnsafeFreezeMD+ , primSizeD+ , primSizeMD)+ +where+import Data.Array.Parallel.Unlifted.Distributed.Primitive+import Data.Array.Parallel.Unlifted.Sequential.Vector+import Data.Array.Parallel.Base+import Control.Monad+import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as V+import qualified Data.Vector.Unboxed.Mutable as MV+import Prelude as P++-- 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 => String -> Dist a -> Int -> a+primIndexD str = (V.index (str P.++ "/primIndexD")) . unDPrim+{-# INLINE primIndexD #-}+++-- | Create a new distributed value, having as many members as threads+-- in the given 'Gang'.+primNewMD :: DPrim a => Gang -> ST s (MDist a s)+primNewMD = liftM mkMDPrim . MV.new . gangSize+{-# INLINE primNewMD #-}+++-- | Read the member of a distributed value corresponding to the given thread index.+primReadMD :: DPrim a => MDist a s -> Int -> ST s a+primReadMD = MV.read . unMDPrim+{-# INLINE primReadMD #-}+++-- | Write the member of a distributed value corresponding to the given thread index.+primWriteMD :: DPrim a => MDist a s -> Int -> a -> ST s ()+primWriteMD = MV.write . unMDPrim+{-# INLINE primWriteMD #-}+++-- | 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)+primUnsafeFreezeMD = liftM mkDPrim . V.unsafeFreeze . unMDPrim+{-# INLINE primUnsafeFreezeMD #-}+++-- | 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+primSizeD = V.length . unDPrim+{-# INLINE primSizeD #-}+++-- | 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+primSizeMD = MV.length . unMDPrim+{-# INLINE primSizeMD #-}+
+ Data/Array/Parallel/Unlifted/Distributed/Primitive/DT.hs view
@@ -0,0 +1,108 @@++{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+module Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+ ( -- * Distributable Types+ DT(..)+ + -- * Checking+ , checkGangD+ , checkGangMD++ -- * General Operations+ , newD+ , debugD)+where+import Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang+import Data.Array.Parallel.Base+import Data.List++here :: String -> String+here s = "Data.Array.Parallel.Unlifted.Distributed.Primitive.DT" ++ s+++-- Distributed Types ----------------------------------------------------------+infixl 9 `indexD`++-- | 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 :: String -> 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)++ -- | Ensure a distributed value is fully evaluated.+ deepSeqD :: a -> b -> b+ deepSeqD = seq+++ -- Debugging ------------------------+ -- | Number of elements in the distributed value.+ -- + -- * For debugging only, as code shouldn't be sensitive to the return value.+ sizeD :: Dist a -> Int++ -- | Number of elements in the mutable distributed value.+ -- + -- * For debugging only, as code shouldn't be sensitive to the return value.+ sizeMD :: MDist a s -> Int++ -- | Show a distributed value.+ --+ -- * For debugging only.+ measureD :: a -> String+ measureD _ = "None"+++-- Show -----------------------------------------------------------------------+-- Show instance (for debugging only) --+instance (Show a, DT a) => Show (Dist a) where+ show d = show (Prelude.map (indexD (here "show") 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 mkInit =+ runST (do+ mdt <- newMD g+ mkInit mdt+ unsafeFreezeMD mdt)++-- | Show all members of a distributed value.+debugD :: DT a => Dist a -> String+debugD d = "["+ ++ intercalate "," [measureD (indexD (here "debugD") d i) + | i <- [0 .. sizeD d-1]]+ ++ "]"
+ Data/Array/Parallel/Unlifted/Distributed/Primitive/DistST.hs view
@@ -0,0 +1,205 @@++{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# 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.+--+-- TODO: +--+module Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST + ( DistST++ -- * Primitives.+ , stToDistST+ , distST_, distST+ , runDistST, runDistST_seq+ , myIndex+ , myD+ , readMyMD, writeMyMD++ -- * Monadic combinators+ , mapDST_, mapDST, zipWithDST_, zipWithDST)+where+import qualified Data.Array.Parallel.Unlifted.Distributed.What as W+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang+import Data.Array.Parallel.Unlifted.Distributed.Data.Tuple+import Data.Array.Parallel.Base (ST, runST)+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+++-- Primitives -----------------------------------------------------------------+-- | Yields the index of the current thread within its gang.+myIndex :: DistST s Int+myIndex = DistST return+{-# INLINE myIndex #-}+++-- | Lifts an 'ST' computation into the 'DistST' monad.+-- The lifted computation should be data parallel.+stToDistST :: ST s a -> DistST s a+stToDistST p = DistST $ \_ -> p+{-# INLINE stToDistST #-}+++-- | Yields the 'Dist' element owned by the current thread.+myD :: DT a => Dist a -> DistST s a+myD dt = liftM (indexD "myD" dt) myIndex+{-# NOINLINE myD #-}+++-- | Yields the 'MDist' element owned by the current thread.+readMyMD :: DT a => MDist a s -> DistST s a+readMyMD mdt + = do i <- myIndex+ stToDistST $ readMD mdt i+{-# NOINLINE readMyMD #-}+++-- | Writes the 'MDist' element owned by the current thread.+writeMyMD :: DT a => MDist a s -> a -> DistST s ()+writeMyMD mdt x + = do i <- myIndex+ stToDistST $ writeMD mdt i x+{-# NOINLINE writeMyMD #-}++++-- Running --------------------------------------------------------------------+-- | Run a data-parallel computation, yielding the distributed result.+runDistST :: DT a => W.Comp -> Gang -> (forall s. DistST s a) -> Dist a+runDistST comp g p + = runST $ distST comp g p+{-# NOINLINE runDistST #-}+++runDistST_seq + :: forall a. DT a + => Gang -> (forall s. DistST s a) -> Dist a+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 ()+{-# NOINLINE runDistST_seq #-}+++-- | Execute a data-parallel computation, yielding the distributed result.+distST :: DT a + => W.Comp -> Gang + -> DistST s a -> ST s (Dist a)+distST comp g p + = do md <- newMD g++ distST_ comp g + $ writeMyMD md =<< p++ unsafeFreezeMD md+{-# INLINE distST #-}+++-- | Execute a data-parallel computation on a 'Gang'.+-- The same DistST comutation runs on each thread.+distST_ :: W.Comp -> Gang -> DistST s () -> ST s ()+distST_ comp gang proc+ = gangST gang + (show comp) + (workloadOfComp comp)+ $ unDistST proc+{-# INLINE distST_ #-}++workloadOfComp :: W.Comp -> Workload+workloadOfComp cc+ = case cc of+ W.CDist w -> workloadOfWhat w+ _ -> WorkUnknown++workloadOfWhat :: W.What -> Workload+workloadOfWhat ww+ = case ww of+ W.WJoinCopy elems -> WorkCopy elems + _ -> WorkUnknown++-- Combinators ----------------------------------------------------------------+-- 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, DT b) + => W.What -> Gang -> (a -> DistST s b) -> Dist a -> ST s (Dist b)+mapDST what g p !d + = mapDST' what g (\x -> x `deepSeqD` p x) d+{-# INLINE mapDST #-}+++mapDST_ :: DT a => W.What -> Gang -> (a -> DistST s ()) -> Dist a -> ST s ()+mapDST_ what g p !d + = mapDST_' what g (\x -> x `deepSeqD` p x) d+{-# INLINE mapDST_ #-}+++mapDST' :: (DT a, DT b) => W.What -> Gang -> (a -> DistST s b) -> Dist a -> ST s (Dist b)+mapDST' what g p !d + = distST (W.CDist what) g (myD d >>= p)+{-# INLINE mapDST' #-}+++mapDST_' + :: DT a + => W.What -> Gang -> (a -> DistST s ()) -> Dist a -> ST s ()+mapDST_' what g p !d + = distST_ (W.CDist what) g (myD d >>= p)+{-# INLINE mapDST_' #-}+++zipWithDST + :: (DT a, DT b, DT c)+ => W.What + -> Gang+ -> (a -> b -> DistST s c) -> Dist a -> Dist b -> ST s (Dist c)+zipWithDST what g p !dx !dy + = mapDST what g (uncurry p) (zipD dx dy)+{-# INLINE zipWithDST #-}+++zipWithDST_ + :: (DT a, DT b)+ => W.What -> Gang -> (a -> b -> DistST s ()) -> Dist a -> Dist b -> ST s ()+zipWithDST_ what g p !dx !dy + = mapDST_ what g (uncurry p) (zipD dx dy)+{-# INLINE zipWithDST_ #-}+++
+ Data/Array/Parallel/Unlifted/Distributed/Primitive/Gang.hs view
@@ -0,0 +1,296 @@+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE CPP #-}++-- If a work request is sent to the gang while another is already running+-- then just run it sequentially instead of dying.+#define SEQ_IF_GANG_BUSY 1++-- Trace all work requests sent to the gang.+#define TRACE_GANG 1++-- | Gang primitives.+module Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang+ ( Gang+ , Workload (..)+ , seqGang+ , forkGang+ , gangSize+ , gangIO, gangST)+where+import GHC.IO+import GHC.ST+import Control.Concurrent (forkOn)+import Control.Concurrent.MVar+import Control.Exception (assert)+import Control.Monad++#if TRACE_GANG+import Debug.Trace (traceEventIO)+import System.Time ( ClockTime(..), getClockTime )+#endif +++-- 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 _ 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.+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 traceWorker threadId $ "ready."+ req <- takeMVar varReq+ + case req of+ ReqDo action varDone+ -> do traceWorker threadId $ " begin."+ start <- getGangTime+ action threadId+ end <- getGangTime+ traceWorker threadId $ " end (" ++ diffTime start end ++ ")."+ + putMVar varDone ()+ gangWorker threadId varReq++ ReqShutdown varDone+ -> do traceWorker threadId $ " shutting down."+ putMVar varDone ()++traceWorker :: Int -> String -> IO ()+traceWorker threadId str+ = traceGang + $ "Worker " ++ show threadId + ++ " "+ ++ replicate (threadId * 10) ' '+ ++ str+++-- | Finaliser for worker threads.+-- We want to shutdown the corresponding thread when it's MVar becomes+-- unreachable. Without this the program can compilain about +-- "Blocked indefinitely on an MVar" because worker threads are still+-- blocked on the request MVars when the program ends. Whether this finalizer+-- is called or not is very racey. It can happen 1 in 10 times, or less often.+-- +-- 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+++-- | O(1). Yield the number of threads in the 'Gang'.+gangSize :: Gang -> Int+gangSize (Gang n _ _) = n+++-------------------------------------------------------------------------------+data Workload+ -- | Unknown workload. Just run it in parallel.+ = WorkUnknown++ -- | Memory bound copy-like workload, + -- of the given number of bytes.+ | WorkCopy Int+ deriving (Eq, Show)+++-- | Decide whether a workload is too small to bother running in parallel.+-- TODO: We want to determine this based on similar workloads that +-- we have run before. The gang should know what its minumum latency is.+workloadIsSmall :: Workload -> Bool+workloadIsSmall ww+ = case ww of+ WorkUnknown -> False+ WorkCopy bytes -> bytes < 1000++++-- | 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. +gangIO :: Gang+ -> String + -> Workload+ -> (Int -> IO ())+ -> IO ()++-- Hrm. Gang hasn't been created yet. +-- Just run the requests in the main thread.+gangIO (Gang n [] _) _what _workload p + = mapM_ p [0 .. n-1]++#if SEQ_IF_GANG_BUSY+gangIO (Gang n mvs busy) what workload p + = do let !small = workloadIsSmall workload+ if small + then do+ traceGang $ "Issuing small " ++ what+ mapM_ p [0 .. n-1]+ else do+ isBusy <- swapMVar busy True+ if isBusy + then do + traceGang $ "WARNING: Gang was already busy, running sequentially: " ++ what+ mapM_ p [0 .. n-1]+ else do+ traceGangSplit $ "Issuing par " ++ what+ parIO what n mvs p+ _ <- swapMVar busy False+ return ()+#else+gangIO (Gang n mvs busy) what _workload p + = parIO n mvs p+#endif+++-- | Issue some requests to the worker threads and wait for them to complete.+parIO :: String+ -> 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 what n mvs p + = do start <- getGangTime+ reqs <- sequence . replicate n $ newReq p++ zipWithM_ putMVar mvs reqs++ traceGang $ "Running."+ mapM_ waitReq reqs+ end <- getGangTime++ traceGangSplit $ "Complete par " ++ what ++ " in " ++ diffTime start end ++ "us."+++-- | Same as 'gangIO' but in the 'ST' monad.+gangST :: Gang -> String -> Workload -> (Int -> ST s ()) -> ST s ()+gangST gang what workload p + = unsafeIOToST + $ gangIO gang what workload+ $ unsafeSTToIO . p+++-- Tracing -------------------------------------------------------------------+#if TRACE_GANG+getGangTime :: IO Integer+getGangTime+ = do TOD sec pico <- getClockTime+ let !micro = pico `div` 1000000+ return (micro + sec * 1000000)++diffTime :: Integer -> Integer -> String+diffTime x y = show (y-x)++-- | Emit a GHC event for debugging, but don't mind if it gets truncated+traceGang :: String -> IO ()+traceGang s+ = do traceEventIO $ "GANG " ++ s++-- | Emit a GHC event for debugging. Split across multiple events if necessary.+traceGangSplit :: String -> IO ()+traceGangSplit s+ = do let xs = chunks 500 s+ let max' = show $ length xs+ mapM_ (\(x,i) -> traceEventIO + $ "GANG[" ++ show i ++ "/" ++ max' ++ "] " ++ x) + (xs `zip` [1 :: Int ..])+ where+ chunks _ [] = []+ chunks len str+ = let (f,r) = splitAt len str+ in f : chunks len r++#else+getGangTime :: IO ()+getGangTime = return ()++diffTime :: () -> () -> String+diffTime _ _ = ""++-- | Emit a GHC event for debugging.+traceGang :: String -> IO ()+traceGang _ = return ()++traceGangSplit :: String -> IO ()+traceGangSplit _ = return ()+#endif
+ Data/Array/Parallel/Unlifted/Distributed/Primitive/Operators.hs view
@@ -0,0 +1,127 @@++-- Primitive Gang Operators are fundamental computations that run on the gang.+-- At runtime we can record how long each one runs using GHC events.+-- As Gang Operators are not inlined, fusion between them is done via GHC rewrite rules.+--+{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}+{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE CPP #-}+#include "fusion-phases.h"++-- | Standard combinators for distributed types.+module Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators + ( generateD+ , generateD_cheap+ , imapD'+ , foldD+ , scanD)+where+import Data.Array.Parallel.Base ( ST, runST)+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST+import Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+import Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang+import qualified Data.Array.Parallel.Unlifted.Distributed.What as W+import Debug.Trace++here s = "Data.Array.Parallel.Unlifted.Distributed.Combinators." ++ s+++-- | Create a distributed value, given a function to create the instance+-- for each thread.+generateD + :: DT a + => W.What -- ^ What is the worker function doing.+ -> Gang + -> (Int -> a) + -> Dist a++generateD what gang f + = runDistST (W.CGen False what) + gang + (myIndex >>= return . f)+{-# NOINLINE generateD #-}+++-- | Create a distributed value, but do it sequentially.+-- +-- This function is used when we want to operate on a distributed value, but+-- there isn't much data involved. For example, if we want to distribute +-- a single integer to each thread, then there's no need to fire up the +-- gang for this.+-- +generateD_cheap + :: DT a + => W.What -- ^ What is the worker function doing.+ -> Gang + -> (Int -> a) + -> Dist a++generateD_cheap what g f + = traceEvent (show $ W.CGen True what) + $ runDistST_seq g (myIndex >>= return . f)+{-# NOINLINE generateD_cheap #-}+++-- Mapping --------------------------------------------------------------------+-- | Map a function across all elements of a distributed value.+-- The worker function also gets the current thread index.+imapD' :: (DT a, DT b) + => W.What -> Gang -> (Int -> a -> b) -> Dist a -> Dist b+imapD' what gang f !d + = runDistST (W.CMap what) gang + $ do i <- myIndex+ x <- myD d+ let result = f i x+ deepSeqD result (return ())+ return result+{-# NOINLINE imapD' #-}+++-- Folding --------------------------------------------------------------------+-- | Fold all the instances of a distributed value.+foldD :: DT a => W.What -> Gang -> (a -> a -> a) -> Dist a -> a+foldD what gang f !d + = traceEvent (show (W.CFold what))+ $ checkGangD ("here foldD") gang d + $ fold 1 (indexD (here "foldD") d 0)+ where+ !n = gangSize gang++ fold i x + | i == n = x+ | otherwise = fold (i+1) (f x $ indexD (here "foldD") d i)+{-# NOINLINE foldD #-}+++-- Scanning -------------------------------------------------------------------+-- | Prefix sum of the instances of a distributed value.+scanD :: forall a. DT a => W.What -> Gang -> (a -> a -> a) -> a -> Dist a -> (Dist a, a)+scanD what gang f z !d+ = traceEvent (show (W.CScan what))+ $ checkGangD (here "scanD") gang d + $ runST (do+ md <- newMD gang+ s <- scan md 0 z+ d' <- unsafeFreezeMD md+ return (d',s))+ where+ !n = gangSize gang+ + 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 $ indexD (here "scanD") d i)+{-# NOINLINE scanD #-}+++-- | Emit a GHC event for debugging, in the `ST` monad.+{-+traceGangST :: String -> ST s ()+traceGangST s = unsafeIOToST (traceGang s)++traceDistIO :: String -> IO ()+traceDistIO s+ = do traceEventIO $ "DIST " ++ s+-}
− Data/Array/Parallel/Unlifted/Distributed/Scalars.hs
@@ -1,37 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}---- | Operations on 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 - ( scalarD- , orD, andD- , sumD)-where-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Unlifted.Distributed.Types-import Data.Array.Parallel.Unlifted.Distributed.Combinators----- | Distribute a scalar.--- Each thread gets its own copy of the same value.--- Example: scalarD theGangN4 10 = [10, 10, 10, 10] -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
@@ -1,22 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-module Data.Array.Parallel.Unlifted.Distributed.TheGang - (theGang)-where-import Data.Array.Parallel.Unlifted.Distributed.Gang -import Control.Concurrent (getNumCapabilities)-import System.IO.Unsafe (unsafePerformIO)---- | DPH programs use this single, shared gang of threads.--- The gang exists at top level, and is initialised at program start.--- --- The vectoriser guarantees that the gang is only used by a single--- computation at a time. This is true because the program produced--- by the vector only uses flat parallelism, so parallel computations--- don't invoke further parallel computations. If the vectorised program--- tries to use nested parallelism then there is a bug in the vectoriser,--- and the code will run sequentially.----theGang :: Gang-theGang = unsafePerformIO (getNumCapabilities >>= forkGang)-{-# NOINLINE theGang #-}-
− Data/Array/Parallel/Unlifted/Distributed/Types.hs
@@ -1,21 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Distributed types.-module Data.Array.Parallel.Unlifted.Distributed.Types - ( module Data.Array.Parallel.Unlifted.Distributed.Types.Vector- , module Data.Array.Parallel.Unlifted.Distributed.Types.Maybe- , module Data.Array.Parallel.Unlifted.Distributed.Types.Tuple- , module Data.Array.Parallel.Unlifted.Distributed.Types.Prim- , module Data.Array.Parallel.Unlifted.Distributed.Types.Unit- , module Data.Array.Parallel.Unlifted.Distributed.Types.Base)-where-import Data.Array.Parallel.Unlifted.Distributed.Types.Vector-import Data.Array.Parallel.Unlifted.Distributed.Types.Maybe-import Data.Array.Parallel.Unlifted.Distributed.Types.Tuple-import Data.Array.Parallel.Unlifted.Distributed.Types.Prim-import Data.Array.Parallel.Unlifted.Distributed.Types.Unit-import Data.Array.Parallel.Unlifted.Distributed.Types.Base--
− Data/Array/Parallel/Unlifted/Distributed/Types/Base.hs
@@ -1,108 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-module Data.Array.Parallel.Unlifted.Distributed.Types.Base - ( -- * Distributable Types- DT(..)- - -- * Checking- , checkGangD- , checkGangMD-- -- * General Operations- , newD- , debugD)-where-import Data.Array.Parallel.Unlifted.Distributed.Gang (Gang, gangSize)-import Data.Array.Parallel.Base-import Data.List (intercalate)--here :: String -> String-here s = "Data.Array.Parallel.Unlifted.Distributed.Types.Base." ++ s----- Distributed Types -----------------------------------------------------------infixl 9 `indexD`---- | 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 :: String -> 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)-- -- | Ensure a distributed value is fully evaluated.- deepSeqD :: a -> b -> b- deepSeqD = seq--- -- Debugging ------------------------- -- | Number of elements in the distributed value.- -- - -- * For debugging only, as code shouldn't be sensitive to the return value.- sizeD :: Dist a -> Int-- -- | Number of elements in the mutable distributed value.- -- - -- * For debugging only, as code shouldn't be sensitive to the return value.- sizeMD :: MDist a s -> Int-- -- | Show a distributed value.- --- -- * For debugging only.- measureD :: a -> String- measureD _ = "None"----- Show -------------------------------------------------------------------------- Show instance (for debugging only) ---instance (Show a, DT a) => Show (Dist a) where- show d = show (Prelude.map (indexD (here "show") 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 mkInit =- runST (do- mdt <- newMD g- mkInit mdt- unsafeFreezeMD mdt)---- | Show all members of a distributed value.-debugD :: DT a => Dist a -> String-debugD d = "["- ++ intercalate "," [measureD (indexD (here "debugD") d i) - | i <- [0 .. sizeD d-1]]- ++ "]"-
− Data/Array/Parallel/Unlifted/Distributed/Types/Maybe.hs
@@ -1,45 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}---- | Distribution of Maybes.-module Data.Array.Parallel.Unlifted.Distributed.Types.Maybe where-import Data.Array.Parallel.Unlifted.Distributed.Types.Prim ()-import Data.Array.Parallel.Unlifted.Distributed.Types.Base-import Control.Monad--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 str (DMaybe bs as) i- | indexD (str ++ "/indexD[Maybe]") bs i- = Just $ indexD (str ++ "/indexD[Maybe]" ++ str) as 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 _) i Nothing - = writeMD bs i False-- writeMD (MDMaybe bs as) i (Just x)- = do 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 ++ ")"
− Data/Array/Parallel/Unlifted/Distributed/Types/Prim.hs
@@ -1,250 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}---- | Distribution of values of primitive types.-module Data.Array.Parallel.Unlifted.Distributed.Types.Prim - ( DPrim (..)- , DT (..)- , Dist (..))-where-import Data.Array.Parallel.Unlifted.Distributed.Types.Base-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Unlifted.Sequential.Vector-import Data.Array.Parallel.Base-import Data.Array.Parallel.Pretty-import Data.Word-import Control.Monad-import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as V-import qualified Data.Vector.Unboxed.Mutable as MV-import Prelude as P---- 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 => String -> Dist a -> Int -> a-primIndexD str = (V.index (str P.++ "/primIndexD")) . unDPrim-{-# INLINE primIndexD #-}----- | Create a new distributed value, having as many members as threads--- in the given 'Gang'.-primNewMD :: DPrim a => Gang -> ST s (MDist a s)-primNewMD = liftM mkMDPrim . MV.new . gangSize-{-# INLINE primNewMD #-}----- | Read the member of a distributed value corresponding to the given thread index.-primReadMD :: DPrim a => MDist a s -> Int -> ST s a-primReadMD = MV.read . unMDPrim-{-# INLINE primReadMD #-}----- | Write the member of a distributed value corresponding to the given thread index.-primWriteMD :: DPrim a => MDist a s -> Int -> a -> ST s ()-primWriteMD = MV.write . unMDPrim-{-# INLINE primWriteMD #-}----- | 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)-primUnsafeFreezeMD = liftM mkDPrim . V.unsafeFreeze . unMDPrim-{-# INLINE primUnsafeFreezeMD #-}----- | 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-primSizeD = V.length . unDPrim-{-# INLINE primSizeD #-}----- | 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-primSizeMD = MV.length . unMDPrim-{-# INLINE primSizeMD #-}----- 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----- Ordering ------------------------------------------------------------------------instance DPrim Ordering where- mkDPrim = DOrdering- unDPrim (DOrdering a) = a-- mkMDPrim = MDOrdering- unMDPrim (MDOrdering a) = a---instance DT Ordering where- data Dist Ordering = DOrdering !(V.Vector Ordering)- data MDist Ordering s = MDOrdering !(MV.STVector s Ordering)-- indexD = primIndexD- newMD = primNewMD- readMD = primReadMD- writeMD = primWriteMD- unsafeFreezeMD = primUnsafeFreezeMD- sizeD = primSizeD- sizeMD = primSizeMD----- Integer -------------------------------------------------------------------------- FIXME: fake instances-instance DPrim Integer-instance DT Integer----- 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 " P.++ show n--instance PprPhysical (Dist Int) where- pprp (DInt xs)- = text "DInt" <+> text (show $ V.toList xs)----- 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--
− Data/Array/Parallel/Unlifted/Distributed/Types/Tuple.hs
@@ -1,145 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Distribution of Tuples-module Data.Array.Parallel.Unlifted.Distributed.Types.Tuple - ( -- * Pairs- zipD, unzipD, fstD, sndD- - -- * Triples- , zip3D, unzip3D)-where-import Data.Array.Parallel.Unlifted.Distributed.Types.Base-import Data.Array.Parallel.Base-import Data.Array.Parallel.Pretty-import Control.Monad--here :: String -> String-here s = "Data.Array.Parallel.Unlifted.Distributed.Types.Tuple." ++ s----- 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 str d i- = ( indexD (str ++ "/indexD[Tuple2]") (fstD d) i- , indexD (str ++ "/indexD[Tuple2]") (sndD d) 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)- = do 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 ++ ")"---instance (PprPhysical (Dist a), PprPhysical (Dist b)) - => PprPhysical (Dist (a, b)) where- pprp (DProd xs ys)- = text "DProd"- $$ (nest 8 $ vcat- [ pprp xs- , pprp ys ])----- | 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)-zipD !x !y - = checkEq (here "zipDT") "Size mismatch" (sizeD x) (sizeD y) - $ DProd x y-{-# INLINE [0] zipD #-}----- | Unpairing of distributed values.-unzipD :: (DT a, DT b) => Dist (a,b) -> (Dist a, Dist b)-unzipD (DProd dx dy) = (dx,dy)-{-# INLINE_DIST unzipD #-}----- | Extract the first elements of a distributed pair.-fstD :: (DT a, DT b) => Dist (a,b) -> Dist a-fstD = fst . unzipD-{-# INLINE_DIST fstD #-}----- | Extract the second elements of a distributed pair.-sndD :: (DT a, DT b) => Dist (a,b) -> Dist b-sndD = snd . unzipD-{-# INLINE_DIST sndD #-}----- Triples ---------------------------------------------------------------------instance (DT a, DT b, DT c) => DT (a,b,c) where- data Dist (a,b,c) = DProd3 !(Dist a) !(Dist b) !(Dist c)- data MDist (a,b,c) s = MDProd3 !(MDist a s) !(MDist b s) !(MDist c s)-- indexD str (DProd3 xs ys zs) i- = ( indexD (here $ "indexD[Tuple3]/" ++ str) xs i- , indexD (here $ "indexD[Tuple3]/" ++ str) ys i- , indexD (here $ "indexD[Tuple3]/" ++ str) zs i)-- newMD g- = liftM3 MDProd3 (newMD g) (newMD g) (newMD g)-- readMD (MDProd3 xs ys zs) i- = liftM3 (,,) (readMD xs i) (readMD ys i) (readMD zs i)-- writeMD (MDProd3 xs ys zs) i (x,y,z)- = do writeMD xs i x- writeMD ys i y- writeMD zs i z-- unsafeFreezeMD (MDProd3 xs ys zs)- = liftM3 DProd3 (unsafeFreezeMD xs)- (unsafeFreezeMD ys)- (unsafeFreezeMD zs)-- {-# INLINE deepSeqD #-}- deepSeqD (x,y,z) k - = deepSeqD x (deepSeqD y (deepSeqD z k))-- sizeD (DProd3 x _ _) = sizeD x- sizeMD (MDProd3 x _ _) = sizeMD x-- measureD (x,y,z)- = "Triple " - ++ "(" ++ measureD x ++ ") "- ++ "(" ++ measureD y ++ ") "- ++ "(" ++ measureD z ++ ")"----- | Pairing of distributed values.--- /The two values must belong to the same/ 'Gang'.-zip3D :: (DT a, DT b, DT c) => Dist a -> Dist b -> Dist c -> Dist (a,b,c)-zip3D !x !y !z- = checkEq (here "zip3DT") "Size mismatch" (sizeD x) (sizeD y) - $ checkEq (here "zip3DT") "Size mismatch" (sizeD x) (sizeD z) - $ DProd3 x y z-{-# INLINE [0] zip3D #-}----- | Unpairing of distributed values.-unzip3D :: (DT a, DT b, DT c) => Dist (a,b,c) -> (Dist a, Dist b, Dist c)-unzip3D (DProd3 dx dy dz) = (dx,dy,dz)-{-# INLINE_DIST unzip3D #-}
− Data/Array/Parallel/Unlifted/Distributed/Types/USSegd.hs
@@ -1,130 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Distribution of Segment Descriptors-module Data.Array.Parallel.Unlifted.Distributed.Types.USSegd - ( lengthD- , takeLengthsD- , takeIndicesD- , takeElementsD- , takeStartsD- , takeSourcesD- , takeUSegdD)-where-import Data.Array.Parallel.Unlifted.Distributed.Types.Base-import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd)-import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)-import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)-import Data.Array.Parallel.Pretty-import Control.Monad-import Prelude as P-import qualified Data.Array.Parallel.Unlifted.Distributed.Types.USegd as DUSegd-import qualified Data.Array.Parallel.Unlifted.Distributed.Types.Vector as DV-import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd---instance DT USSegd where- data Dist USSegd - = DUSSegd !(Dist (Vector Int)) -- segment starts- !(Dist (Vector Int)) -- segment sources- !(Dist USegd) -- distributed usegd-- data MDist USSegd s - = MDUSSegd !(MDist (Vector Int) s) -- segment starts- !(MDist (Vector Int) s) -- segment sources- !(MDist USegd s) -- distributed usegd-- indexD str (DUSSegd starts sources usegds) i- = USSegd.mkUSSegd- (indexD (str ++ "/indexD[USSegd]") starts i)- (indexD (str ++ "/indexD[USSegd]") sources i)- (indexD (str ++ "/indexD[USSegd]") usegds i)-- newMD g- = liftM3 MDUSSegd (newMD g) (newMD g) (newMD g)-- readMD (MDUSSegd starts sources usegds) i- = liftM3 USSegd.mkUSSegd (readMD starts i) (readMD sources i) (readMD usegds i)-- writeMD (MDUSSegd starts sources usegds) i ussegd- = do writeMD starts i (USSegd.takeStarts ussegd)- writeMD sources i (USSegd.takeSources ussegd)- writeMD usegds i (USSegd.takeUSegd ussegd)-- unsafeFreezeMD (MDUSSegd starts sources usegds)- = liftM3 DUSSegd (unsafeFreezeMD starts)- (unsafeFreezeMD sources)- (unsafeFreezeMD usegds)-- deepSeqD ussegd z- = deepSeqD (USSegd.takeStarts ussegd)- $ deepSeqD (USSegd.takeSources ussegd)- $ deepSeqD (USSegd.takeUSegd ussegd) z-- sizeD (DUSSegd _ _ usegd) = sizeD usegd- sizeMD (MDUSSegd _ _ usegd) = sizeMD usegd-- measureD ussegd - = "USSegd " P.++ show (USSegd.takeStarts ussegd)- P.++ " " P.++ show (USSegd.takeSources ussegd)- P.++ " " P.++ measureD (USSegd.takeUSegd ussegd)---instance PprPhysical (Dist USSegd) where- pprp (DUSSegd starts sources usegds)- = text "DUSSegd"- $$ (nest 7 $ vcat- [ text "starts: " <+> pprp starts- , text "sources: " <+> pprp sources- , text "usegds: " <+> pprp usegds])----- | O(1). Yield the overall number of segments.-lengthD :: Dist USSegd -> Dist Int-lengthD (DUSSegd starts _ _) - = DV.lengthD starts-{-# INLINE_DIST lengthD #-}----- | O(1). Yield the lengths of the individual segments.-takeLengthsD :: Dist USSegd -> Dist (Vector Int)-takeLengthsD (DUSSegd _ _ usegds)- = DUSegd.takeLengthsD usegds-{-# INLINE_DIST takeLengthsD #-}----- | O(1). Yield the segment indices.-takeIndicesD :: Dist USSegd -> Dist (Vector Int)-takeIndicesD (DUSSegd _ _ usegds)- = DUSegd.takeIndicesD usegds-{-# INLINE_DIST takeIndicesD #-}----- | O(1). Yield the number of data elements.-takeElementsD :: Dist USSegd -> Dist Int-takeElementsD (DUSSegd _ _ usegds)- = DUSegd.takeElementsD usegds-{-# INLINE_DIST takeElementsD #-}----- | O(1). Yield the starting indices.-takeStartsD :: Dist USSegd -> Dist (Vector Int)-takeStartsD (DUSSegd starts _ _)- = starts-{-# INLINE_DIST takeStartsD #-}- ---- | O(1). Yield the source ids-takeSourcesD :: Dist USSegd -> Dist (Vector Int)-takeSourcesD (DUSSegd _ sources _)- = sources-{-# INLINE_DIST takeSourcesD #-}----- | O(1). Yield the USegd-takeUSegdD :: Dist USSegd -> Dist USegd-takeUSegdD (DUSSegd _ _ usegd)- = usegd-{-# INLINE_DIST takeUSegdD #-}-
− Data/Array/Parallel/Unlifted/Distributed/Types/USegd.hs
@@ -1,113 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Distribution of Segment Descriptors-module Data.Array.Parallel.Unlifted.Distributed.Types.USegd - ( mkDUSegd- , lengthD- , takeLengthsD- , takeIndicesD- , takeElementsD)-where-import Data.Array.Parallel.Unlifted.Distributed.Types.Base-import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)-import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)-import Data.Array.Parallel.Pretty-import Control.Monad-import qualified Data.Array.Parallel.Unlifted.Distributed.Types.Vector as DV-import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd-import Prelude as P---instance DT USegd where- data Dist USegd - = DUSegd !(Dist (Vector Int)) -- segment lengths- !(Dist (Vector Int)) -- segment indices- !(Dist Int) -- number of elements in this chunk-- data MDist USegd s - = MDUSegd !(MDist (Vector Int) s) -- segment lengths- !(MDist (Vector Int) s) -- segment indices- !(MDist Int s) -- number of elements in this chunk-- indexD str (DUSegd lens idxs eles) i- = USegd.mkUSegd- (indexD (str ++ "/indexD[USegd]") lens i)- (indexD (str ++ "/indexD[USegd]") idxs i)- (indexD (str ++ "/indexD[USegd]") eles i)-- newMD g- = liftM3 MDUSegd (newMD g) (newMD g) (newMD g)-- readMD (MDUSegd lens idxs eles) i- = liftM3 USegd.mkUSegd (readMD lens i) (readMD idxs i) (readMD eles i)-- writeMD (MDUSegd lens idxs eles) i segd- = do writeMD lens i (USegd.takeLengths segd)- writeMD idxs i (USegd.takeIndices segd)- writeMD eles i (USegd.takeElements segd)-- unsafeFreezeMD (MDUSegd lens idxs eles)- = liftM3 DUSegd (unsafeFreezeMD lens)- (unsafeFreezeMD idxs)- (unsafeFreezeMD eles)-- deepSeqD segd z- = deepSeqD (USegd.takeLengths segd)- $ deepSeqD (USegd.takeIndices segd)- $ deepSeqD (USegd.takeElements segd) z-- sizeD (DUSegd _ _ eles) = sizeD eles- sizeMD (MDUSegd _ _ eles) = sizeMD eles-- measureD segd - = "Segd " P.++ show (USegd.length segd)- P.++ " " P.++ show (USegd.takeElements segd)---instance PprPhysical (Dist USegd) where- pprp (DUSegd lens indices elements)- = text "DUSegd"- $$ (nest 7 $ vcat- [ text "lengths: " <+> pprp lens- , text "indices: " <+> pprp indices- , text "elements:" <+> pprp elements])----- | O(1). Construct a distributed segment descriptor-mkDUSegd - :: Dist (Vector Int) -- ^ segment lengths- -> Dist (Vector Int) -- ^ segment indices- -> Dist Int -- ^ number of elements in each chunk- -> Dist USegd--mkDUSegd = DUSegd----- | O(1). Yield the overall number of segments.-lengthD :: Dist USegd -> Dist Int-lengthD (DUSegd lens _ _) - = DV.lengthD lens-{-# INLINE_DIST lengthD #-}----- | O(1). Yield the lengths of the individual segments.-takeLengthsD :: Dist USegd -> Dist (Vector Int)-takeLengthsD (DUSegd lens _ _ )- = lens-{-# INLINE_DIST takeLengthsD #-}----- | O(1). Yield the segment indices of a segment descriptor.-takeIndicesD :: Dist USegd -> Dist (Vector Int)-takeIndicesD (DUSegd _ idxs _)- = idxs-{-# INLINE_DIST takeIndicesD #-}----- | O(1). Yield the number of data elements.-takeElementsD :: Dist USegd -> Dist Int-takeElementsD (DUSegd _ _ dns)- = dns-{-# INLINE_DIST takeElementsD #-}
− Data/Array/Parallel/Unlifted/Distributed/Types/UVSegd.hs
@@ -1,128 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Distribution of Virtual Segment Descriptors-module Data.Array.Parallel.Unlifted.Distributed.Types.UVSegd - ( lengthD- , takeLengthsD- , takeIndicesD- , takeElementsD- , takeStartsD- , takeSourcesD- , takeVSegidsD- , takeUSSegdD)-where-import Data.Array.Parallel.Unlifted.Distributed.Types.Base-import Data.Array.Parallel.Unlifted.Sequential.UVSegd (UVSegd)-import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd)-import Data.Array.Parallel.Unlifted.Sequential.Vector-import Data.Array.Parallel.Pretty-import Control.Monad-import Prelude as P-import qualified Data.Array.Parallel.Unlifted.Sequential.UVSegd as UVSegd-import qualified Data.Array.Parallel.Unlifted.Distributed.Types.USSegd as DUSegd---instance DT UVSegd where- data Dist UVSegd - = DUVSegd !(Dist (Vector Int)) -- vsegids- !(Dist USSegd) -- distributed ussegd-- data MDist UVSegd s - = MDUVSegd !(MDist (Vector Int) s) -- vsegids- !(MDist USSegd s) -- distributed ussegd-- indexD str (DUVSegd vsegids ussegds) i- = UVSegd.mkUVSegd- (indexD (str P.++ "/indexD[UVSegd]") vsegids i)- (indexD (str P.++ "/indexD[UVSegd]") ussegds i)-- newMD g- = liftM2 MDUVSegd (newMD g) (newMD g)-- readMD (MDUVSegd vsegids ussegds) i- = liftM2 UVSegd.mkUVSegd (readMD vsegids i) (readMD ussegds i)-- writeMD (MDUVSegd vsegids ussegds) i uvsegd- = do writeMD vsegids i (UVSegd.takeVSegids uvsegd)- writeMD ussegds i (UVSegd.takeUSSegd uvsegd)-- unsafeFreezeMD (MDUVSegd vsegids ussegds)- = liftM2 DUVSegd (unsafeFreezeMD vsegids)- (unsafeFreezeMD ussegds)-- deepSeqD uvsegd z- = deepSeqD (UVSegd.takeVSegids uvsegd)- $ deepSeqD (UVSegd.takeUSSegd uvsegd) z-- sizeD (DUVSegd _ ussegd) = sizeD ussegd- sizeMD (MDUVSegd _ ussegd) = sizeMD ussegd-- measureD uvsegd - = "UVSegd " P.++ show (UVSegd.takeVSegids uvsegd)- P.++ " " P.++ measureD (UVSegd.takeUSSegd uvsegd)---instance PprPhysical (Dist UVSegd) where- pprp (DUVSegd vsegids ussegds)- = text "DUVSegd"- $$ (nest 7 $ vcat- [ text "vsegids: " <+> pprp vsegids- , text "ussegds: " <+> pprp ussegds])----- | O(1). Yield the overall number of segments.-lengthD :: Dist UVSegd -> Dist Int-lengthD (DUVSegd _ ussegd) - = DUSegd.lengthD ussegd-{-# INLINE_DIST lengthD #-}----- | O(1). Yield the lengths of the individual segments.-takeLengthsD :: Dist UVSegd -> Dist (Vector Int)-takeLengthsD (DUVSegd _ ussegd)- = DUSegd.takeLengthsD ussegd-{-# INLINE_DIST takeLengthsD #-}----- | O(1). Yield the segment indices.-takeIndicesD :: Dist UVSegd -> Dist (Vector Int)-takeIndicesD (DUVSegd _ ussegd)- = DUSegd.takeIndicesD ussegd-{-# INLINE_DIST takeIndicesD #-}----- | O(1). Yield the number of data elements.-takeElementsD :: Dist UVSegd -> Dist Int-takeElementsD (DUVSegd _ ussegd)- = DUSegd.takeElementsD ussegd-{-# INLINE_DIST takeElementsD #-}----- | O(1). Yield the starting indices.-takeStartsD :: Dist UVSegd -> Dist (Vector Int)-takeStartsD (DUVSegd _ ussegd)- = DUSegd.takeStartsD ussegd-{-# INLINE_DIST takeStartsD #-}- - --- | O(1). Yield the source ids-takeSourcesD :: Dist UVSegd -> Dist (Vector Int)-takeSourcesD (DUVSegd _ ussegd)- = DUSegd.takeSourcesD ussegd-{-# INLINE_DIST takeSourcesD #-}----- | O(1). Yield the vsegids-takeVSegidsD :: Dist UVSegd -> Dist (Vector Int)-takeVSegidsD (DUVSegd vsegids _)- = vsegids-{-# INLINE_DIST takeVSegidsD #-}----- | O(1). Yield the USSegd-takeUSSegdD :: Dist UVSegd -> Dist USSegd-takeUSSegdD (DUVSegd _ ussegd)- = ussegd-{-# INLINE_DIST takeUSSegdD #-}
− Data/Array/Parallel/Unlifted/Distributed/Types/Unit.hs
@@ -1,45 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Distribution of unit values.-module Data.Array.Parallel.Unlifted.Distributed.Types.Unit - (unitD)-where-import Data.Array.Parallel.Unlifted.Distributed.Types.Base-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Base--here :: String -> String-here s = "Data.Array.Parallel.Unlifted.Distributed.Types.Unit." ++ s--instance DT () where- data Dist () = DUnit !Int- data MDist () s = MDUnit !Int-- indexD str (DUnit n) i- = check (str ++ "/indexD[Unit]") 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 $ here "sizeD undefined"- sizeMD = error $ here "sizeMD undefined"----- | Yield a distributed unit.-unitD :: Gang -> Dist ()-unitD = DUnit . gangSize-{-# INLINE_DIST unitD #-}
− Data/Array/Parallel/Unlifted/Distributed/Types/Vector.hs
@@ -1,61 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}---- | Distribution of Vectors.-module Data.Array.Parallel.Unlifted.Distributed.Types.Vector- (lengthD)-where-import qualified Data.Array.Parallel.Base as B-import Data.Array.Parallel.Unlifted.Distributed.Types.Prim-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Pretty-import Data.Array.Parallel.Unlifted.Sequential.Vector as V-import qualified Data.Vector as BV-import qualified Data.Vector.Mutable as MBV-import Prelude as P-import Control.Monad---here :: String -> String-here s = "Data.Array.Parallel.Unlifted.Distributed.Types.Vector." P.++ s---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 str (DVector _ a) i- = B.check (here ("indexD[Vector]/" P.++ str)) (BV.length 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 " P.++ show (V.length xs)---instance (Unbox a, Show a) => PprPhysical (Dist (V.Vector a)) where- pprp (DVector (DInt lengths) chunks)- = text "DVector"- $$ (nest 8 $ vcat- [ text "lengths:" <+> (text $ show $ V.toList lengths)- , text "chunks: " <+> (text $ show $ BV.toList $ BV.map V.toList chunks) ])----- | Yield the distributed length of a distributed array.-lengthD :: Unbox a => Dist (Vector a) -> Dist Int-lengthD (DVector l _) = l
− Data/Array/Parallel/Unlifted/Distributed/USSegd.hs
@@ -1,245 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Operations on Distributed Segment Descriptors-module Data.Array.Parallel.Unlifted.Distributed.USSegd - (splitSSegdOnElemsD)-where-import Data.Array.Parallel.Unlifted.Distributed.Arrays-import Data.Array.Parallel.Unlifted.Distributed.Combinators-import Data.Array.Parallel.Unlifted.Distributed.Types-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd)-import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector)-import Data.Array.Parallel.Base-import Data.Bits (shiftR)-import Control.Monad (when)-import Data.Array.Parallel.Unlifted.Distributed.Types.USSegd ()-import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd-import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd-import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as Seq--here :: String -> String-here s = "Data.Array.Parallel.Unlifted.Distributed.USSegd." ++ s------------------------------------------------------------------------------------ | Split a segment descriptor across the gang, element wise.--- We try to put the same number of elements on each thread, which means--- that segments are sometimes split across threads.------ Each thread gets a slice of segment descriptor, the segid of the first --- slice, and the offset of the first slice in its segment.--- --- Example:--- In this picture each X represents 5 elements, and we have 5 segements in total.------ @ segs: ----------------------- --- ------- --------------- ---------------------- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|--- | thread1 | thread2 | thread3 | thread4 |--- segid: 0 0 3 4--- offset: 0 45 0 5------ pprp $ splitSegdOnElemsD theGang --- $ lengthsToUSegd $ fromList [60, 10, 20, 40, 50 :: Int]------ segd: DUSegd lengths: DVector lengths: [1,3,2,1]--- chunks: [[45],[15,10,20],[40,5],[45]]--- indices: DVector lengths: [1,3,2,1]--- chunks: [[0], [0,15,25], [0,40],[0]]--- elements: DInt [45,45,45,45]------ segids: DInt [0,0,3,4] (segment id of first slice on thread)--- offsets: DInt [0,45,0,5] (offset of that slice in its segment)--- @----splitSSegdOnElemsD :: Gang -> USSegd -> Dist ((USSegd,Int),Int)-splitSSegdOnElemsD g !segd - = {-# SCC "splitSSegdOnElemsD" #-}- imapD g mk (splitLenIdxD g (USegd.takeElements $ USSegd.takeUSegd segd))- where - -- Number of threads in gang.- !nThreads = gangSize g--- -- Build a USSegd from just the lengths, starts and sources fields.- -- The indices and elems fields of the contained USegd are - -- generated from the lengths.- buildUSSegd :: Vector Int -> Vector Int -> Vector Int -> USSegd- buildUSSegd lengths starts sources- = USSegd.mkUSSegd starts sources- $ USegd.fromLengths lengths-- -- Determine what elements go on a thread- mk :: Int -- Thread index.- -> (Int, Int) -- Number of elements on this thread,- -- and starting offset into the flat array.- -> ((USSegd, Int), Int) -- Segd for this thread, segid of first slice,- -- and offset of first slice.-- mk i (nElems, ixStart) - = case chunk segd ixStart nElems (i == nThreads - 1) of- (# lengths, starts, sources, l, o #) - -> ((buildUSSegd lengths starts sources, l), o)--{-# NOINLINE splitSSegdOnElemsD #-}--- NOINLINE because it's complicated and won't fuse with anything.--- This function has a large body of code and we don't want to blow up--- the client modules by inlining it everywhere.------------------------------------------------------------------------------------- | Determine what elements go on a thread.--- The 'chunk' refers to the a chunk of the flat array, and is defined--- by a set of segment slices. ------ Example:--- In this picture each X represents 5 elements, and we have 5 segements in total.------ @ segs: ----------------------- --- ------- --------------- ---------------------- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|--- | thread1 | thread2 | thread3 | thread4 |--- segid: 0 0 3 4--- offset: 0 45 0 5--- k: 0 1 3 5--- k': 1 3 5 5--- left: 0 15 0 45--- right: 45 20 5 0--- left_len: 0 1 0 1--- left_off: 0 45 0 5--- n': 1 3 2 1--- @-chunk :: USSegd -- ^ Segment descriptor of entire array.- -> Int -- ^ Starting offset into the flat array for the first- -- slice on this thread.- -> Int -- ^ Number of elements in this thread.- -> Bool -- ^ Whether this is the last thread in the gang.- -> (# Vector Int -- Lengths of segment slices, - , Vector Int -- Starting index of data in its vector- , Vector Int -- Source id- , Int -- segid of first slice- , Int #) -- offset of first slice.--chunk !ussegd !nStart !nElems is_last- = (# lengths', starts', sources', k-left_len, left_off #)- where- -- Lengths of all segments.- -- eg: [60, 10, 20, 40, 50]- lengths = USSegd.takeLengths ussegd-- -- Indices indices of all segments.- -- eg: [0, 60, 70, 90, 130]- indices = USSegd.takeIndices ussegd-- -- Starting indices for all segments.- starts = USSegd.takeStarts ussegd-- -- Source ids for all segments.- sources = USSegd.takeSources ussegd- - -- Total number of segments defined by segment descriptor.- -- eg: 5- n = Seq.length lengths-- -- Segid of the first seg that starts after the left of this chunk.- k = search nStart indices-- -- Segid of the first seg that starts after the right of this chunk.- k' | is_last = n- | otherwise = search (nStart + nElems) indices-- -- The length of the left-most slice of this chunk.- left | k == n = nElems- | otherwise = min ((Seq.index (here "chunk") indices k) - nStart) nElems-- -- The length of the right-most slice of this chunk.- length_right - | k' == k = 0- | otherwise = nStart + nElems - (Seq.index (here "chunk") indices (k'-1))-- -- Whether the first element in this chunk is an internal element of- -- of a segment. Alternatively, indicates that the first element of - -- the chunk is not the first element of a segment. - left_len | left == 0 = 0- | otherwise = 1-- -- If the first element of the chunk starts within a segment, - -- then gives the index within that segment, otherwise 0.- left_off | left == 0 = 0- | otherwise = nStart - (Seq.index (here "chunk") indices (k-1))-- -- How many segments this chunk straddles.- n' = left_len + (k'-k)-- -- Create the lengths for this chunk by first copying out the lengths- -- from the original segment descriptor. If the slices on the left- -- and right cover partial segments, then we update the corresponding- -- lengths.- (!lengths', !starts', !sources')- = runST (do- -- Create a new array big enough to hold all the lengths for this chunk.- mlengths' <- Seq.newM n'- msources' <- Seq.newM n'- mstarts' <- Seq.newM n'-- -- If the first element is inside a segment, - -- then update the length to be the length of the slice.- when (left /= 0) - $ do Seq.write mlengths' 0 left- Seq.write mstarts' 0 (Seq.index (here "chunk") starts (k - left_len) + left_off)- Seq.write msources' 0 (Seq.index (here "chunk") sources (k - left_len))-- -- Copy out array lengths for this chunk.- Seq.copy (Seq.mdrop left_len mlengths') (Seq.slice (here "chunk") lengths k (k'-k))- Seq.copy (Seq.mdrop left_len mstarts') (Seq.slice (here "chunk") starts k (k'-k))- Seq.copy (Seq.mdrop left_len msources') (Seq.slice (here "chunk") sources k (k'-k))-- -- If the last element is inside a segment, - -- then update the length to be the length of the slice.- when (length_right /= 0)- $ do Seq.write mlengths' (n' - 1) length_right-- clengths' <- Seq.unsafeFreeze mlengths'- cstarts' <- Seq.unsafeFreeze mstarts'- csources' <- Seq.unsafeFreeze msources'- return (clengths', cstarts', csources'))--{- = trace - (render $ vcat- [ text "CHUNK"- , pprp segd- , text "nStart: " <+> int nStart- , text "nElems: " <+> int nElems- , text "k: " <+> int k- , text "k': " <+> int k'- , text "left: " <+> int left- , text "right: " <+> int right- , text "left_len:" <+> int left_len- , text "left_off:" <+> int left_off- , text "n': " <+> int n'- , text ""]) lens'--}--{-# INLINE chunk #-}--- INLINE even though it should be inlined into splitSSegdOnElemsD anyway--- because that function contains the only use.------------------------------------------------------------------------------------- O(log n).--- Given a monotonically increasing vector of `Int`s,--- find the first element that is larger than the given value.--- --- eg search 75 [0, 60, 70, 90, 130] = 90--- search 43 [0, 60, 70, 90, 130] = 60----search :: Int -> Vector Int -> Int-search !x ys = go 0 (Seq.length ys)- where- go i n | n <= 0 = i- | Seq.index (here "search") ys mid < x- = go (mid + 1) (n - half - 1)- | otherwise = go i half- where- half = n `shiftR` 1- mid = i + half
− Data/Array/Parallel/Unlifted/Distributed/USegd.hs
@@ -1,349 +0,0 @@-{-# OPTIONS -Wall -fno-warn-orphans -fno-warn-missing-signatures #-}-{-# LANGUAGE CPP #-}-#include "fusion-phases.h"---- | Operations on Distributed Segment Descriptors-module Data.Array.Parallel.Unlifted.Distributed.USegd - ( splitSegdOnSegsD- , splitSegdOnElemsD- , splitSD- , joinSegdD- , glueSegdD)-where-import Data.Array.Parallel.Unlifted.Distributed.Arrays-import Data.Array.Parallel.Unlifted.Distributed.Combinators-import Data.Array.Parallel.Unlifted.Distributed.Types-import Data.Array.Parallel.Unlifted.Distributed.Gang-import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)-import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector, Unbox)-import Data.Array.Parallel.Base-import Data.Bits (shiftR)-import Control.Monad (when)-import qualified Data.Array.Parallel.Unlifted.Distributed.Types.USegd as DUSegd-import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd-import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as Seq--here :: String -> String-here s = "Data.Array.Parallel.Unlifted.Distributed.USegd." ++ s------------------------------------------------------------------------------------ | Split a segment descriptor across the gang, segment wise.--- Whole segments are placed on each thread, and we try to balance out--- the segments so each thread has the same number of array elements.------ We don't split segments across threads, as this would limit our ability--- to perform intra-thread fusion of lifted operations. The down side--- of this is that if we have few segments with an un-even size distribution--- then large segments can cause the gang to become unbalanced.------ In the following example the segment with size 100 dominates and--- unbalances the gang. There is no reason to put any segments on the--- the last thread because we need to wait for the first to finish anyway.------ @ > pprp $ splitSegdOnSegsD theGang--- $ lengthsToUSegd $ fromList [100, 10, 20, 40, 50 :: Int]--- --- DUSegd lengths: DVector lengths: [ 1, 3, 1, 0]--- chunks: [[100],[10,20,40],[50],[]]--- --- indices: DVector lengths: [1,3,1,0]--- chunks: [[0], [0,10,30], [0], []]------ elements: DInt [100,70,50,0]--- @------ NOTE: This splitSegdOnSegsD function isn't currently used.----splitSegdOnSegsD :: Gang -> USegd -> Dist USegd-splitSegdOnSegsD g !segd - = mapD g USegd.fromLengths- $ splitAsD g d lens- where- !d = snd- . mapAccumLD g chunks 0- . splitLenD g- $ USegd.takeElements segd-- n = USegd.length segd- lens = USegd.takeLengths segd-- chunks !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 = Seq.index (here "splitSegdOnSegsD") lens i-{-# NOINLINE splitSegdOnSegsD #-}------------------------------------------------------------------------------------- | Split a segment descriptor across the gang, element wise.--- We try to put the same number of elements on each thread, which means--- that segments are sometimes split across threads.------ Each thread gets a slice of segment descriptor, the segid of the first --- slice, and the offset of the first slice in its segment.--- --- Example:--- In this picture each X represents 5 elements, and we have 5 segements in total.------ @ segs: ----------------------- --- ------- --------------- ---------------------- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|--- | thread1 | thread2 | thread3 | thread4 |--- segid: 0 0 3 4--- offset: 0 45 0 5------ pprp $ splitSegdOnElemsD theGang4--- $ lengthsToUSegd $ fromList [60, 10, 20, 40, 50 :: Int]------ segd: DUSegd lengths: DVector lengths: [1,3,2,1]--- chunks: [[45],[15,10,20],[40,5],[45]]--- indices: DVector lengths: [1,3,2,1]--- chunks: [[0], [0,15,25], [0,40],[0]]--- elements: DInt [45,45,45,45]------ segids: DInt [0,0,3,4] (segment id of first slice on thread)--- offsets: DInt [0,45,0,5] (offset of that slice in its segment)--- @----splitSegdOnElemsD :: Gang -> USegd -> Dist ((USegd,Int),Int)-splitSegdOnElemsD g !segd - = {-# SCC "splitSegdOnElemsD" #-} - imapD g mk (splitLenIdxD g (USegd.takeElements segd))- where - -- Number of threads in gang.- !nThreads = gangSize g-- -- Determine what elements go on a thread- mk :: Int -- Thread index.- -> (Int, Int) -- Number of elements on this thread,- -- and starting offset into the flat array.- -> ((USegd, Int), Int) -- Segd for this thread, segid of first slice,- -- and offset of first slice.-- mk i (nElems, ixStart) - = case getChunk segd ixStart nElems (i == nThreads - 1) of- (# lens, l, o #) -> ((USegd.fromLengths lens, l), o)--{-# NOINLINE splitSegdOnElemsD #-}--- NOINLINE because this function has a large body of code and we don't want--- to blow up the client modules by inlining it everywhere.------------------------------------------------------------------------------------- | Determine what elements go on a thread.--- The 'chunk' refers to the a chunk of the flat array, and is defined--- by a set of segment slices. ------ Example:--- In this picture each X represents 5 elements, and we have 5 segements in total.------ @--- segs: ----------------------- --- ------- --------------- ---------------------- elems: |X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|X X X X X X X X X|--- | thread1 | thread2 | thread3 | thread4 |--- segid: 0 0 3 4--- offset: 0 45 0 5--- k: 0 1 3 5--- k': 1 3 5 5--- left: 0 15 0 45--- right: 45 20 5 0--- left_len: 0 1 0 1--- left_off: 0 45 0 5--- n': 1 3 2 1--- @-getChunk- :: USegd -- ^ Segment descriptor of entire array.- -> Int -- ^ Starting offset into the flat array for the first- -- slice on this thread.- -> Int -- ^ Number of elements in this thread.- -> Bool -- ^ Whether this is the last thread in the gang.- -> (# Vector Int -- Lengths of segment slices, - , Int -- segid of first slice,- , Int #) -- offset of first slice.--getChunk !segd !nStart !nElems is_last- = (# lens'', k-left_len, left_off #)- where- -- Lengths of all segments.- -- eg: [60, 10, 20, 40, 50]- !lens = USegd.takeLengths segd-- -- Indices indices of all segments.- -- eg: [0, 60, 70, 90, 130]- !idxs = USegd.takeIndices segd- - -- Total number of segments defined by segment descriptor.- -- eg: 5- !n = Seq.length lens-- -- Segid of the first seg that starts after the left of this chunk.- !k = search nStart idxs-- -- Segid of the first seg that starts after the right of this chunk.- !k' | is_last = n- | otherwise = search (nStart + nElems) idxs-- -- The length of the left-most slice of this chunk.- !left | k == n = nElems- | otherwise = min ((Seq.index (here "getChunk") idxs k) - nStart) nElems-- -- The length of the right-most slice of this chunk.- !right | k' == k = 0- | otherwise = nStart + nElems - (Seq.index (here "getChunk") idxs (k'-1))-- -- Whether the first element in this chunk is an internal element of- -- of a segment. Alternatively, indicates that the first element of - -- the chunk is not the first element of a segment. - !left_len | left == 0 = 0- | otherwise = 1-- -- If the first element of the chunk starts within a segment, - -- then gives the index within that segment, otherwise 0.- !left_off | left == 0 = 0- | otherwise = nStart - (Seq.index (here "getChunk") idxs (k-1))-- -- How many segments this chunk straddles.- !n' = left_len + (k'-k)-- -- Create the lengths for this chunk by first copying out the lengths- -- from the original segment descriptor. If the slices on the left- -- and right cover partial segments, then we update the corresponding- -- lengths.- !lens' - = runST (do- -- Create a new array big enough to hold all the lengths for this chunk.- !mlens' <- Seq.newM n'-- -- If the first element is inside a segment, - -- then update the length to be the length of the slice.- when (left /= 0) - $ Seq.write mlens' 0 left-- -- Copy out array lengths for this chunk.- Seq.copy (Seq.mdrop left_len mlens')- (Seq.slice "getChunk" lens k (k'-k))-- -- If the last element is inside a segment, - -- then update the length to be the length of the slice.- when (right /= 0)- $ Seq.write mlens' (n' - 1) right-- Seq.unsafeFreeze mlens')-- !lens'' = lens'-{- = trace - (render $ vcat- [ text "CHUNK"- , pprp segd- , text "nStart: " <+> int nStart- , text "nElems: " <+> int nElems- , text "k: " <+> int k- , text "k': " <+> int k'- , text "left: " <+> int left- , text "right: " <+> int right- , text "left_len:" <+> int left_len- , text "left_off:" <+> int left_off- , text "n': " <+> int n'- , text ""]) lens'--}--{-# INLINE getChunk #-}--- INLINE even though it should be inlined into splitSSegdOnElemsD anyway--- because that function contains the only use.------------------------------------------------------------------------------------- O(log n). Given a monotonically increasing vector of `Int`s,--- find the first element that is larger than the given value.--- --- eg search 75 [0, 60, 70, 90, 130] = 90--- search 43 [0, 60, 70, 90, 130] = 60----search :: Int -> Vector Int -> Int-search !x ys = go 0 (Seq.length ys)- where- go i n | n <= 0 = i-- | Seq.index (here "search") ys mid < x - = go (mid + 1) (n - half - 1)-- | otherwise = go i half- where- half = n `shiftR` 1- mid = i + half------------------------------------------------------------------------------------- | time O(segs)--- Join a distributed segment descriptor into a global one.--- This simply joins the distributed lengths and indices fields, but does--- not reconstruct the original segment descriptor as it was before splitting.--- --- @ > pprp $ joinSegdD theGang4 --- $ fstD $ fstD $ splitSegdOnElemsD theGang--- $ lengthsToUSegd $ fromList [60, 10, 20, 40, 50]--- --- USegd lengths: [45,15,10,20,40,5,45]--- indices: [0,45,60,70,90,130,135]--- elements: 180--- @--- --- TODO: sequential runtime is O(segs) due to application of lengthsToUSegd--- -joinSegdD :: Gang -> Dist USegd -> USegd-joinSegdD gang- = USegd.fromLengths- . joinD gang unbalanced- . mapD gang USegd.takeLengths-{-# INLINE_DIST joinSegdD #-}------------------------------------------------------------------------------------- | Glue a distributed segment descriptor back into the original global one.--- Prop: glueSegdD gang $ splitSegdOnElems gang usegd = usegd------ NOTE: This is runs sequentially and should only be used for testing purposes.----glueSegdD :: Gang -> Dist ((USegd, Int), Int) -> Dist USegd-glueSegdD gang bundle- = let !usegd = fstD $ fstD $ bundle- !lengths = DUSegd.takeLengthsD usegd- - !firstSegOffsets = sndD bundle-- -- | Whether the last segment in this chunk extends into the next chunk.- segSplits :: Dist Bool- !segSplits- = generateD_cheap gang $ \ix - -> if ix >= sizeD lengths - 1- then False- else indexD (here "glueSegdD") firstSegOffsets (ix + 1) /= 0-- !lengths' = fst $ carryD gang (+) 0 segSplits lengths- !dusegd' = mapD gang USegd.fromLengths lengths'-- in dusegd'-{-# INLINE_DIST glueSegdD #-}-----------------------------------------------------------------------------------splitSD :: Unbox a => Gang -> Dist USegd -> Vector a -> Dist (Vector a)-splitSD g dsegd xs- = splitAsD g (DUSegd.takeElementsD dsegd) xs-{-# INLINE_DIST splitSD #-}--{-# 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/What.hs view
@@ -0,0 +1,34 @@++module Data.Array.Parallel.Unlifted.Distributed.What+ ( Comp (..)+ , What (..))+where+ +-- | What distributed computation we are doing.+data Comp+ = CGen Bool -- cheap+ What++ | CMap What+ | CFold What+ | CScan What+ | CDist What+ deriving (Eq,Ord,Read,Show)++-- | What sort of thing is being computed.+data What+ = What String+ | WScalar + | WZip+ | WSlice+ | WLength+ | WLengthIdx+ | WBpermute++ -- Copy due to a join instruction.+ | WJoinCopy Int -- number elements++ | WFMapMap What What+ | WFMapGen What What+ | WFZipMap What What+ deriving (Eq,Ord,Read,Show)
Data/Array/Parallel/Unlifted/Parallel.hs view
@@ -43,15 +43,20 @@ -- * Segmented , replicateRSUP , appendSUP+ , appendSUPV , foldRUP , sumRUP -- * Index and Extracts , indexsFromVector , indexsFromVectorsUPVSegd+ , indexsFromVectorsUPVSegdP+ , extractsFromNestedUPSSegd , extractsFromVectorsUPSSegd+ , extractsFromVectorsUPVSegd+ , extractsFromVectorsUPVSegdP -- * Subarrays , dropUP
Data/Array/Parallel/Unlifted/Parallel/Basics.hs view
@@ -13,6 +13,7 @@ where import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Distributed.What import Data.Array.Parallel.Unlifted.Parallel.Combinators (mapUP) import Data.Array.Parallel.Unlifted.Parallel.Enum (enumFromToUP) import Data.Array.Parallel.Unlifted.Parallel.Permute (bpermuteUP)@@ -29,7 +30,7 @@ replicateUP :: Unbox e => Int -> e -> Vector e replicateUP n !e = joinD theGang balanced- . mapD theGang (\n' ->Seq.replicate n' e)+ . mapD (What "replicateUP/replicate") theGang (\n' ->Seq.replicate n' e) $ splitLenD theGang n {-# INLINE_UP replicateUP #-} @@ -62,7 +63,7 @@ interleaveUP :: Unbox e => Vector e -> Vector e -> Vector e interleaveUP xs ys = joinD theGang unbalanced- $ zipWithD theGang Seq.interleave+ $ zipWithD (What "interleaveUP/interleave") theGang Seq.interleave (splitD theGang balanced xs) (splitD theGang balanced ys) {-# INLINE_UP interleaveUP #-}@@ -73,9 +74,13 @@ 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+ sizes arr + = fst + $ scanD (What "indexedUP/length") theGang (+) 0 $ lengthD arr++ indexedFn + = \arr -> zipWithD (What "indexedUP.map") theGang + (\o -> Seq.map (\(x,y) -> (x + o, y)))+ (sizes arr) + $ mapD (What "indexedUP/indexed") theGang Seq.indexed arr {-# INLINE_UP indexedUP #-}
Data/Array/Parallel/Unlifted/Parallel/Combinators.hs view
@@ -13,6 +13,7 @@ where import Data.Array.Parallel.Base import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Distributed.What import Data.Array.Parallel.Unlifted.Parallel.UPSel import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq @@ -23,7 +24,8 @@ -- | Apply a worker to all elements of an array. mapUP :: (Unbox a, Unbox b) => (a -> b) -> Vector a -> Vector b mapUP f xs - = splitJoinD theGang (mapD theGang (Seq.map f)) xs+ = splitJoinD theGang + (mapD (What "mapUP/map") theGang (Seq.map f)) xs {-# INLINE_UP mapUP #-} @@ -31,7 +33,7 @@ filterUP :: Unbox a => (a -> Bool) -> Vector a -> Vector a filterUP f = joinD theGang unbalanced- . mapD theGang (Seq.filter f)+ . mapD (What "filterUP/filter") theGang (Seq.filter f) . splitD theGang unbalanced {-# INLINE_UP filterUP #-} @@ -75,7 +77,7 @@ ("tags length /= sum of args length") (Seq.length tags) (Seq.length xs + Seq.length ys) $ joinD theGang balanced- $ zipWithD theGang go rep+ $ zipWithD (What "combine2UP/go") theGang go rep $ splitD theGang balanced tags where go ((i,j), (m,n)) ts = Seq.combine2ByTag ts @@ -89,7 +91,7 @@ => (a -> b -> c) -> Vector a -> Vector b -> Vector c zipWithUP f xs ys = splitJoinD theGang - (mapD theGang (Seq.map (uncurry f))) + (mapD (What "zipWithUP/map") theGang (Seq.map (uncurry f))) (Seq.zip xs ys) {-# INLINE_UP zipWithUP #-} @@ -109,9 +111,9 @@ -- foldUP :: (Unbox a, DT a) => (a -> a -> a) -> a -> Vector a -> a foldUP f !z xs- = foldD theGang f- (mapD theGang (Seq.fold f z)- (splitD theGang unbalanced xs))+ = foldD (What "foldUP/f") theGang f+ $ mapD (What "foldUP/fold") theGang (Seq.fold f z)+ $ splitD theGang unbalanced xs {-# INLINE_UP foldUP #-} @@ -150,8 +152,8 @@ foldl1UP :: (DT a, Unbox a) => (a -> a -> a) -> Vector a -> a foldl1UP f arr = (maybe z (f z)- . foldD theGang combine'- . mapD theGang (Seq.foldl1Maybe f)+ . foldD (What "fold1UP/foldD") theGang combine'+ . mapD (What "fold1UP/fold1Maybe") theGang (Seq.foldl1Maybe f) . splitD theGang unbalanced) arr where z = Seq.index (here "fold1UP") arr 0@@ -171,9 +173,14 @@ -- scanUP :: (DT a, Unbox a) => (a -> a -> a) -> a -> Vector a -> Vector a 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+ = splitJoinD theGang go+ where go xs + = let (ds,zs) = unzipD + $ mapD (What "scanUP/scanRes") theGang (Seq.scanRes f z) xs++ zs' = fst + $ scanD (What "scanUP/scan") theGang f z zs++ in zipWithD (What "scanUP/map") theGang (Seq.map . f) zs' ds {-# INLINE_UP scanUP #-}
Data/Array/Parallel/Unlifted/Parallel/Enum.hs view
@@ -10,6 +10,7 @@ where import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Distributed.What import Data.Array.Parallel.Unlifted.Parallel.Combinators (mapUP) import GHC.Base (divInt) @@ -45,7 +46,7 @@ enumFromStepLenUP :: Int -> Int -> Int -> Vector Int enumFromStepLenUP start delta len = joinD theGang balanced- (mapD theGang gen+ (mapD (What "enumFromStepLenUP/gen") theGang gen (splitLenIdxD theGang len)) where gen (n,i) = Seq.enumFromStepLen (i * delta + start) delta n@@ -56,7 +57,7 @@ :: Int -> Vector Int -> Vector Int -> Vector Int -> Vector Int enumFromStepLenEachUP _n starts steps lens = joinD theGang unbalanced- $ mapD theGang enum+ $ mapD (What "enumFromStepLenEachUP/enum") theGang enum $ splitD theGang unbalanced (Seq.zip (Seq.zip starts steps) lens) where enum ps = let (qs, llens) = Seq.unzip ps
Data/Array/Parallel/Unlifted/Parallel/Extracts.hs view
@@ -5,25 +5,33 @@ module Data.Array.Parallel.Unlifted.Parallel.Extracts ( -- * Scattered indexing indexsFromVector+ , indexsFromVectorsUPVSegdP , indexsFromVectorsUPVSegd -- * Scattered extracts , extractsFromNestedUPSSegd , extractsFromVectorsUPSSegd- , extractsFromVectorsUPVSegd)++ , extractsFromVectorsUPVSegdP+ , extractsFromVectorsUPVSegd+ , extractsFromVectorsUPSSegdSegmap) where+import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Distributed.What import Data.Array.Parallel.Unlifted.Parallel.UPSSegd (UPSSegd) import Data.Array.Parallel.Unlifted.Parallel.UPVSegd (UPVSegd)-import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq import Data.Array.Parallel.Unlifted.Vectors (Vectors)+import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq import qualified Data.Array.Parallel.Unlifted.Parallel.UPSSegd as UPSSegd import qualified Data.Array.Parallel.Unlifted.Parallel.UPVSegd as UPVSegd import qualified Data.Array.Parallel.Unlifted.Sequential.UVSegd as UVSegd+import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd+import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd import qualified Data.Array.Parallel.Unlifted.Vectors as US import qualified Data.Array.Parallel.Unlifted.Stream as US import qualified Data.Array.Parallel.Unlifted.Sequential as Seq import qualified Data.Vector as V-+import Prelude as P -- Indexvs -------------------------------------------------------------------- -- | Lookup elements from a `Vector`.@@ -41,6 +49,19 @@ -- -- TODO: make this parallel. --+indexsFromVectorsUPVSegdP + :: (Unbox a, US.Unboxes a)+ => Vectors a -> UPVSegd -> Vector (Int, Int) -> Vector a++indexsFromVectorsUPVSegdP vectors upvsegd vsrcixs+ = splitJoinD theGang + (mapD (What "indexsFromVectorsUPVSegdP") theGang+ (indexsFromVectorsUPVSegd vectors upvsegd))+ vsrcixs+{-# INLINE_UP indexsFromVectorsUPVSegdP #-}+++-- | Lookup elements from some Vectors through a `UPVSegd` indexsFromVectorsUPVSegd :: (Unbox a, US.Unboxes a) => Vectors a -> UPVSegd -> Vector (Int, Int) -> Vector a@@ -51,7 +72,7 @@ !vsegids = UPVSegd.takeVSegidsRedundant upvsegd !upssegd = UPVSegd.takeUPSSegdRedundant upvsegd !ussegd = UPSSegd.takeUSSegd upssegd- in Seq.unstream+ in Seq.unstream $ US.streamElemsFromVectors vectors $ US.streamSrcIxsThroughUSSegd ussegd $ US.streamSrcIxsThroughVSegids vsegids@@ -72,7 +93,6 @@ (UPSSegd.takeUSSegd upssegd) {-# INLINE_U extractsFromNestedUPSSegd #-} - -- | TODO: make this parallel. extractsFromVectorsUPSSegd :: (Unbox a, US.Unboxes a)@@ -87,7 +107,50 @@ {-# INLINE_UP extractsFromVectorsUPSSegd #-} --- | TODO: make this parallel.++-- From UPVSegd ---------------------------------------------------------------+-- | Parallel extracts from UPVSegd and Segmap+-- TODO: This just distributes the segmap over the gang, and will be unbalanced+-- if there aren't many segments, or they have varying sizes.+extractsFromVectorsUPVSegdP+ :: (Unbox a, US.Unboxes a)+ => UPVSegd+ -> Vectors a+ -> Vector a++extractsFromVectorsUPVSegdP !upvsegd !vectors+ = let !segs = UPVSegd.takeDistributed upvsegd+ !vsegids = UPVSegd.takeVSegidsRedundant upvsegd+ !ussegd = UPSSegd.takeUSSegd+ $ UPVSegd.takeUPSSegdRedundant upvsegd+ in joinD theGang balanced+ $ mapD (What "extractsFromVectorsUPVSegdP")+ theGang+ (extractsFromVectorsUPSSegd_split+ ussegd+ vsegids+ vectors)+ segs++{-# INLINE_UP extractsFromVectorsUPVSegdP #-}++-- | Sequential extracts from USSegd and Segmap+extractsFromVectorsUPSSegd_split+ :: (Unbox a, US.Unboxes a)+ => USSegd.USSegd+ -> Vector Int+ -> Vectors a+ -> ((USegd.USegd,Int),Int)+ -> Vector a++extractsFromVectorsUPSSegd_split !ussegd !vsegids !vectors !which+ = Seq.unstream + $ US.streamSegsFromVectorsUSSegd_split vectors+ ussegd vsegids which+{-# INLINE_UP extractsFromVectorsUPSSegd_split #-}+++-- | Sequential extracts from UPVSegd. extractsFromVectorsUPVSegd :: (Unbox a, US.Unboxes a) => UPVSegd@@ -102,3 +165,18 @@ (UPSSegd.takeUSSegd $ UPVSegd.takeUPSSegdRedundant upvsegd) {-# INLINE_UP extractsFromVectorsUPVSegd #-} ++-- | Sequential extracts from USSegd and Segmap+extractsFromVectorsUPSSegdSegmap+ :: (Unbox a, US.Unboxes a)+ => UPSSegd+ -> Vectors a+ -> Vector Int+ -> Vector a++extractsFromVectorsUPSSegdSegmap upssegd vectors segmap+ = Seq.unstream + $ US.streamSegsFromVectorsUSSegdSegmap vectors+ (UPSSegd.takeUSSegd upssegd)+ segmap+{-# INLINE_UP extractsFromVectorsUPSSegdSegmap #-}
Data/Array/Parallel/Unlifted/Parallel/Segmented.hs view
@@ -5,21 +5,30 @@ module Data.Array.Parallel.Unlifted.Parallel.Segmented ( replicateRSUP , appendSUP+ , appendSUP_old+ , appendSUPV , foldRUP , sumRUP) where import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Distributed.What import Data.Array.Parallel.Unlifted.Parallel.Basics import Data.Array.Parallel.Unlifted.Parallel.UPSegd (UPSegd) import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd) import Data.Array.Parallel.Unlifted.Sequential.Vector as Seq+import qualified Data.Array.Parallel.Unlifted.Vectors as Vs import qualified Data.Array.Parallel.Unlifted.Parallel.UPSegd as UPSegd import qualified Data.Array.Parallel.Unlifted.Sequential as Seq import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd+import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd+import qualified Data.Array.Parallel.Unlifted.Parallel.UPVSegd as UPVSegd+import qualified Data.Array.Parallel.Unlifted.Parallel.UPSSegd as UPSSegd import Data.Vector.Fusion.Stream.Monadic ( Stream(..), Step(..) ) import Data.Vector.Fusion.Stream.Size ( Size(..) ) import qualified Data.Vector.Fusion.Stream as S +import GHC.Exts -- for unboxed primops+ here :: String -> String here s = "Data.Array.Parallel.Unlifted.Parallel.Segmented." Prelude.++ s @@ -42,14 +51,16 @@ -- | Segmented append. appendSUP :: Unbox a- => UPSegd- -> UPSegd -> Vector a- -> UPSegd -> Vector a+ => UPSegd -- ^ segment descriptor of result+ -> UPSegd -- ^ left-hand segd+ -> Vector a -- ^ left-hand data+ -> UPSegd -- ^ right-hand segd+ -> Vector a -- ^ right-hand data -> Vector a appendSUP segd !xd !xs !yd !ys = joinD theGang balanced- . mapD theGang append+ . mapD (What "appendSUP/append") theGang append $ UPSegd.takeDistributed segd where append ((segd',seg_off),el_off) = Seq.unstream@@ -67,12 +78,331 @@ -> Vector a -- ^ Data of first array -> USegd -- ^ Segment descriptor of second array. -> Vector a -- ^ Data of second array.+ -> Int -- ^ How many elements to return+ -> Int -- ^ Segment offset+ -> Int -- ^ Element offset+ -> S.Stream a++appendSegS !xd !xs !yd !ys !n seg_off el_off+ = Stream next state (Exact n)+ where+ !xlens = USegd.takeLengths xd+ !ylens = USegd.takeLengths yd++ -- Two index functions because of monomorphism restriction+ {-# INLINE index1 #-}+ index1 = Seq.index (here "appendSegS")++ {-# INLINE index2 #-}+ index2 = Seq.index (here "appendSegS")++ {-# INLINE unbox #-}+ unbox (I# i) = i+ + state+ -- Nothing to return+ | n == 0 = ASDo+ { as_takefrom = 0#+ , as_seg_off = 0#+ , as_xs_index = 0#+ , as_ys_index = 0#+ , as_next_swap= 0#+ , as_remain = 0# }++ -- Start returning data from xs+ | el_off < xlens `index1` seg_off+ = let xi = (USegd.takeIndices xd `index1` seg_off) + el_off+ yi = USegd.takeIndices yd `index1` seg_off+ swap = (USegd.takeLengths xd `index1` seg_off) - el_off+ in ASDo+ -- start reading from xs, then read from ys at end of this xs segment+ { as_takefrom = 0#+ , as_seg_off = unbox seg_off+ , as_xs_index = unbox xi+ , as_ys_index = unbox yi+ , as_next_swap= unbox swap+ , as_remain = unbox n }++ -- Start with ys+ | otherwise+ = let -- NOTE: *not* indicesUSegd xd ! (seg_off+1) since seg_off+1+ -- might be out of bounds+ xi = (USegd.takeIndices xd `index1` seg_off) + (USegd.takeLengths xd `index1` seg_off)+ el_off' = el_off - USegd.takeLengths xd `index1` seg_off+ yi = (USegd.takeIndices yd `index1` seg_off) + el_off'+ swap = (USegd.takeLengths yd `index1` seg_off) - el_off'+ in ASDo+ { as_takefrom = 1#+ , as_seg_off = unbox seg_off+ , as_xs_index = unbox xi+ , as_ys_index = unbox yi+ , as_next_swap= unbox swap+ , as_remain = unbox n }++ {-# INLINE next #-}+ next ASDo{as_remain=0#} = return Done++ -- Reading from xs+ next s@ASDo{as_takefrom=0#}+ -- Done reading xs, so read the rest of this segment from ys.+ | as_next_swap s ==# 0#+ = return $ Skip (s{as_takefrom=1#, as_next_swap= unbox (ylens `index1` I# (as_seg_off s))})++ -- Grab a value from xs+ | otherwise = return $ Yield (xs `index2` I# (as_xs_index s)) (inc s)++ -- Reading from ys; takefrom nonzero+ next s+ -- Done reading ys, so we need to look at the next segment's xs+ | as_next_swap s ==# 0#+ = let seg' = as_seg_off s +# 1#+ in return $ Skip (s {as_takefrom=0#, as_seg_off=seg', as_next_swap= unbox (xlens `index1` I# seg')})++ -- Grab a value from ys+ | otherwise = return $ Yield (ys `index2` I# (as_ys_index s)) (inc s)++ {-# INLINE inc #-}+ -- Move data pointer forward, and decrease remaining and swap+ inc s@ASDo{as_takefrom=0#, as_xs_index=xi, as_next_swap=swap, as_remain=n'}+ = s{as_xs_index=xi +# 1#, as_next_swap=swap -# 1#, as_remain=n' -# 1#}++ -- Takefrom is nonzero: reading from ys+ inc s@ASDo{as_ys_index=yi, as_next_swap=swap, as_remain=n'}+ = s{as_ys_index=yi +# 1#, as_next_swap=swap -# 1#, as_remain=n' -# 1#}+{-# INLINE_STREAM appendSegS #-}++data AppendState+ = ASDo+ { as_takefrom :: Int# -- ^ 0 = xs, nonzero = ys+ , as_seg_off :: Int# -- ^ current segment+ , as_xs_index :: Int# -- ^ pointer into xs data+ , as_ys_index :: Int# -- ^ pointer into ys data+ , as_next_swap:: Int# -- ^ toggle takefrom in this many elements+ , as_remain :: Int# -- ^ how many left+ }++-- virtual scattered append+appendSUPV+ :: (Vs.Unboxes a, Unbox a)+ => UPSegd -- ^ segment descriptor of result+ -> UPVSegd.UPVSegd -- ^ left-hand segd+ -> Vs.Vectors a -- ^ left-hand data+ -> UPVSegd.UPVSegd -- ^ right-hand segd+ -> Vs.Vectors a -- ^ right-hand data+ -> Vector a++appendSUPV segd !xd !xs !yd !ys+ = joinD theGang balanced+ . mapD (What "appendSUPV/append") theGang append+ $ UPSegd.takeDistributed segd+ where append ((segd',seg_off),el_off)+ = Seq.unstream+ $ appendUPVSegS xd xs+ yd ys+ (USegd.takeElements segd')+ seg_off el_off+{-# INLINE_UP appendSUPV #-}+appendUPVSegS+ :: Vs.Unboxes a+ => UPVSegd.UPVSegd -- ^ Segment descriptor of first array.+ -> Vs.Vectors a -- ^ Data of first array+ -> UPVSegd.UPVSegd -- ^ Segment descriptor of second array.+ -> Vs.Vectors a -- ^ Data of second array.+ -> Int -- ^ How many elements to return+ -> Int -- ^ Segment offset+ -> Int -- ^ Element offset+ -> S.Stream a++appendUPVSegS !xd !xs !yd !ys !n seg_off el_off+ = Stream next state (Exact n)+ where+ !xvsegs= UPVSegd.takeVSegidsRedundant xd+ !yvsegs= UPVSegd.takeVSegidsRedundant yd++ !xssegd= UPSSegd.takeUSSegd $ UPVSegd.takeUPSSegdRedundant xd+ !yssegd= UPSSegd.takeUSSegd $ UPVSegd.takeUPSSegdRedundant yd++ !xsegd = USSegd.takeUSegd xssegd+ !ysegd = USSegd.takeUSegd yssegd++ -- get physical segment id+ {-#INLINE xpseg #-}+ xpseg s = index1 xvsegs "xpseg" s+ {-#INLINE ypseg #-}+ ypseg s = index1 yvsegs "ypseg" s++ !xseglens = USegd.takeLengths xsegd+ !yseglens = USegd.takeLengths ysegd++ !xsrc = USSegd.takeSources xssegd+ !ysrc = USSegd.takeSources yssegd++ !xstrt = USSegd.takeStarts xssegd+ !ystrt = USSegd.takeStarts yssegd++ -- physical lengths+ {-#INLINE xplen #-}+ xplen s = index1 xseglens "xplen1" (xpseg s)+ {-#INLINE yplen #-}+ yplen s = index1 yseglens "yplen1" (ypseg s)++ -- get actual data+ {-# INLINE gdata #-}+ gdata gs st+ = let !src = avs_ssrc st+ !strt = avs_sstart st+ !ix = avs_index st+ in index2 gs (I# src) (I# (strt +# ix))++ -- get scattered segment source and starts+ {-# INLINE getscatter #-}+ getscatter gpseg gsrcs gstrts segid+ = let !phys = gpseg segid in+ let !src = index1 gsrcs "src" phys in+ let !strt = index1 gstrts "strt" phys in+ (src, strt)++ {-# INLINE index1 #-}+ --index1 v i = Seq.index (here "appendUVSegS") v i++ index1 v h i = Seq.index (here $ "appendUVSegS:" Prelude.++ h) v i++ {-# INLINE index2 #-}+ index2 v i1 i2 = Vs.index2 (here "appendUVSegS") v i1 i2+++ {-# INLINE unbox #-}+ unbox (I# i) = i++ state+ -- Nothing to return+ | n == 0 = ASUPVDo+ { avs_takefrom = 0#+ , avs_seg_off = 0#+ , avs_index = 0#+ , avs_next_swap= 0#+ , avs_remain = 0#+ , avs_sstart = 0#+ , avs_ssrc = 0# }++ -- Start returning data from xs+ | el_off < xplen seg_off+ = let (src,strt) = getscatter xpseg xsrc xstrt seg_off+ swap = (xplen seg_off) - el_off+ in ASUPVDo+ -- start reading from xs, then read from ys at end of this xs segment+ { avs_takefrom = 0#+ , avs_seg_off = unbox seg_off+ , avs_index = unbox el_off+ , avs_next_swap= unbox swap+ , avs_remain = unbox n+ , avs_sstart = unbox strt+ , avs_ssrc = unbox src }++ -- Start with ys+ | otherwise+ = let (src,strt) = getscatter ypseg ysrc ystrt seg_off+ el_off' = el_off - xplen seg_off+ swap = (yplen seg_off) - el_off'+ in ASUPVDo+ { avs_takefrom = 1#+ , avs_seg_off = unbox seg_off+ , avs_index = unbox el_off'+ , avs_next_swap= unbox swap+ , avs_remain = unbox n+ , avs_sstart = unbox strt+ , avs_ssrc = unbox src }++ {-# INLINE next #-}+ next ASUPVDo{avs_remain=0#} = return Done++ -- Reading from xs+ next s@ASUPVDo{avs_takefrom=0#}+ -- Done reading xs, so read the rest of this segment from ys.+ | avs_next_swap s ==# 0# =+ let seg' = I# (avs_seg_off s)+ (src,strt) = getscatter ypseg ysrc ystrt seg'+ in return $ Skip $+ s {+ avs_takefrom = 1#+ , avs_index = 0#+ , avs_next_swap = unbox (yplen seg')+ , avs_sstart = unbox strt+ , avs_ssrc = unbox src }+ -- Grab a value from xs+ | otherwise = return $ Yield (gdata xs s) (inc s)++ -- Reading from ys, so avs_takefrom=1#+ next s+ -- Done reading ys, so we need to look at the next segment's xs+ | avs_next_swap s ==# 0#+ = let seg' = I# (avs_seg_off s +# 1#)+ (src,strt) = getscatter xpseg xsrc xstrt seg'+ in return $ Skip $+ s {+ avs_takefrom = 0#+ , avs_seg_off = unbox seg'+ , avs_index = 0#+ , avs_next_swap = unbox (xplen seg')+ , avs_sstart = unbox strt+ , avs_ssrc = unbox src }++ -- Grab a value from ys+ | otherwise = return $ Yield (gdata ys s) (inc s)++ {-# INLINE inc #-}+ -- Move data pointer forward and decrease remaining and swap+ inc s@ASUPVDo{avs_index=ix, avs_next_swap=swap, avs_remain=n'}+ = s{avs_index=ix +# 1#, avs_next_swap=swap -# 1#, avs_remain=n' -# 1#}+{-# INLINE_STREAM appendUPVSegS #-}++data AppendUPVState+ = ASUPVDo+ { avs_takefrom :: Int# -- ^ 0 = xs, 1/else = ys+ , avs_seg_off :: Int# -- ^ current segment+ , avs_index :: Int# -- ^ index into current segment+ , avs_next_swap:: Int# -- ^ toggle takefrom in this many elements+ , avs_remain :: Int# -- ^ how many left+ , avs_sstart :: Int# -- ^ scattered segment start+ , avs_ssrc :: Int# -- ^ scattered segment source+ }++-- Append ---------------------------------------------------------------------+-- | Segmented append.+-- -old+appendSUP_old+ :: Unbox a+ => UPSegd+ -> UPSegd -> Vector a+ -> UPSegd -> Vector a+ -> Vector a++appendSUP_old segd !xd !xs !yd !ys+ = joinD theGang balanced+ . mapD (What "appendSUP_old/append") theGang append+ $ UPSegd.takeDistributed segd+ where append ((segd',seg_off),el_off)+ = Seq.unstream+ $ appendSegS_old (UPSegd.takeUSegd xd) xs+ (UPSegd.takeUSegd yd) ys+ (USegd.takeElements segd')+ seg_off el_off+{-# INLINE_UP appendSUP_old #-}+++-- append ---------------------------------------------------------------------+appendSegS_old+ :: 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 -appendSegS !xd !xs !yd !ys !n seg_off el_off+appendSegS_old !xd !xs !yd !ys !n seg_off el_off = Stream next state (Exact n) where !xlens = USegd.takeLengths xd@@ -116,7 +446,7 @@ in return $ Skip (Just (False, seg', i, j, xlens `index1` seg', n')) | otherwise = return $ Yield (ys `index2` j) (Just (True, seg, i, j+1, k-1, n'-1))-{-# INLINE_STREAM appendSegS #-}+{-# INLINE_STREAM appendSegS_old #-} -- foldR ----------------------------------------------------------------------@@ -124,9 +454,9 @@ foldRUP :: (Unbox a, Unbox b) => (b -> a -> b) -> b -> Int -> Vector a -> Vector b foldRUP f z !segSize xs = joinD theGang unbalanced- (mapD theGang + (mapD (What "foldRUP/foldRU") theGang (Seq.foldlRU f z segSize)- (splitAsD theGang (mapD theGang (*segSize) dlen) xs))+ (splitAsD theGang (mapD (What "foldRUP/segSize") theGang (*segSize) dlen) xs)) where noOfSegs = Seq.length xs `div` segSize dlen = splitLenD theGang noOfSegs
Data/Array/Parallel/Unlifted/Parallel/UPSSegd.hs view
@@ -38,12 +38,13 @@ where import Data.Array.Parallel.Pretty hiding (empty) import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Distributed.What import Data.Array.Parallel.Unlifted.Parallel.UPSegd (UPSegd) import Data.Array.Parallel.Unlifted.Sequential.USSegd (USSegd) import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector, MVector, Unbox) import Data.Array.Parallel.Unlifted.Vectors (Vectors, Unboxes) import qualified Data.Array.Parallel.Unlifted.Parallel.UPSegd as UPSegd-import qualified Data.Array.Parallel.Unlifted.Distributed.USSegd as DUSSegd+import qualified Data.Array.Parallel.Unlifted.Distributed.Data.USSegd as DUSSegd import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as US import qualified Data.Array.Parallel.Unlifted.Sequential as Seq@@ -283,7 +284,8 @@ where (dcarry,drs) = unzipD- $ mapD theGang partial (takeDistributed segd)+ $ mapD (What "UPSSegd.foldSegsWithP/partial") theGang + partial (takeDistributed segd) partial ((ssegd, k), off) = let rs = fSeg ssegd xss
Data/Array/Parallel/Unlifted/Parallel/UPSegd.hs view
@@ -35,8 +35,9 @@ , foldSegsWithP) where import Data.Array.Parallel.Unlifted.Distributed-import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)-import qualified Data.Array.Parallel.Unlifted.Distributed.USegd as USegd+import Data.Array.Parallel.Unlifted.Distributed.What+import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd)+import qualified Data.Array.Parallel.Unlifted.Distributed.Data.USegd as USegd import qualified Data.Array.Parallel.Unlifted.Sequential as Seq import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as US import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd@@ -183,7 +184,7 @@ indicesP :: UPSegd -> Vector Int indicesP = joinD theGang balanced- . mapD theGang indices+ . mapD (What "UPSegd.indicesP/indices") theGang indices . takeDistributed where indices ((segd,_k),off) = Seq.indicesSU' off segd@@ -200,7 +201,7 @@ replicateWithP :: Unbox a => UPSegd -> Vector a -> Vector a replicateWithP segd !xs = joinD theGang balanced- . mapD theGang rep+ . mapD (What "UPSegd.replicateWithP/replicateSU") theGang rep $ takeDistributed segd where rep ((dsegd,di),_)@@ -254,7 +255,7 @@ where (dcarry,drs) = unzipD- $ mapD theGang partial+ $ mapD (What "UPSegd.foldSegsWithP/partial") theGang partial $ zipD (takeDistributed segd) (splitD theGang balanced xs)
Data/Array/Parallel/Unlifted/Parallel/UPSel.hs view
@@ -23,6 +23,7 @@ import Data.Array.Parallel.Unlifted.Sequential.Vector as US import Data.Array.Parallel.Unlifted.Sequential.USel import Data.Array.Parallel.Unlifted.Distributed+import Data.Array.Parallel.Unlifted.Distributed.What import Data.Array.Parallel.Base (Tag, tagToInt) @@ -88,11 +89,11 @@ mkUPSelRep2 :: Vector Tag -> UPSelRep2 mkUPSelRep2 tags = zipD idxs lens where- lens = mapD theGang count+ lens = mapD (What "UPSelRep2.mkUPSelRep2/count") theGang count $ splitD theGang balanced tags idxs = fst- $ scanD theGang add (0,0) lens+ $ scanD (What "UPSelRep2.mkUPSelRep2/add") theGang add (0,0) lens count bs = let ones = US.sum (US.map tagToInt bs) in (US.length bs - ones,ones)@@ -104,7 +105,7 @@ indicesUPSelRep2 :: Vector Tag -> UPSelRep2 -> Vector Int indicesUPSelRep2 tags rep = joinD theGang balanced- $ zipWithD theGang indices+ $ zipWithD (What "UPSel.indicesUPSelRep2/split") theGang indices (splitD theGang balanced tags) rep where
Data/Array/Parallel/Unlifted/Parallel/UPVSegd.hs view
@@ -29,6 +29,7 @@ , length , takeVSegids, takeVSegidsRedundant , takeUPSSegd, takeUPSSegdRedundant+ , takeDistributed , takeLengths , getSeg @@ -43,14 +44,18 @@ , appendWith , combine2) where+import Data.Array.Parallel.Unlifted.Distributed import Data.Array.Parallel.Unlifted.Parallel.Permute import Data.Array.Parallel.Unlifted.Parallel.UPSel (UPSel2) import Data.Array.Parallel.Unlifted.Parallel.UPSSegd (UPSSegd) import Data.Array.Parallel.Unlifted.Parallel.UPSegd (UPSegd)+import Data.Array.Parallel.Unlifted.Sequential.USegd (USegd) import Data.Array.Parallel.Unlifted.Sequential.Vector (Vector) import Data.Array.Parallel.Pretty hiding (empty) import Prelude hiding (length)+import qualified Data.Array.Parallel.Unlifted.Distributed.Data.USegd as USegd import qualified Data.Array.Parallel.Unlifted.Sequential.Vector as US+import qualified Data.Array.Parallel.Unlifted.Sequential.USegd as USegd import qualified Data.Array.Parallel.Unlifted.Sequential.USSegd as USSegd import qualified Data.Array.Parallel.Unlifted.Parallel.UPSel as UPSel import qualified Data.Array.Parallel.Unlifted.Parallel.UPSegd as UPSegd@@ -78,6 +83,13 @@ -- are layed out in memory. , upvsegd_upssegd_redundant :: UPSSegd -- LAZY FIELD , upvsegd_upssegd_culled :: UPSSegd -- LAZY FIELD++ -- | Segment descriptors distributed over the gang,+ -- along with logical segment id and element offsets.+ -- Note that segment ids here refer to a vsegid,+ -- not a physical or scattered segment id.+ -- See `splitSegdOfElemsD` for an example.+ , upvsegd_dsegd :: Dist ((USegd,Int),Int) -- LAZY FIELD -- IMPORTANT: -- When vsegids are transformed due to a segmented replication operation, @@ -106,7 +118,7 @@ -- | Pretty print the physical representation of a `UVSegd` instance PprPhysical UPVSegd where- pprp (UPVSegd _ _ vsegids _ upssegd)+ pprp (UPVSegd _ _ vsegids _ upssegd _) = vcat [ text "UPVSegd" $$ (nest 7 $ text "vsegids: " <+> (text $ show $ US.toList vsegids)) , pprp upssegd ]@@ -115,7 +127,7 @@ -- | O(1). Check the internal consistency of a virutal segmentation descriptor. --- -- * TODO: this doesn't do any checks yet.---+-- valid :: UPVSegd -> Bool valid UPVSegd{} = True {-# NOINLINE valid #-}@@ -136,8 +148,19 @@ mkUPVSegd vsegids ussegd = UPVSegd False vsegids vsegids ussegd ussegd+ -- Lazy, so doesn't contribute to big-O+ $ mkDist vsegids ussegd {-# INLINE_UP mkUPVSegd #-} +-- | O(segs). Distribute the logical over the gang.+mkDist+ :: Vector Int -- ^ vsegids+ -> UPSSegd -- ^ Scattered segment descriptor+ -> Dist ((USegd,Int),Int)+mkDist vsegids ussegd+ = USegd.splitSegdOnElemsD theGang+ $ USegd.fromLengths+ $ takeLengths' vsegids ussegd -- | O(segs). Promote a `UPSSegd` to a `UPVSegd`. -- The result contains one virtual segment for every physical segment@@ -149,6 +172,7 @@ fromUPSSegd upssegd = let vsegids = US.enumFromTo 0 (UPSSegd.length upssegd - 1) in UPVSegd True vsegids vsegids upssegd upssegd+ (mkDist vsegids upssegd) {-# INLINE_UP fromUPSSegd #-} @@ -168,6 +192,7 @@ = let vsegids = US.empty upssegd = UPSSegd.empty in UPVSegd True vsegids vsegids upssegd upssegd+ (mkDist vsegids upssegd) {-# INLINE_UP empty #-} @@ -179,6 +204,7 @@ = let vsegids = US.singleton 0 upssegd = UPSSegd.singleton n in UPVSegd True vsegids vsegids upssegd upssegd+ (mkDist vsegids upssegd) {-# INLINE_UP singleton #-} @@ -273,21 +299,36 @@ {-# INLINE takeUPSSegdRedundant #-} +-- | O(1) or O(segs). Yield `USegd`s distributed over a logical view of this `UPVSegd`.+-- The complexity is only O(1) if this has already been evaluated.+takeDistributed :: UPVSegd -> Dist ((USegd,Int),Int)+takeDistributed = upvsegd_dsegd+{-# INLINE takeDistributed #-}++ -- | O(segs). Yield the lengths of the segments described by a `UPVSegd`. --- -- * TODO: This is slow and sequential. -- takeLengths :: UPVSegd -> Vector Int-takeLengths (UPVSegd manifest _ vsegids _ upssegd)+takeLengths (UPVSegd manifest _ vsegids _ upssegd _) | manifest = UPSSegd.takeLengths upssegd | otherwise - = let !lengths = (UPSSegd.takeLengths upssegd)- in US.map (US.index (here "takeLengths") lengths) vsegids+ = takeLengths' vsegids upssegd {-# NOINLINE takeLengths #-} -- NOINLINE because we don't want a case expression due to the test on the -- manifest flag to appear in the core program. +takeLengths'+ :: Vector Int -- ^ Vsegids+ -> UPSSegd -- ^ Scattered segment descriptor+ -> Vector Int+takeLengths' vsegids upssegd+ = let !lengths = (UPSSegd.takeLengths upssegd)+ in US.map (US.index (here "takeLengths") lengths) vsegids+{-# INLINE takeLengths' #-} + -- | O(1). Get the length, starting index, and source id of a segment. --- -- NOTE: We don't return the segment index field from the `USSegd` as this refers@@ -343,7 +384,7 @@ -- * TODO: if the upvsegd is manifest and contiguous this can be O(1). -- unsafeDemoteToUPSegd :: UPVSegd -> UPSegd-unsafeDemoteToUPSegd (UPVSegd _ _ vsegids _ upssegd)+unsafeDemoteToUPSegd (UPVSegd _ _ vsegids _ upssegd _) = {-# SCC "unsafeDemoteToUPSegd" #-} UPSegd.fromLengths $ bpermuteUP (UPSSegd.takeLengths upssegd) vsegids@@ -365,7 +406,7 @@ -- It runs the sequential 'cull' then reconstructs the UPSSegd. -- updateVSegs :: (Vector Int -> Vector Int) -> UPVSegd -> UPVSegd-updateVSegs fUpdate (UPVSegd _ vsegids _ upssegd _)+updateVSegs fUpdate (UPVSegd _ vsegids _ upssegd _ _) = let -- When we transform the vsegids, we don't know whether they all -- made it into the result. @@ -383,6 +424,7 @@ in UPVSegd False vsegids_redundant vsegids_culled upssegd upssegd_culled+ (mkDist vsegids_redundant upssegd) {-# NOINLINE updateVSegs #-} -- NOINLINE because we want to see this happening in core. @@ -399,9 +441,10 @@ -- like segmented fold will have the wrong work complexity. -- updateVSegsReachable :: (Vector Int -> Vector Int) -> UPVSegd -> UPVSegd-updateVSegsReachable fUpdate (UPVSegd _ _ vsegids _ upssegd)+updateVSegsReachable fUpdate (UPVSegd _ _ vsegids _ upssegd _) = let vsegids' = fUpdate vsegids in UPVSegd False vsegids' vsegids' upssegd upssegd+ (mkDist vsegids' upssegd) {-# NOINLINE updateVSegsReachable #-} -- NOINLINE because we want to see this happening in core. @@ -419,8 +462,8 @@ -> UPVSegd appendWith- (UPVSegd _ _ vsegids1 _ upssegd1) pdatas1- (UPVSegd _ _ vsegids2 _ upssegd2) pdatas2+ (UPVSegd _ _ vsegids1 _ upssegd1 _) pdatas1+ (UPVSegd _ _ vsegids2 _ upssegd2 _) pdatas2 = let -- vsegids releative to appended psegs vsegids1' = vsegids1@@ -435,6 +478,7 @@ upssegd2 pdatas2 in UPVSegd False vsegids' vsegids' upssegd' upssegd'+ (mkDist vsegids' upssegd') {-# NOINLINE appendWith #-} -- NOINLINE because it doesn't need to be specialised -- and we're worried about code explosion.@@ -455,8 +499,8 @@ combine2 upsel2- (UPVSegd _ _ vsegids1 _ upssegd1) pdatas1- (UPVSegd _ _ vsegids2 _ upssegd2) pdatas2+ (UPVSegd _ _ vsegids1 _ upssegd1 _) pdatas1+ (UPVSegd _ _ vsegids2 _ upssegd2 _) pdatas2 = let -- vsegids relative to combined psegs vsegids1' = vsegids1@@ -472,6 +516,7 @@ upssegd2 pdatas2 in UPVSegd False vsegids' vsegids' upssegd' upssegd'+ (mkDist vsegids' upssegd') {-# NOINLINE combine2 #-} -- NOINLINE because it doesn't need to be specialised -- and we're worried about code explosion.
dph-prim-par.cabal view
@@ -1,5 +1,5 @@ Name: dph-prim-par-Version: 0.6.1.1+Version: 0.7.0.1 License: BSD3 License-File: LICENSE Author: The DPH Team@@ -15,34 +15,37 @@ Library Exposed-Modules:- Data.Array.Parallel.Unlifted.Distributed- Data.Array.Parallel.Unlifted.Distributed.Gang- Data.Array.Parallel.Unlifted.Distributed.TheGang- Data.Array.Parallel.Unlifted.Parallel- Data.Array.Parallel.Unlifted.Parallel.UPSegd- Data.Array.Parallel.Unlifted.Parallel.UPSSegd- Data.Array.Parallel.Unlifted.Parallel.UPVSegd- Data.Array.Parallel.Unlifted.Parallel.UPSel- Data.Array.Parallel.Unlifted+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DT+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DPrim+ Data.Array.Parallel.Unlifted.Distributed.Primitive.DistST+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Gang+ Data.Array.Parallel.Unlifted.Distributed.Primitive.Operators+ Data.Array.Parallel.Unlifted.Distributed.Primitive - Other-Modules:- 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.Data.Bool+ Data.Array.Parallel.Unlifted.Distributed.Data.Maybe+ Data.Array.Parallel.Unlifted.Distributed.Data.Ordering+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar.Base+ Data.Array.Parallel.Unlifted.Distributed.Data.Scalar+ Data.Array.Parallel.Unlifted.Distributed.Data.Tuple+ Data.Array.Parallel.Unlifted.Distributed.Data.Unit+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.DT+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Base+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd.Split+ Data.Array.Parallel.Unlifted.Distributed.Data.USegd+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.DT+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Base+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd.Split+ Data.Array.Parallel.Unlifted.Distributed.Data.USSegd+ Data.Array.Parallel.Unlifted.Distributed.Data.UVSegd+ Data.Array.Parallel.Unlifted.Distributed.Data.Vector+ Data.Array.Parallel.Unlifted.Distributed.Arrays- Data.Array.Parallel.Unlifted.Distributed.USegd- Data.Array.Parallel.Unlifted.Distributed.USSegd Data.Array.Parallel.Unlifted.Distributed.Basics- Data.Array.Parallel.Unlifted.Distributed.Types.USegd- Data.Array.Parallel.Unlifted.Distributed.Types.USSegd- Data.Array.Parallel.Unlifted.Distributed.Types.UVSegd- Data.Array.Parallel.Unlifted.Distributed.Types.Vector- Data.Array.Parallel.Unlifted.Distributed.Types.Maybe- Data.Array.Parallel.Unlifted.Distributed.Types.Tuple- Data.Array.Parallel.Unlifted.Distributed.Types.Prim- Data.Array.Parallel.Unlifted.Distributed.Types.Unit- Data.Array.Parallel.Unlifted.Distributed.Types.Base+ Data.Array.Parallel.Unlifted.Distributed.Combinators+ Data.Array.Parallel.Unlifted.Distributed.What+ Data.Array.Parallel.Unlifted.Distributed+ Data.Array.Parallel.Unlifted.Parallel.Basics Data.Array.Parallel.Unlifted.Parallel.Combinators Data.Array.Parallel.Unlifted.Parallel.Enum@@ -52,7 +55,14 @@ Data.Array.Parallel.Unlifted.Parallel.Subarrays Data.Array.Parallel.Unlifted.Parallel.Sums Data.Array.Parallel.Unlifted.Parallel.Text+ Data.Array.Parallel.Unlifted.Parallel.UPSegd+ Data.Array.Parallel.Unlifted.Parallel.UPSSegd+ Data.Array.Parallel.Unlifted.Parallel.UPVSegd+ Data.Array.Parallel.Unlifted.Parallel.UPSel+ Data.Array.Parallel.Unlifted.Parallel + Data.Array.Parallel.Unlifted+ Exposed: False Extensions: TypeFamilies, GADTs, RankNTypes,@@ -65,10 +75,10 @@ -fcpr-off -Wall Build-Depends: - base == 4.5.*,+ base == 4.6.*, random == 1.0.*,- vector == 0.9.*,+ vector == 0.10.*, old-time == 1.1.*,- dph-base == 0.6.1.*,- dph-prim-interface == 0.6.1.*,- dph-prim-seq == 0.6.1.*+ dph-base == 0.7.*,+ dph-prim-interface == 0.7.*,+ dph-prim-seq == 0.7.*