diff options
| author | Tom Smeding <tom@tomsmeding.com> | 2026-08-02 21:19:10 +0200 |
|---|---|---|
| committer | Tom Smeding <tom@tomsmeding.com> | 2026-08-02 21:21:54 +0200 |
| commit | 6e8c2b8df3898d0f50463f3151831f5d7c9306a2 (patch) | |
| tree | 0964eb6afd22d8c3739cbf3c66eb6c3798195c0a /src | |
| parent | faff2c8249b1cc0342867dd22774301ae980647c (diff) | |
Index logs in parallel
Diffstat (limited to 'src')
| -rw-r--r-- | src/Index.hs | 3 | ||||
| -rw-r--r-- | src/Parallel.hs | 30 |
2 files changed, 32 insertions, 1 deletions
diff --git a/src/Index.hs b/src/Index.hs index 81eb5f2..ee3cdaf 100644 --- a/src/Index.hs +++ b/src/Index.hs @@ -56,6 +56,7 @@ import Cache import Config (Channel(..), prettyChannel) import ImmutGrowVector qualified as IGV import Mmap +import Parallel import Util import ZNC.Parser @@ -151,7 +152,7 @@ initIndex basedir toimport = do let nw = T.unpack nwT ch = T.unpack chT files <- listDirectory (basedir </> nw </> ch) - days <- fmap sort . forM (sort files) $ \fn -> do + days <- fmap sort . parallelForM (sort files) $ \fn -> do -- atomicPrintS $ " -> " ++ fn let path = basedir </> nw </> ch </> fn date = case parseFileName fn of diff --git a/src/Parallel.hs b/src/Parallel.hs new file mode 100644 index 0000000..cc6eb22 --- /dev/null +++ b/src/Parallel.hs @@ -0,0 +1,30 @@ +module Parallel where + +import Control.Concurrent +import Control.Monad (replicateM_) +import Data.IORef + + +-- | Does not return results in-order. +parallelForM :: [a] -> (a -> IO b) -> IO [b] +parallelForM inputList action = do + nthread <- getNumCapabilities + + listref <- newIORef inputList + outref <- newIORef [] + donechan <- newChan + replicateM_ nthread $ forkIO $ + let loop = do + mitem <- atomicModifyIORef' listref (\case l@[] -> (l, Nothing) + item : l -> (l, Just item)) + case mitem of + Just item -> do + res <- action item + atomicModifyIORef' outref (\l -> (res : l, ())) + loop + Nothing -> do + writeChan donechan () + in loop + + replicateM_ nthread $ readChan donechan + reverse <$> readIORef outref |
