time-manager 0.2.4 → 0.4.0
raw patch · 7 files changed
Files
- ChangeLog.md +57/−0
- README.md +13/−0
- System/ThreadManager.hs +95/−36
- System/TimeManager.hs +370/−155
- System/TimeManager/Internal.hs +117/−0
- test/Spec.hs +225/−0
- time-manager.cabal +47/−25
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