transient 0.5.9.2 → 0.6.0.0
raw patch · 8 files changed
+533/−361 lines, 8 filesdep ~base
Dependency ranges changed: base
Files
- src/Transient/Base.hs +3/−3
- src/Transient/EVars.hs +1/−1
- src/Transient/Indeterminism.hs +7/−11
- src/Transient/Internals.hs +264/−203
- src/Transient/Logged.hs +73/−60
- src/Transient/Parse.hs +100/−0
- tests/TestSuite.hs +1/−1
- transient.cabal +84/−82
src/Transient/Base.hs view
@@ -246,10 +246,10 @@ -- * Task Creation -- $taskgen , StreamData(..)-,parallel, async, waitEvents, sample, spawn, react+,parallel, async, waitEvents, sample, spawn, react, abduce -- * State management-,setData, getSData, getData, delData, modifyData, try, setState, getState, delState, getRState,setRState, modifyState+,setData, getSData, getData, delData, modifyData, modifyData', try, setState, getState, delState, getRState,setRState, modifyState -- * Thread management , threads,addThreads, freeThreads, hookedThreads,oneThread, killChilds@@ -257,7 +257,7 @@ -- * Exceptions -- $exceptions -,onException, cutExceptions, continue, catcht, throwt+,onException, onException', cutExceptions, continue, catcht, throwt -- * Utilities ,genId
src/Transient/EVars.hs view
@@ -6,7 +6,7 @@ import Control.Applicative import Control.Concurrent.STM-import Control.Monad.IO.Class+import Control.Monad.State
src/Transient/Indeterminism.hs view
@@ -11,7 +11,7 @@ -- | see <https://www.fpcomplete.com/user/agocorona/beautiful-parallel-non-determinism-transient-effects-iii> -- ------------------------------------------------------------------------------{-# LANGUAGE ScopedTypeVariables #-}+{-# LANGUAGE ScopedTypeVariables, CPP #-} module Transient.Indeterminism ( choose, choose', collect, collect', group, groupByTime ) where@@ -21,15 +21,18 @@ import Data.IORef import Control.Applicative import Data.Monoid-import Control.Concurrent+import Control.Concurrent import Data.Typeable import Control.Monad.State import GHC.Conc import Data.Time.Clock-import Control.Exception-import Data.Atomics +import Control.Exception +#ifndef ETA_VERSION+import Data.Atomics+#endif + -- | Converts a list of pure values into a transient task set. You can use the -- 'threads' primitive to control the parallelism. --@@ -90,13 +93,6 @@ Just xs -> return xs ---- XXX Collect might return before collecting @n@ results if the runtime--- detects that there is absolutely no possibility of more tasks to be--- collected (using 'BlockedIndefinitelyOnMVar'). This is inconsistent with the--- behavior of group. We should perhaps return 'stop' in the failure case to--- keep it consistent with group.--- -- | Collect the results of the first @n@ tasks. Synchronizes concurrent tasks -- to collect the results safely and kills all the non-free threads before -- returning the results. Results are returned in the thread where 'collect'
src/Transient/Internals.hs view
@@ -24,6 +24,7 @@ {-# LANGUAGE InstanceSigs #-} {-# LANGUAGE ConstraintKinds #-} + module Transient.Internals where import Control.Applicative@@ -35,34 +36,39 @@ import Control.Exception hiding (try,onException) import qualified Control.Exception (try) import Control.Concurrent-import GHC.Conc(unsafeIOToSTM)-import Control.Concurrent.STM hiding (retry)-import qualified Control.Concurrent.STM as STM (retry)+-- import GHC.Real+--import GHC.Conc(unsafeIOToSTM)+-- import Control.Concurrent.STM hiding (retry)+-- import qualified Control.Concurrent.STM as STM (retry) import System.Mem.StableName import Data.Maybe- import Data.List import Data.IORef import System.Environment import System.IO import qualified Data.ByteString.Char8 as BS-import Data.Atomics import Data.Typeable +#ifndef ETA_VERSION+import Data.Atomics+#endif+ #ifdef DEBUG -import Data.Monoid import Debug.Trace import System.Exit +tshow :: Show a => a -> x -> x+tshow= Debug.Trace.traceShow {-# INLINE (!>) #-} (!>) :: Show a => b -> a -> b (!>) x y = trace (show y) x infixr 0 !> #else-+tshow :: Show a => a -> x -> x+tshow _ y= y {-# INLINE (!>) #-} (!>) :: a -> b -> a (!>) = const@@ -70,8 +76,14 @@ #endif +#ifdef ETA_VERSION+atomicModifyIORefCAS = atomicModifyIORef+#endif+ type StateIO = StateT EventF IO ++ newtype TransIO a = Transient { runTrans :: StateIO (Maybe a) } type SData = ()@@ -114,7 +126,7 @@ -- ^ Label the thread with its lifecycle state and a label string } deriving Typeable -+ instance MonadState EventF TransIO where get = Transient $ get >>= return . Just@@ -254,7 +266,7 @@ hasWait _ = False appf k = Transient $ do- Log rec _ full <- getData `onNothing` return (Log False [] [])+ Log rec _ full _ <- getData `onNothing` return (Log False [] [] 0) (liftIO $ writeIORef rf (Just k,full)) -- !> ( show $ unsafePerformIO myThreadId) ++"APPF" @@ -268,7 +280,7 @@ return $ Just k <*> x appg x = Transient $ do- Log rec _ full <- getData `onNothing` return (Log False [] [])+ Log rec _ full _ <- getData `onNothing` return (Log False [] [] 0) liftIO $ writeIORef rg (Just x, full) -- !> ( show $ unsafePerformIO myThreadId) ++ "APPG"@@ -292,7 +304,7 @@ was <- getData `onNothing` return NoRemote when (was == WasParallel) $ setData NoRemote - Log recovery _ full <- getData `onNothing` return (Log False [] [])+ Log recovery _ full _ <- getData `onNothing` return (Log False [] [] 0) @@ -315,7 +327,7 @@ x <- runTrans g -- !> ( show $ unsafePerformIO myThreadId) ++ "RUN g" - Log recovery _ full' <- getData `onNothing` return (Log False [] [])+ Log recovery _ full' _ <- getData `onNothing` return (Log False [] [] 0) liftIO $ writeIORef rg (x,full') restoreStack fs k'' <- if was== WasParallel@@ -375,19 +387,33 @@ justx -> return justx -readWithErr :: (Typeable a, Read a) => String -> IO [(a, String)]-readWithErr line =+readWithErr :: (Typeable a, Read a) => Int -> String -> IO [(a, String)]+readWithErr n line = (v `seq` return [(v, left)])- `catch` (\(e :: SomeException) ->+ `catch` (\(_ :: SomeException) -> error $ "read error trying to read type: \"" ++ show (typeOf v) ++ "\" in: " ++ " <" ++ show line ++ "> ")- where [(v, left)] = readsPrec 0 line+ where (v, left):_ = readsPrec n line -readsPrec' _ = unsafePerformIO . readWithErr+read' s= case readsPrec' 0 s of+ [(x,"")] -> x+ _ -> error $ "reading " ++ s+ +readsPrec' n = unsafePerformIO . readWithErr n -- | Constraint type synonym for a value that can be logged. type Loggable a = (Show a, Read a, Typeable a) +-- data Serializable a where+-- serialize :: a .ç-> BS.ByteString+-- deserialize :: BS.ByteString -> a ++-- instance Serialize a => Serialie [a] where+-- serialize (x:xs)= +-- let s= serialize x+-- l= length s+-- in makeByteString (#(#l,s#),serialize xs #)+ -- | Dynamic serializable data for logging. data IDynamic = IDyns String@@ -403,11 +429,12 @@ type Recover = Bool type CurrentPointer = [LogElem] type LogEntries = [LogElem]+type Hash = Int data LogElem = Wait | Exec | Var IDynamic deriving (Read, Show) -data Log = Log Recover CurrentPointer LogEntries+data Log = Log Recover CurrentPointer LogEntries Hash deriving (Typeable, Show) data RemoteStatus = WasRemote | WasParallel | NoRemote@@ -419,9 +446,9 @@ stop :: Alternative m => m stopped stop = empty ---instance (Num a,Eq a,Fractional a) =>Fractional (TransIO a)where--- mf / mg = (/) <$> mf <*> mg--- fromRational (x:%y) = fromInteger x % fromInteger y+instance (Num a,Eq a,Fractional a) =>Fractional (TransIO a)where+ mf / mg = (/) <$> mf <*> mg+ fromRational r = return $ fromRational r instance (Num a, Eq a) => Num (TransIO a) where@@ -617,18 +644,35 @@ -- | Return the state variable of the type desired with which a thread, identified by his number in the treee was initiated showState :: (Typeable a, MonadIO m, Alternative m) => String -> EventF -> m (Maybe a) showState th top = resp- where resp = do+ where + resp = do let thstring = drop 9 . show $ threadId top if thstring == th then getstate top else do sts <- liftIO $ readMVar $ children top foldl (<|>) empty $ map (showState th) sts- getstate st =+ getstate st = case M.lookup (typeOf $ typeResp resp) $ mfData st of Just x -> return . Just $ unsafeCoerce x Nothing -> return Nothing- typeResp :: m (Maybe x) -> x+ typeResp :: m (Maybe x) -> x+ typeResp = undefined++-- | return all the states of the type desired that are created by direct child threads+processStates :: Typeable a => (a-> TransIO ()) -> EventF -> TransIO()+processStates display st = do+ + getstate st >>= display+ liftIO $ print $ threadId st+ sts <- liftIO $ readMVar $ children st+ mapM_ (processStates display) sts+ where + getstate st =+ case M.lookup (typeOf $ typeResp display) $ mfData st of+ Just x -> return $ unsafeCoerce x+ Nothing -> empty+ typeResp :: (a -> TransIO()) -> a typeResp = undefined -- | Add n threads to the limit of threads. If there is no limit, the limit is set.@@ -721,15 +765,15 @@ typeResp :: m (Maybe x) -> x typeResp = undefined ++ -- | Retrieve a previously stored data item of the given data type from the--- monad state. The data type to retrieve is implicitly determined from the--- requested type context.--- If the data item is not found, an 'empty' value (a void event) is returned.--- Remember that an empty value stops the monad computation. If you want to--- print an error message or a default value in that case, you can use an--- 'Alternative' composition. For example:+-- monad state. The data type to retrieve is implicitly determined by the data type.+-- If the data item is not found, empty is executed, so the alternative computation will be executed, if any, or+-- Otherwise, the computation will stop.. +-- If you want to print an error message or a default value, you can use an 'Alternative' composition. For example: ----- > getSData <|> error "no data"+-- > getSData <|> error "no data of the type desired" -- > getInt = getSData <|> return (0 :: Int) getSData :: Typeable a => TransIO a getSData = Transient getData@@ -775,6 +819,19 @@ Just x -> Just $ unsafeCoerce x Nothing -> Nothing +-- | Either modify according with the first parameter or insert according with the second, depending on if the data exist or not.+--+-- > runTransient $ do modifyData1 (\h -> h ++ " world") "hello new" ; r <- getSData ; liftIO $ putStrLn r -- > "hello new"+-- > runTransient $ do setData "hello" ; modifyData1 (\h -> h ++ " world") "hello new" ; r <- getSData ; liftIO $ putStrLn r -- > "hello world"+modifyData' :: (MonadState EventF m, Typeable a) => (a -> a) -> a ->m a+modifyData' f v= do + st <- get+ let (ma,nmap)= M.insertLookupWithKey alterf t (unsafeCoerce v) (mfData st) + put st { mfData =nmap}+ return $ if isNothing ma then v else unsafeCoerce $ fromJust ma+ where t = typeOf v+ alterf _ _ x = unsafeCoerce $ f $ unsafeCoerce x+ -- | Same as modifyData modifyState :: (MonadState EventF m, Typeable a) => (Maybe a -> Maybe a) -> m () modifyState = modifyData@@ -801,11 +858,11 @@ -- Initialized the first time it is set. setRState:: Typeable a => a -> TransIO () setRState x= do- Ref ref <- getSData- liftIO $ atomicModifyIORefCAS ref $ const (x,())+ Ref ref <- getSData+ liftIO $ atomicModifyIORefCAS ref $ const (x,()) <|> do- ref <- liftIO (newIORef x)- setData $ Ref ref+ ref <- liftIO (newIORef x)+ setData $ Ref ref getRState :: Typeable a => TransIO a getRState= do@@ -829,6 +886,12 @@ sd <- gets mfData mx <*** modify (\s ->s { mfData = sd}) + -- | generates an identifier that is unique within the current program execution+genGlobalId :: MonadIO m => m Int+genGlobalId= liftIO $ atomicModifyIORefCAS rglobalId $ \n -> (n +1,n)++rglobalId= unsafePerformIO $ newIORef (0 :: Int)+ -- | Generator of identifiers that are unique within the current monadic -- sequence They are not unique in the whole program. genId :: MonadState EventF m => m Int@@ -842,8 +905,9 @@ getPrevId = gets mfSequence instance Read SomeException where+ readsPrec n str = [(SomeException $ ErrorCall s, r)]- where [(s , r)] = read str+ where [(s , r)] = readsPrec n str -- | 'StreamData' represents a task in a task stream being generated. data StreamData a =@@ -853,7 +917,7 @@ | SError SomeException -- ^ An error occurred deriving (Typeable, Show,Read) --- | An task stream generator that produces an infinite stream of tasks by+-- | A task stream generator that produces an infinite stream of tasks by -- running an IO computation in a loop. A task is triggered carrying the output -- of the computation. See 'parallel' for notes on the return value. waitEvents :: IO a -> TransIO a@@ -888,7 +952,7 @@ spawn :: IO a -> TransIO a spawn = freeThreads . waitEvents --- | An task stream generator that produces an infinite stream of tasks by+-- | A task stream generator that produces an infinite stream of tasks by -- running an IO computation periodically at the specified time interval. The -- task carries the result of the computation. A new task is generated only if -- the output of the computation is different from the previous one. See@@ -987,7 +1051,7 @@ $ \me -> do case me of- Left e -> exceptBack cont e >> return () !> "exceptBack 2" + Left e -> exceptBack cont e >> return () -- !> "exceptBack 2" @@ -1011,8 +1075,8 @@ mask $ \restore -> forkIO $ Control.Exception.try (restore action) >>= and_then free th env= do- -- return () !> ("freeing",th,"in",threadId env)- let sibling= children env+-- return () !> ("freeing",th,"in",threadId env)+ let sibling= children env (sbs',found) <- modifyMVar sibling $ \sbs -> do let (sbs', found) = drop [] th sbs@@ -1063,7 +1127,7 @@ mapM_ (killChildren . children) ths - mapM_ (killThread . threadId) ths -- !> ("KILL", map threadId ths )+ mapM_ (killThread . threadId) ths !> ("KILL", map threadId ths ) @@ -1088,11 +1152,11 @@ -> IO response -> TransIO eventdata react setHandler iob= Transient $ do- cont <- get+ cont <- get case event cont of Nothing -> do liftIO $ setHandler $ \dat ->do- runStateT (runCont cont) cont{event= Just $ unsafeCoerce dat}+ runStateT (runCont cont) cont{event= Just $ unsafeCoerce dat} `catch` exceptBack cont iob was <- getData `onNothing` return NoRemote when (was /= WasRemote) $ setData WasParallel@@ -1102,133 +1166,147 @@ put cont{event=Nothing} return $ unsafeCoerce j --- | Runs a computation asynchronously without generating any events. Returns--- 'empty' in an 'Alternative' composition.-+-- | Runs the rest of the computation in a new thread. Returns 'empty' to the current thread abduce = async $ return () +-- * Non-blocking keyboard input +-- getLineRef= unsafePerformIO $ newTVarIO Nothing --- * non-blocking keyboard input -getLineRef= unsafePerformIO $ newTVarIO Nothing+-- | listen stdin and triggers a new task every time the input data+-- matches the first parameter. The value contained by the task is the matched+-- value i.e. the first argument itself. The second parameter is a message to the user for+-- the user. The label is displayed in the console when the option match.+option :: (Typeable b, Show b, Read b, Eq b) =>+ b -> String -> TransIO b+option = optionf False +-- Implements the same functionality than `option` but only wait for one input+option1 :: (Typeable b, Show b, Read b, Eq b) =>+ b -> String -> TransIO b+option1= optionf True -roption= unsafePerformIO $ newMVar [] --- | Waits on stdin in a loop and triggers a new task every time the input data--- matches the first parameter. The value contained by the task is the matched--- value i.e. the first argument itself. The second parameter is a label for--- the option. The label is displayed on the console when the option is--- activated.------ Note that if two independent invocations of 'option' are expecting the same--- input, only one of them gets it and triggers a task. It cannot be--- predicted which one gets it.----option :: (Typeable b, Show b, Read b, Eq b) =>- b -> String -> TransIO b-option ret message= do- let sret= show ret+optionf :: (Typeable b, Show b, Read b, Eq b) =>+ Bool -> b -> String -> TransIO b+optionf flag ret message = do + let sret= if typeOf ret == typeOf "" then unsafeCoerce ret else show ret+ liftIO $ putStrLn $ "Enter "++sret++"\tto: " ++ message+ inputf flag sret message Nothing ( == sret)+ liftIO $ putStr "\noption: " >> putStrLn sret+ -- abduce+ return ret+ +inputf flag ident message mv cond= do+ r <- react (addListener ident) (return ()) - liftIO $ putStrLn $ "Enter "++sret++"\tto: " ++ message- liftIO $ modifyMVar_ roption $ \msgs-> return $ sret:msgs- waitEvents $ getLine' (==ret)- liftIO $ putStr "\noption: " >> putStrLn (show ret)- return ret+ when (null r) $ liftIO $ writeIORef rconsumed True + let rr= read1 r + when flag $ liftIO $ delListener ident+ case rr of+ Just x -> if cond x + then do + liftIO $ do+ writeIORef rconsumed True + print x; + return x + else do liftIO $ when (isJust mv) $ putStrLn ""; returnm mv+ _ -> do liftIO $ when (isJust mv) $ putStrLn ""; returnm mv + where+ returnm (Just x)= return x + returnm _ = empty+ + read1 s= x where+ x= if typeOf(typeOfr x) == typeOf "" + then Just $ unsafeCoerce s+ else unsafePerformIO $ do+ (let r= read s in r `seq` return (Just r)) `catch` \(e :: SomeException) -> (return Nothing)+ typeOfr :: Maybe a -> a+ typeOfr = undefined --- | Waits on stdin and triggers a task when a console input matches the+-- | Waits on stdin and return a value when a console input matches the -- predicate specified in the first argument. The second parameter is a string -- to be displayed on the console before waiting.----input cond prompt= input' Nothing cond prompt+input :: (Typeable a, Read a,Show a) => (a -> Bool) -> String -> TransIO a+input= input' Nothing input' :: (Typeable a, Read a,Show a) => Maybe a -> (a -> Bool) -> String -> TransIO a-input' mv cond prompt= Transient . liftIO $do- putStr prompt >> hFlush stdout- atomically $ do- mr <- readTVar getLineRef- case mr of- Nothing -> STM.retry- Just r ->- case reads2 r of- (s,_):_ -> if cond s !> show (cond s)- then do- unsafeIOToSTM $ print s- writeTVar getLineRef Nothing !>"match"- return $ Just s+input' mv cond prompt= do+ liftIO $ putStr prompt >> hFlush stdout+ inputf True "input" prompt mv cond+ - else return mv- _ -> return mv !> "return "+rcb= unsafePerformIO $ newIORef M.empty :: IORef (M.Map String (String -> IO())) - where- reads2 s= x where- x= if typeOf(typeOfr x) == typeOf "" - then unsafeCoerce[(s,"")] - else unsafePerformIO $ return (reads s) `catch` \(e :: SomeException) -> (return [])+addListener :: String -> (String -> IO ()) -> IO ()+addListener name cb= atomicModifyIORef rcb $ \cbs -> (M.insert name cb cbs,()) - typeOfr :: [(a,String)] -> a- typeOfr = undefined+delListener :: String -> IO ()+delListener name= atomicModifyIORef rcb $ \cbs -> (M.delete name cbs,())+ --- | Non blocking `getLine` with a validator-getLine' cond= do- atomically $ do- mr <- readTVar getLineRef- case mr of- Nothing -> STM.retry- Just r ->- case reads1 r of -- !> ("received " ++ show r ++ show (unsafePerformIO myThreadId)) of- (s,_):_ -> if cond s -- !> show (cond s)- then do- writeTVar getLineRef Nothing -- !>"match"- return s - else STM.retry- _ -> STM.retry- reads1 s=x where x= if typeOf(typeOfr x) == typeOf "" then unsafeCoerce[(s,"")] else readsPrec' 0 s typeOfr :: [(a,String)] -> a typeOfr = undefined - inputLoop= do- r <- getLine- -- XXX hoping that the previous value has been consumed by now.- -- otherwise its just lost by overwriting.- atomically $ writeTVar getLineRef Nothing- processLine r- inputLoop+ x <- getLine + processLine x+ putStr "> "; hFlush stdout+ inputLoop -processLine r= do+rconsumed = unsafePerformIO $ newIORef False - let rs = breakSlash [] r+processLine r = do+ let rs = breakSlash [] r+ return () !> rs+ mapM' invoke rs - -- XXX this blocks forever if an input is not consumed by any consumer.- -- e.g. try this "xxx/xxx" on the stdin- liftIO $ mapM_ (\ r ->- atomically $ do--- threadDelay 1000000- t <- readTVar getLineRef- when (isJust t) STM.retry- writeTVar getLineRef $ Just r ) rs+ where+ invoke x= do+ mbs <- readIORef rcb+ mapM (\cb -> cb x) $ M.elems mbs + + mapM' f []= return ()+ mapM' f (xss@(x:xs)) =do + f x + r <- readIORef rconsumed+ if r + then do+ writeIORef riterloop 0+ writeIORef rconsumed False+ mapM' f xs + else do+ threadDelay 1000+ n <- atomicModifyIORef riterloop $ \n -> (n+1,n)+ if n==100+ then do+ when (not $ null x) $ putStr x >> putStrLn ": can't read, skip"+ writeIORef riterloop 0+ writeIORef rconsumed False+ mapM' f xs+ else mapM' f xss - where+ riterloop= unsafePerformIO $ newIORef 0 breakSlash :: [String] -> String -> [String] breakSlash [] ""= [""] breakSlash s ""= s breakSlash res ('\"':s)=- let (r,rest) = span(/= '\"') s- in breakSlash (res++[r]) $ tail1 rest-+ let (r,rest) = span(/= '\"') s+ in breakSlash (res++[r]) $ tail1 rest+ breakSlash res s=- let (r,rest) = span(\x -> x /= '/' && x /= ' ') s- in breakSlash (res++[r]) $ tail1 rest-- tail1 []=[]+ let (r,rest) = span(\x -> x /= '/' && x /= ' ') s+ in breakSlash (res++[r]) $ tail1 rest+ + tail1 []= [] tail1 x= tail x+ @@ -1277,18 +1355,23 @@ -- keep :: Typeable a => TransIO a -> IO (Maybe a) keep mx = do- liftIO $ hSetBuffering stdout LineBuffering rexit <- newEmptyMVar forkIO $ do -- liftIO $ putMVar rexit $ Right Nothing runTransient $ do+ onException $ \(e :: SomeException ) -> liftIO $ putStr "keep block: " >> print e st <- get setData $ Exit rexit- (async (return ()) >> labelState "input" >> liftIO inputLoop)-+ (abduce >> labelState "input" >> liftIO inputLoop >> empty) <|> do+ option "options" "show all options"+ mbs <- liftIO $ readIORef rcb+ liftIO $ mapM_ (\c ->do putStr c; putStr "|") $ M.keys mbs+ liftIO $ putStrLn ""+ empty+ <|> do option "ps" "show threads" liftIO $ showThreads st empty@@ -1296,14 +1379,17 @@ option "log" "inspect the log of a thread" th <- input (const True) "thread number>" ml <- liftIO $ showState th st- liftIO $ print $ fmap (\(Log _ _ log) -> reverse log) ml+ liftIO $ print $ fmap (\(Log _ _ log _) -> reverse log) ml empty <|> do option "end" "exit"+ liftIO $ putStrLn "exiting..."+ abduce killChilds liftIO $ putMVar rexit Nothing empty- <|> mx++ <|> mx return () threadDelay 10000 execCommandLine@@ -1335,17 +1421,17 @@ forkIO $ execCommandLine stay rexit - execCommandLine= do- args <- getArgs- let mindex = findIndex (\o -> o == "-p" || o == "--path" ) args- when (isJust mindex) $ do- let i= fromJust mindex +1- when (length args >= i) $ do- let path= args !! i- putStr "Executing: " >> print path- processLine path+ args <- getArgs+ let mindex = findIndex (\o -> o == "-p" || o == "--path" ) args+ when (isJust mindex) $ do+ let i= fromJust mindex +1+ when (length args >= i) $ do+ let path= args !! i+ putStr "Executing: " >> print path+ processLine path + -- | Exit the main thread, and thus all the Transient threads (and the -- application if there is no more code) exit :: Typeable a => a -> TransIO a@@ -1473,29 +1559,31 @@ -- returned. -- back :: (Typeable b, Show b) => b -> TransientIO a-back reason = Transient $ do+back reason = do bs <- getData `onNothing` return (backStateOf reason) - goBackt bs !>"GOBACK"+ goBackt bs -- !>"GOBACK" where+ runClosure :: EventF -> TransIO a+ runClosure EventF { xcomp = x } = unsafeCoerce x+ runContinuation :: EventF -> a -> TransIO b+ runContinuation EventF { fcomp = fs } = (unsafeCoerce $ compose $ fs) - goBackt (Backtrack _ [] )= return Nothing + goBackt (Backtrack _ [] )= empty goBackt (Backtrack b (stack@(first : bs)) )= do setData $ Backtrack (Just reason) stack - mr <- runClosure first -- !> ("RUNCLOSURE",length stack)-+ x <- runClosure first -- !> ("RUNCLOSURE",length stack)+ return () !> "runclosure" Backtrack back _ <- getData `onNothing` return (backStateOf reason) -- !> "END RUNCLOSURE" - case mr of- Nothing -> return Nothing -- !> "END EXECUTION"- Just x -> case back of+ case back of Nothing -> runContinuation first x -- !> "FORWARD EXEC" justreason ->do setData $ Backtrack justreason bs- goBackt $ Backtrack justreason bs !> ("BACK AGAIN",back)- return Nothing + goBackt $ Backtrack justreason bs -- !> ("BACK AGAIN",back)+ empty backStateOf :: (Show a, Typeable a) => a -> Backtrack a backStateOf reason= Backtrack (Nothing `asTypeOf` (Just reason)) []@@ -1557,7 +1645,7 @@ -- | Install an exception handler. Handlers are executed in reverse (i.e. last in, first out) order when such exception happens in the -- continuation. Note that multiple handlers can be installed for the same exception type. ----- The semantic is thus very different than the one of `Control.Exception.Base.onException`+-- The semantic is, thus, very different than the one of `Control.Exception.Base.onException` onException :: Exception e => (e -> TransIO ()) -> TransIO () onException exc= return () `onException'` exc @@ -1570,19 +1658,16 @@ where onAnyException :: TransIO a -> (SomeException ->TransIO a) -> TransIO a onAnyException mx f= ioexp f `onBack` f-- ioexp f = Transient $ do+ ioexp f = Transient $ do st <- get+ (mx,st') <- liftIO $ (runStateT (do- case event st of Nothing -> do- r <- runTrans mx -+ r <- runTrans mx modify $ \s -> s{event= Just $ unsafeCoerce r}- -+ return () !> "MX" runCont st was <- getData `onNothing` return NoRemote when (was /= WasRemote) $ setData WasParallel@@ -1591,25 +1676,26 @@ Just r -> do modify $ \s -> s{event=Nothing} + return () !> "JUSTTTTTTTTTTT" return $ unsafeCoerce r) st) `catch` exceptBack st put st' return mx exceptBack st = \(e ::SomeException) -> do -- recursive catch itself- runStateT ( runTrans $ back e ) st- `catch` exceptBack (stex st)+ runStateT ( runTrans $ back e ) st !> "EXCEPTBACK"+ `catch` exceptBack st - where- -- drop the current exception handler from the stack- stex st = - let list = mfData st- emptyback=backStateOf(undefined :: SomeException)- Backtrack b stack = case M.lookup (typeOf emptyback) list of- Just x -> unsafeCoerce x- Nothing -> emptyback+ -- where+ -- -- drop the current exception handler from the stack+ -- stex st = + -- let list = mfData st+ -- emptyback= backStateOf(undefined :: SomeException)+ -- Backtrack b stack = case M.lookup (typeOf emptyback) list of+ -- Just x -> unsafeCoerce x+ -- Nothing -> emptyback - in st { mfData = M.insert (typeOf emptyback) (unsafeCoerce $ Backtrack b $ tail stack) (mfData st) }+ -- in st { mfData = M.insert (typeOf emptyback) (unsafeCoerce $ Backtrack b $ tail stack) (mfData st) } -- | Delete all the exception handlers registered till now. cutExceptions :: TransIO ()@@ -1618,7 +1704,7 @@ -- | Use it inside an exception handler. it stop executing any further exception -- handlers and resume normal execution from this point on. continue :: TransIO ()-continue = forward (undefined :: SomeException) !> "CONTINUE"+continue = forward (undefined :: SomeException) -- !> "CONTINUE" -- | catch an exception in a Transient block --@@ -1643,28 +1729,3 @@ throwt :: Exception e => e -> TransIO a throwt= back . toException ---- catcht1 :: Exception e => TransIO a -> (e ->TransIO a) -> TransIO a--- catcht1 mx exc= Transient $ do --- st <- get- --- case event st of --- Nothing -> do--- (mx,st') <- liftIO $ (runStateT(runTrans mx) st ) `catch` \(e ::SomeException) -> --- runStateT ( runTrans $ f e ) st--- put st'--- modify $ \s -> s{event= Just $ unsafeCoerce mx}--- runCont st--- was <- getData `onNothing` return NoRemote--- when (was /= WasRemote) $ setData WasParallel---- return Nothing---- Just r -> do--- modify $ \s -> s{event=Nothing} --- return $ unsafeCoerce r- --- where--- f e=case fromException e of--- Nothing -> empty--- Just e' -> exc e'
src/Transient/Logged.hs view
@@ -35,25 +35,26 @@ -- liftIO $ print (\"C",r) -- @ ------------------------------------------------------------------------------{-# LANGUAGE CPP,ExistentialQuantification, FlexibleInstances, ScopedTypeVariables, UndecidableInstances #-}+{-# LANGUAGE CPP, ExistentialQuantification, FlexibleInstances, ScopedTypeVariables, UndecidableInstances #-} module Transient.Logged(-Loggable, logged, received, param-+Loggable, logged, received, param, #ifndef ghcjs_HOST_OS-, suspend, checkpoint, restore+ suspend, checkpoint, restore, #endif+-- * low level+fromIDyn,maybeFromIDyn,toIDyn ) where import Data.Typeable import Unsafe.Coerce import Transient.Base-import Transient.Internals(Loggable)+import Transient.Internals(Loggable,read') import Transient.Indeterminism(choose) import Transient.Internals -- (onNothing,reads1,IDynamic(..),Log(..),LogElem(..),RemoteStatus(..),StateIO) import Control.Applicative-import Control.Monad.IO.Class+import Control.Monad.State import System.Directory-import Control.Exception+import Control.Exception import Control.Monad import Control.Concurrent.MVar @@ -76,7 +77,7 @@ liftIO $ createDirectory logs `catch` (\(e :: SomeException) -> return ()) list <- liftIO $ getDirectoryContents logs `catch` (\(e::SomeException) -> return [])- if length list== 2 then proc else do+ if null list || length list== 2 then proc else do let list'= filter ((/=) '.' . head) list file <- choose list' -- !> list'@@ -103,7 +104,7 @@ -- suspend :: Typeable a => a -> TransIO a suspend x= do- Log recovery _ log <- getData `onNothing` return (Log False [] [])+ Log recovery _ log _ <- getData `onNothing` return (Log False [] [] 0) if recovery then return x else do logAll log exit x@@ -112,22 +113,32 @@ -- does not exit. checkpoint :: TransIO () checkpoint = do- Log recovery _ log <- getData `onNothing` return (Log False [] [])+ Log recovery _ log _ <- getData `onNothing` return (Log False [] [] 0) if recovery then return () else logAll log -logAll log= do- newlogfile <- liftIO $ (logs ++) <$> replicateM 7 (randomRIO ('a','z'))- liftIO $ writeFile newlogfile $ show log+logAll log= liftIO $do+ newlogfile <- (logs ++) <$> replicateM 7 (randomRIO ('a','z'))+ logsExist <- doesDirectoryExist logs+ when (not logsExist) $ createDirectory logs+ writeFile newlogfile $ show log -- :: TransIO () #endif +maybeFromIDyn :: Loggable a => IDynamic -> Maybe a+maybeFromIDyn (IDynamic x)= r+ where+ r= if typeOf (Just x) == typeOf r then Just $ unsafeCoerce x else Nothing +maybeFromIDyn (IDyns s) = case reads s of+ [] -> Nothing+ [(x,"")] -> Just x+ fromIDyn :: Loggable a => IDynamic -> a fromIDyn (IDynamic x)=r where r= unsafeCoerce x -- !> "coerce" ++ " to type "++ show (typeOf r) -fromIDyn (IDyns s)=r `seq`r where r= read s -- !> "read " ++ s ++ " to type "++ show (typeOf r)+fromIDyn (IDyns s)=r `seq`r where r= read' s -- !> "read " ++ s ++ " to type "++ show (typeOf r) @@ -145,48 +156,50 @@ -- discarded. -- logged :: Loggable a => TransIO a -> TransIO a-logged mx = Transient $ do- Log recover rs full <- getData `onNothing` return ( Log False [][])- runTrans $- case (recover ,rs) of -- !> ("logged enter",recover,rs) of- (True, Var x: rs') -> do- setData $ Log True rs' full- return $ fromIDyn x--- !> ("Var:", x)-- (True, Exec:rs') -> do- setData $ Log True rs' full- mx--- !> "Exec"-- (True, Wait:rs') -> do- setData (Log True rs' full) -- !> "Wait"- empty-- _ -> do--- let add= Exec: full- setData $ Log False (Exec : rs) (Exec: full) -- !> ("setLog False", Exec:rs)-- r <- mx <** ( do -- when p1 <|> p2, to avoid the re-execution of p1 at the- -- recovery when p1 is asynchronous- r <- getSData <|> return NoRemote- case r of- WasParallel ->--- let add= Wait: full- setData $ Log False (Wait: rs) (Wait: full)- _ -> return ())-- Log recoverAfter lognew _ <- getData `onNothing` return ( Log False [][])- let add= Var (toIDyn r): full- if recoverAfter && (not $ null lognew) -- !> ("recoverAfter", recoverAfter)- then (setData $ Log True lognew (reverse lognew ++ add) )- -- !> ("recover",reverse lognew ,add)- else if recoverAfter && (null lognew) then- setData $ Log False [] add- else- (setData $ Log False (Var (toIDyn r):rs) add) -- !> ("restore", (Var (toIDyn r):rs))- return r-+logged mx = Transient $ do+ Log recover rs full hash <- getData `onNothing` return ( Log False [][] 0)+ runTrans $+ case (recover ,rs) of -- !> ("logged enter",recover,rs,reverse full) of+ (True, Var x: rs') -> do+ return () -- !> ("Var:", x)+ setData $ Log True rs' full (hash+ 10000000)+ return $ fromIDyn x+ + + (True, Exec:rs') -> do+ setData $ Log True rs' full (hash + 1000)+ mx -- !> "Exec"+ + (True, Wait:rs') -> do+ setData $ Log True rs' full (hash + 100000) + empty !> "Wait"+ + _ -> do+ -- let add= Exec: full+ setData $ Log False (Exec : rs) (Exec: full) (hash + 1000) -- !> ("setLog False", Exec:rs)+ + r <- mx <** ( do -- when p1 <|> p2, to avoid the re-execution of p1 at the+ -- recovery when p1 is asynchronous+ r <- getSData <|> return NoRemote+ case r of+ WasParallel ->+ setData $ Log False (Wait: rs) (Wait: full) (hash+ 100000) + _ -> return ())+ + Log recoverAfter lognew _ _ <- getData `onNothing` return ( Log False [][] 0)+ let add= Var (toIDyn r): full+ if recoverAfter && (not $ null lognew) -- !> ("recoverAfter", recoverAfter)+ then do+ (setData $ Log True lognew (reverse lognew ++ add) (hash + 10000000) )+ -- !> ("recover",reverse (reverse lognew ++add))+ else if recoverAfter && (null lognew) then do + -- showClosure+ setData $ Log False [] add (hash + 10000000) -- !> ("recover2",reverse add)+ else do+ -- showClosure+ (setData $ Log False (Var (toIDyn r):rs) add (hash +10000000)) -- !> ("restore", reverse $ (Var (toIDyn r):rs))+ + return r @@ -194,12 +207,12 @@ received :: Loggable a => a -> TransIO () received n=Transient $ do- Log recover rs full <- getData `onNothing` return ( Log False [][])+ Log recover rs full hash <- getData `onNothing` return ( Log False [][] 0) case rs of [] -> return Nothing Var (IDyns s):t -> if s == show1 n then do- setData $ Log recover t full+ setData $ Log recover t full hash return $ Just () else return Nothing _ -> return Nothing@@ -209,11 +222,11 @@ param :: Loggable a => TransIO a param= res where res= Transient $ do- Log recover rs full <- getData `onNothing` return ( Log False [][])+ Log recover rs full hash<- getData `onNothing` return ( Log False [][] 0) case rs of [] -> return Nothing Var (IDynamic v):t ->do- setData $ Log recover t full+ setData $ Log recover t full hash return $ cast v Var (IDyns s):t -> do let mr = reads1 s `asTypeOf` type1 res
+ src/Transient/Parse.hs view
@@ -0,0 +1,100 @@+{-#LANGUAGE ExistentialQuantification, OverloadedStrings, TypeSynonymInstances, FlexibleInstances #-} +module Transient.Parse where +import Transient.Internals +import Data.String +import Data.Typeable +import Control.Applicative +import Data.Char +import Data.Monoid +import System.IO.Unsafe +import Control.Monad +import Control.Monad.State + +import qualified Data.ByteString.Lazy.Char8 as BS + +-- | set a stream of strings to be parsed +setParseStream :: (Typeable str,IsString str) => IO str -> TransIO () +setParseStream iox= setState $ ParseContext iox "" + +-- | set a string to be parsed +setParseString :: (Typeable str,IsString str) => str -> TransIO () +setParseString x = setState $ ParseContext (error "end of parse string") x + +-- | The parse context contains either the string to be parsed or a computation that gives an stream of +-- strings or both. First, the string is parsed. If it is empty, the stream is pulled for more. +data ParseContext str = IsString str => ParseContext (IO str) str deriving Typeable + + +string :: BS.ByteString -> TransIO BS.ByteString +string s=withData $ \str -> do + let len= BS.length s + ret@(s',_) = BS.splitAt len str + if s == s' + then return ret + else empty + +-- many p= do +-- r <- p +-- rs <- many p <|> return [] +-- return $ r:rs + +manyTill p end = scan + where + scan = do{ end; return [] } + <|> + do{ x <- p; xs <- scan; return (x:xs) } + + +dropSpaces= parse $ \str ->((),BS.dropWhile isSpace str) + +dropChar= parse $ \r -> ((), BS.tail r) + +endline c= c== '\r' || c =='\n' + +--tGetLine= tTakeWhile . not . endline + +parseString= do + dropSpaces + tTakeWhile (not . isSpace) + +tTakeWhile :: (Char -> Bool) -> TransIO BS.ByteString +tTakeWhile cond= parse (BS.span cond) + +tTakeWhile' :: (Char -> Bool) -> TransIO BS.ByteString +tTakeWhile' cond= parse ((\(h,t) -> (h, if BS.null t then t else BS.tail t)) . BS.span cond) + +tTake n= parse ( BS.splitAt n) + +tChar= parse $ \s -> (BS.head s,BS.tail s) + +parse :: (BS.ByteString -> (b, BS.ByteString)) -> TransIO b +parse split= withData $ \str -> + if str== mempty then empty + else return $ split str + + +-- | bring the data of a parse context as a lazy byteString to a parser +-- and actualize the parse context with the result +withData :: (BS.ByteString -> TransIO (a,BS.ByteString)) -> TransIO a +withData parser= Transient $ do + ParseContext readMore s <- getData `onNothing` error "parser: no context" + let loop = unsafeInterleaveIO $ do + r <- readMore + (r <>) `liftM` loop + str <- liftIO $ (s <> ) `liftM` loop + mr <- runTrans $ parser str + case mr of + Nothing -> return Nothing + Just (v,str') -> do + setData $ ParseContext readMore str' + return $ Just v + +-- | bring the data of the parse context as a lazy byteString +giveData= noTrans $ do + ParseContext readMore s <- getData `onNothing` error "parser: no context" + :: StateIO (ParseContext BS.ByteString) -- change to strict BS + + let loop = unsafeInterleaveIO $ do + r <- readMore + (r <>) `liftM` loop + liftIO $ (s <> ) `liftM` loop
tests/TestSuite.hs view
@@ -25,7 +25,7 @@ mappend = (+) main= do- keep $ do+ keep' $ do let genElem :: a -> TransIO a genElem x= do isasync <- liftIO randomIO
transient.cabal view
@@ -1,82 +1,84 @@-name: transient-version: 0.5.9.2-author: Alberto G. Corona-extra-source-files:- ChangeLog.md README.md-maintainer: agocorona@gmail.com-cabal-version: >=1.10-build-type: Simple-license: MIT-license-file: LICENSE-homepage: http://www.fpcomplete.com/user/agocorona-bug-reports: https://github.com/agocorona/transient/issues-synopsis: composing programs with multithreading, events and distributed computing-description: See <http://github.com/agocorona/transient>- Distributed primitives are in the transient-universe package. Web primitives are in the axiom package.-category: Control, Concurrency-data-dir: ""--flag debug- description: Enable debugging outputs- default: False- manual: True--library- -- Note: `stack sdist/upload` will add missing bounds (via "pvp-bounds: both") in `build-depends`- -- support GHC 7.10.3 and later; lower bounds below denote GHC 7.10.3's bundled versions- build-depends: base >= 4.8.1 && < 5- , containers >= 0.5.6- , transformers >= 0.4.2- , time >= 1.5- , directory >= 1.2.2- , bytestring >= 0.10.6-- -- libraries not bundled w/ GHC- , mtl- , stm- , random- , atomic-primops--- exposed-modules: Transient.Backtrack- Transient.Base- Transient.EVars- Transient.Indeterminism- Transient.Internals- Transient.Logged-- exposed: True- buildable: True- exposed: True- default-language: Haskell2010- hs-source-dirs: src .- if flag(debug)- cpp-options: -DDEBUG--source-repository head- type: git- location: https://github.com/agocorona/transient--test-suite test-transient-- if !impl(ghcjs >=0.1)- build-depends:- base >= 4.8.1 && < 5- , containers >= 0.5.6- , transformers >= 0.4.2- , time >= 1.5- , directory >= 1.2.2- , bytestring >= 0.10.6-- -- libraries not bundled w/ GHC- , mtl- , stm- , random- , atomic-primops- type: exitcode-stdio-1.0- main-is: TestSuite.hs- build-depends:- base >4- default-language: Haskell2010- hs-source-dirs: tests src .- ghc-options: -threaded -rtsopts+name: transient +version: 0.6.0.0 +author: Alberto G. Corona +extra-source-files: + ChangeLog.md README.md +maintainer: agocorona@gmail.com +cabal-version: >=1.10 +build-type: Simple +license: MIT +license-file: LICENSE +homepage: https://github.com/transient-haskell/transient +bug-reports: https://github.com/transient-haskell/transient/issues +synopsis: composing programs with multithreading, events and distributed computing +description: See <http://github.com/agocorona/transient> + Distributed primitives are in the transient-universe package. Web primitives are in the axiom package. +category: Control, Concurrency +data-dir: "" + +flag debug + description: Enable debugging outputs + default: False + manual: True + +library + -- Note: `stack sdist/upload` will add missing bounds (via "pvp-bounds: both") in `build-depends` + -- support GHC 7.10.3 and later; lower bounds below denote GHC 7.10.3's bundled versions + build-depends: base >= 4.8.0 && < 5 + , containers >= 0.5.6 + , transformers >= 0.4.2 + , time >= 1.5 + , directory >= 1.2.2 + , bytestring >= 0.10.6 + + -- libraries not bundled w/ GHC + , mtl + , stm + , random + ,atomic-primops + + exposed-modules: Transient.Backtrack + Transient.Base + Transient.EVars + Transient.Indeterminism + Transient.Internals + Transient.Logged + Transient.Parse + + exposed: True + buildable: True + default-language: Haskell2010 + hs-source-dirs: src . + + -- eta-options: -ddump-stg -ddump-to-file + + if flag(debug) + cpp-options: -DDEBUG + +source-repository head + type: git + location: https://github.com/agocorona/transient + +test-suite test-transient + + if !impl(ghcjs >=0.1) + build-depends: + base >= 4.8.1 && < 5 + , containers >= 0.5.6 + , transformers >= 0.4.2 + , time >= 1.5 + , directory >= 1.2.2 + , bytestring >= 0.10.6 + + -- libraries not bundled w/ GHC + , mtl + , stm + , random + , atomic-primops + type: exitcode-stdio-1.0 + main-is: TestSuite.hs + build-depends: + base >4 + default-language: Haskell2010 + hs-source-dirs: tests src . + ghc-options: -threaded -rtsopts