packages feed

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 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