| ... |
... |
@@ -1787,9 +1787,10 @@ parDfsBuild base_map roots key expand = ReaderT $ \ds_env -> do |
|
1787
|
1787
|
visited_var <- newTVarIO (fromMaybe Map.empty base_map)
|
|
1788
|
1788
|
pending <- newTVarIO Set.empty
|
|
1789
|
1789
|
worklist <- newTQueueIO
|
|
|
1790
|
+ threads <- newTVarIO []
|
|
1790
|
1791
|
|
|
1791
|
1792
|
coord_tid <- forkIO $
|
|
1792
|
|
- coordinator ds_env exc_var visited_var worklist pending
|
|
|
1793
|
+ coordinator ds_env exc_var visited_var worklist pending threads
|
|
1793
|
1794
|
`MC.catch` \case
|
|
1794
|
1795
|
(e::MC.SomeException)
|
|
1795
|
1796
|
-- exit cleanly when killed
|
| ... |
... |
@@ -1798,11 +1799,12 @@ parDfsBuild base_map roots key expand = ReaderT $ \ds_env -> do |
|
1798
|
1799
|
-- signal the exc_var for the main thread to throw it
|
|
1799
|
1800
|
| otherwise -> atomically (modifyTVar' exc_var (<|> Just e))
|
|
1800
|
1801
|
|
|
1801
|
|
- mapM_ (atomically . writeTQueue worklist) roots
|
|
|
1802
|
+ atomically $ mapM_ (writeTQueue worklist) roots
|
|
1802
|
1803
|
|
|
1803
|
1804
|
mb_exc <- wait_done exc_var worklist pending
|
|
1804
|
1805
|
`MC.finally` do
|
|
1805
|
1806
|
killThread coord_tid
|
|
|
1807
|
+ mapM_ killThread =<< readTVarIO threads
|
|
1806
|
1808
|
|
|
1807
|
1809
|
case mb_exc of
|
|
1808
|
1810
|
Just e -> throwIO e
|
| ... |
... |
@@ -1820,7 +1822,7 @@ parDfsBuild base_map roots key expand = ReaderT $ \ds_env -> do |
|
1820
|
1822
|
unless (empty_worklist && empty_pending) retry
|
|
1821
|
1823
|
return Nothing
|
|
1822
|
1824
|
|
|
1823
|
|
- coordinator ds_env exc_var visvar worklist pendvar = forever $ do
|
|
|
1825
|
+ coordinator ds_env exc_var visvar worklist pendvar threads = forever $ do
|
|
1824
|
1826
|
mb_node_to_expand <- atomically $ do
|
|
1825
|
1827
|
node <- readTQueue worklist
|
|
1826
|
1828
|
let k = key node
|
| ... |
... |
@@ -1839,27 +1841,32 @@ parDfsBuild base_map roots key expand = ReaderT $ \ds_env -> do |
|
1839
|
1841
|
|
|
1840
|
1842
|
case mb_node_to_expand of
|
|
1841
|
1843
|
Nothing -> return ()
|
|
1842
|
|
- Just (k, node) -> void $ do
|
|
1843
|
|
- withLocalTmpFSMake (ds_make_env ds_env) $ \make_env ->
|
|
1844
|
|
- forkIOWithUnmask $ \unmask ->
|
|
|
1844
|
+ Just (k, node) -> do
|
|
|
1845
|
+ tid <- withLocalTmpFSMake (ds_make_env ds_env) $ \make_env ->
|
|
|
1846
|
+ MC.mask_ $ forkIOWithUnmask $ \unmask ->
|
|
1845
|
1847
|
unmask (worker ds_env{ds_make_env = make_env} visvar worklist pendvar k node)
|
|
1846
|
|
- -- write exception in the worker to the main thread
|
|
1847
|
|
- `MC.catch` \(e::MC.SomeException) ->
|
|
1848
|
|
- atomically (modifyTVar' exc_var (<|> Just e))
|
|
|
1848
|
+ `MC.catch` \case
|
|
|
1849
|
+ e | Just (_ :: SomeAsyncException) <- fromException e
|
|
|
1850
|
+ -> throwIO e -- async exceptions like KillThread get thrown
|
|
|
1851
|
+ | otherwise -- exceptions in workers are written for main thread
|
|
|
1852
|
+ -> atomically (modifyTVar' exc_var (<|> Just e))
|
|
|
1853
|
+
|
|
|
1854
|
+ atomically $ modifyTVar' threads (tid:)
|
|
1849
|
1855
|
|
|
1850
|
1856
|
worker ds_env@DownsweepEnv{..} visvar worklist pendvar k node =
|
|
1851
|
1857
|
withAbstractSem (compile_sem ds_make_env) $ do
|
|
1852
|
1858
|
r <- runDownsweepM ds_env $
|
|
1853
|
1859
|
expand node -- do the main work!
|
|
1854
|
1860
|
|
|
1855
|
|
- case r of
|
|
1856
|
|
- NSkip ->
|
|
1857
|
|
- atomically $ modifyTVar' visvar (Map.insert k NSkip)
|
|
1858
|
|
- NSuccess (v,ns) -> do
|
|
1859
|
|
- atomically $ modifyTVar' visvar (Map.insert k (NSuccess v))
|
|
1860
|
|
- mapM_ (atomically . writeTQueue worklist) ns
|
|
|
1861
|
+ atomically $ do
|
|
|
1862
|
+ case r of
|
|
|
1863
|
+ NSkip ->
|
|
|
1864
|
+ modifyTVar' visvar (Map.insert k NSkip)
|
|
|
1865
|
+ NSuccess (v,ns) -> do
|
|
|
1866
|
+ modifyTVar' visvar (Map.insert k (NSuccess v))
|
|
|
1867
|
+ mapM_ (writeTQueue worklist) ns
|
|
1861
|
1868
|
|
|
1862
|
|
- atomically $ modifyTVar' pendvar (Set.delete k)
|
|
|
1869
|
+ modifyTVar' pendvar (Set.delete k)
|
|
1863
|
1870
|
|
|
1864
|
1871
|
{-
|
|
1865
|
1872
|
Note [Downsweep Control Flow and Caching]
|