packages feed

time-manager 0.1.3 → 0.4.0

raw patch · 7 files changed

Files

ChangeLog.md view
@@ -1,5 +1,88 @@ # 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`.+* Providing emptyHandle.++## 0.2.3++* Exporting defaultManager.++## 0.2.2++* `initialize` with non positive integer creates a time manager+  which does not maintain timeout.+  [#1017](https://github.com/yesodweb/wai/pull/1017)++## 0.2.1++* Export KilledByThreadManager exception+  [#1016](https://github.com/yesodweb/wai/pull/1016)++## 0.2.0++* Providing `System.ThreadManager`.+* `withHandle` catches `TimeoutThread` internally.+  It returns `Nothing` on timeout.+ ## 0.1.3  * Providing `withHandle` and `withHandleKillThread`.
+ 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
@@ -0,0 +1,277 @@+{-# LANGUAGE CPP #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE ScopedTypeVariables #-}++-- | 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+    forkManaged,+    forkManagedFinally,+    forkManagedUnmask,+    forkManagedTimeout,+    forkManagedTimeoutFinally,++    -- * Synchronization+    waitUntilAllGone,+    isAllGone,++    -- * Re-exports+    T.Manager,+    withHandle,+    T.Handle,+    T.tickle,+    T.pause,+    T.resume,+) where++import Control.Concurrent+import Control.Concurrent.STM+import Control.Exception (Exception (..), SomeException (..))+import qualified Control.Exception as E+import Control.Monad (unless, void)+import Data.Foldable (forM_)+import Data.IORef (IORef, atomicModifyIORef', newIORef)+import Data.Map.Strict (Map)+import qualified Data.Map.Strict as Map+import Data.Word (Word64)+import GHC.Conc.Sync (labelThread)+#if __GLASGOW_HASKELL__ >= 908+import GHC.Conc.Sync (fromThreadId)+#endif+import System.Mem.Weak (Weak, deRefWeak)+import qualified System.TimeManager as T++----------------------------------------------------------------++-- | Manager to manage the thread and the timer.+data ThreadManager = ThreadManager T.Manager (TVar ManagedThreads)++type Key = Word64+type ManagedThreads = Map Key ManagedThread++----------------------------------------------------------------++-- 'IORef' prevents race between WAI TimeManager (TimeoutThread)+-- and stopAfter (KilledByThreadManager).+-- It is initialized with 'False' and turned into 'True' when locked.+-- The winner can throw an asynchronous exception.+data ManagedThread = ManagedThread (Weak ThreadId) (IORef Bool)++----------------------------------------------------------------++-- | 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++----------------------------------------------------------------++-- | An exception used internally to kill a managed thread.+data KilledByThreadManager = KilledByThreadManager (Maybe SomeException)+    deriving (Show)++instance Exception KilledByThreadManager where+    toException = E.asyncExceptionToException+    fromException = E.asyncExceptionFromException++-- | 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 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+            m0 <- readTVar var+            writeTVar var Map.empty+            return m0+        let ths = Map.elems m+            er = either Just (const Nothing) ma+            ex = KilledByThreadManager er+        forM_ ths $ \(ManagedThread wtid ref) -> lockAndKill wtid ref ex+        cleanup ma++----------------------------------------------------------------++-- | Fork a managed thread.+--+-- This guarantees that the thread ID is added to the manager's queue before+-- the thread starts, and is removed again when the thread terminates+-- (normally or abnormally).+forkManaged+    :: ThreadManager+    -> String+    -- ^ Thread name+    -> IO ()+    -- ^ Action+    -> IO ()+forkManaged mgr label io =+    forkManagedUnmask mgr label $ \unmask -> unmask io++-- | Like 'forkManaged', but run action with exceptions masked+forkManagedUnmask+    :: 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+    -- ^ Thread name+    -> (T.Handle -> IO ())+    -- ^ Action with timeout handle+    -> IO ()+forkManagedTimeout (ThreadManager timmgr var) label io =+    void $ forkIO $ do+        labelMe label+        E.bracket (setup var) (clear var) $ \(_n, wtid, ref) ->+            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+    -- ^ 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+    -- ^ 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++setup :: TVar (Map Key ManagedThread) -> IO (Key, Weak ThreadId, IORef Bool)+setup var = do+    (wtid, n) <- myWeakThreadId+    ref <- newIORef False+    let ent = ManagedThread wtid ref+    -- asking to throw KilledByThreadManager to me+    atomically $ modifyTVar' var $ Map.insert n ent+    return (n, wtid, ref)++lockAndKill :: Exception e => Weak ThreadId -> IORef Bool -> e -> IO ()+lockAndKill wtid ref e = do+    alreadyLocked <- atomicModifyIORef' ref (\b -> (True, b)) -- try to lock+    unless alreadyLocked $ do+        mtid <- deRefWeak wtid+        case mtid of+            Nothing -> return ()+            Just tid -> E.throwTo tid e++clear+    :: TVar (Map Key ManagedThread)+    -> (Key, Weak ThreadId, IORef Bool)+    -> IO ()+clear var (n, _, _) = atomically $ modifyTVar' var $ Map.delete n++ignore :: KilledByThreadManager -> IO ()+ignore (KilledByThreadManager _) = return ()++-- | Wait until all managed threads are finished.+waitUntilAllGone :: ThreadManager -> IO ()+waitUntilAllGone tm =+    atomically $+        isAllGone tm >>= check++-- | STM action that checks if all managed threads are finished.+isAllGone :: ThreadManager -> STM Bool+isAllGone (ThreadManager _timmgr var) =+    Map.null <$> readTVar var++----------------------------------------------------------------++myWeakThreadId :: IO (Weak ThreadId, Key)+myWeakThreadId = do+    tid <- myThreadId+    wtid <- mkWeakThreadId tid+    let n = fromThreadId tid+    return (wtid, n)++labelMe :: String -> IO ()+labelMe l = do+    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 a+withHandle (ThreadManager timmgr _) = T.withHandle timmgr++#if __GLASGOW_HASKELL__ < 908+fromThreadId :: ThreadId -> Word64+fromThreadId tid = read (drop 9 $ show tid)+#endif
System/TimeManager.hs view
@@ -1,11 +1,26 @@-{-# 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,+    defaultManager,     TimeoutAction,     Handle,+    emptyHandle,      -- ** Manager     initialize,@@ -18,138 +33,420 @@     withHandle,     withHandleKillThread, -    -- ** Control+    -- ** Control timeout     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.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 GHC.Weak (deRefWeak)+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-type Manager = Reaper [Handle] Handle+-- | A manager whose timeout value is 0 (no callbacks are fired).+defaultManager :: Manager+defaultManager = Manager 0 --- | An action to be performed on timeout.-type TimeoutAction = IO ()+---------------------------------------------------------------- --- | A handle used by 'Manager'-data Handle = Handle Manager !(IORef TimeoutAction) !(IORef State)+-- | Dummy 'Handle'.+emptyHandle :: Handle+emptyHandle =+    Handle+        { 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 micro seconds---   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 =-    mkReaper-        defaultReaperSettings-            { reaperAction = mkListAction prune-            , reaperDelay = timeout-            , reaperThreadName = "WAI timeout manager (Reaper)"-            }-  where-    prune m@(Handle _ actionRef stateRef) = do-        state <- I.atomicModifyIORef' stateRef (\x -> (inactivate x, x))-        case state of-            Inactive -> do-                onTimeout <- I.readIORef actionRef-                onTimeout `E.catch` ignoreAll-                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 mgr = E.mask_ (reaperStop mgr >>= mapM_ fire)-  where-    fire (Handle _ actionRef _) = do-        onTimeout <- I.readIORef actionRef-        onTimeout `E.catch` ignoreAll--ignoreAll :: E.SomeException -> IO ()-ignoreAll _ = return ()+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 = reaperKill+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. withHandle :: Manager -> TimeoutAction -> (Handle -> IO a) -> IO a-withHandle mgr onTimeout action =-    E.bracket (register mgr onTimeout) cancel action+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 handler $ 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-    handler 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 mgr !onTimeout = do-    actionRef <- I.newIORef onTimeout-    stateRef <- I.newIORef Active-    let h = Handle mgr actionRef 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 mgr _ stateRef) = do-    _ <- reaperModify mgr filt-    return ()+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 _ _ stateRef') : hs)-        | stateRef == stateRef' = 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" @@ -159,57 +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 _ _ stateRef) = I.writeIORef stateRef Active---- | Setting the state to paused.---   'Manager' does not change the value.-pause :: Handle -> IO ()-pause (Handle _ _ stateRef) = I.writeIORef stateRef 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+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,21 +1,48 @@-Name:                time-manager-Version:             0.1.3-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.-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-  Default-Language:  Haskell2010-  Exposed-modules:   System.TimeManager-  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