lvish 1.1.1.3 → 1.1.2
raw patch · 23 files changed
+563/−1806 lines, 23 filesdep ~containers
Dependency ranges changed: containers
Files
- Control/LVish/BulkRetry.hs +0/−163
- Control/LVish/DeepFrz/Internal.hs +1/−1
- Control/LVish/Logging.hs +107/−105
- Control/LVish/SchedIdempotent.hs +173/−140
- Control/LVish/SchedIdempotentInternal.hs +45/−50
- Control/LVish/Types.hs +1/−0
- Data/LVar/AddRemoveSet.hs +0/−91
- Data/LVar/CycGraph.hs +0/−576
- Data/LVar/MaxCounter.hs +0/−79
- Data/LVar/Memo.hs +0/−110
- Data/LVar/NatArray.hs +0/−284
- Data/LVar/NatArray/Unsafe.hs +0/−27
- Data/LVar/PNCounter.hs +0/−73
- Data/LVar/PureMap.hs +1/−1
- Data/LVar/PureSet.hs +3/−2
- lvish.cabal +42/−42
- tests/AddRemoveSetTests.hs +114/−0
- tests/ArrayTests.hs +12/−8
- tests/LVishAndIVar.hs +9/−6
- tests/Main.hs +2/−2
- tests/MaxCounterTests.hs +0/−45
- tests/MaxPosIntTests.hs +45/−0
- tests/TestHelpers.hs +8/−1
− Control/LVish/BulkRetry.hs
@@ -1,163 +0,0 @@-{-# LANGUAGE ScopedTypeVariables, BangPatterns #-}-{-# LANGUAGE DataKinds #-}---- | EXPERIMENTAL version which eventually should be made generic across Par monads--- (i.e. a BulkRetryT transformer), and should thus be extended to transparently--- catch any attempts by a thread to block, not just the special non-blocking calls--- provided by *this* library.--module Control.LVish.BulkRetry-- where--import qualified Data.Bits.Atomic as B-import Foreign.Storable (sizeOf, Storable)-import Control.Monad (unless) -import Control.LVish-import Control.LVish.Internal (unsafeDet)-import Control.Par.Class (LVarSched(returnToSched))--- import Data.LVar.NatArray-import Data.LVar.NatArray.Unsafe (NatArray, unsafePeek)--import Data.Par.Splittable (pforEach)-import Data.Par.Range (range)-import Data.Par.Set () -- Instances only.--import qualified Data.Foldable as F-import qualified Data.Set as S--- import Data.LVar.PureSet as IS-import Data.LVar.SLSet as IS-import Data.LVar.Generic (freeze)---- import Data.Par.Range-------------------------------------------------------------------------------------- | The point where users send abort messages.-data RetryHub s = RetryHub (ISet s Int) -- ^ This stores the iterations that fail.- Int -- ^ This is the current iteration---- -- | Non-blocking get on a `NatArray`.--- getNB :: forall s d elt . (Storable elt, B.AtomicBits elt, Num elt) =>--- RetryHub s -> NatArray s elt -> Int -> Par d s elt--- -- LVarSched (Par d s) --- getNB (RetryHub fails) arr ind = do--- x <- unsafePeek arr ind--- -- if empty, don't block, do this:--- case x of--- Nothing -> do logDbgLn 4 $ " [dbg-lvish] getNB: iteration failed, enqueue for retry: "++show ind--- insert ind fails--- returnToSched--- Just res -> return res----- | Non-blocking get on a `NatArray`. In this prototype we require that the user--- manually CPS the computation, so that the delimited continuation between this get--- and the end of the loop iteration is passed explicitly as an argument.------ The current reason for this compromise is that the HandlerPool mechanism is not--- robust to us dropping the current continuation with `returnToSched`. We would--- need a version of HandlerPool's that interoperates with a user-level callCC, that is--- we would need something like bracket/dynamic-wind for our continuation monad.-getNB_cps :: forall s d elt . (Storable elt, B.AtomicBits elt, Num elt) =>- RetryHub s- -> NatArray s elt -- ^ Array to dereference- -> Int -- ^ Which index to get- -> (elt -> Par d s ()) -- ^ Delimited continuation.- -> Par d s ()--- LVarSched (Par d s) -getNB_cps (RetryHub fails thisiter) arr ind cont = do- x <- unsafePeek arr ind- -- if empty, don't block, do this:- case x of- Nothing -> do logDbgLn 4 $ " [dbg-lvish] getNB: iteration "++ show thisiter- ++" failed, due to get on index "++show ind- insert thisiter fails- return ()- Just res -> do logDbgLn 4 $ " [dbg-lvish] getNB: result available, calling continuation (iter "++show thisiter++")"- cont res-{-# INLINE getNB_cps #-}--desired_tasks :: Int-desired_tasks = 16 -- FIXME: num procs * overpartition---- | A parallel for-loop which aborts and retries failed iterations in bulk, rather--- than allowing them to "block" and suffering the overhead of capturing and storing--- their continuations.--- --- `forSpeculative` continues retrying until ALL iterations have completed. It is--- thus a *synchronous* parallel for loop.-forSpeculative :: (Int, Int) -- ^ Inclusive/Exclusive range to run.- -> (RetryHub s -> Int -> Par d s ()) -- ^ Body of the loop- -> Par d s ()--- forSpeculative :: (Int, Int) -> (RetryHub s -> Int -> Par QuasiDet s ()) -> Par QuasiDet s ()--- TODO: Requires idempotency!!-forSpeculative (st,end) bodyfn = do- logDbgLn 2 $ " [dbg-lvish] Begin forSpeculative, bounds "++show (st,end)- let sz = end - st- -- Even in a trivial loop, 2000 iters per task should be enough:- prefix = min sz (2000 * desired_tasks)- -- TODO: automatic strategies for tuning the input prefix size would be helpful.- -- One approach that might make sense would be to auto-tune based on the- -- time/iteration observed. That is, gradually increase to try to approximate a- -- minimum reasonable task size and no bigger. -- body' = bodyfn- -- body' retry ix = bodyfn retry ix- - let flush leftover fails = - -- unless (S.null leftover) $ do- -- TODO: need parallel fold, this is sequential...- F.foldlM (\ () ix -> do- logDbgLn 3 $ " [dbg-lvish] forSpeculative: flushing iter "++show ix- body' (RetryHub fails ix) ix)- () leftover- let flushLoop leftover = do- fails <- newEmptySet- -- FIXME: Add parallelism- flush leftover fails -- Sequential... - snap <- unsafeDet $ freeze fails- logDbgLn 3 $ " [dbg-lvish] forSpeculative: did one sequential flush, remaining: "++show snap- unless (S.null snap) $- -- error$ "forSpeculative: failures not flushed with a sequential run!:\n "++show snap- flushLoop snap- - -- Outer loop of "rounds", in which we try a prefix of the iteration space. - let loop !round leftover offset 0 = do- logDbgLn 3 $ " [dbg-lvish] forSpeculative: got to the end, only failures left."- flushLoop leftover- - loop !round leftover offset remain = do- logDbgLn 3 $ " [dbg-lvish] forSpeculative starting round "++- show round++": offset "++show offset++", remaining "++show remain- -- Set of iterations that failed in THIS upcoming round:- fails <- newEmptySet- let chunkend = offset + (min prefix remain) -- hp <- newPool- -- Here we keep the failed iterations "to the left" of the new batch, i.e. we- -- fork them first.- - -- FINISHME: need Split instance.- logDbgLn 4 $ " [dbg-lvish] forSpeculative RElaunching failures: "++show leftover- -- This version is poor because it forks on a per-iteration basis upon retry:- -- F.foldrM (\ ix () -> forkHP (Just hp) (body' (RetryHub fails ix) ix)) () leftover- -- F.foldrM (\ ix () -> body' (RetryHub fails ix) ix) () leftover- -- pforEach leftover $ bodyfn (RetryHub fails)- asyncForEachHP (Just hp) leftover $ \ ix -> bodyfn (RetryHub fails ix) ix- - -- TODO: if we keep failing it's better to expand the prefix. That way we- -- end up with a logarithmic number of retries for each iterate in the worst- -- case, rather than linear (making the whole loop unnecessarily quadratic).-- logDbgLn 4 $ " [dbg-lvish] forSpeculative launching new batch: "++show (offset,chunkend)- asyncForEachHP (Just hp) (range offset chunkend) $ \ ix -> - body' (RetryHub fails ix) ix- logDbgLn 4 $ " [dbg-lvish] forSpeculative: return from par for-loop; now quiesce."- quiesce hp- logDbgLn 4 $ " [dbg-lvish] forSpeculative: quiesce finished, next freeze failed set."- snap <- unsafeDet $ freeze fails- logDbgLn 4 $ " [dbg-lvish] forSpeculative finish round; failed iterates: "++show snap- loop (round+1) snap chunkend (remain - (chunkend - offset))- loop 0 S.empty 0 sz - -- After the last quiesce, we're done.
Control/LVish/DeepFrz/Internal.hs view
@@ -17,7 +17,7 @@ -- by the user, however. Rather, it is the final step in a -- `runParThenFreeze` invocation. --- An instance of DeepFrz is a valid return valud for `runParThenFreeze`+-- An instance of DeepFrz is a valid return value for `runParThenFreeze` class DeepFrz a where -- | This type function is public. It maps pre-frozen types to -- frozen ones. It should be idempotent.
Control/LVish/Logging.hs view
@@ -28,10 +28,12 @@ -- * New logger interface newLogger, logOn, Logger(closeIt, flushLogs),- WaitMode(..), LogMsg(..), OutDest(..),+ WaitMode(..), LogMsg(..), mapMsg, OutDest(..), -- * General utilities- forkWithExceptions+ forkWithExceptions,++ Backoff(totalWait), newBackoff, backoff ) where @@ -87,16 +89,14 @@ -- | Several different ways we know to wait for quiescence in the concurrent mutator -- before proceeding.-data WaitMode = WaitTids [ThreadId] (IO Bool)- -- ^ Wait until a certain set of threads is blocked before proceeding.- -- If that conditional holds ALSO make sure the provided polling action- -- returns True as well.- | WaitDynamic -- ^ UNFINISHED: Dynamically track tasks/workers. The+data WaitMode = WaitDynamic -- ^ UNFINISHED: Dynamically track tasks/workers. The -- num workers starts at 1 and then is modified -- with `incrTasks` and `decrTasks`. | WaitNum { numThreads :: Int, -- ^ How many threads total must check in? downThreads :: IO Int -- ^ Poll how many threads won't participate this round.+ -- After all productive threads have checked in + -- this number must grow to eventually include all other threads. } -- ^ A fixed set of threads must check-in each round before proceeding. | DontWait -- ^ In this mode, logging calls are non-blocking and return -- immediately, rather than waiting on a central coordinator.@@ -114,10 +114,21 @@ -- distinction is not that important, because only *thunks* should be logged; the -- thread printing the logs should deal with forcing those thunks. data LogMsg = StrMsg { lvl::Int, body::String }+ | OffTheRecord { lvl :: Int, obod :: String }+ -- ^ This sort of message is chatter and NOT meant + -- to participate in the scheduler-testing framework. -- | ByteStrMsg { lvl::Int, }+ deriving (Show,Eq,Ord,Read) -toString x@(StrMsg{}) = body x+mapMsg :: (String -> String) -> LogMsg -> LogMsg+mapMsg f (StrMsg l s) = StrMsg l (f s)+mapMsg f (OffTheRecord l s) = OffTheRecord l (f s) +toString :: LogMsg -> String+toString x = case x of + StrMsg {body} -> body+ OffTheRecord _ s -> s+ maxWait :: Int maxWait = 10*1000 -- 10ms @@ -152,66 +163,54 @@ parent <- myThreadId let flushLogs = atomicModifyIORef' logged $ \ ls -> ([],reverse ls) - let -- When all threads are quiescent, we can flush the remaining messagers from- -- the channel to get the whole set of waiting tasks. Return in chronological order.- flushChan !acc = do- x <- tryReadSmplChan checkPoint- case x of- Just h -> flushChan (h:acc)- Nothing -> return $ reverse acc- - -- This is the format we use for debugging messages- formatMessage extra Writer{msg} = "|"++show (lvl msg)++ "| "++extra++ toString msg- -- One of these message reports how many tasks are in parallel with it:- messageInContext pos len wr = formatMessage ("#"++show (1+pos)++" of "++show len ++": ") wr- printOne str (OutputTo h) = hPrintf h "%s\n" str- printOne str OutputEvents = traceEventIO str- printOne str OutputInMemory =- -- This needs to be atomic because other messages might be calling "flush"- -- at the same time.- atomicModifyIORef' logged $ \ ls -> (str:ls,())- printAll str = mapM_ (printOne str) loutDests-- shutdownFlag <- newIORef False -- When true, time to shutdown.- shutdownComplete <- newEmptyMVar+ shutdownFlag <- newIORef False -- When true, time to start shutdown. -- Here's the new thread that corresponds to this logger:- coordinator <- A.async $ E.handle (catchAll parent) $- -- BEGIN defs for the async task:- --------------------------------------------------------------------------------- -- Proceed in rounds, gather the set of actions that may happen in parallel, then- -- pick one. We log the series of decisions we make for reproducability.- let schedloop :: Int -> Int -- ^ length of list `waiting`- -> [Writer] -> Backoff -> IO ()- schedloop !iters !num !waiting !bkoff = do+ coordinator <- A.async $ E.handle (catchAll parent) $ do+ runCoordinator waitWorkers shutdownFlag checkPoint logged loutDests+ let closeIt = do+ atomicModifyIORef' shutdownFlag (\_ -> (True,())) -- Declare that it's time to shutdown:+ A.wait coordinator -- Gently wait for it to be done.+ return $! Logger { coordinator, checkPoint, closeIt, loutDests,+ logged, flushLogs,+ waitWorkers, minLvl, maxLvl }++--------------------------------------------------------------------------------++-- | Run a logging coordinator thread until completion/shutdown.+runCoordinator :: WaitMode -> IORef Bool -> IORef (Seq.Seq Writer) -> IORef [String] -> [OutDest] -> IO ()+runCoordinator waitWorkers shutdownFlag checkPoint logged loutDests = + case waitWorkers of+ DontWait -> printLoop =<< newBackoff maxWait+ _ -> schedloop (0::Int) [] =<< newBackoff maxWait -- Kick things off.+ where+ -- Proceed in rounds, gather the set of actions that may happen in parallel, then+ -- pick one. We log the series of decisions we make for reproducability.+ schedloop :: Int + -> [Writer] -- ^ Waiting threads, reverse chronological (newest first)+ -> Backoff -> IO ()+ schedloop !iters !waiting !bkoff = do when (iters > 0 && iters `mod` 500 == 0) $- putStrLn $ "Warning: logger has spun for "++show iters++" iterations, "++show num++" are waiting."+ putStrLn $ "Warning: logger has spun for "++show iters++" iterations, "++show (length waiting)++" are waiting." hFlush stdout fl <- readIORef shutdownFlag if fl then flushLoop else do - let keepWaiting = do b <- backoff bkoff- schedloop (iters+1) num waiting b- waitMore = do w <- readSmplChan checkPoint -- Blocking! (or spinning)- b <- newBackoff maxWait -- We got something, reset this.- schedloop (iters+1) (num+1) (w:waiting) b+ let keepWaiting w = do b <- backoff bkoff+ schedloop (iters+1) w b case waitWorkers of DontWait -> error "newLogger: internal invariant broken." WaitNum target extra -> do+ waiting2 <- flushChan waiting+ let numWait = length waiting2 n <- extra -- Atomically check how many extra workers are blocked.- if (num + n >= target)- then pickAndProceed waiting- else waitMore- WaitTids tids poll -> do- -- FIXME: This is not watertight... it will work with high probability but can't be trusted:- andM [checkTids tids, poll, checkTids tids, poll]- (do ls <- flushChan waiting- case ls of- [] -> do chatter " [Logger] Warning: No active tasks?"- bk2 <- backoff bkoff- schedloop (iters+1) 0 [] bk2- _ -> pickAndProceed ls)- keepWaiting+ -- putStrLn $ "TEMP: schedloop/WaitNum: polled for waiting/extra workers: "+ -- ++show (numWait,n)++" target "++show target+ if (numWait + n >= target)+ then if numWait > 0 + then pickAndProceed waiting2+ else keepWaiting waiting2 -- This sounds like a shutdown is happening, all are idle.+ else keepWaiting waiting2 -- We don't know if we're waiting for idles to arrive or blocked waiters. -- | Keep printing messages until there is (transiently) nothing left. flushLoop = do @@ -221,13 +220,24 @@ flushLoop Nothing -> return () + flushChan !acc = do+ x <- tryReadSmplChan checkPoint+ case x of+ Just h -> case msg h of + StrMsg {} -> flushChan (h:acc)+ OffTheRecord {} -> do printAll (formatMessage "" h) + flushChan acc+ Nothing -> return acc+ -- | A simpler alternative schedloop that only does printing (e.g. for DontWait mode).- printLoop = do+ printLoop bk = do fl <- readIORef shutdownFlag if fl then flushLoop- else do wr <- readSmplChan checkPoint- printAll (formatMessage "" wr)- printLoop+ else do mwr <- tryReadSmplChan checkPoint+ case mwr of + Nothing -> do printLoop =<< backoff bk + Just wr -> do printAll (formatMessage "" wr)+ printLoop =<< newBackoff (cap bk) -- Take the set of logically-in-parallel tasks, choose one, execute it, and -- then return to the main scheduler loop.@@ -247,45 +257,32 @@ let pick = sorted !! pos (pref,suf) = splitAt pos sorted rst = pref ++ tail suf+ -- putStrLn$ "TEMP: pickAndProceed, unblocking "++show (pos,len,msg pick) unblockTask pos len pick -- The task will asynchronously run when it can. yield -- If running on one thread, give it a chance to run. -- Return to the scheduler to wait for the next quiescent point: bnew <- newBackoff maxWait- schedloop 0 (length rst) rst bnew+ schedloop 0 rst bnew unblockTask pos len wr@Writer{continue} = do printAll (messageInContext pos len wr) putMVar continue () -- Signal that the thread may continue. - -- Check whether the worker threads are all quiesced - checkTids [] = return True- checkTids (tid:rst) = do - st <- threadStatus tid- case st of- ThreadRunning -> return False- ThreadFinished -> checkTids rst- -- WARNING: this design is flawed because it is possible when compiled- -- with -threaded that IO will spuriously showed up as BlockedOnMVar:- ThreadBlocked BlockedOnMVar -> checkTids rst- ThreadBlocked _ -> return False- ThreadDied -> checkTids rst -- Should this be an error condition!?- in -- Main body of async task:- do case waitWorkers of- DontWait -> printLoop - _ -> schedloop (0::Int) (0::Int) [] =<< newBackoff maxWait -- Kick things off.- putMVar shutdownComplete ()- return () -- End: async thread- -- END async task.- --------------------------------------------------------------------------------+ -- This is the format we use for debugging messages+ formatMessage extra Writer{msg} = "|"++show (lvl msg)++ "| "++extra++ toString msg+ -- One of these message reports how many tasks are in parallel with it:+ messageInContext pos len wr = formatMessage ("#"++show (1+pos)++" of "++show len ++": ") wr+ printOne str (OutputTo h) = hPrintf h "%s\n" str+ printOne str OutputEvents = traceEventIO str+ printOne str OutputInMemory =+ -- This needs to be atomic because other messages might be calling "flush"+ -- at the same time.+ atomicModifyIORef' logged $ \ ls -> (str:ls,())+ printAll str = mapM_ (printOne str) loutDests - let closeIt = do- atomicModifyIORef' shutdownFlag (\_ -> (True,()))- readMVar shutdownComplete- A.cancel coordinator -- Just to make sure its completely done.- return $! Logger { coordinator, checkPoint, closeIt, loutDests,- logged, flushLogs,- waitWorkers, minLvl, maxLvl } ++ chatter :: String -> IO () -- chatter = hPrintf stderr -- chatter = printf "%s\n"@@ -301,15 +298,17 @@ -- message falls into the range accepted by the given `Logger`, -- otherwise, the message is ignored. logOn :: Logger -> LogMsg -> IO ()-logOn Logger{checkPoint,minLvl,maxLvl,waitWorkers} msg- | (minLvl <= lvl msg) && (lvl msg <= maxLvl) = do +logOn Logger{checkPoint,minLvl,maxLvl,waitWorkers} msg = do + + if (minLvl <= lvl msg) && (lvl msg <= maxLvl) then do + -- putStrLn$ "TEMP: "++show (minLvl,maxLvl)++" attempt to log msg: "++show msg case waitWorkers of -- In this mode we are non-blocking: DontWait -> writeSmplChan checkPoint Writer{who="",continue=dummyMVar,msg} _ -> do continue <- newEmptyMVar writeSmplChan checkPoint Writer{who="",continue,msg} takeMVar continue -- Block until we're given permission to proceed.- | otherwise = return ()+ else return () {-# NOINLINE dummyMVar #-} dummyMVar :: MVar ()@@ -321,24 +320,27 @@ -- | The state for an exponential backoff. data Backoff = Backoff { current :: !Int , cap :: !Int -- ^ Maximum nanoseconds to wait.+ , totalWait :: !Int } deriving Show --newBackoff :: Int -> IO Backoff-newBackoff cap = return Backoff{cap,current=0}+-- | Create an object used for exponentential backoff; see `backoff`.+newBackoff :: Int -- ^ Maximum delay, nanoseconds+ -> IO Backoff+newBackoff cap = return Backoff{cap,current=0,totalWait=0} +-- | Perform the backoff, possibly delaying the thread. backoff :: Backoff -> IO Backoff--- backoff b = do yield; return b-backoff Backoff{current,cap} = - case current of- -- Yield once before we start delaying:- 0 -> do yield- return Backoff{cap,current=1}- n -> do let next = min cap (2*n)- threadDelay n- return Backoff{cap,current=next}- +backoff Backoff{current,cap,totalWait} = do+ if current < 1 then + -- Yield before we start delaying:+ do yield+ return Backoff{cap,current=current+1,totalWait}+ else+ do let nxt = min cap (2*current)+ threadDelay current+ return Backoff{cap,current=nxt,totalWait=totalWait+current}+ ---------------------------------------------------------------------------------------------------- -- Simple channels: we need non-blocking reads so we can't use -- Control.Concurrent.Chan. We could use TChan, but I don't want to bring STM into
Control/LVish/SchedIdempotent.hs view
@@ -49,6 +49,7 @@ import Control.Concurrent hiding (yield) import qualified Control.Concurrent as Conc import qualified Control.Exception as E+import qualified Control.Concurrent.Async as A import Control.DeepSeq import Control.Applicative import Control.LVish.MonadToss@@ -57,9 +58,11 @@ import Data.IORef import Data.Atomics import Data.Typeable+import qualified Data.Atomics.Counter as C2 import qualified Data.Concurrent.Counter as C import qualified Data.Concurrent.Bag as B import GHC.Conc hiding (yield)+import qualified GHC.Conc import System.IO import System.IO.Unsafe (unsafePerformIO) import System.Environment(getEnvironment)@@ -67,13 +70,10 @@ import Prelude hiding (mapM, sequence, head, tail) import qualified Prelude import System.Random (random)--#ifdef DEBUG_LVAR -import Text.Printf (printf)-#endif+import Text.Printf (printf, hPrintf) -- import Control.Compose ((:.), unO)-import Data.Traversable +import Data.Traversable hiding (forM) import Control.LVish.Types import qualified Control.LVish.SchedIdempotentInternal as Sched@@ -198,19 +198,29 @@ logStrLn lvl str = when (dbgLvl >= 1) $ do lgr <- getLogger num <- getWorkerNum- liftIO$ L.logOn lgr (L.StrMsg lvl ("(wrkr"++show num ++") "++ str))+ if lvl < 0+ then liftIO$ logHelper (Just lgr) num (L.OffTheRecord (-lvl) str)+ else liftIO$ logHelper (Just lgr) num (L.StrMsg lvl str) #else logStrLn _ _ = return () #endif -logWith :: Sched.State a s -> Int -> String -> IO ()+logHelper :: Maybe Logger -> Int -> LogMsg -> IO ()+logHelper lgr num msg = when (dbgLvl >= 1) $ do+ let msg' = L.mapMsg (("wrkr"++show num++" ")++) msg+ case lgr of + Just lgr -> L.logOn lgr msg'+ Nothing -> hPutStrLn stderr ("WARNING/nologger:"++show msg')++logWith :: Sched.State a s -> Int -> String -> IO ()+logOffRecord :: Sched.State a s -> Int -> String -> IO () #ifdef DEBUG_LVAR -- Only when the debug level is 1 or higher is the logger even initialized:-logWith q lvl str = when (dbgLvl >= 1) $ do- Just lgr <- readIORef (Sched.logger q)- L.logOn lgr (L.StrMsg lvl str)+logWith q lvl str = logHelper (Sched.logger q) (Sched.no q) (L.StrMsg lvl str)+logOffRecord q lvl str = logHelper (Sched.logger q) (Sched.no q) (L.OffTheRecord lvl str) #else logWith _ _ _ = return ()+logOffRecord _ _ _ = return () #endif ------------------------------------------------------------------------------@@ -219,7 +229,8 @@ -- | Create an LVar. newLV :: IO a -> Par (LVar a d)-newLV init = mkPar $ \k q -> do+newLV init = mkPar $ \k q -> do + logOffRecord q 7$ " [dbg-lvish] newLV: allocating... " state <- init listeners <- B.new status <- newIORef $ Active listeners@@ -251,31 +262,21 @@ -- continuation immediately Nothing -> do -- /transiently/ not past the threshhold; block - -#if GET_ONCE- execFlag <- newIORef False-#endif- ++ execFlag <- newDedupCheck let onUpdate d = unblockWhen $ deltaThresh d onFreeze = unblockWhen $ globalThresh state True- + {-# INLINE unblockWhen #-} unblockWhen thresh tok q = do let uniqsuf = ", lv "++(show$ unsafeName state)++" on worker "++(show$ Sched.no q) logWith q 7$ " [dbg-lvish] getLV (active): callback: check thresh"++uniqsuf tripped <- thresh whenJust tripped $ \b -> do B.remove tok-#if GET_ONCE- logWith q 8$ " [dbg-lvish] getLV (active): read execFlag for dedup"++uniqsuf- ticket <- readForCAS execFlag- unless (peekTicket ticket) $ do- (winner, _) <- do logWith q 8$ " [dbg-lvish] getLV (active): CAS execFlag dedup"++uniqsuf- casIORef execFlag ticket True- when winner $ Sched.pushWork q (k b) -#else - Sched.pushWork q (k b) -#endif- logWith q 4$ " [dbg-lvish] getLV: blocking on LVar, registering listeners"++uniqsuf+ winnerCheck execFlag q (Sched.pushWork q (k b)) (return ())++ logWith q 8$ " [dbg-lvish] getLV "++show(unsafeName execFlag)+++ ": blocking on LVar, registering listeners..." -- add listener, i.e., move the continuation to the waiting bag tok <- B.put listeners $ Listener onUpdate onFreeze @@ -290,10 +291,13 @@ Just b -> do logWith q 7$ " [dbg-lvish] getLV (active): second globalThresh tripped, remove tok"++uniqsuf B.remove tok -- remove the listener we just added, and- exec (k b) q -- execute the continuation. this work might be- -- redundant, but by idempotence that's OK++ winnerCheck execFlag q (exec (k b) q) (sched q)+ -- execute the continuation. this work might be+ -- redundant, but in idempotence-mode that's OK Nothing -> sched q + -------------------------------------------------------------------------------- -- Freezing or Frozen: _ -> do logWith q 7$ " [dbg-lvish] getLV (frozen): about to check globalThresh"++uniqsuf@@ -306,6 +310,52 @@ -- Shouldn't this be an ERROR? (blocked-indefinitely) -- Depends on our semantics for runPar quiescence / errors states. ++{-# INLINE newDedupCheck #-}+{-# INLINE winnerCheck #-}+winnerCheck :: DedupCell -> Sched.State a s -> IO () -> IO () -> IO ()+newDedupCheck :: IO DedupCell++#if GET_ONCE++# if 0+type DedupCell = IORef Bool+newDedupCheck = newIORef False -- True means someone has already won.+winnerCheck execFlag q tru fal = do + ticket <- readForCAS execFlag+ if (peekTicket ticket) + then do logWith q 8 $ " [dbg-lvish] getLV winnerCheck failed.."+ fal+ else do+ (winner, _) <- casIORef execFlag ticket True+ logWith q 8 $ " [dbg-lvish] getLV "++show(unsafeName execFlag)+ ++" on worker "++ (show$ Sched.no q) ++": winner check? " ++show winner+ ++ ", ticks " ++ show (ticket, peekTicket ticket)+ if winner then tru else fal+# else++type DedupCell = C2.AtomicCounter+newDedupCheck = C2.newCounter 0 +winnerCheck execFlag q tru fal = do+ cnt <- C2.incrCounter 1 execFlag+ logWith q 8 $ " [dbg-lvish] getLV "++show(unsafeName execFlag)+ ++" on worker "++ (show$ Sched.no q) ++": winner check? " ++show (cnt==1)+ ++ ", counter val " ++ show cnt+ if cnt==1 then tru else fal++# endif+#else+type DedupCell = ()+newDedupCheck = return ()+winnerCheck _ _ tr _ = tr+#endif+++++++ -- | Update an LVar. putLV_ :: LVar a d -- ^ the LVar -> (a -> Par (Maybe d, b)) -- ^ how to do the put, and whether the LVar's@@ -449,10 +499,6 @@ onUpdate d _ q = spawnWhen (updateThresh d) q onFreeze _ _ = return () - runWhen thresh q = do- tripped <- thresh- whenJust tripped $ \cb -> - exec (close cb nullCont) q in mkPar $ \k q -> do curStatus <- readIORef status case curStatus of@@ -464,10 +510,7 @@ logWith q 4 " [dbg-lvish] addHandler: calling globalCB.." -- At registration time, traverse (globally) over the previously inserted items -- to launch any required callbacks.- exec (close (globalCB state) nullCont) q- exec (k ()) q --nullCont = (\() -> ClosedPar (\_ -> return ()))+ exec (close (globalCB state) k) q -- | Block until a handler pool is quiescent. quiesce :: HandlerPool -> Par ()@@ -476,20 +519,22 @@ -- tradeoff: we assume that the pool is not yet quiescent, and thus enroll as -- a blocked thread prior to checking for quiescence tok <- B.put bag (k ())+ hpMsg q " [dbg-lvish] quiesce: poll count" hp quiescent <- C.poll cnt if quiescent then do+ hpMsg q " [dbg-lvish] already quiesced, remove token from bag" hp B.remove tok- hpMsg q " [dbg-lvish] -> Quiescent already!" hp exec (k ()) q else do - hpMsg q " [dbg-lvish] -> Not quiescent yet, back to sched" hp+ logOffRecord q 4 " [dbg-lvish] -> Not quiescent yet, back to sched" sched q -- | A global barrier. quiesceAll :: Par () quiesceAll = mkPar $ \k q -> do+ logWith q 1 " [dbg-lvish] quiesceAll: initiating global barrier." sched q- logWith q 1 " [dbg-lvish] Return from global barrier."+ logWith q 1 " [dbg-lvish] quiesceAll: Past global barrier." exec (k ()) q -- | Freeze an LVar after a given handler quiesces.@@ -534,8 +579,8 @@ -- | IF compiled with debugging support, this will return the Logger used by the -- current Par session, otherwise it will simply throw an exception. getLogger :: Par L.Logger-getLogger = mkPar $ \k q -> do- Just lgr <- readIORef (Sched.logger q)+getLogger = mkPar $ \k q -> + let Just lgr = Sched.logger q in exec (k lgr) q -- | Return the worker that we happen to be running on. (NONDETERMINISTIC.)@@ -584,106 +629,84 @@ -> Int -- ^ How many worker threads to use. -> Par a -- ^ The computation to run. -> IO ([String], Either E.SomeException a)-runParDetailed DbgCfg {dbgRange, dbgDests, dbgScheduling } numWrkrs comp = do- queues <- Sched.new numWrkrs noName- +runParDetailed cfg@DbgCfg{dbgRange, dbgDests, dbgScheduling } numWrkrs comp = do+ (lgr,queues) <- Sched.new cfg numWrkrs noName + -- We create a thread on each CPU with forkOn. The CPU on which -- the current thread is running will host the main thread; the -- other CPUs will host worker threads. main_cpu <- Sched.currentCPU answerMV <- newEmptyMVar- wrkrtids <- newIORef [] - -- Debugging: spin the main thread (not beginning work) until we can fully- -- initialize the logging data structure.- --- -- TODO: This would be easier to deal with if we used the current thread directly- -- as the main worker thread...- let setLogger = do- ls <- readIORef wrkrtids- if length ls == numWrkrs- then Sched.initLogger queues ls (minLvl,maxLvl) dbgDests dbgScheduling- else do Conc.yield- setLogger- (minLvl, maxLvl) = case dbgRange of- Just b -> b- Nothing -> (0,dbgLvl)- -- Option 1: forkWithExceptions version:- ---------------------------------------------------------------------------------- -#if 1- let forkit = forM_ (zip [0..] queues) $ \(cpu, q) -> do - tid <- L.forkWithExceptions (forkOn cpu) "worker thread" $ do- if cpu == main_cpu - then let k x = ClosedPar $ \q -> do - sched q -- ensure any remaining, enabled threads run to - putMVar answerMV x -- completion prior to returning the result- -- [TODO: ^ perhaps better to use a binary notification tree to signal the workers to stop...]- in do -#ifdef DEBUG_LVAR- -- This is painful, we may need to spin and wait for everybody to be forked:- when (maxLvl >= 1) setLogger-#endif- exec (close comp k) q- -- Note: The above is important: it is sketchy to leave any workers running after- -- the main thread exits. Subsequent exceptions on child threads, even if- -- forwarded asynchronously, can arrive much later at the main thread- -- (e.g. after it has exited, or set up a new handler, etc).- else sched q- atomicModifyIORef_ wrkrtids (tid:)- -- logWith (Prelude.head queues) " [dbg-lvish] About to fork workers..." - ans <- E.catch (forkit >> fmap Right (takeMVar answerMV))- (\ (e :: E.SomeException) -> do - tids <- readIORef wrkrtids- logWith (Prelude.head queues) 1 $ " [dbg-lvish] Killing off workers due to exception: "++show tids- mapM_ killThread tids- -- if length tids < length queues then do -- TODO: we could try to chase these down in the idle list.- mytid <- myThreadId- -- when (maxLvl >= 1) printLog -- Unfortunately this races with the log printing thread.- -- E.throw$ LVarSpecificExn ("EXCEPTION in runPar("++show mytid++"): "++show e)- return $! Left e- )- logWith (Prelude.head queues) 1 " [dbg-lvish] parent thread escaped unscathed"- mlgr <- readIORef (Sched.logger (Prelude.head queues))- logs <- case mlgr of - Nothing -> return []- Just lgr -> do L.closeIt lgr- L.flushLogs lgr -- If in-memory logging is off, this will be empty.- return $! (logs,ans)-#else--- Option 2: This was an experiment to use Control.Concurrent.Async to deal with exceptions:------------------------------------------------------------------------------------ let runWorker (cpu, q) = do + let grabLogs = do + logOffRecord (Prelude.head queues) 1 " [dbg-lvish] parent thread escaped unscathed. Optionally closing logger."+ case lgr of + Nothing -> return []+ Just lgr -> do L.closeIt lgr+ L.flushLogs lgr -- If in-memory logging is off, this will be empty.+ mlog s = case lgr of + Nothing -> return ()+ Just l -> L.logOn l (L.OffTheRecord 4 s)++ -- Use Control.Concurrent.Async to deal with exceptions:+ ----------------------------------------------------------------------------------+ let runWorker :: (Int,Sched.State ClosedPar LVarID) -> IO ()+ runWorker (cpu, q) = do if (cpu /= main_cpu)- then sched q- else let k x = ClosedPar $ \q -> do + then do logOffRecord q 3 $ " [dbg-lvish] Auxillary worker #"++show cpu++" starting."+ sched q+ logOffRecord q 3 $ " [dbg-lvish] Auxillary worker #"++show cpu++" exitting."+ else let k x = ClosedPar $ \q -> do + logOffRecord q 3 " [dbg-lvish] Final continuation of main worker: reenter sched to cleanup." sched q -- ensure any remaining, enabled threads run to - putMVar answerMV x -- completion prior to returning the result- in exec (close comp k) q+ -- completion prior to returning the result+ -- FIXME: this continuation gets duplicated.+ logOffRecord q 3 " [dbg-lvish] Main worker: past global barrier, putting answer."+ b <- tryPutMVar answerMV x+#ifdef GET_ONCE+ unless b $ error "Final continuation of Par computation was duplicated, in spite of GET_ONCE!"+#endif+ return ()+ in do logOffRecord q 3 " [dbg-lvish] Main worker thread starting."+ exec (close comp k) q -- Here we want a traditional, fork-join parallel loop with proper exception handling:- let loop [] asyncs = mapM_ wait asyncs+ let loop [] asyncs = do tid <- myThreadId+ mlog $ " [dbg-lvish] (tid "++show tid++") Wait on at least one async to complete.."+ (_,x) <- A.waitAnyCatch asyncs+ -- We could do a binary tree of waitBoth here, but this should work for now:+ case x of+ Left e -> return $! Left e+ Right () -> waitloop asyncs -- If one finishes, all are trying to. loop ((cpu,q):tl) asyncs = --- withAsync (runWorker state)- withAsyncOn cpu (runWorker (cpu,q))- (\a -> loop tl (a:asyncs))+ A.withAsyncOn cpu (runWorker (cpu,q))+ (\a -> loop tl (a:asyncs))+ waitloop [] = do + mlog " [dbg-lvish] All asyncs complete, read final answer MVar."+ fmap Right (dbgTakeMVar [] "runPar/final answer" answerMV)+-- fmap Right (takeMVar answerMV)+ waitloop (hd:tl) = do mlog " [dbg-lvish] Waiting for one async.."+ me <- A.waitCatch hd+ case me of + Left e -> return $! Left e+ Right () -> waitloop tl+ ----------------------------------------+ -- (1) There was a BUG in 'loop' at some point:+ -- "thread blocked indefinitely in an STM transaction"+ ans <- loop (zip [0..] queues) []+ ----------------------------------------+ -- (2) This has the same problem as 'loop':+ -- ls <- mapM (\ pr@(cpu,_) -> Async.asyncOn cpu (runWorker pr)) (zip [0..] queues)+ -- mapM_ wait ls+ ----------------------------------------+ -- (3) Using this FOR NOW, but it does NOT pin to the right processors:+ -- A.mapConcurrently runWorker (zip [0..] queues)+ ----------------------------------------+ logs <- grabLogs+ return $! (logs,ans) -------------------------------------------- (1) There is a BUG in 'loop' presently:--- "thread blocked indefinitely in an STM transaction"--- loop (zip [0..] queues) []-------------------------------------------- (2) This has the same problem as 'loop':--- ls <- mapM (\ pr@(cpu,_) -> Async.asyncOn cpu (runWorker pr)) (zip [0..] queues)--- mapM_ wait ls-------------------------------------------- (3) Using this FOR NOW, but it does NOT pin to the right processors:- mapConcurrently runWorker (zip [0..] queues)------------------------------------------ -- Now that child threads are done, it's safe for the main thread- -- to call it quits.- takeMVar answerMV -#endif + defaultRun :: Par b -> IO b defaultRun = fmap (fromRight . snd) . runParDetailed cfg numCapabilities@@ -753,21 +776,31 @@ -- | For debugging purposes. This can help us figure out (by an ugly -- process of elimination) which MVar reads are leading to a "Thread -- blocked indefinitely" exception.-{--busyTakeMVar :: String -> MVar a -> IO a-busyTakeMVar msg mv = try (10 * 1000 * 1000)+busyTakeMVar :: [ThreadId] -> String -> MVar a -> IO a+busyTakeMVar tids msg mv = + do b <- L.newBackoff maxWait+ try b where- try 0 = do- when dbg $ do+ maxWait = 10000 -- nanoseconds+ timeOut = (3 * 1000 * 1000) -- three seconds, only for debugging.+ try bkoff | totalWait bkoff >= timeOut = do+ error "OVER WAIT"+ -- when dbg $ do tid <- myThreadId -- After we've failed enough times, start complaining:- printf "%s not getting anywhere, msg: %s\n" (show tid) msg- try (100 * 1000)- try n = do+ hPrintf stderr "%s not unblocked yet, for: %s\n" (show tid) msg+ stats <- Prelude.mapM threadStatus tids + hPrintf stderr $ "Worker statuses: " ++ show (zip tids stats) ++"\n"+ try =<< L.backoff bkoff + try bkoff = do x <- tryTakeMVar mv case x of Just y -> return y- Nothing -> do yield; try (n-1)--}+ Nothing -> try =<< L.backoff bkoff +#ifdef DEBUG_LVAR+dbgTakeMVar = busyTakeMVar+#else+dbgTakeMVar _ _ = takeMVar+#endif
Control/LVish/SchedIdempotentInternal.hs view
@@ -4,7 +4,7 @@ {-# LANGUAGE RecursiveDo #-} module Control.LVish.SchedIdempotentInternal (- State(logger, no), initLogger,+ State(logger, no), new, number, next, pushWork, nullQ, yieldWork, currentCPU, setStatus, await, prng ) where @@ -21,6 +21,7 @@ import Text.Printf import qualified Control.LVish.Logging as L+import Control.LVish.Types (DbgCfg(..)) #ifdef CHASE_LEV #warning "Compiling with Chase-Lev work-stealing deque"@@ -90,10 +91,8 @@ workpool :: Deque a, -- ^ The thread-local work deque idle :: IORef [MVar Bool], -- ^ global list of idle workers states :: [State a s], -- ^ global list of all worker states.- logger :: IORef (Maybe L.Logger)- -- ^ The Logger object used by the current Par session. (This should not- -- change during runtime, it is mutable only to support deferred- -- initialization.)+ logger :: Maybe L.Logger+ -- ^ The Logger object used by the current Par session, if debugging is activated. } -- | Process the next item on the work queue or, failing that, go into@@ -116,8 +115,8 @@ -- This function does NOT return until the complete runPar session is complete (all -- workers idle). steal :: State a s -> IO (Maybe a)-steal State{ idle, states, no=my_no, numWorkers } = do- chatter $ printf "!cpu %d stealing\n" my_no+steal State{ idle, states, no=my_no, numWorkers, logger } = do+ chatter logger $ "!cpu "++show my_no++" stealing" go states where -- After a failed sweep, go idle:@@ -125,18 +124,18 @@ r <- atomicModifyIORef idle $ \is -> (m:is, is) if length r == numWorkers - 1 then do- chatter$ printf "!cpu %d initiating shutdown\n" my_no- mapM_ (\m -> putMVar m True) r+ chatter logger $ printf "!cpu %d initiating shutdown\n" my_no+ mapM_ (\m -> putMVar m True) r -- Signal to all but us. return Nothing else do- chatter $ printf "!cpu %d going idle...\n" my_no+ chatter logger $ printf "!cpu %d going idle...\n" my_no done <- takeMVar m if done then do- chatter $ printf "!cpu %d shutting down\n" my_no+ chatter logger $ printf "!cpu %d shutting down\n" my_no return Nothing else do- chatter $ printf "!cpu %d woken up\n" my_no+ chatter logger $ printf "!cpu %d woken up\n" my_no go states go (x:xs) | no x == my_no = go xs@@ -144,7 +143,7 @@ r <- popOther (workpool x) case r of Just t -> do- -- printf "cpu %d got work from cpu %d\n" my_no (no x)+ chatter logger $ printf "cpu %d got work from cpu %d\n" my_no (no x) return r Nothing -> go xs @@ -152,7 +151,8 @@ pushWork :: State a s -> a -> IO () -- TODO: If we're really going to do wakeup on every push we could consider giving -- the formerly-idle worker the work item directly and thus avoid touching the deque.-pushWork State { workpool, idle } t = do+pushWork State { workpool, idle, logger, no } t = do+ chatter logger $ "Starting pushWork on worker "++show no pushMine workpool t idles <- readIORef idle when (not (null idles)) $ do@@ -166,44 +166,33 @@ pushYield workpool t -- AJT: should this also wake an idle thread? -- | Create a new set of scheduler states.-new :: Int -> s -> IO [State a s]-new numWorkers s = do- idle <- newIORef []- logger <- newIORef Nothing+new :: DbgCfg -> Int -> s -> IO (Maybe L.Logger,[State a s])+new DbgCfg{dbgDests,dbgRange,dbgScheduling} numWorkers s = do+ idle <- newIORef [] -- Shared by all workers.+ let (minLvl, maxLvl) = case dbgRange of+ Just b -> b+ Nothing -> (0,L.dbgLvl)+ let mkLogger = do + lgr <- L.newLogger (minLvl,maxLvl) dbgDests+ (if dbgScheduling + then L.WaitNum numWorkers countIdle + else L.DontWait)+ -- L.logOn lgr (L.StrMsg 1 " [dbg-lvish] Initialized Logger... ")+ return lgr+ -- Atomically count how many workers are currently registered as idle:+ countIdle = do ls <- readIORef idle+ return $! length ls+ logger <- if L.dbgLvl > 0 + then fmap Just $ mkLogger+ else return Nothing let mkState states i = do workpool <- newDeque status <- newIORef s prng <- newIORef $ mkStdGen i return State { no = i, workpool, idle, status, states, prng, logger, numWorkers } rec states <- forM [0..(numWorkers-1)] $ mkState states- return states+ return (logger,states) --- | Takes a full set of worker states and correspoding threadIds and initializes the--- loggers.-initLogger :: [State a s] -> [ThreadId] -> (Int,Int) -> [L.OutDest] -> Bool -> IO ()-initLogger [] _ _ _ _ = error "initLogger: cannot take empty list of workers"-initLogger queues@(hd:_) tids bounds outDests debugScheduling- | len1 /= len2 = error "initLogger: length of arguments did not match"- | otherwise = do- lgr <- L.newLogger bounds outDests- (if debugScheduling then waitAll else L.DontWait)- -- lgr <- L.newLogger Nothing (L.WaitNum len1 countIdle)- L.logOn lgr (L.StrMsg 1 " [dbg-lvish] Initializing Logger... ")- -- Setting one of them sets all of them -- this field is shared:- writeIORef (logger hd) (Just lgr)- -- TODO: ASSERT that they are all actually the same IORef?- return ()- where- waitAll = (L.WaitTids tids (pollDeques queues))- - len1 = length queues- len2 = length tids- countIdle = do ls <- readIORef (idle hd)- return $! length ls- pollDeques [] = return True- pollDeques (h:t) = do b <- nullQ (workpool h)- if b then pollDeques t- else return False number :: State a s -> Int number State { no } = no@@ -217,8 +206,7 @@ let awaitOne state@(State { status, no=no2 }) = do cur <- readIORef status unless (p cur) $ do- mlgr <- readIORef logger- case mlgr of+ case logger of Nothing -> return () Just lgr -> L.logOn lgr (L.StrMsg 7 (" [dbg-lvish] busy-waiting on worker "++show no1++ ", for status to change on worker "++show no2))@@ -247,7 +235,14 @@ return 0 #endif --chatter :: String -> IO ()+-- | Local chatter function for this module+chatter :: Maybe L.Logger -> String -> IO () -- chatter s = putStrLn s-chatter _ = return ()+-- chatter _ s = printf "%s\n" s+-- chatter _ _ = return ()++-- We should NOT do this if dbgScheduling is on.+chatter mlg s = do + case mlg of + Nothing -> return ()+ Just lg -> L.logOn lg (L.OffTheRecord 7 s)
Control/LVish/Types.hs view
@@ -33,6 +33,7 @@ -- ^ Inclusive range of debug messages to accept -- (i.e. filter on priority level). If Nothing, use the default level, -- which is (0,N) where N is controlled by the DEBUG environment variable.+ -- The convention is to use Just (0,0) to disable logging. , dbgDests :: [OutDest] -- ^ Destinations for debug log messages. , dbgScheduling :: Bool -- ^ In additional to logging debug messages, control
− Data/LVar/AddRemoveSet.hs
@@ -1,91 +0,0 @@-{-# LANGUAGE BangPatterns #-}--{-|--This module provides sets that allow both addition and removal of-elements. This is possible because, under the hood, it's represented-with two monotonically growing sets, one for additions and one for-removals. It is inspired by /2P-Sets/ from the literature on-/conflict-free replicated data types/.-- -}-module Data.LVar.AddRemoveSet- (- AddRemoveSet,- newEmptySet, newSet, newFromList,- insert, waitAddedElem, waitAddedSize,- remove, waitRemovedElem, waitRemovedSize,-- freezeSet- - ) where-import qualified Data.Set as S-import Control.LVish-import Control.LVish.Internal-import qualified Data.LVar.PureSet as PS-import Control.Applicative---- | The set datatype.-data AddRemoveSet s a =- AddRemoveSet !(PS.ISet s a)- !(PS.ISet s a)---- | Create a new, empty `AddRemoveSet`.-newEmptySet :: Ord a => Par d s (AddRemoveSet s a)-newEmptySet = newSet S.empty---- | Create a new `AddRemoveSet` populated with initial elements.-newSet :: Ord a => S.Set a -> Par d s (AddRemoveSet s a)--- Here we're creating two new PureSets, one from the provided initial--- elements (the "add" set) and one empty (the "remove" set), and--- then, since both of those return `Par` computations, we're using--- our friends `<$>` and `<*>`.-newSet set = AddRemoveSet <$> (PS.newSet set) <*> PS.newEmptySet--- Alternate version that works if we import `Control.Monad`:--- newSet set = ap (fmap AddRemoveSet (PS.newSet set)) PS.newEmptySet- --- | A simple convenience function. Create a new 'ISet' drawing--- initial elements from an existing list.-newFromList :: Ord a => [a] -> Par d s (AddRemoveSet s a)-newFromList ls = newSet (S.fromList ls)---- | Put a single element in the set. (WHNF) Strict in the element--- being put in the set.-insert :: Ord a => a -> AddRemoveSet s a -> Par d s ()--- Because the two sets inside an AddRemoveSet are already PureSets,--- we really just have to call the provided `insert` method for--- PureSet. We don't need to call `putLV` or anything like that!-insert !elm (AddRemoveSet added removed) = PS.insert elm added---- | Wait for the set to contain a specified element.-waitAddedElem :: Ord a => a -> AddRemoveSet s a -> Par d s ()--- And similarly here, we don't have to call `getLV` ourselves.-waitAddedElem !elm (AddRemoveSet added removed) = PS.waitElem elm added---- | Wait on the size of the set of added elements.-waitAddedSize :: Int -> AddRemoveSet s a -> Par d s ()--- You get the idea...-waitAddedSize !sz (AddRemoveSet added removed) = PS.waitSize sz added---- | Remove a single element from the set.-remove :: Ord a => a -> AddRemoveSet s a -> Par d s ()--- We remove an element by adding it to the `removed` set!-remove !elm (AddRemoveSet added removed) = PS.insert elm removed---- | Wait for a single element to be removed from the set.-waitRemovedElem :: Ord a => a -> AddRemoveSet s a -> Par d s ()-waitRemovedElem !elm (AddRemoveSet added removed) = PS.waitElem elm removed---- | Wait on the size of the set of removed elements.-waitRemovedSize :: Int -> AddRemoveSet s a -> Par d s ()-waitRemovedSize !sz (AddRemoveSet added removed) = PS.waitSize sz removed---- | Get the exact contents of the set. As with any--- quasi-deterministic operation, using `freezeSet` may cause your--- program to exhibit a limited form of nondeterminism: it will never--- return the wrong answer, but it may include synchronization bugs--- that can (nondeterministically) cause exceptions.-freezeSet :: Ord a => AddRemoveSet s a -> QPar s (S.Set a)--- Freezing takes the set difference of added and removed elements.-freezeSet (AddRemoveSet added removed) =- liftA2 S.difference (PS.freezeSet added) (PS.freezeSet removed)
− Data/LVar/CycGraph.hs
@@ -1,576 +0,0 @@-{-# LANGUAGE ScopedTypeVariables, DataKinds #-}-{-# LANGUAGE KindSignatures, EmptyDataDecls #-}-{-# LANGUAGE NamedFieldPuns, ParallelListComp #-}-{-# LANGUAGE BangPatterns, CPP #-}-{-# LANGUAGE FlexibleInstances #-}--- {-# LANGUAGE UndecidableInstances #-}-{-# OPTIONS_GHC -O2 #-}--{-|--In contrast with "Data.LVar.Memo", this module provides a way to run a computation-for each node of a graph WITH support for cycles. Cycles are explicitly recognized-and then may be handled in an application specific fashion.-- -}--module Data.LVar.CycGraph- (- -- * An idiom for fixed point computations- exploreGraph_seq,- Response(..),-- -- * A parallel version- exploreGraph, NodeValue(..), NodeAction,-- -- * Debugging aides- ShortShow(..), shortTwo- )- where--- Standard:-import Data.Set (Set)-import Control.Monad-import qualified Data.Set as S-import qualified Data.Map as M-import Data.IORef-import Data.Char (ord)-import Data.List (intersperse)-import Data.Int-import qualified Data.Foldable as F-import System.IO.Unsafe-import Debug.Trace---- LVish:-import Control.LVish-import qualified Control.LVish.Internal as LV-import qualified Control.LVish.SchedIdempotent as LI-import Data.LVar.PureSet as IS-import Data.LVar.IVar as IV-import qualified Data.Concurrent.SkipListMap as SLM-import qualified Data.Set as S-import qualified Data.LVar.PureMap as IM--- import qualified Data.LVar.SLMap as IM--- import qualified Data.LVar.PureSet as S------- For debugging: -----#ifdef DEBUG_MEMO -import System.Environment (getEnvironment)-import Data.Graph.Inductive.Graph as G-import Data.Graph.Inductive.PatriciaTree as G-import Data.GraphViz as GV-import qualified Data.GraphViz.Attributes.Complete as GA-import qualified Data.GraphViz.Attributes.Colors as GC-import Data.Text.Lazy (pack)-#endif------------------------------------------------------------------------------------ Simple atomic Set accumulators------------------------------------------------------------------------------------- | Could use a more scalable structure here... but we need union as well as--- elementwise insertion.-type SetAcc a = IORef (S.Set a)---- Here @SetAcc@s are LINKED to downstream SetAcc's which must receive all the same--- inserts that they do.--- newtype SetAcc a = SetAcc (IORef (S.Set a, [SetAcc a]))--newSetAcc :: Par d s (SetAcc a)-newSetAcc = LV.WrapPar $ LI.liftIO $ newIORef S.empty-readSetAcc :: (SetAcc a) -> Par d s (S.Set a)-readSetAcc r = LV.WrapPar $ LI.liftIO $ readIORef r-insertSetAcc :: Ord a => a -> SetAcc a -> Par d s (S.Set a)-insertSetAcc x ref = LV.WrapPar $ LI.liftIO $- atomicModifyIORef' ref (\ s -> let ss = S.insert x s in (ss,ss))-unionSetAcc :: Ord a => Set a -> SetAcc a -> Par d s (S.Set a)-unionSetAcc x ref = LV.WrapPar $ LI.liftIO $- atomicModifyIORef' ref (\ s -> let ss = S.union x s in (ss,ss))------------------------------------------------------------------------------------- Types------------------------------------------------------------------------------------- | A Memo-table that stores cached results of executing a `Par` computation.--- --- This, enhanced, version of the Memo-table also is required to track all the keys--- that are reachable from each key (for cycle-detection).-data Memo (d::Determinism) s k v =- -- Here we keep both a Ivars of return values, and a set of keys whose computations- -- have traversed through THIS key. If we see a cycle there, we can catch it.--- !(IM.IMap k s (SetAcc k, IVar s v))- - Memo !(IS.ISet s k)- -- EXPENSIVE version:- !(IM.IMap k s (NodeRecord s k v))- -- ^ Store all the keys that we know *can reach this key*---- | All the information associated with one node in the graph of keys.-data NodeRecord s k v = NodeRecord- { mykey :: k- , chldrn :: [k]- , reachme :: !(IS.ISet s k) -- ^ Which keys are upstream of me in the graph- , in_cycle :: !(IVar s Bool) -- ^ Does this node participate in any cycle?- , result :: !(IVar s v) -- ^ The result of the per-node computation.- } deriving (Eq)------------------------------------------------------------------------------------- Cycle-detecting mapping of a computation over graph neighborhoods------------------------------------------------------------------------------------- | A means of building a dynamic graph. The node computation returns a response--- which may either be a final value, or a request to explore more nodes (together--- with a continuation for the resulting value).------ Note that because only one key is requested at a time, this cannot express--- parallel graph traversals.-data Response par key ans =- Done !ans- | Request !key (RequestCont par key ans)- -type RequestCont par key ans = (ans -> par (Response par key ans))------------------------------------------------------------------------------------- Sequential version:---- | This supercombinator does a parallel depth-first search of a dynamic graph, with--- detection of cycles.--- --- Each node in the graph is a computation whose input is the `key` (the vertex ID).--- Each such computation dynamically computes which other keys it depends on and--- requests the values associated with those keys.------ This implementation uses a sequential depth-first-search (DFS), starting from the--- initially requested key. One can picture this search as a directed tree radiating--- from the starting key. When a cycle is detected at any leaf of this tree, an--- alternate cycle handler is called instead of running the normal computation for--- that key.-exploreGraph_seq :: forall d s k v . (Ord k, Eq v, Show k, Show v) =>- (k -> Par d s (Response (Par d s) k v)) -- ^ The computation to perform for new requests- -> (k -> Par d s v) -- ^ Handler for a cycle on @k@. The- -- value it returns is in lieu of running- -- the main computation at this- -- particular node in the graph.- -> k -- ^ Key to lookup.- -> Par d s v-exploreGraph_seq initCont cycHndlr initKey = do- -- Start things off:- resp <- initCont initKey- v <- loop initKey (S.singleton initKey) resp return- return v- where- loop :: k -> S.Set k -> (Response (Par d s) k v) -> (v -> Par d s v) -> Par d s v- loop current hist resp kont = do- dbgPr (" [MemoFixedPoint] going around loop, key "++showID current++", hist size "++show (S.size hist))- case resp of- Done ans -> do dbgPr (" !! Final result, answer "++show ans)- kont ans- Request key2 newCont- -- Here we have hit a cycle, and label it as such for the CURRENT node.- | S.member key2 hist -> do- dbgPr (" Stopping before hitting a cycle on "++showID key2++", call cycHndlr on "++showID current)- ans <- cycHndlr current- kont ans- | otherwise -> do- dbgPr (" Requesting child computation with key "++showWID key2)- resp' <- initCont key2- loop key2 (S.insert key2 hist) resp' $ \ ans2 -> do- dbgPr (" DONE blocking on child key, cont invoked with answer: "++show ans2)- resp'' <- newCont ans2- -- Popping back to processing the current key, which may not be finished.- loop current hist resp'' kont- --- -- if wasloop then do--- if False then do --- -- Here the child computation ended up being processed as a cycle, so we must be as well:--- dbgPr (" Child comp "++showID key2++" of "++showID current++" hit a cycle...")--- ans3 <- cycHndlr current--- kont (True,ans3)-- -----------------------------------------------------------------------------------type IsCycle = Bool---- | The handler at a particular node (key) in the graph. This takes as argument a--- key, along with a boolean indicating whether the current node has been found to--- be part of a cycle.--- --- Also, for each child node, this handler is provided a way to demand the--- resulting value of that child node, plus an indication of whether the child node--- participates in a cycle.------ Finally, this handler is expected to produce a value which becomes associated--- with the key.-type NodeAction d s k v =--- Bool -> k -> [(Bool,Par d s v)] -> Par d s v- IsCycle -> k -> [(k,IsCycle,IV.IVar s v)] -> Par d s (NodeValue k v)- -- One thing that's missing here is WHICH child node(s) puts us in a cycle.---- | At the end of the handler execution, the value of a node is either ready, or it--- is instead deferred to be exactly the value provided by another key.-data NodeValue k v = FinalValue !v | Defer k - deriving (Show,Eq,Ord)----- | This combinator provides parallel exploration of a graph that contains cycles.--- The limitation is that the work to be performed at each node (`NodeAction`) is not--- invoked until the graph is fully traversed, i.e. after a barrier. Thus the graph--- explored is not a "dynamic graph" in the sense of being computed on the fly by the--- `NodeAction`.------ The algorithm used in this function is fairly expensive. For each node, it uses a--- monotonic data structure to track the full set of other nodes that can reach it in--- the graph.-#ifdef DEBUG_MEMO-exploreGraph :: forall s k v . (Ord k, Eq v, ShortShow k, Show v) =>-#else-exploreGraph :: forall s k v . (Ord k, Eq v, Show k, Show v) =>-#endif- (k -> Par QuasiDet s [k]) -- ^ Sketch the graph: map a key onto its children.- -> NodeAction QuasiDet s k v -- ^ The computation to run at each graph node.- -> k -- ^ The initial node (key) from which to explore.- -> Par QuasiDet s v-exploreGraph keyNbrs nodeHndlr initKey = do-- -- First: propogate key requests.- -- This will not diverge because the Set here suppressed duplicate callbacks:- set <- IS.newEmptySet - -- The map stores results:- mp <- IM.newEmptyMap-- keywalkHP <- newPool-- IS.forEachHP (Just keywalkHP) set $ \ key0 -> do- dbgPr ("![MemoFixedPoint] Start new key "++show key0)- -- Make some empty space for results:- key0_res <- IV.new- key0_cycle <- IV.new - key0_reach <- IS.newEmptySet- -- Next fetch the child node identities:- child_keys <- keyNbrs key0 - IM.insert key0 (NodeRecord key0 child_keys key0_reach key0_cycle key0_res) mp- dbgPr (" Computed nbrs of "++showID key0++" to be: "++ (showIDs child_keys))-- case child_keys of- [] -> return () -- IV.put_ key0_cycle False- _ -> do - -- Spawn traversals of child nodes:- forM_ child_keys (`IS.insert` set)- - -- Establish the (expensive) cycle-checker handler:- IS.forEachHP (Just keywalkHP) key0_reach $ \ key1 ->- when (key1 == key0) $ do- dbgPr (" !! Cycle detected on key "++showID key0)- IV.put_ key0_cycle True-- -- Now we must wait for records to come up, and establish ourselves as upstream- -- of each child:- chldrecs <- forM child_keys $ \child -> do - nrec@NodeRecord{reachme} <- IM.getKey child mp- IS.insert key0 reachme -- Child is reachable from us.- -- Further, what reaches us, reaches the child:- copyTo keywalkHP key0_reach reachme- dbgPr (" Inserted ourselves ("++showID key0++") in reachme list of child: "++showID child)- return nrec-- -- If all our children are do not participate in a cycle, neither do we.- -- fork $ let loop [] = IV.put_ key0_cycle False- -- loop (NodeRecord{in_cycle}:tl) = do- -- bl <- IV.get in_cycle- -- case bl of- -- True -> return ()- -- False -> loop tl- -- in loop chldrecs - -- FINISHME: If we have some cycle children and some leafish ones....- -- then we may need to do an unsafe peek at our reachme set, no?- return ()-- IS.insert initKey set- quiesce keywalkHP- -- fset <- IS.freezeSet set- frmap <- IM.freezeMap mp-- dbgPr ("Froze map: "++show (M.keys frmap))- - -- TODO: need parallel traversable:- let getcyc vr = do mb <- IV.freezeIVar vr- if mb == Just True- then return True- else return False- showCyc bl = if bl then "cycle" else "Nocyc"- fn NodeRecord{mykey, chldrn, reachme,in_cycle=mecyc,result=myres} () = fork$ do- bl <- getcyc mecyc- bls <- mapM (getcyc . in_cycle . (frmap #)) chldrn- dbgPr (" !! Invoking node handler at key "++showID mykey++" "++- showCyc bl ++" chldrn "++concat (intersperse " "$ map showCyc bls))- x <- nodeHndlr bl mykey [ (k, b, result (frmap # k)) | b <- bls- | k <- chldrn ]- case x of- FinalValue vv -> do - dbgPr (" !! Writing result into key "++showID mykey++" value: "++show x)- IV.put_ myres vv- Defer tokey -> do dbgPr (" !! No result yet on key "++showID mykey++", DEFERing to key "++showID tokey)- fork $ do kv <- IV.get (result(frmap # tokey))- dbgPr (" .. Delegated key "++showID tokey++", of key "++showID mykey++" produced result: "++show kv)- IV.put_ myres kv- F.foldrM fn () frmap-- let NodeRecord{result} = frmap # initKey- final <- IV.get result- ------------------------------------------------------------- -- TEMP: Debugging- -------------------------------------------------------------#ifdef DEBUG_MEMO - when (dbg_lvl >= 4) $ do- dbgPr ("| START creating dot graph...")- dg <- debugVizMemoGraph True initKey frmap- unsafePerformIO (GV.runGraphviz dg GV.Pdf "MemoCyc_short.pdf")- `seq` return () - dg <- debugVizMemoGraph False initKey frmap- unsafePerformIO (GV.runGraphviz dg GV.Pdf "MemoCyc.pdf")- `seq` return ()- dbgPr ("| DONE creating dot graph...") -#endif - ------------------------------------------------------------ - return final--- return $! Memo set mp --{------ | This version watches for, and catches, cyclic requests to the memotable that--- would normally diverge. Once caught, the user specifies what to do with these--- cycles by providing a handler. The handler is called on the key which formed the--- cycle. That is, computing the invocation spawned by that key results in a demand--- for that key. -makeMemoCyclic :: (MemoTable d s a b -> a -> Par d s b) -> (a -> Par d s b) -> Par d s (MemoTable d s a b)-makeMemoCyclic normalFn ifCycle = undefined--- FIXME: Are there races where more than one cycle can be hit? Can we guarantee--- that all are hit? ------ | Cancel an outstanding speculative computation. This recursively attempts to--- cancel any downstream computations in this or other memo-tables that are children--- of the given `MemoFuture`.-cancel :: MemoFuture Det s b -> Par Det s ()--- FIXME: Det needs to be replaced here with "GetOnly".-cancel fut = undefined---}------------------------------------------------------------------------------------- Misc Helpers and Utilities-----------------------------------------------------------------------------------(#) :: (Ord a1, Show a1) => M.Map a1 a -> a1 -> a-m # k = case M.lookup k m of- Nothing -> error$ "Key was missing from map: "++show k- Just x -> x--showMapContents :: (Eq t1, Show a, Show a1) => IM.IMap a1 s (IORef (Set a), IV.IVar t t1) -> IO String-showMapContents (IM.IMap lv) = do- mp <- readIORef (LV.state lv)- let lst = M.toList mp- return$ " Map Contents: (length "++ show (length lst) ++")\n" ++- concat [ " "++fullempt++" "++showWID k++" -> "++vals++"\n"- | (k,(v,IV.IVar ivr)) <- lst--- , let vals = "hello"- , let lst = S.toList $ unsafePerformIO (readIORef v)- , let vals = "#"++show (length lst)++"["++ (concat $ intersperse ", " $ map showID lst) ++"]"- , let fullempt = if Nothing == unsafePerformIO (readIORef (LV.state ivr))- then "[empty]"- else "[full]"- ]--showMapContents2 :: (Eq t3, Show t1, Show a) => IM.IMap a s (ISet t t1, IV.IVar t2 t3) -> IO String-showMapContents2 (IM.IMap lv) = do- mp <- readIORef (LV.state lv)- let lst = M.toList mp- return$ " Map Contents: (length "++ show (length lst) ++")\n" ++- concat [ " "++fullempt++" "++showWID k++" -> "++vals++"\n"- | (k,(IS.ISet setlv, IV.IVar ivr)) <- lst--- , let vals = "hello"- , let lst = S.toList $ unsafePerformIO (readIORef (LV.state setlv))- , let vals = "#"++show (length lst)++"["++ (concat $ intersperse ", " $ map showID lst) ++"]"- , let fullempt = if Nothing == unsafePerformIO (readIORef (LV.state ivr))- then "[empty]"- else "[full]"- ]---- | Variant of `union` that optionally ties the handlers in the resulting set to the same--- handler pool as those in the two input sets.-copyTo :: Ord a => HandlerPool -> IS.ISet s a -> IS.ISet s a -> Par d s ()-copyTo hp sfrom sto = do- IS.forEachHP (Just hp) sfrom (`insert` sto)--{-# INLINE dbgPr #-}-dbgPr :: Monad m => String -> m ()-#ifdef DEBUG_MEMO-dbgPr s | dbg_lvl >= 1 = trace s (return ())- | otherwise = return ()-#else-dbgPr _ = return ()-#endif--showWID :: Show a => a -> String-showWID x = let str = (show x) in- if length str < 10- then str- else showID x++"__"++str--showID :: Show a => a -> String-showID x = let str = (show x) in- if length str < 10 then str- else (show (length str))++"-"++ show (checksum str)--showIDs ls = ("{"++(concat$ intersperse ", " $ map showID ls)++"}")--checksum :: String -> Int-checksum str = sum (map ord str)-------------------------------------------------------------------------------------- DEBUGGING------------------------------------------------------------------------------------- | A show class that tries to stay under a budget.-class Show t => ShortShow t where- shortShow :: Int -> t -> String- shortShow n x = take n (show x)--instance ShortShow Bool where- shortShow 1 True = "t"- shortShow 1 False = "f"- shortShow 2 True = "#t"- shortShow 2 False = "#f" - shortShow n b = take n (show b)--instance ShortShow Integer where shortShow = shortShowNum-instance ShortShow Int where shortShow = shortShowNum-instance ShortShow Int8 where shortShow = shortShowNum-instance ShortShow Int16 where shortShow = shortShowNum-instance ShortShow Int32 where shortShow = shortShowNum-instance ShortShow Int64 where shortShow = shortShowNum --shortShowNum :: Show a => Int -> a -> String-shortShowNum n num =- let str = show num- len = length str in- if len > n then- (take (n-2) str)++".."- else str- -instance ShortShow String where- shortShow n str =- let len = length str in- if len > 2 && n ==2- then ".."- else if len > 1 && n == 1- then "?"- else take n str--instance (ShortShow a, ShortShow b) => ShortShow (a,b) where- shortShow 1 _ = "?"- shortShow 2 _ = ".."- shortShow n (a,b) = let (l,r) = shortTwo (n-3) a b - in "("++ l ++","++ r ++")"---- | Combine two things within a given size budget.-shortTwo :: (ShortShow t, ShortShow t1) => Int -> t -> t1 -> (String, String)--- this could be better...-shortTwo n a b = (left, shortShow (half+remain) b)- where- remain = abs (half - length left)- left = shortShow half a- (q,r) = quotRem (abs(n-3)) 2 - half = q + r------------------------------------------------------------------------------------#ifdef DEBUG_MEMO---- | Debugging flag shared by all accelerate-backend-kit modules.--- This is activated by setting the environment variable DEBUG=1..5-dbg_lvl :: Int-dbg_lvl = case lookup "DEBUG" theEnv of- Nothing -> defaultDbg- Just "" -> defaultDbg- Just "0" -> defaultDbg- Just s ->- trace (" ! Responding to env Var: DEBUG="++s)$- case reads s of- ((n,_):_) -> n- [] -> error$"Attempt to parse DEBUG env var as Int failed: "++show s--theEnv :: [(String, String)]-theEnv = unsafePerformIO getEnvironment--defaultDbg :: Int-defaultDbg = 0--debugVizMemoGraph :: forall s t t1 t2 . (Ord t1, ShortShow t1, Show t2, F.Foldable t) =>- Bool -- ^ Use shorter `showID` for keys.- -> t1 -- ^ The inital key.- -> t (NodeRecord s t1 t2) -- ^ A frozen map of graph nodes.--- Par d s (Gr (Bool,String) ())- -> Par QuasiDet s (GV.DotGraph G.Node)-debugVizMemoGraph idOnly initKey frmap = do- let showKey = if idOnly then showID- else shortShow 40- let gcons :: NodeRecord s t1 t2- -> (M.Map t1 G.Node, G.Gr (Bool,t1,t2) ())- -> Par QuasiDet s (M.Map t1 G.Node, G.Gr (Bool,t1,t2) ())- gcons NodeRecord{mykey, in_cycle,result}- (labmap, gracc) = do- dbgPr (" .. About to wait for node result, key "++show mykey)- res <- IV.get result- dbgPr (" .. About to wait for node in_cycle, key "++show mykey)- cyc <- IV.freezeIVar in_cycle- let num = 1 + G.noNodes gracc - gr' = G.insNode (num, (cyc == Just True,mykey,res)) $ - gracc- labmap' = M.insert mykey num labmap- return (labmap',gr')- - gedges :: NodeRecord s t1 t2- -> (M.Map t1 G.Node, G.Gr (Bool,t1,t2) ())- -> Par d s (M.Map t1 G.Node, G.Gr (Bool,t1,t2) ())- gedges NodeRecord{mykey, chldrn }- (labmap, gracc) = do - let chldnodes = map (labmap #) chldrn- num = labmap # mykey- gr' = G.insEdges [ (num,cnd::Int,()) | cnd <- chldnodes ] $- gracc- labmap' = M.insert mykey num labmap- return (labmap',gr')- - dbgPr (" !! Creating graphviz graph from MemoCyc map of size "++show (F.foldr (\ _ n -> 1+n) 0 frmap))--- dbgPr (" !! All keys "++show frmap)- - -- Two passes, first add nodes, then edges:- (lm,graph0) <- F.foldrM gcons (M.empty, G.empty) frmap- dbgPr (" .. Added all nodes to the graph...") - (_,graph) <- F.foldrM gedges (lm, graph0) frmap- dbgPr (" .. Added all edges to the graph...") - let -- dg = graphToDot nonClusteredParams graph- myparams :: GV.GraphvizParams G.Node (Bool,t1,t2) () () (Bool,t1,t2)- myparams = GV.defaultParams { GV.fmtNode= nodeAttrs }-- nodeAttrs :: (Int, (Bool,t1,t2)) -> [GA.Attribute]--- nodeAttrs :: (Int, String) -> [GA.Attribute] - nodeAttrs (_num, (cyc,key,res)) =- let lbl = showKey key++"\n=> "++ show res in- [ GA.Label$ GA.StrLabel $ pack lbl ] ++- (if key == initKey - then [GA.Color [weighted$ GA.X11Color GV.Red]]- else []) ++- (if cyc then []- else [GA.Shape GA.BoxShape])-- dg = GV.graphToDot myparams graph -- (G.nmap uid graph)- return dg--weighted c = GC.WC {GC.wColor=c, GC.weighting=Nothing}--#endif--- End DEBUG_MEMO
− Data/LVar/MaxCounter.hs
@@ -1,79 +0,0 @@-{-# LANGUAGE Trustworthy #-}--{-# LANGUAGE DataKinds, BangPatterns, MagicHash #-}-{-# LANGUAGE TypeSynonymInstances, FlexibleInstances, MultiParamTypeClasses, TypeFamilies #-}---- | A counter that contains the maximum value of all puts.---- TODO: Add 'Min', 'Or', 'And' and other idempotent ops...--module Data.LVar.MaxCounter- ( MaxCounter,- newMaxCounter, put, waitThresh, freezeMaxCounter- ) where--import Control.LVish hiding (freeze, put)-import Control.LVish.Internal (state)-import Control.LVish.DeepFrz.Internal-import Data.IORef-import Data.LVar.Generic-import Data.LVar.Internal.Pure as P-import Algebra.Lattice-import System.IO.Unsafe (unsafeDupablePerformIO)-import GHC.Prim (unsafeCoerce#)-------------------------------------------------------------------------------------- | A @MaxCounter@ is really a constant-space ongoing @fold max@ operation.--- --- A @MaxCounter@ is an example of a `PureLVar`. It is implemented simply as a--- pure value in a mutable box.-type MaxCounter s = PureLVar s MC--newtype MC = MC Int- deriving (Eq, Show, Ord, Read)--instance JoinSemiLattice MC where - join (MC !a) (MC !b) = MC (a `max` b)--instance BoundedJoinSemiLattice MC where- bottom = MC minBound---- | Create a new counter with the given initial value.-newMaxCounter :: Int -> Par d s (MaxCounter s)-newMaxCounter n = newPureLVar (MC n)---- | Incorporate a new value in the max-fold. If the previous maximum is less than--- the new value, increase it.-put :: MaxCounter s -> Int -> Par d s ()-put lv n = putPureLVar lv (MC n)---- | Wait until the maximum observed value reaches some threshold, then return.-waitThresh :: MaxCounter s -> Int -> Par d s ()-waitThresh lv n = waitPureLVar lv (MC n)---- | Observe what the final value of the counter was.-freezeMaxCounter :: MaxCounter s -> Par QuasiDet s Int-freezeMaxCounter lv = do- MC n <- freezePureLVar lv- return n---- | Once frozen, for example by `runParThenFreeze`, a MaxCounter can be converted--- directly into an Int.-fromMaxCounter :: MaxCounter Frzn -> Int-fromMaxCounter (PureLVar lv) =- case unsafeDupablePerformIO (readIORef (state lv)) of- MC n -> n--instance DeepFrz MC where- type FrzType MC = MC---- Don't need this because there is an instance for `PureLVar`:-{---- | @MaxCounter@ values can be returned in the results of a--- `runParThenFreeze`. Hence they need a `DeepFrz` instance.--- @DeepFrz@ is just a type-coercion. No bits flipped at runtime.-instance DeepFrz (MaxCounter s) where- type FrzType (MaxCounter s) = (MaxCounter Frzn)- frz = unsafeCoerce#--}
− Data/LVar/Memo.hs
@@ -1,110 +0,0 @@-{-# LANGUAGE ScopedTypeVariables, DataKinds #-}-{-# LANGUAGE KindSignatures, EmptyDataDecls #-}-{-# LANGUAGE BangPatterns #-}-{-# OPTIONS_GHC -O2 #-}--{-|--This basic version of memotables is implemented on top of existing LVars without-breaking any rules.--The problem is that it cannot do cycle detection, because that requires tracking-extra information (where we've been) which is NOT exposed to the user and NOT used -- -}-module Data.LVar.Memo- (- -- * Memo tables and defered lookups - Memo, MemoFuture, makeMemo,- - -- * Memo table operations- getLazy, getMemo, force- ) where-import Debug.Trace--import Control.LVish-import qualified Data.Set as S--- import qualified Data.LVar.SLMap as IM--- import Data.LVar.SLSet as IS-import qualified Data.LVar.PureMap as IM-import Data.LVar.PureSet as IS-import Data.LVar.IVar as IV------------------------------------------------------------------------------------- Types------------------------------------------------------------------------------------- | A Memo-table that stores cached results of executing a `Par` computation.-data Memo (d::Determinism) s a b =- Memo !(IS.ISet s a)- !(IM.IMap a s b)---- | A result from a lookup in a Memo-table, unforced.--- The two-stage `getLazy`/`force` lookup is useful to separate--- spawning the work from demanding its result.-newtype MemoFuture (d :: Determinism) s b = MemoFuture (Par d s b)-------------------------------------------------------------------------------------- | Reify a function in the `Par` monad as an explicit memoization table.-makeMemo :: (Ord a, Eq b, Show a, Show b) =>- (a -> Par d s b) -> Par d s (Memo d s a b)-makeMemo fn = do- st <- newEmptySet- mp <- IM.newEmptyMap- IS.forEach st $ \ elm -> do- res <- fn elm- trace ("makeMemo, about to insert result: "++show (show elm, show res)) $ - IM.insert elm res mp- return $! Memo st mp--- TODO: this version may want to have access to the memo-table within the handler as--- well....----- | Read from the memo-table. If the value must be computed, do that right away and--- block until its complete.-getMemo :: (Ord a, Eq b) => Memo d s a b -> a -> Par d s b -getMemo tab key =- do fut <- getLazy tab key- force fut---- | Begin to read from the memo-table. Initiate the computation if the key is not--- already present. Don't block on the computation being complete, rather, return a--- future.-getLazy :: (Ord a, Eq b) => Memo d s a b -> a -> Par d s (MemoFuture d s b)-getLazy (Memo st mp) key = do - IS.insert key st- return $! MemoFuture (IM.getKey key mp)----- | This will throw exceptions that were raised during the computation, INCLUDING--- multiple put.-force :: MemoFuture d s b -> Par d s b -force (MemoFuture pr) = pr--- FIXME!!! Where do errors in the memoized function (e.g. multiple put) surface?--- We must pick a determined, consistent place.--- --- Multiple put errors may not be able to wait until this point to get--- thrown. Otherwise we'd have to be at least quasideterministic here. If you have--- a MemoFuture you never force, it and an outside computation may be racing to do a--- put. If the outside one wins the MemoFuture is the one that gets the exception--- (and hides it), otherwise the exception is exposed. Quasideterminism.---- It may be fair to distinguish between internal problems with the MemoFuture--- (deferred exceptions), and problematic interactions with the outside world (double--- put) which would then not be deferred. Such futures can't be canceled anyway, so--- there's really no need to defer the exceptions.----{------ | Cancel an outstanding speculative computation. This recursively attempts to--- cancel any downstream computations in this or other memo-tables that are children--- of the given `MemoFuture`.-cancel :: MemoFuture Det s b -> Par Det s ()--- FIXME: Det needs to be replaced here with "GetOnly".-cancel fut = undefined---}
− Data/LVar/NatArray.hs
@@ -1,284 +0,0 @@-{-# LANGUAGE Trustworthy #-}-{-# LANGUAGE BangPatterns #-}-{-# LANGUAGE NamedFieldPuns #-}-{-# LANGUAGE RankNTypes #-}-{-# LANGUAGE ScopedTypeVariables #-}-{-# LANGUAGE TypeFamilies #-}-{-# LANGUAGE CPP #-}-{-# LANGUAGE FlexibleContexts #-}-{-# LANGUAGE MultiParamTypeClasses #-}-{-# LANGUAGE FlexibleInstances #-}-{-# LANGUAGE DataKinds #-}-{-# LANGUAGE InstanceSigs #-}-{-# LANGUAGE GADTs #-}--{-|--An I-structure (array) of /positive/ numbers. A `NatArray` cannot store zeros.--This particular implementation makes a trade-off between expressiveness (monomorphic-in array contents) and efficiency. The efficiency gained of course is that the array-may be unboxed, and we don't need extra bits to store empty/full status.--/However/, relative to "Data.LVar.IStructure", there is a performance disadvantage as-well. As of [2013.09.28] and their initial release, `NatArray`s are implemented as a-/single/ `LVar`, which means they share a single wait-list of blocked computations.-If there are many computations blocking on different elements within a `NatArray`,-scalability will be much worse than with other `IStructure` implementations.--The holy grail is to get unboxed arrays and scalable blocking, but we don't have this-yet.--Finally, note that this data-structure has an EXPERIMENTAL status and may be removed-in future releases as we find better ways to support unboxed array structures with-per-element synchronization.---}--module Data.LVar.NatArray- (- -- * Basic operations- NatArray,- newNatArray, put, get,-- -- * Iteration and callbacks- forEach, forEachHP-- -- -- * Quasi-deterministic operations- -- freezeSetAfter, withCallbacksThenFreeze, freezeSet,-- -- -- * Higher-level derived operations- -- copy, traverseSet, traverseSet_, union, intersection,- -- cartesianProd, cartesianProds, -- -- -- * Alternate versions of derived ops that expose HandlerPools they create.- -- forEachHP, traverseSetHP, traverseSetHP_,- -- cartesianProdHP, cartesianProdsHP- ) where---- import qualified Data.Vector.Unboxed as U--- import qualified Data.Vector.Unboxed.Mutable as M---import Data.LVar.NatArray.Unsafe--import qualified Data.Vector.Storable as U-import qualified Data.Vector.Storable.Mutable as M-import Foreign.Marshal.MissingAlloc (callocBytes)-import Foreign.Marshal.Alloc (finalizerFree)-import Foreign.Storable (sizeOf, Storable)-import Foreign.ForeignPtr (newForeignPtr, withForeignPtr)-import qualified Foreign.Ptr as P-import qualified Data.Bits.Atomic as B-import Data.Bits ((.&.))--import Control.Monad (void)-import Control.Exception (throw)-import Data.IORef-import Data.Maybe (fromMaybe)-import qualified Data.Set as S-import qualified Data.LVar.IVar as IV-import qualified Data.Foldable as F-import qualified Data.Traversable as T-import Data.LVar.Generic--import Control.LVish as LV hiding (addHandler, put,get)-import Control.LVish.DeepFrz.Internal as DF-import Control.LVish.Internal as LI-import Control.LVish.SchedIdempotent (newLV, putLV, getLV, freezeLV,- freezeLVAfter, liftIO)-import qualified Control.LVish.SchedIdempotent as L-import System.IO.Unsafe (unsafeDupablePerformIO)-import Data.LVar.NatArray.Unsafe (NatArray(..))----------------------------------------------------------------------------------- Toggles--#define USE_CALLOC--- A low-level optimization below.----------------------------------------------------------------------------------unNatArray (NatArray lv) = lv---- | Physical identity, just as with IORefs.--- instance Eq (NatArray s v) where--- NatArray lv1 == NatArray lv2 = state lv1 == state lv2 ---- | Create a new, empty, monotonically growing 'NatArray' of a given size.--- All entries start off as zero, which must be BOTTOM.-newNatArray :: forall elt d s . (Storable elt, Num elt) =>- Int -> Par d s (NatArray s elt)-newNatArray len = WrapPar $ fmap (NatArray . WrapLVar) $ newLV $ do-#ifdef USE_CALLOC- let bytes = sizeOf (undefined::elt) * len- mem <- callocBytes bytes- fp <- newForeignPtr finalizerFree mem- return $! M.unsafeFromForeignPtr0 fp len-#else- M.replicate len 0-#endif---- | /O(1)/ Freeze operation that directly returns a nice, usable, representation of--- the array data.-freezeNatArray :: Storable a => NatArray s a -> LV.Par QuasiDet s (U.Vector a)-freezeNatArray (NatArray lv) = do--- freezeLV --- U.unsafeFreeze (state lv))- error "FINISHME -- freezeNatArray "- -- LI.liftIO $ U.unsafeFreeze (LI.state lv)------------------------------------------------------------------------------------- Instances:---- FIXME: there is a tension here.. should NatArray really be a generic LVarData1 at all?--- Can it really store anything in Storable!?!? Or do we need to fix it to numbers--- to ensure the zero-trick makes sense?--{---instance DeepFrz a => DeepFrz (NatArray s a) where- type FrzType (NatArray s a) = NatArray Frzn (FrzType a)- frz = unsafeCoerceLVar---- | /O(1)/: Convert from a frozen `NatArray` to a plain vector.--- This is only permitted when the `NatArray` has already been frozen.--- This is useful for processing the result of `Control.LVish.DeepFrz.runParThenFreeze`.-fromNatArray :: NatArray Frzn a -> U.Vector a-fromNatArray (NatArray lv) = unsafeDupablePerformIO (readIORef (state lv))---}------------------------------------------------------------------------------------{-# INLINE forEachHP #-}--- | Add an (asynchronous) callback that listens for all new elements added to--- the array, optionally enrolled in a handler pool.-forEachHP :: (Storable a, Eq a, Num a) =>- Maybe HandlerPool -- ^ pool to enroll in, if any- -> NatArray s a -- ^ array to listen to- -> (Int -> a -> Par d s ()) -- ^ callback- -> Par d s ()-forEachHP hp (NatArray (WrapLVar lv)) callb = WrapPar $ do- L.addHandler hp lv globalCB deltaCB- return ()- where- deltaCB (ix,x) = return$ Just$ unWrapPar$ callb ix x- globalCB vec = unWrapPar$- -- FIXME / TODO: need a better (parallel) for loop:- forVec vec $ \ ix elm ->- -- FIXME: When it starts off, it is SPARSE... there must be a good way to- -- avoid testing each position for zero.- if elm == 0- then return () - else forkHP hp $ callb ix elm--{-# INLINE forVec #-}--- | Simple for-each loops over vector elements.-forVec :: Storable a =>- M.IOVector a -> (Int -> a -> Par d s ()) -> Par d s ()-forVec vec fn = loop 0 - where- len = M.length vec- loop i | i == len = return ()- | otherwise = do elm <- LI.liftIO$ M.unsafeRead vec i- fn i elm- loop (i+1)--{-# INLINE forEach #-}--- | Add an (asynchronous) callback that listens for all new elements added to--- the set-forEach :: (Num a, Storable a, Eq a) =>- NatArray s a -> (Int -> a -> Par d s ()) -> Par d s ()-forEach = forEachHP Nothing---{-# INLINE put #-}--- | Put a single element in the array. That slot must be previously empty. (WHNF)--- Strict in the element being put in the set.-put :: forall s d elt . (Storable elt, B.AtomicBits elt, Num elt, Show elt) =>- NatArray s elt -> Int -> elt -> Par d s ()-put _ !ix 0 = throw (LVarSpecificExn$ "NatArray: violation! Attempt to put zero to index: "++show ix)-put (NatArray (WrapLVar lv)) !ix !elm = WrapPar$ putLV lv (putter ix)- where putter ix vec@(M.MVector _len fptr) =- withForeignPtr fptr $ \ ptr -> do - let offset = sizeOf (undefined::elt) * ix- -- ARG, if it weren't for the idempotency requirement we could use fetchAndAdd here:- -- orig <- B.fetchAndAdd (P.plusPtr ptr offset) elm - orig <- B.compareAndSwap (P.plusPtr ptr offset) 0 elm- case orig of- 0 -> return (Just (ix, elm))- i | i == elm -> return Nothing -- Allow repeated, equal puts.- | otherwise -> throw$ ConflictingPutExn$ "Multiple puts to index of a NatArray: "++- show ix++" new/old : "++show elm++"/"++show orig--{-# INLINE get #-}--- | Wait for an indexed entry to contain a non-zero value.--- --- Warning: this is inefficient if it needs to block, because the deltaThresh must--- monitor EVERY new addition.-get :: forall s d elt . (Storable elt, B.AtomicBits elt, Num elt) =>- NatArray s elt -> Int -> Par d s elt-get (NatArray (WrapLVar lv)) !ix = WrapPar $- getLV lv globalThresh deltaThresh- where- globalThresh ref _frzn = do - elm <- M.read ref ix - if elm == 0- then return Nothing- else return (Just elm)- -- FIXME: we don't actually want to call the deltaThresh on every element...- -- We want more locality than that...- deltaThresh (ix2,e2) | ix == ix2 = return$! Just e2- | otherwise = return Nothing ----- | A sequential for-loop with a catch. The body of the loop gets access to a--- special get function. This getter will not block subsequent iterations of the--- loop. Parallelism will be introduced minimally, only as neccessary to avoid--- blocking.-seqLoopNonblocking :: Int -> Int ->- ((NatArray s elt -> Int -> Par d s elt) -> Int -> Par d s ()) ->- Par d s ()-seqLoopNonblocking start end fn = do- error "TODO - FINISHME: seqLoopNonblocking optimization"- where- par =- L.Par $ \k -> L.ClosedPar $ \q -> do- -- tripped <- globalThresh state False--- case tripped of- -- Just b -> exec (k b) q -- already past the threshold; invoke the--- forkHP mh child = mkPar $ \k q -> do--- closed <- closeInPool mh child--- Sched.pushWork q (k ()) -- "Work-first" policy.--- -- hpMsg " [dbg-lvish] incremented and pushed work in forkInPool, now running cont" hp --- exec closed q - undefined--{--parFor :: (ParFuture iv p) => InclusiveRange -> (Int -> p ()) -> p ()-parFor (InclusiveRange start end) body =- do- let run (x,y) = for_ x (y+1) body- range_segments = splitInclusiveRange (4*numCapabilities) (start,end)-- vars <- M.forM range_segments (\ pr -> spawn_ (run pr))- M.mapM_ get vars- return ()--splitInclusiveRange :: Int -> (Int, Int) -> [(Int, Int)]-splitInclusiveRange pieces (start,end) =- map largepiece [0..remain-1] ++- map smallpiece [remain..pieces-1]- where- len = end - start + 1 -- inclusive [start,end]- (portion, remain) = len `quotRem` pieces- largepiece i =- let offset = start + (i * (portion + 1))- in (offset, offset + portion)- smallpiece i =- let offset = start + (i * portion) + remain- in (offset, offset + portion - 1)--data InclusiveRange = InclusiveRange Int Int--}
− Data/LVar/NatArray/Unsafe.hs
@@ -1,27 +0,0 @@-{-# LANGUAGE GADTs #-}---- | Unsafe operations on NatArray. NOT for end-user applications.--module Data.LVar.NatArray.Unsafe- ( NatArray(..), unsafePeek )- where-import qualified Data.Vector.Storable.Mutable as M-import Foreign.Storable (sizeOf, Storable)--- import System.IO.Unsafe (unsafeDupablePerformIO)-import Control.LVish.Internal as LI------------------------------------------------------------------------------------------------ | An array of bit-fields with a monotonic OR operation. This can be used to model--- a set of Ints by setting the vector entries to zero or one, but it can also--- model other finite lattices for each index.--- newtype NatArray s a = NatArray (LVar s (M.IOVector a) (Int,a))-data NatArray s a = Storable a => NatArray !(LVar s (M.IOVector a) (Int,a))--unsafePeek :: (Num a, Eq a) => NatArray s a -> Int -> Par d s (Maybe a)-unsafePeek (NatArray lv) ix = do- peek <- LI.liftIO $ M.read (LI.state lv) ix- case peek of- -- TODO: generalize:- 0 -> return Nothing- x -> return $! Just x
− Data/LVar/PNCounter.hs
@@ -1,73 +0,0 @@-{-# LANGUAGE BangPatterns #-}--{-|--This module provides a /PN-Counter/, a counter that allows both-increment and decrement operations. This is possible because, under-the hood, it's represented with two monotonically growing counters,-one for increments and one for decrements. The name "PN-Counter"-comes from the literature on /conflict-free replicated data types/.-- -}-module Data.LVar.PNCounter- (- PNCounter,- newCounter, newCounterWithValue,- increment, waitForIncrements,- decrement, waitForDecrements,-- freezeCounter- - ) where-import Control.LVish-import Control.LVish.Internal-import qualified Data.Atomics.Counter.Reference as AC--- LK: FIXME: it can't be okay to use SchedIdempotent if we're using bump, can it?!-import Control.LVish.SchedIdempotent (newLV)-import Data.IORef----- | The counter datatype.---- LK: LVar around the outside, or PureLVar? What's the difference?-data PNCounter s = LVar s (AC.AtomicCounter, AC.AtomicCounter)- --- | Create a new `PNCounter` set to zero.-newCounter :: Par d s (PNCounter s)-newCounter = newCounterWithValue 0---- | Create a new `PNCounter` with the specified initial value.-newCounterWithValue :: Int -> Par d s (PNCounter s)--- LK: hm, how do I create IORefs and then return a Par? I think what--- I'm supposed to be doing here is wrapping an unsafe internal Par--- computation (that's allowed to do IO) in a safe one that I return.-newCounterWithValue n = undefined--- FIXME...- -- do- -- incs <- newIORef (Just n)- -- decs <- newIORef Nothing---- | Increment the `PNCounter`.-increment :: PNCounter s -> Par d s ()-increment = undefined---- | Wait for the number of increments to reach a given number.-waitForIncrements :: Int -> PNCounter s -> Par d s ()-waitForIncrements = undefined---- | Decrement the `PNCounter`.-decrement :: PNCounter s -> Par d s ()-decrement = undefined---- | Wait for the number of decrements to reach a given number.-waitForDecrements :: Int -> PNCounter s -> Par d s ()-waitForDecrements = undefined---- | Get the exact contents of the counter. As with any--- quasi-deterministic operation, using `freezeCounter` may cause your--- program to exhibit a limited form of nondeterminism: it will never--- return the wrong answer, but it may include synchronization bugs--- that can (nondeterministically) cause exceptions.-freezeCounter :: PNCounter s -> QPar s Int--- Freezing takes the difference of increments and decrements.-freezeCounter = undefined
Data/LVar/PureMap.hs view
@@ -227,7 +227,7 @@ True -> return (Just ()) False -> return (Nothing) -- Here's an example of a situation where we CANNOT TELL if a delta puts it over- -- the threshold.a+ -- the threshold. deltaThresh _ = globalThresh (L.state lv) False -- | Get the exact contents of the map. As with any
Data/LVar/PureSet.hs view
@@ -255,8 +255,9 @@ -- | Wait on the /size/ of the set, not its contents. waitSize :: Int -> ISet s a -> Par d s ()-waitSize !sz (ISet lv) = WrapPar$- getLV (unWrapLVar lv) globalThresh deltaThresh+waitSize !sz (ISet lv) = do+ logDbgLn (-2) "PureSet.waitSize: about to (potentially) block:"+ WrapPar$ getLV (unWrapLVar lv) globalThresh deltaThresh where globalThresh ref _frzn = do set <- readIORef ref
lvish.cabal view
@@ -10,7 +10,7 @@ -- PVP summary: +-+------- breaking API changes -- | | +----- non-breaking API additions -- | | | +--- code changes with no API change-version: 1.1.1.3+version: 1.1.2 -- Changelog: -- 0.2 -- switch SLMap over to O(1) freeze@@ -21,8 +21,8 @@ -- 1.1.0.2 -- add verifyFiniteJoin -- 1.1.0.3 -- expose BulkRetry prototype -- 1.1.1.0 -- expose logging routines--- 1.1.1.1 -- restrict exports for interm hackage release--- 1.1.1.3 -- hide modules more aggressively, not just by flag+-- 1.1.1.5 -- various bugfixes+-- 1.1.2 -- bugfixes and small additions, work-around for -fbeta problems synopsis: Parallel scheduler, LVar data structures, and infrastructure to build more. @@ -42,6 +42,8 @@ Change Log: . * 1.0.0.6 - tighten up dependencies; remove unused flags; very minor doc fixes.+ . + * 1.1.2 - many bugfixes, plus a new logging infrastructure, and verification of finite lattices license: BSD3 license-file: LICENSE@@ -62,9 +64,9 @@ description: Use the Chase-Lev work-stealing deque default: False -flag newcontainers- description: Use a pre-release version of containers to enable splitting.- default: False+-- flag newcontainers+-- description: Use a pre-release version of containers to enable splitting.+-- default: False flag getonce description: Ensure that continuations of get run at most once @@ -74,11 +76,11 @@ -- We won't really support this until LVish 2.0: flag generic description: Use (forthcoming) generic interfaces for Par monads.- default: False+ default: True -flag beta- description: These features are in beta and not fully supported yet.- default: False+-- flag beta+-- description: These features are in beta and not fully supported yet.+-- default: True -------------------------------------------------------------------------------- library@@ -108,30 +110,29 @@ -- These are only for developing new LVars: Data.LVar.Internal.Pure Data.LVar.Generic.Internal+ Data.LVar.PureMap.Unsafe+ Data.LVar.SLMap.Unsafe Control.LVish.Internal Control.LVish.DeepFrz.Internal -- This is also not recommended for general use yet. Data.Concurrent.SkipListMap- - if flag(beta)- other-modules:- -- Not quite ready for prime-time yet:- Data.LVar.NatArray- Data.LVar.NatArray.Unsafe- Data.LVar.MaxCounter- Data.LVar.AddRemoveSet- Data.LVar.PNCounter- -------------------------------------------- -- New / Experimental:- Data.LVar.Memo - Data.LVar.CycGraph-- Data.LVar.PureMap.Unsafe- Data.LVar.SLMap.Unsafe+ -- if flag(beta)+ -- exposed-modules: + -- -- Not quite ready for prime-time yet:+ -- Data.LVar.NatArray+ -- Data.LVar.NatArray.Unsafe+ -- Data.LVar.MaxPosInt+ -- Data.LVar.Counter+ -- Data.LVar.AddRemoveSet+ -- Data.LVar.PNCounter+ -- -------------------------------------------+ -- -- New / Experimental:+ -- Data.LVar.Memo + -- Data.LVar.CycGraph - if flag(beta) && flag(newcontainers)- other-modules: - Control.LVish.BulkRetry+ -- if flag(beta) && flag(newcontainers)+ -- exposed-modules: + -- Control.LVish.BulkRetry -- Modules included in this library but not exported. other-modules:@@ -160,12 +161,12 @@ random, transformers, ghc-prim,- async- if flag(beta) - -- Used in NatArray:- build-depends: bits-atomic, missing-foreign- if flag(newcontainers) { build-depends: containers >= 0.5.3.2 } - else { build-depends: containers >= 0.5 }+ async,+ -- Used in NatArray:+ bits-atomic, missing-foreign+ -- if flag(newcontainers) { build-depends: containers >= 0.5.3.2 } + -- else { build-depends: containers >= 0.5 }+ build-depends: containers >= 0.5 if flag(generic) cpp-options: -DGENERIC_PAR build-depends: par-classes >= 1.0 && < 2.0,@@ -200,10 +201,10 @@ PureMapTests, SLMapTests, SetTests,- MaxCounterTests+ MaxPosIntTests,+ AddRemoveSetTests --- ghc-options: -O2 -threaded -rtsopts -with-rtsopts=-N4- ghc-options: -O2 -rtsopts+ ghc-options: -O2 -threaded -rtsopts -with-rtsopts=-N4 -- DUPLICATED: from above build-depends: base >= 4.6 && <= 4.8, @@ -215,10 +216,9 @@ random, transformers, ghc-prim,- async- if flag(beta) - -- Used in NatArray:- build-depends: bits-atomic, missing-foreign+ async,+ -- Used in NatArray:+ bits-atomic, missing-foreign if flag(generic) cpp-options: -DGENERIC_PAR build-depends: par-classes >= 1.0 && < 2.0,
+ tests/AddRemoveSetTests.hs view
@@ -0,0 +1,114 @@+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE DataKinds #-}++-- | Tests for the Data.LVar.AddRemoveSet module.++module AddRemoveSetTests(tests, runTests) where++import Control.Concurrent+import Test.Framework.Providers.HUnit +import Test.Framework (Test, defaultMain, testGroup)+import Test.HUnit (Assertion, assertEqual, assertBool, Counts(..))+import Test.Framework.TH (testGroupGenerator)+import qualified Test.HUnit as HU+import TestHelpers2 as T++import qualified Data.Set as S++import qualified Data.LVar.AddRemoveSet as ARS++import Control.LVish+import Control.LVish.DeepFrz (DeepFrz(..), Frzn, Trvrsbl, runParThenFreeze, runParThenFreezeIO)+import Control.LVish.Internal (liftIO)++--------------------------------------------------------------------------------++tests :: Test+tests = $(testGroupGenerator)++runTests :: IO ()+runTests = T.defaultMainSeqTests [tests]++--------------------------------------------------------------------------------++case_v1 :: Assertion+case_v1 = v1 >>= assertEqual "freeze with 3 elements"+ (S.fromList [1..3] :: S.Set Int)++-- If you have a computation that does freezing, you have to run it with runParIO.+v1 :: IO (S.Set Int)+v1 = runParIO $+ do s <- ARS.newEmptySet+ ARS.insert 1 s+ ARS.insert 2 s+ ARS.insert 3 s+ ARS.waitAddedSize 3 s+ ARS.freezeSet s++case_v2 :: Assertion+case_v2 = v2 >>= assertEqual "freeze with 10 elements added, asynchronously"+ (S.fromList [1.. v2size] :: S.Set Int)++v2 :: IO (S.Set Int)+v2 = runParIO $+ do s <- ARS.newEmptySet+ mapM_ (\n -> fork $ do+ -- liftIO$ threadDelay 5000 + logDbgLn 3$ " [AR-v2] Doing one insert: "++show n+ ARS.insert n s) [1.. v2size]+ logDbgLn 3$ " [AR-v2] now waiting.."+ ARS.waitAddedSize v2size s+ logDbgLn 3$ " [AR-v2] now freezing.."+ ARS.freezeSet s++v2size =+ case numElems of+ Just x -> x+ Nothing -> 10++case_v3 :: Assertion+case_v3 = stressTest T.stressTestReps 15 v3 (\()->True)++-- "freeze with 3 elements added, asynchronously"+-- If we're doing a guaranteed-deterministic computation we can't+-- actually read out the contents of the set.+v3 :: Par d s ()+v3 = + do s <- ARS.newEmptySet+ mapM_ (\n -> fork $ ARS.insert n s) [1..10]+ ARS.waitAddedSize 10 s++-- Getting occasional failures here with -N2, don't know what's+-- wrong. :(+case_v4 :: Assertion+case_v4 = stressTest T.stressTestReps 30 v4 (== (S.fromList [1..10] :: S.Set Int))++-- "additions and removals"+v4 :: Par QuasiDet s (S.Set Int)+v4 = + do s <- ARS.newEmptySet+ mapM_ (\n -> fork $ ARS.insert n s) [1..15]+ mapM_ (\n -> fork $ ARS.remove n s) [11..15]+ ARS.waitAddedSize 15 s + ARS.waitRemovedSize 5 s + ARS.freezeSet s++-- This one is intentionally undersynchronized.+case_i1 :: Assertion+case_i1 = do+ allowSomeExceptions ["Attempt to change a frozen LVar"] $+ do x <- i1+ assertEqual "additions and removals, undersynchronized"+ (S.fromList [1..10] :: S.Set Int) x+ return ()++-- Unblock too early, leaving a put-after-freeze possibility.+i1 :: IO (S.Set Int)+i1 = runParIO $+ do s <- ARS.newEmptySet+ mapM_ (\n -> fork $ ARS.insert n s) [1..15]+ mapM_ (\n -> fork $ ARS.remove n s) [11..15]+ -- If we don't wait for 15 additions, they might not all be+ -- there when we check.+ ARS.waitRemovedSize 5 s + ARS.freezeSet s
tests/ArrayTests.hs view
@@ -178,14 +178,14 @@ -- | Here's the same test as v9e with an actual array of IVars. -- This one is reliable, but takes about 0.20-0.30 seconds.-case_v9f1_ivarArr :: Assertion+case_v9f1_fillIvarArr :: Assertion -- [2013.08.05] RRN: Actually I'm seeing the same non-deterministic -- thread-blocked-indefinitely problem here. -- [2013.12.13] It can even happen at NUMELEMS=1000 (with debug messages slowing it) -- Could this possibly be a GHC bug? -- [2013.12.13] Runaway duplication of callbacks is ALSO possible on this test. -- Bafflingly that happens on DEBUG=2 but not 5.-case_v9f1_ivarArr = assertEqual "Array of ivars, compare effficiency:" out9e =<< v9f+case_v9f1_fillIvarArr = assertEqual "Array of ivars, compare effficiency:" out9e =<< v9f v9f :: IO Word64 v9f = runParIO$ do let size = in9e@@ -203,9 +203,9 @@ loop (acc+v) (ix+1) loop 0 0 --- | A variation of the previous. In this version-case_v9f2 :: Assertion-case_v9f2 = assertEqual "Array of ivars, compare effficiency:" out9e =<< runParIO (do +-- | A variation of the previous, change the order work is spawned to tickle the scheduler differently.+case_v9f2_seq_fillIvarArray :: Assertion+case_v9f2_seq_fillIvarArray = assertEqual "Array of ivars, compare effficiency:" out9e =<< runParIO (do let size = in9e news = V.replicate size IV.new arr <- V.sequence news@@ -213,15 +213,19 @@ logDbgLn 1 " [v9f2] Beginning putter loop.." forM_ [0..size-1] $ \ix -> IV.put_ (arr V.! ix) (fromIntegral ix + 1)- logDbgLn 1 " [v9f2] After fork."+ logDbgLn 1 " [v9f2] After puts are complete." let loop !acc ix | ix == size = return acc | otherwise = do v <- IV.get (arr V.! ix) when (ix `mod` 1000 == 0) $ logDbgLn 2 $ " [v9f2] get completed at: "++show ix++" -> "++show v loop (acc+v) (ix+1) fut <- spawn (loop 0 0)- putters- IV.get fut)+ putters + res <- IV.get fut -- Parallel+ logDbgLn 1 " [v9f2] Test is DONE."+ return res+ -- putters; loop 0 0 -- Sequential+ ) --------------------------------------------------------------------------------
tests/LVishAndIVar.hs view
@@ -7,7 +7,7 @@ -- | Core tests for the LVish scheduler and basic futures/IVars. -module LVishAndIVar(tests, runTests) where+module LVishAndIVar(tests,runTests, runParStress, lotsaRunPar) where import Test.Framework.Providers.HUnit import Test.Framework (Test, defaultMain, testGroup)@@ -64,7 +64,9 @@ -- | This stress test does nothing but run runPar again and again. case_runParStress :: HU.Assertion-case_runParStress = stressTest T.stressTestReps 15 (return ()) (\()->True)+case_runParStress = runParStress+runParStress :: HU.Assertion+runParStress = stressTest T.stressTestReps 15 (return ()) (\()->True) -- TEMP: another version that uses the simplest possible method to run lots of runPars. -- Nothing else that could POSSIBLY get in the way.@@ -82,14 +84,16 @@ -- can't make the runtime use more capabilities than we fork par worker threads. -- This could be a GHC runtime bug relating to thread migration? case_lotsaRunPar :: Assertion-case_lotsaRunPar = loop iters+case_lotsaRunPar = lotsaRunPar+lotsaRunPar = loop iters where iters = 5000+ threads = 15 -- numCapabilities loop 0 = putStrLn "" loop i = do -- We need to do runParIO to make sure the compiler does the runPar each time. -- runParIO (return ()) -- Can't crash this one.- runParDetailed (DbgCfg Nothing [] True) 15 (return ()) + runParDetailed (DbgCfg Nothing [] False) threads (return ()) -- This version can start going RIDICULOUSLY slowly with -N20. It will use <20% CPU while it does it. -- But it won't use much memory either... what is it doing? With -N4 it goes light years faster, and with -N2 -- faster yet. Extra capabilities result in a crazy slowdown here.@@ -101,11 +105,10 @@ -- -qm seems to EXACERBATE the problem, making it happen from the start and consistently. -- (even then, it is fine with -N15, the mismatch is the problem) -- Playing around with -C, -qb -qg -qi doesn't seem to do anything.- traceEventIO ("Finish iteration "++show (iters-i))+ -- traceEventIO ("Finish iteration "++show (iters-i)) -- For debugging I put in this traceEvent and ran with +RTS -N18 -qm -la putStr "."; hFlush stdout loop (i-1)- -- Disabling thread-variation due to below bug:
tests/Main.hs view
@@ -18,7 +18,7 @@ import qualified PureMapTests import qualified SLMapTests import qualified SetTests-import qualified MaxCounterTests+import qualified MaxPosIntTests import qualified AddRemoveSetTests #ifdef GENERIC_PAR@@ -34,7 +34,7 @@ , ArrayTests.tests , MemoTests.tests , LogicalTests.tests- , MaxCounterTests.tests+ , MaxPosIntTests.tests , SetTests.tests , PureMapTests.tests #ifdef FAILING_TESTS
− tests/MaxCounterTests.hs
@@ -1,45 +0,0 @@-{-# LANGUAGE TemplateHaskell #-}-{-# LANGUAGE DataKinds #-}---- | Tests for the Data.LVar.MaxCounter module.--module MaxCounterTests(tests, runTests) where--import Test.Framework.Providers.HUnit -import Test.Framework (Test, defaultMain, testGroup)-import Test.HUnit (Assertion, assertEqual, assertBool, Counts(..))-import Test.Framework.TH (testGroupGenerator)-import qualified Test.HUnit as HU-import TestHelpers as T--import Control.Concurrent (killThread, myThreadId)--import Data.LVar.MaxCounter-import Control.LVish hiding (put)-import Control.LVish.DeepFrz (DeepFrz(..), Frzn, Trvrsbl, runParThenFreeze, runParThenFreezeIO)-import qualified Control.LVish.Internal as I------------------------------------------------------------------------------------tests :: Test-tests = $(testGroupGenerator)--runTests :: IO ()-runTests = defaultMainSeqTests [tests]------------------------------------------------------------------------------------case_mc1 :: Assertion--- Spuriously failing currently:--- case_mc1 = assertEqual "mc1" (Just ()) $ timeOutPure 0.3 $ runPar $ do-case_mc1 = assertEqual "mc1" () $ runPar $ do- num <- newMaxCounter 0- fork $ put num 3- fork $ put num 4- waitThresh num 4--case_mc2 :: Assertion-case_mc2 = assertEqual "mc2" () $ runPar $ do- num <- newMaxCounter 0- fork $ put num 3- fork $ put num 4
+ tests/MaxPosIntTests.hs view
@@ -0,0 +1,45 @@+{-# LANGUAGE TemplateHaskell #-}+{-# LANGUAGE DataKinds #-}++-- | Tests for the Data.LVar.MaxPosInt module.++module MaxPosIntTests(tests, runTests) where++import Test.Framework.Providers.HUnit +import Test.Framework (Test, defaultMain, testGroup)+import Test.HUnit (Assertion, assertEqual, assertBool, Counts(..))+import Test.Framework.TH (testGroupGenerator)+import qualified Test.HUnit as HU+import TestHelpers as T++import Control.Concurrent (killThread, myThreadId)++import Data.LVar.MaxPosInt+import Control.LVish hiding (put)+import Control.LVish.DeepFrz (DeepFrz(..), Frzn, Trvrsbl, runParThenFreeze, runParThenFreezeIO)+import qualified Control.LVish.Internal as I++--------------------------------------------------------------------------------++tests :: Test+tests = $(testGroupGenerator)++runTests :: IO ()+runTests = defaultMainSeqTests [tests]++--------------------------------------------------------------------------------++case_mc1 :: Assertion+-- Spuriously failing currently:+-- case_mc1 = assertEqual "mc1" (Just ()) $ timeOutPure 0.3 $ runPar $ do+case_mc1 = assertEqual "mc1" () $ runPar $ do+ num <- newMaxPosInt 0+ fork $ put num 3+ fork $ put num 4+ waitThresh num 4++case_mc2 :: Assertion+case_mc2 = assertEqual "mc2" () $ runPar $ do+ num <- newMaxPosInt 0+ fork $ put num 3+ fork $ put num 4
tests/TestHelpers.hs view
@@ -368,7 +368,8 @@ res <- try (case x of Left err -> error$ "defaultMainSeqTests: "++err Right (opts,_) -> do let opts' = ((mempty{ ropt_threads= Just 1- , ropt_test_options = Just (mempty{ topt_timeout=(Just$ Just$ 3*1000*1000)})})+ , ropt_test_options = Just (mempty{ + topt_timeout=(Just$ Just defaultTestTimeout)})}) `mappend` opts) putStrLn $ " [*] Using "++ show (ropt_threads opts')++ " worker threads for testing." defaultMainWithOpts tests opts'@@ -381,3 +382,9 @@ threadDelay (30 * 1000) putStrLn " [*] Main thread exiting." exitWith e++-- | In nanoseconds.+defaultTestTimeout :: Int+-- defaultTestTimeout = 3*1000*1000+defaultTestTimeout = 10*1000*1000+-- defaultTestTimeout = 100*1000*1000