packages feed

time-manager 0.2.4 → 0.4.0

raw patch · 7 files changed

Files

ChangeLog.md view
@@ -1,5 +1,62 @@ # ChangeLog for time-manager +## 0.4.0++* CHANGES IN BEHAVIOUR:+  * `tickle` is rate-limited/debounced. The renewal is skipped unless a quarter+    of the timeout (capped at one second) has passed since the timeout was last+    registered or updated. This does mean a timeout _might_ run a bit earlier+    than the last `tickle` would indicate, but never more than the maximum+    debounce period.+  * `cancel` completely stops the timeout, making it un`resume`able.+    `resume` will only resume a timeout that has been `pause`d.+  * Prior to this major version, the `Handle` could be reused to run more+    timeout actions. Now, a timeout action will only ever run, at most, once.+    After a timeout action has run, the `Handle` is turned off and won't be+    `resume`able, necessitating a call to `register` to start a new timeout.++  [#1109](https://github.com/yesodweb/wai/pull/1109)++## 0.3.2++* Add `stopAfterWithResult`. [#1069](https://github.com/yesodweb/wai/pull/1069)++## 0.3.1.1++* Added `README.md` file and improved documentation of modules and functions.+  [#1057](https://github.com/yesodweb/wai/pull/1057)++## 0.3.1++BUGFIXES:++* [#1055](https://github.com/yesodweb/wai/pull/1055)+  * `resume` now acts as a `tickle` if the `Handle` isn't paused.+    This is the same behaviour as before version `0.3.0`.+  * `registerKillThread` now throws the `TimeoutThread` in a separate+    thread, so as to not block the GHC's System TimerManager.+    This does mean that `TimeoutThread` exceptions could technically+    be thrown "out of order", but they will be more prompt.+++## 0.3.0++* [#1048](https://github.com/yesodweb/wai/pull/1048)+  * New architecture. The backend is switched from the designated thread+    to GHC's System TimerManager. From this version, this library is+    just wrapper APIs of GHC's System TimerManager. Unlike v0.2 or+    earlier, callbacks are executed at the exact time. System+    TimerManager uses a PSQ (a tree) while v0.2 or earlier uses a list.+    So, this version hopefully scales better.+  * Deprecated functions: `stopManager`, `killManager` and `withManager'`.+  * `tickle` sets the specified timeout from now.+  * `pause` is now identical to `cancel`.+  * `resume` is now re-registration of timeout.+  * The signature of `withHandle` is changed.++This change also means that using this package only works with the threaded runtime.+The moment a timeout is registered on a non-threaded runtime, an exception will be thrown.+ ## 0.2.4  * Providing `isAllGone`.
+ README.md view
@@ -0,0 +1,13 @@+## time-manager++This package provides module to let you run actions with resettable timeouts+(i.e. `System.TimeManager`) and run actions that will make sure that all+threads forked with the given `ThreadManager` will be killed when the action+finishes.++## WARNINGS++Since version `0.3.0`, **the timeout manager relies on GHC internals**.+This change also means that using **this package only works with the threaded**+**runtime**. The moment a timeout is registered on a non-threaded runtime, an+exception will be thrown.
System/ThreadManager.hs view
@@ -4,10 +4,22 @@  -- | A thread manager including a time manager. --   The manager has responsibility to kill managed threads.+--+-- Because this is based on the accompanying "System.TimeManager" module,+-- the same caveats apply:+--+--   * Only works for GHC.+--   * Only works with a threaded runtime.+--   * Users of older versions should check the current semantics.+--   * Using 32-bit systems means the max timeout is @'maxBound' :: Int@+--     (2147483647) microseconds, which is less than 36 minutes.+--   * Using the same 'Handle' in different threads might cause issues in some+--     edge cases. (i.e. using cancel/pause in one thread, and resume in another) module System.ThreadManager (     ThreadManager,     newThreadManager,     stopAfter,+    stopAfterWithResult,     KilledByThreadManager (..),      -- * Fork@@ -36,7 +48,7 @@ import qualified Control.Exception as E import Control.Monad (unless, void) import Data.Foldable (forM_)-import Data.IORef+import Data.IORef (IORef, atomicModifyIORef', newIORef) import Data.Map.Strict (Map) import qualified Data.Map.Strict as Map import Data.Word (Word64)@@ -65,10 +77,12 @@  ---------------------------------------------------------------- --- | Starting a thread manager.---   Its action is initially set to 'return ()' and should be set---   by 'setAction'. This allows that the action can include---   the manager itself.+-- | Create a thread manager.+--+-- To create a 'ThreadManager', you will first have to create a+-- 'T.Manager' from the "System.TimeManager" module.+--+-- You can use either 'System.TimeManager.initialize' or 'System.TimeManager.withManager'. newThreadManager :: T.Manager -> IO ThreadManager newThreadManager timmgr = ThreadManager timmgr <$> newTVarIO Map.empty @@ -84,12 +98,34 @@  -- | Stopping the manager. --+-- @+-- stopAfter threadManager action cleanup+-- @+-- -- The action is run in the scope of an exception handler that catches all -- exceptions (including asynchronous ones); this allows the cleanup handler -- to cleanup in all circumstances. If an exception is caught, it is rethrown -- after the cleanup is complete. stopAfter :: ThreadManager -> IO a -> (Maybe SomeException -> IO ()) -> IO a-stopAfter (ThreadManager _timmgr var) action cleanup = do+stopAfter mgr action cleanup =+    stopAfterWithResult mgr action $ \mResult ->+        case mResult of+            Left err -> do+                cleanup (Just err)+                E.throwIO err+            Right result -> do+                cleanup Nothing+                return result++-- | Generalization of 'stopAfter' where the cleanup handler is allowed to construct the final result+--+-- Unlike in 'stopAfter', if an exception is thrown, it is not re-thrown after the cleanup+-- handler completes; the cleanup handler itself can decide to rethrow it or compute a result.+--+-- @since 0.3.2+stopAfterWithResult+    :: ThreadManager -> IO a -> (Either SomeException a -> IO b) -> IO b+stopAfterWithResult (ThreadManager _timmgr var) action cleanup = do     E.mask $ \unmask -> do         ma <- E.try $ unmask action         m <- atomically $ do@@ -100,9 +136,7 @@             er = either Just (const Nothing) ma             ex = KilledByThreadManager er         forM_ ths $ \(ManagedThread wtid ref) -> lockAndKill wtid ref ex-        case ma of-            Left err -> cleanup (Just err) >> E.throwIO err-            Right a -> cleanup Nothing >> return a+        cleanup ma  ---------------------------------------------------------------- @@ -123,42 +157,65 @@  -- | Like 'forkManaged', but run action with exceptions masked forkManagedUnmask-    :: ThreadManager -> String -> ((forall x. IO x -> IO x) -> IO ()) -> IO ()+    :: ThreadManager+    -> String+    -- ^ Thread name+    -> ((forall x. IO x -> IO x) -> IO ())+    -- ^ Action with unmask argument+    -> IO () forkManagedUnmask (ThreadManager _timmgr var) label io =     void $ E.mask_ $ forkIOWithUnmask $ \unmask -> E.handle ignore $ do         labelMe label         E.bracket (setup var) (clear var) $ \_ -> io unmask  -- | Fork a managed thread with a handle created by a timeout manager.-forkManagedTimeout :: ThreadManager -> String -> (T.Handle -> IO ()) -> IO ()+forkManagedTimeout+    :: ThreadManager+    -> String+    -- ^ Thread name+    -> (T.Handle -> IO ())+    -- ^ Action with timeout handle+    -> IO () forkManagedTimeout (ThreadManager timmgr var) label io =-    void $ forkIO $ E.handle ignore $ do+    void $ forkIO $ do         labelMe label         E.bracket (setup var) (clear var) $ \(_n, wtid, ref) ->-            -- 'TimeoutThread' is ignored by 'withHandle'.-            void $ T.withHandle timmgr (lockAndKill wtid ref T.TimeoutThread) io+            E.handle ignore $ T.withHandle timmgr (lockAndKill wtid ref ex) io+  where+    ex = KilledByThreadManager Nothing  -- | Fork a managed thread with a cleanup function.-forkManagedFinally :: ThreadManager -> String -> IO () -> IO () -> IO ()-forkManagedFinally mgr label io final = E.mask $ \restore ->-    forkManaged-        mgr-        label-        (E.try (restore io) >>= \(_ :: Either E.SomeException ()) -> final)+forkManagedFinally+    :: ThreadManager+    -> String+    -- ^ Thread name+    -> IO ()+    -- ^ Action+    -> IO ()+    -- ^ Cleanup function+    -> IO ()+forkManagedFinally mgr label io final =+    forkManagedUnmask mgr label $ \restore ->+        E.try (restore io) >>= \(_ :: Either E.SomeException ()) -> final  -- | Fork a managed thread with a handle created by a timeout manager -- and with a cleanup function. forkManagedTimeoutFinally-    :: ThreadManager -> String -> (T.Handle -> IO ()) -> IO () -> IO ()+    :: ThreadManager+    -> String+    -- ^ Thread name+    -> (T.Handle -> IO ())+    -- ^ Action with timeout handle+    -> IO ()+    -- ^ Cleanup function+    -> IO () forkManagedTimeoutFinally mgr label io final = E.mask $ \restore ->-    forkManagedTimeout-        mgr-        label-        (\th -> E.try (restore $ io th) >>= \(_ :: Either E.SomeException ()) -> final)+    forkManagedTimeout mgr label $ \th ->+        E.try (restore $ io th) >>= \(_ :: Either E.SomeException ()) -> final  setup :: TVar (Map Key ManagedThread) -> IO (Key, Weak ThreadId, IORef Bool) setup var = do-    (wtid, n) <- myWeakThradId+    (wtid, n) <- myWeakThreadId     ref <- newIORef False     let ent = ManagedThread wtid ref     -- asking to throw KilledByThreadManager to me@@ -183,21 +240,21 @@ ignore :: KilledByThreadManager -> IO () ignore (KilledByThreadManager _) = return () --- | Wait until all managed thread are finished.+-- | Wait until all managed threads are finished. waitUntilAllGone :: ThreadManager -> IO ()-waitUntilAllGone (ThreadManager _timmgr var) = atomically $ do-    m <- readTVar var-    check (Map.size m == 0)+waitUntilAllGone tm =+    atomically $+        isAllGone tm >>= check +-- | STM action that checks if all managed threads are finished. isAllGone :: ThreadManager -> STM Bool-isAllGone (ThreadManager _timmgr var) = do-    m <- readTVar var-    return (Map.size m == 0)+isAllGone (ThreadManager _timmgr var) =+    Map.null <$> readTVar var  ---------------------------------------------------------------- -myWeakThradId :: IO (Weak ThreadId, Key)-myWeakThradId = do+myWeakThreadId :: IO (Weak ThreadId, Key)+myWeakThreadId = do     tid <- myThreadId     wtid <- mkWeakThreadId tid     let n = fromThreadId tid@@ -208,8 +265,10 @@     tid <- myThreadId     labelThread tid l +-- | Registering a 'T.TimeoutAction' and unregister its 'T.Handle'+--   when the body action is finished. withHandle-    :: ThreadManager -> T.TimeoutAction -> (T.Handle -> IO a) -> IO (Maybe a)+    :: ThreadManager -> T.TimeoutAction -> (T.Handle -> IO a) -> IO a withHandle (ThreadManager timmgr _) = T.withHandle timmgr  #if __GLASGOW_HASKELL__ < 908
System/TimeManager.hs view
@@ -1,7 +1,19 @@-{-# LANGUAGE BangPatterns #-}-{-# LANGUAGE DeriveDataTypeable #-}+{-# LANGUAGE CPP #-}+{-# LANGUAGE NumericUnderscores #-} {-# LANGUAGE RecordWildCards #-} +-- | Timeout manager. Since @v0.3.0@, timeout manager is a wrapper of+-- GHC System TimerManager.+--+-- Some caveats of using this package:+--+--   * Only works for GHC+--   * Only works with a threaded runtime+--   * Users of older versions should check the current semantics.+--   * Using 32-bit systems means the max timeout is @'maxBound' :: Int@+--     (2147483647) microseconds, which is less than 36 minutes.+--   * Using the same 'Handle' in different threads might cause issues in some+--     edge cases. (i.e. using 'cancel'/'pause' in one thread, and 'resume' in another) module System.TimeManager (     -- ** Types     Manager,@@ -25,174 +37,416 @@     tickle,     pause,     resume,+    cancel,      -- ** Low level     register,     registerKillThread,-    cancel,      -- ** Exceptions     TimeoutThread (..), ) where -import Control.Concurrent (mkWeakThreadId, myThreadId)+import Control.Concurrent (forkIO, mkWeakThreadId, myThreadId, newMVar) import qualified Control.Exception as E-import Control.Monad (void)-import Control.Reaper-import Data.IORef (IORef)+import Control.Monad (void, when)+import Data.Bits (shiftR) import qualified Data.IORef as I-import Data.Typeable (Typeable)-import System.IO.Unsafe+import Data.Word (Word64)+import GHC.Clock (getMonotonicTimeNSec)+import System.IO.Unsafe (unsafePerformIO) import System.Mem.Weak (deRefWeak)+import System.TimeManager.Internal -----------------------------------------------------------------+#if defined(mingw32_HOST_OS)+import qualified GHC.Event.Windows as EV+#else+import qualified GHC.Event as EV+#endif --- | A timeout manager-data Manager = Manager (Reaper [Handle] Handle) | NoManager+---------------------------------------------------------------- --- | No manager.+-- | A manager whose timeout value is 0 (no callbacks are fired). defaultManager :: Manager-defaultManager = NoManager---- | An action to be performed on timeout.-type TimeoutAction = IO ()---- | A handle used by a timeout manager.-data Handle = Handle-    { handleManager :: Manager-    , handleActionRef :: IORef TimeoutAction-    , handleStateRef :: IORef State-    }--{-# NOINLINE emptyAction #-}-emptyAction :: IORef TimeoutAction-emptyAction = unsafePerformIO $ I.newIORef (return ())+defaultManager = Manager 0 -{-# NOINLINE emptyState #-}-emptyState :: IORef State-emptyState = unsafePerformIO $ I.newIORef Inactive+---------------------------------------------------------------- +-- | Dummy 'Handle'. emptyHandle :: Handle emptyHandle =     Handle-        { handleManager = NoManager-        , handleActionRef = emptyAction-        , handleStateRef = emptyState+        { handleTimeout = 0+        , handleAction = pure ()+        , handleTimerManager = mutError "handleTimerManager"+        , handleState = mutError "handleState"+        , handleLastRenewed = mutError "handleLastRenewed"+        , handleMinRenewGap = 0+        , handleLock = emptyLock         }+  where+    mutError s = error $ "time-manager: Handle." <> s <> " not set" -data State-    = Active -- Manager turns it to Inactive.-    | Inactive -- Manager removes it with timeout action.-    | Paused -- Manager does not change it.+emptyLock :: Lock+emptyLock = unsafePerformIO $ newMVar ()+{-# NOINLINE emptyLock #-}  ---------------------------------------------------------------- --- | Creating timeout manager which works every N microseconds---   where N is the first argument.+-- | Creating timeout manager with a timeout value in microseconds.+--+--   Setting the timeout to zero or lower @(<= 0)@ will produce a+--   `defaultManager`.+--+--   __WARNING for Windows users:__ /the precision of extending timeouts/+--   /is only full "seconds". The provided microseconds will be floored/+--   /to the first full second. (i.e. @initialize 2_500_000@ will get/+--   /extended by 2 seconds on a 'tickle')/+--   /This also means timeouts of less than one second will not be extended/+--   /when using 'tickle'./ initialize :: Int -> IO Manager-initialize timeout-    | timeout <= 0 = return NoManager-initialize timeout =-    Manager-        <$> mkReaper-            defaultReaperSettings-                { -- Data.Set cannot be used since 'partition' cannot be used-                  -- with 'readIORef`. So, let's just use a list.-                  reaperAction = mkListAction prune-                , reaperDelay = timeout-                , reaperThreadName = "WAI timeout manager (Reaper)"-                }-  where-    prune m@Handle{..} = do-        state <- I.atomicModifyIORef' handleStateRef (\x -> (inactivate x, x))-        case state of-            Inactive -> do-                onTimeout <- I.readIORef handleActionRef-                onTimeout `E.catch` ignoreSync-                return Nothing-            _ -> return $ Just m--    inactivate Active = Inactive-    inactivate x = x+initialize = pure . Manager . max 0  ---------------------------------------------------------------- --- | Stopping timeout manager with onTimeout fired.+-- | Obsoleted since version 0.3.0+--   Is now equivalent to @pure ()@. stopManager :: Manager -> IO ()-stopManager NoManager = return ()-stopManager (Manager mgr) = E.mask_ (reaperStop mgr >>= mapM_ fire)-  where-    fire Handle{..} = do-        onTimeout <- I.readIORef handleActionRef-        onTimeout `E.catch` ignoreSync+stopManager _ = pure ()+{-# DEPRECATED stopManager "This function does nothing since version 0.3.0" #-} --- | Killing timeout manager immediately without firing onTimeout.+-- | Obsoleted since version 0.3.0+--   Is now equivalent to @pure ()@. killManager :: Manager -> IO ()-killManager NoManager = return ()-killManager (Manager mgr) = reaperKill mgr+killManager _ = pure ()+{-# DEPRECATED killManager "This function does nothing since version 0.3.0" #-}  ----------------------------------------------------------------  -- | Registering a timeout action and unregister its handle --   when the body action is finished.---   'Nothing' is returned on timeout.-withHandle :: Manager -> TimeoutAction -> (Handle -> IO a) -> IO (Maybe a)-withHandle mgr onTimeout action =-    E.handle ignore $ E.bracket (register mgr onTimeout) cancel $ \th ->-        Just <$> action th-  where-    ignore TimeoutThread = return Nothing+withHandle :: Manager -> TimeoutAction -> (Handle -> IO a) -> IO a+withHandle mgr onTimeout action+    | isNoManager mgr = action emptyHandle+    | otherwise = E.bracket (register mgr onTimeout) cancel action  -- | Registering a timeout action of killing this thread and --   unregister its handle when the body action is killed or finished. withHandleKillThread :: Manager -> TimeoutAction -> (Handle -> IO ()) -> IO ()-withHandleKillThread mgr onTimeout action =-    E.handle ignore $ E.bracket (registerKillThread mgr onTimeout) cancel action+withHandleKillThread mgr onTimeout action+    | isNoManager mgr = action emptyHandle+    | otherwise =+        E.handle ignore $ E.bracket (registerKillThread mgr onTimeout) cancel action   where-    ignore TimeoutThread = return ()+    ignore TimeoutThread = pure ()  ---------------------------------------------------------------- +-- ============== NOTE ABOUT THREAD SAFETY ==============+--+-- The use of 'IORef's are fine in the current situation where+-- the 'TimeManager' is supposed to be used in a single thread.+--+-- The triggered action, though, is run by the Timer Manager+-- outside of the thread it was registered in.+-- This will potentially cause race conditions if we implement+-- anything that depends on the 'Handle's state.+--+-- Given the following:+--   - If run in one thread: 'register/tickle/pause/resume/cancel' never+--     overlap, making them devoid of race conditions in the general sense.+--   - We want to hit the Timer Manager as little as possible.+--   - We want to keep the 'resume/pause' surface functionality intact, while+--     not hitting the Timer Manager when we don't have to. This means not+--     cancelling the timeout on a pause, but rather mark the timeout paused.+--   - Not actually stopping the timeout on 'pause' introduces race+--     conditions, because the registered action will need to check the 'Handle'+--     state to see whether it should actually run (Active) or if it should+--     drop the action (Paused/Stopped).+--   - Not hitting the Timer Manager on a 'pause' will increase performance on+--     hot 'resume/pause' loops, like 'warp' has when using a streaming response.+--   - 'tickle' gets a sort of debounce to avoid repeated updates in hot loops.+--     - The debounce is 1/4 of the timeout, but we cap it to a maximum of 1 second.+--     - This means the registered action might run earlier than the timeout+--       would indicate; that difference going up to a maximum of 'handleMinRenewGap'.+--   - The following can happen:+--     - == The "Surprise Active" issue ==+--        A 'resume' might get called right after a 'Paused' registered action+--        starts running, and sets the state to 'Active' __before__ the action+--        inspects the 'Handle' state.+--     - == The "Dropped Active" issue ==+--        A 'resume' might get called right after a 'Paused' registered action+--        starts running, but inspects the state __before__ the action sets the+--        state to 'Stopped', and the registered action inspects the state+--        __before__ the 'resume' has set it to 'Active'. Essentially missing+--        the 'resume' completely.+--     - == The "Dropped Cancel" issue ==+--       A 'cancel' getting called right after a 'Paused' registered action+--       starts running, and cancelling __after__ the action reads the state+--       will have the registered action overwrite the state to 'Stopped',+--       when it shouldn't register a new action, but stop everything.+--     - A 'pause' should technically not be an issue, as it will only run when+--       the state is 'Active', but it is a function that changes the state, so+--       just to be cautious, we let it grab the lock.+--     - A 'tickle' in the same situation doesn't matter, as a 'tickle'+--       shouldn't activate a 'Paused' state. (and doesn't change any state)+--   - The "Surprise Active" issue can be mitigated by checking the+--     'handleLastRenewed' time and reregistering the timeout action with the+--     remaining amount of microseconds in the case where it has not yet been+--     'handleTimeout' amount of time.+--     - A 'tickle' could also cause this if the state was 'Active' all along,+--       but we'll accept the 'tickle' as being on time to extend the timeout.+--   - The "Dropped Active" issue is a bit more difficult to mitigate. We'll+--     need a lock to guarantee that either the activation of 'resume' is seen+--     by the registered action, or that the termination of the registered+--     action is seen by the 'resume'.+--   - The "Dropped Cancel" issue will also be avoided when using a lock.+--   - The lock will generally never be contested. It is there only for the+--     off-chance that a state-changing function runs JUST after the registered+--     action triggers. So in general, we don't expect the lock to reduce+--     performance noticeably.++----------------------------------------------------------------+ -- | Registering a timeout action. register :: Manager -> TimeoutAction -> IO Handle-register NoManager _ = return emptyHandle-register m@(Manager mgr) !onTimeout = do-    actionRef <- I.newIORef onTimeout-    stateRef <- I.newIORef Active-    let h =-            Handle-                { handleManager = m-                , handleActionRef = actionRef-                , handleStateRef = stateRef-                }-    reaperAdd mgr h-    return h+register mgr@(Manager timeout) onTimeout+    | isNoManager mgr = pure emptyHandle+    | otherwise = do+        -- The system timer manager is stable for the lifetime of the+        -- process (and even if it were replaced, e.g. around a fork,+        -- the key registered below would only be meaningful to the+        -- manager it was registered with). So fetch it once here and+        -- cache it in the 'Handle' instead of re-reading the global+        -- IORef on every tickle/pause/resume.+        sysmgr <- getTimerManager+        stateRef <- I.newIORef Stopped+        lock <- newLock+        lastRenewedRef <- I.newIORef =<< getMonotonicTimeNSec+        let h =+                Handle+                    { handleTimeout = timeout+                    , handleAction = onTimeout+                    , handleTimerManager = sysmgr+                    , handleState = stateRef+                    , handleLastRenewed = lastRenewedRef+                    , handleMinRenewGap = minRenewGap timeout+                    , handleLock = lock+                    }+        -- Just in case the timeout is only 1 microsecond and because of thread+        -- scheduling it runs before we can change the state to 'Active'+        withLock lock $ do+            key <- registerAdjustedTimeout h timeout+            now <- getMonotonicTimeNSec+            I.writeIORef lastRenewedRef now+            I.writeIORef stateRef $ Active key+        pure h --- | Removing the 'Handle' from the 'Manager' immediately.+-- | This function needs a separate 'timeout' argument, because we might not+-- register the full amount of time when continuing a timeout that was started+-- a bit earlier. (cf. "Surprise Active" situation)+registerAdjustedTimeout :: Handle -> Int -> IO EV.TimeoutKey+registerAdjustedTimeout h@Handle{..} timeout = do+    originalKeyRef <-+        I.newIORef $+            error "System.TimeManager.registerAdjustedTimeout: originalKeyRef not filled"+    key <-+        EV.registerTimeout handleTimerManager timeout $+            adjustOnTimeout originalKeyRef h+    I.writeIORef originalKeyRef key+    pure key++-- | Wrapper around a registered action to ensure correct handling.+--+-- We basically need the 'Handle', but this is used before making the handle, so+adjustOnTimeout :: I.IORef EV.TimeoutKey -> Handle -> TimeoutAction+adjustOnTimeout originalKeyRef h@Handle{..} = do+    let writeState = I.atomicWriteIORef handleState+    -- Lock ensures we don't get race conditions.+    -- We return a boolean so that we don't run the (potentially long) action+    -- while holding on to the lock.+    shouldRun <- withLock handleLock $ do+        st <- I.readIORef handleState+        case st of+            -- We can check the @now - handleLastRenewed@ diff+            -- and 'threadDelay' the diff to make the timing better?+            -- @if diff > 'handleTimeout - 'handleMinRenewGap' then runTimeout@+            Active key -> do+                -- set state ref to 'Active'?+                ifSameKey key $ do+                    lastRenewed <- I.readIORef handleLastRenewed+                    now <- getMonotonicTimeNSec+                    let diff = fromIntegral $ now - lastRenewed+                    if diff > handleTimeout+                        -- Valid expiration of the timeout, we run the action+                        then do+                            -- We're going to run the action, so set the state+                            -- so that it won't be resumed.+                            writeState Terminated+                            pure True+                        -- "Surprise Active" situation+                        else do+                            -- We reschedule, but with only the remaining time+                            let remainingTimeout = handleTimeout - diff+                            k <- registerAdjustedTimeout h remainingTimeout+                            writeState $ Active k+                            pure False+            -- We find this action being run after it's been paused. We write+            -- the state to 'Stopped' so that 'resume' knows to reregister the+            -- timeout action.+            Paused key ->+                ifSameKey key $ do+                    writeState Stopped+                    pure False+            -- 'Stopped' and 'Terminated' mean the action shouldn't run.+            _ -> pure False+    when shouldRun handleAction+  where+    -- If the key in the state isn't the same as the one this action+    -- was registered with, then this action shouldn't run.+    -- (Technically, this situation shouldn't happen. but since the registered+    -- action only ever runs once, we can afford to be redundant)+    ifSameKey key f = do+        originalKey <- I.readIORef originalKeyRef+        if key == originalKey+            then f+            else pure False++-- | How long 'tickle' waits before actually renewing the timeout:+--   a quarter of the timeout, capped at one second. Skipping a renewal+--   inside this window only shortens the effective timeout by up to+--   this gap, but turns hot 'tickle' loops (one per chunk sent or+--   received) into a clock read and a comparison.+minRenewGap :: Int -> Word64+minRenewGap timeout =+    -- @shiftR 2 === divide by 4@+    min maxRenewDebounce (microToNano timeout `shiftR` 2)+  where+    microToNano = (* 1_000) . fromIntegral++-- | One second in nanoseconds+maxRenewDebounce :: Word64+maxRenewDebounce = 1_000_000_000++-- | Run 'f' if the minimum renew gap has been crossed.+whenRenew :: Handle -> IO () -> IO ()+whenRenew h f = do+    now <- getMonotonicTimeNSec+    lastRenewed <- I.readIORef $ handleLastRenewed h+    when (now - lastRenewed >= handleMinRenewGap h) f++-- | Unregistering the timeout.+--+-- The timeout can not be 'resume'd. To "resume" the timeout, you need to+-- 'register' again. cancel :: Handle -> IO ()-cancel Handle{..} = case handleManager of-    NoManager -> return ()-    Manager mgr -> void $ reaperModify mgr filt+cancel h@Handle{..} =+    withNonEmptyHandle h $+        -- "Dropped Cancel" remedy+        --+        -- We can eat a potential mutex pause here to avoid race conditions,+        -- because we don't expect 'cancel' to be called in hot loops.+        --+        -- (The race condition being: the 'Terminated' state being overwritten+        -- because the 'cancel' runs JUST after the registered action starts+        -- running, sets the state to 'Terminated', and then the registered+        -- action finishes and overwrites it to 'Stopped')+        withLock handleLock $ do+            withTimeoutKey h $ EV.unregisterTimeout handleTimerManager+            I.atomicWriteIORef handleState Terminated++-- | Extending the timeout.+--+-- To keep frequent callers cheap, the renewal is rate-limited: it is+-- skipped unless at least a quarter of the timeout (capped at one+-- second) has passed since the timeout was last registered or updated.+--+-- Careful: this does NOT reactivate an already 'pause'd 'Handle'!+--+-- __WARNING for Windows users:__ /the precision of extending timeouts/+-- /is only full "seconds". The provided microseconds will be floored/+-- /to the first full second. (i.e. @initialize 2_500_000@ will get/+-- /extended by 2 seconds on a 'tickle')/+-- /This also means timeouts of less than one second will not be extended/+-- /when using 'tickle'./+tickle :: Handle -> IO ()+tickle h@Handle{..} =+    withNonEmptyHandle h $+        whenRenew h $+            withActiveTimeoutKey h $ \key -> do+                updateTheTimeout key+                now <- getMonotonicTimeNSec+                I.atomicWriteIORef handleLastRenewed now   where-    -- It's very important that this function forces the whole workload so we-    -- don't retain old handles, otherwise disasterous leaks occur.-    filt [] = []-    filt (h@(Handle _ _ ref) : hs)-        | handleStateRef == ref = hs-        | otherwise =-            let !hs' = filt hs-             in h : hs'+    -- For some reason the Windows implementation of 'updateTimeout' wants+    -- full seconds, instead of the microseconds that's used when registering...+    updateTheTimeout key =+        EV.updateTimeout handleTimerManager key+#if defined(mingw32_HOST_OS)+            (fromIntegral (handleTimeout `div` 1_000_000))+#else+            handleTimeout+#endif +-- | Pauses the timeout so you can 'resume' it later. Does not stop it entirely.+-- Use 'cancel' if you want to make sure the action will not be resumed.+--+-- To resume a timeout with the same 'Handle', 'resume' MUST be called.+-- Don't call 'tickle' for resumption.+pause :: Handle -> IO ()+pause h@Handle{..} =+    withNonEmptyHandle h $+        withLock handleLock . withActiveTimeoutKey h $+            I.atomicWriteIORef handleState . Paused++-- | Resuming the timeout.+--+-- Works like 'tickle' if the 'Handle' wasn't 'pause'd or 'cancel'ed.+resume :: Handle -> IO ()+resume h@Handle{..} =+    withNonEmptyHandle h $+        -- we ignore the key when paused, because we recheck the state after+        -- grabbing the lock.+        checkStateWith (\_ -> onPausedOrStopped) onPausedOrStopped+  where+    -- "Dropped Active" remedy+    --+    -- Grabbing the lock ensures 'resume' runs either before or after the+    -- registered action changes the state.+    onPausedOrStopped =+        withLock handleLock $ checkStateWith pausedF stoppedF+    checkStateWith onPaused onStopped = do+        state <- I.readIORef handleState+        case state of+            -- 'tickle' doesn't introduce race conditions, so can always be run.+            Active{} -> tickle h+            -- Abort when terminated.+            Terminated -> pure ()+            Paused k -> onPaused k+            Stopped -> onStopped+    pausedF k = do+        -- Set state to 'Active' before 'tickle'ing, because+        -- 'tickle' only runs when the state is 'Active'.+        activateTimeout k+        tickle h+    stoppedF = do+        key <- registerAdjustedTimeout h handleTimeout+        now <- getMonotonicTimeNSec+        I.atomicWriteIORef handleLastRenewed now+        activateTimeout key+    activateTimeout =+        I.atomicWriteIORef handleState . Active+ ----------------------------------------------------------------  -- | The asynchronous exception thrown if a thread is registered via -- 'registerKillThread'. data TimeoutThread = TimeoutThread-    deriving (Typeable)  instance E.Exception TimeoutThread where     toException = E.asyncExceptionToException     fromException = E.asyncExceptionFromException+ instance Show TimeoutThread where     show TimeoutThread = "Thread killed by timeout manager" @@ -202,70 +456,31 @@ --   want to leak the asynchronous exception to GHC RTS. registerKillThread :: Manager -> TimeoutAction -> IO Handle registerKillThread m onTimeout = do-    tid <- myThreadId-    wtid <- mkWeakThreadId tid+    wtid <- myThreadId >>= mkWeakThreadId     -- First run the timeout action in case the child thread is masked.     register m $         onTimeout `E.finally` do             mtid <- deRefWeak wtid             case mtid of-                Nothing -> return ()-                Just tid' -> E.throwTo tid' TimeoutThread---------------------------------------------------------------------- | Setting the state to active.---   'Manager' turns active to inactive repeatedly.-tickle :: Handle -> IO ()-tickle Handle{..} = I.writeIORef handleStateRef Active---- | Setting the state to paused.---   'Manager' does not change the value.-pause :: Handle -> IO ()-pause Handle{..} = I.writeIORef handleStateRef Paused---- | Setting the paused state to active.---   This is an alias to 'tickle'.-resume :: Handle -> IO ()-resume = tickle+                Nothing -> pure ()+                Just tid' -> void . forkIO $ E.throwTo tid' TimeoutThread  ----------------------------------------------------------------  -- | Call the inner function with a timeout manager.---   'stopManager' is used after that. withManager     :: Int     -- ^ timeout in microseconds     -> (Manager -> IO a)     -> IO a-withManager timeout f =-    E.bracket-        (initialize timeout)-        stopManager-        f+withManager timeout f = initialize timeout >>= f  -- | Call the inner function with a timeout manager.---   'killManager' is used after that.+--   This is identical to 'withManager'. withManager'     :: Int     -- ^ timeout in microseconds     -> (Manager -> IO a)     -> IO a-withManager' timeout f =-    E.bracket-        (initialize timeout)-        killManager-        f--------------------------------------------------------------------isAsyncException :: E.Exception e => e -> Bool-isAsyncException e =-    case E.fromException (E.toException e) of-        Just (E.SomeAsyncException _) -> True-        Nothing -> False--ignoreSync :: E.SomeException -> IO ()-ignoreSync se-    | isAsyncException se = E.throwIO se-    | otherwise = return ()+withManager' = withManager+{-# DEPRECATED withManager' "This function is the same as 'withManager' since version 0.3.0" #-}
+ System/TimeManager/Internal.hs view
@@ -0,0 +1,117 @@+{-# LANGUAGE BangPatterns #-}+{-# LANGUAGE CPP #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE StrictData #-}++module System.TimeManager.Internal where++import Control.Concurrent.MVar (MVar, modifyMVar, newMVar)+import Data.IORef (IORef, readIORef)+import Data.Word (Word64)++#if defined(mingw32_HOST_OS)+import qualified GHC.Event.Windows as EV+#else+import qualified GHC.Event as EV+#endif++----------------------------------------------------------------++-- | A timeout manager+newtype Manager = Manager Int++isNoManager :: Manager -> Bool+isNoManager (Manager 0) = True+isNoManager _ = False++----------------------------------------------------------------++-- | An action (callback) to be performed on timeout.+type TimeoutAction = IO ()++-- | A handle used by a timeout manager.+data Handle = Handle+    { handleTimeout :: Int+    , handleAction :: TimeoutAction+    , handleTimerManager :: ~TimerManager+    -- ^ The system timer manager the timeout key was registered with.+    --   Cached so that per-request operations don't re-fetch it.+    , handleLock :: Lock+    -- ^ Used by 'resume', 'pause' and 'cancel' to determine race conditions.+    --+    -- /We intentionally do not use an @MVar HandleState@ for performance reasons./+    -- /The lock only has to be grabbed to avoid race conditions./+    , handleState :: ~(IORef HandleState)+    -- ^ The current state. Used to decide whether a timeout is still going,+    -- paused, or completely terminated.+    --+    -- /We intentionally do not use an @MVar HandleState@ for performance reasons./+    , handleLastRenewed :: ~(IORef Word64)+    -- ^ Monotonic time (in nanoseconds) when the timeout was last+    --   registered or updated.+    , handleMinRenewGap :: Word64+    -- ^ 'tickle' is a no-op unless at least this many nanoseconds have+    --   passed since the last renewal.+    }++-- | Makes sure the function is only run when there's a key to act on.+withTimeoutKey :: Handle -> (EV.TimeoutKey -> IO ()) -> IO ()+withTimeoutKey h keyF = do+    st <- readIORef $ handleState h+    case st of+        Paused key -> keyF key+        Active key -> keyF key+        _ -> pure ()++-- | Makes sure the function is only run when the state is 'Active'.+withActiveTimeoutKey :: Handle -> (EV.TimeoutKey -> IO ()) -> IO ()+withActiveTimeoutKey h keyF = do+    st <- readIORef $ handleState h+    case st of+        Active key -> keyF key+        _ -> pure ()++-- | Used to avoid race conditions in situations when the state has to be changed.+type Lock = MVar ()++newLock :: IO Lock+newLock = newMVar ()++withLock :: Lock -> IO a -> IO a+withLock lock action =+    -- Not sure whether this should be 'modifyMVarMasked' or not.+    modifyMVar lock $ \l -> do+        a <- action+        pure (l, a)++-- | Tracking the state of a handle.+data HandleState+    = -- Timeout is primed to run+      Active EV.TimeoutKey+    | -- Timeout is paused, but still running+      -- ('resume' will set it back to 'Active' and 'tickle')+      Paused EV.TimeoutKey+    | -- Action ran, but timeout was paused, so it is resumable+      -- ('resume' will reregister the action)+      Stopped+    | -- Action was cancelled or run. 'register' is needed to start a new timeout.+      Terminated++isEmptyHandle :: Handle -> Bool+isEmptyHandle Handle{..} = handleTimeout == 0++withNonEmptyHandle :: Handle -> IO () -> IO ()+withNonEmptyHandle h act =+    if isEmptyHandle h then pure () else act++#if defined(mingw32_HOST_OS)+type TimerManager = EV.Manager++getTimerManager :: IO TimerManager+getTimerManager = EV.getSystemManager+#else+type TimerManager = EV.TimerManager++getTimerManager :: IO TimerManager+getTimerManager = EV.getSystemTimerManager+#endif
+ test/Spec.hs view
@@ -0,0 +1,225 @@+{-# LANGUAGE CPP #-}+{-# LANGUAGE NumericUnderscores #-}+{-# LANGUAGE RecordWildCards #-}+{-# LANGUAGE StandaloneDeriving #-}+{-# OPTIONS_GHC -Wno-orphans #-}++module Main where++import Test.Hspec++#if defined(mingw32_HOST_OS)+-- -- Uncomment when reenabling the tests for Windows+-- import qualified GHC.Event.Windows as EV+#else+import qualified GHC.Event as EV+#endif++#if defined(mingw32_HOST_OS)+main :: IO ()+main = hspec $ do+    describe "TimeManager" $ do+        it "tests don't work on windows" $+            pendingWith "requires more testing on a Windows machine"+#else+import Control.Concurrent (threadDelay)+import Control.Monad (forM_, void)+import Data.IORef as I (+    IORef,+    atomicModifyIORef',+    newIORef,+    readIORef,+    writeIORef,+ )+import System.TimeManager+import System.TimeManager.Internal+import Test.HUnit (assertBool)++main :: IO ()+main = hspec $ do+    describe "TimeManager" $ do+        it "defaultManager == no manager" $+            defaultManager `shouldSatisfy` isNoManager++        it "initializes negative manager" $ do+            let check = (`shouldBe` defaultManager)+            initialize (-10) >>= check+            withManager (-5) check++        it "empty handle is correct" $+            handleTimeout emptyHandle `shouldBe` 0++        it "empty handle check is consistent" $ do+            assertBool "emptyHandle not empty" $+                isEmptyHandle emptyHandle++        it "gives emptyHandle when registering defaultManager" $ do+            hndl <- register defaultManager $ pure ()+            assertBool "got non-empty handle" $ isEmptyHandle hndl++        it "throws TimeoutThread exception" $+            throwsTimeoutThread $ do+                mngr <- initialize timeoutAmount+                _hndl <- registerKillThread mngr $ pure ()+                waitLong++        it "defaultManager doesn't kill thread" $ do+            _hndl <- registerKillThread defaultManager $ pure ()+            waitShort++        it "withHandle: registers timeout" $+            withHandleTest mgr1 $ \check _ -> do+                waitShort+                check True++        it "withHandle: doesn't register timeout" $+            withHandleTest defaultManager $ \check _ -> do+                waitShort+                check False++        -- We make a ref on the outside, to check that the ref is indeed+        -- set before the timeout kills the action inside.+        it "withHandleKillThread: registers timeout (and kills)" $ do+            ref <- freshRef+            withHandleKillTest (Just ref) mgr1 $ \_ _ ->+                throwsTimeoutThread waitShort+            ref `refShouldBe` True++        it "withHandleKillThread: doesn't register timeout" $+            withHandleKillTest Nothing defaultManager $ \check _ -> do+                waitShort >> check False++        it "cancel/pause works as expected" $ do+            m <- mkTestManager+            let killUnless f = do+                    hndl <- registerKillThread m (pure ())+                    _ <- f hndl+                    waitLong+            throwsTimeoutThread $ killUnless pure+            killUnless cancel+            killUnless pause++        it "tickle works as expected" $ do+            m <- mkTestManager+            withHandleTest m $ \check hndl -> do+                forM_ [(1 :: Int) .. 20] $ \_ -> do+                    waitShort+                    tickle hndl+                check False++        let runAndWaitForTimeout f =+                runIt $ \hndl -> do+                    void $ f hndl+                    waitLong+        it "resume works as expected (nothing)" $ do+            -- Doing nothing kills the thread+            throwsTimeoutThread . runAndWaitForTimeout $ \_ -> pure ()+        it "resume works as expected (pause)" $ do+            -- Pausing stops the kill+            runAndWaitForTimeout $ \hndl -> do+                waitShort >> pause hndl+        it "resume works as expected (pause/resume)" $ do+            -- Resuming kills the thread again+            throwsTimeoutThread . runAndWaitForTimeout $ \hndl -> do+                waitShort >> pause hndl+                waitLong >> resume hndl+        it "resume works as expected (cancel/resume)" $ do+            -- Cancelling is unresumable+            runAndWaitForTimeout $ \hndl -> do+                waitShort >> cancel hndl+                waitLong >> resume hndl+        it "resume works as expected (cancel/pause/resume)" $ do+            -- Cancelling and then pausing is still unresumable+            runAndWaitForTimeout $ \hndl -> do+                waitShort >> cancel hndl+                waitShort >> pause hndl+                waitLong >> resume hndl+            -- Pausing, then cancelling doesn't change anything+            runAndWaitForTimeout $ \hndl -> do+                waitShort >> pause hndl+                waitShort >> cancel hndl+                waitLong >> resume hndl+        it "finished timeout won't resume" $ do+            -- If the timeout action runs, resume shouldn't work+            counter <- I.newIORef (0 :: Int)+            m <- mkTestManager+            let increase = I.atomicModifyIORef' counter $ \i -> (i + 1, ())+            withHandle m increase $ \h -> do+                let checkCount x = do+                        i <- I.readIORef counter+                        i `shouldBe` x+                    timeoutOnlyRanOnce = waitLong >> checkCount 1++                checkCount 0+                -- waiting lets the timeout+                timeoutOnlyRanOnce+                -- resuming should not influence the counter+                resume h+                timeoutOnlyRanOnce+                -- pausing after it runs also doesn't re-arm the timeout+                pause h+                resume h+                timeoutOnlyRanOnce+                -- cancel also doesn't re-arm the timeout+                cancel h+                pause h+                resume h+                timeoutOnlyRanOnce++        it "resume also works as tickle" $+            testResume resume++        it "resume also works as tickle with pauses" $+            testResume $ \hndl -> do+                resume hndl+                pause hndl+                resume hndl++        it "old resume did NOT work as tickle" $+            throwsTimeoutThread $+                testResume oldResume+  where+    withHandleTest = withTest withHandle Nothing+    withHandleKillTest = withTest withHandleKillThread+    -- Test that starts with a 'False' IORef and on timeout sets it to true+    withTest withF mRef m f = do+        ref <- maybe freshRef pure mRef+        withF m (I.writeIORef ref True) . f $ refShouldBe ref+    -- run with a 20ms timeout and kill+    runIt f = do+        m <- mkTestManager+        void $ f =<< registerKillThread m (pure ())+    timeoutAmount = 20_000+    mkTestManager = initialize timeoutAmount+    -- Waiting a lot less than the timeout takes+    waitShort = threadDelay $ timeoutAmount `div` 5+    -- Waiting a lot longer than the timeout takes+    waitLong = threadDelay $ timeoutAmount * 5+    -- "resuming" every 2.5ms 20 times+    testResume f = do+        runIt $ \hndl -> do+            forM_ [(1 :: Int) .. 20] $ \_ -> waitShort >> f hndl++mgr1 :: Manager+mgr1 = Manager 1++freshRef :: IO (IORef Bool)+freshRef = I.newIORef False++refShouldBe :: IORef Bool -> Bool -> IO ()+refShouldBe ref expected =+    I.readIORef ref >>= (`shouldBe` expected)++throwsTimeoutThread :: IO () -> Expectation+throwsTimeoutThread t = t `shouldThrow` (const True :: TimeoutThread -> Bool)++deriving instance Eq Manager+deriving instance Show Manager++-- copied from time-manager-0.3.0 to check it actually is broken+oldResume :: Handle -> IO ()+oldResume h | isEmptyHandle h = return ()+oldResume Handle{..} = do+    key <- EV.registerTimeout handleTimerManager handleTimeout handleAction+    I.writeIORef handleState $ Active key+#endif
time-manager.cabal view
@@ -1,26 +1,48 @@-Name:                time-manager-Version:             0.2.4-Synopsis:            Scalable timer-License:             MIT-License-file:        LICENSE-Author:              Michael Snoyman and Kazu Yamamoto-Maintainer:          kazu@iij.ad.jp-Homepage:            http://github.com/yesodweb/wai-Category:            System-Build-Type:          Simple-Cabal-Version:       >=1.10-Stability:           Stable-Description:         Scalable timer functions provided by a timer manager-                     and thread management functions to prevent thread-                     leak by a thread manager.-Extra-Source-Files:  ChangeLog.md+cabal-version:      >=1.10+name:               time-manager+version:            0.4.0+license:            MIT+license-file:       LICENSE+maintainer:         kazu@iij.ad.jp+author:             Michael Snoyman and Kazu Yamamoto+stability:          Stable+homepage:           http://github.com/yesodweb/wai+synopsis:           Scalable timer+description:+    Scalable timer functions provided by a timer manager+    and thread management functions to prevent thread+    leak by a thread manager. -Library-  Build-Depends:     base                      >= 4.12       && < 5-                   , auto-update               >= 0.2        && < 0.3-                   , containers-                   , stm-  Default-Language:  Haskell2010-  Exposed-modules:   System.TimeManager-  Exposed-modules:   System.ThreadManager-  Ghc-Options:       -Wall+category:           System+build-type:         Simple+extra-source-files: ChangeLog.md+                    README.md++library+    exposed-modules:+        System.TimeManager+        System.ThreadManager+    other-modules:+        System.TimeManager.Internal++    default-language:   Haskell2010+    default-extensions: Strict StrictData+    ghc-options:        -Wall+    build-depends:+        base >=4.12 && <5,+        containers,+        stm++test-suite spec+    type:             exitcode-stdio-1.0+    main-is:          test/Spec.hs+    other-modules:+        System.TimeManager,+        System.TimeManager.Internal+    build-depends:    base+    default-language: Haskell2010+    ghc-options:      -Wall -threaded+    build-depends:+        base >=4.12 && <5,+        hspec,+        HUnit