Spock-worker 0.2.1.3 → 0.3.0.0
raw patch · 5 files changed
+69/−62 lines, 5 filesdep +errorsdep ~Spockdep ~basePVP ok
version bump matches the API change (PVP)
Dependencies added: errors
Dependency ranges changed: Spock, base
API changes (from Hackage documentation)
- Web.Spock.Worker: instance Eq WorkResult
- Web.Spock.Worker: instance Show WorkResult
- Web.Spock.Worker: type InternalError = String
- Web.Spock.Worker: wc_concurrent :: WorkerConfig -> WorkerConcurrentStrategy
- Web.Spock.Worker: wc_queueLimit :: WorkerConfig -> Int
- Web.Spock.Worker.Internal.Queue: instance (Eq p, Eq v) => Eq (PureQueue p v)
- Web.Spock.Worker.Internal.Queue: instance (Show p, Show v) => Show (PureQueue p v)
- Web.Spock.Worker.Internal.Queue: pq_container :: PureQueue p v -> !(Map p (Vector v))
- Web.Spock.Worker.Internal.Queue: pq_maxSize :: PureQueue p v -> !Int
+ Web.Spock.Worker: InternalError :: a -> InternalError a
+ Web.Spock.Worker: InternalErrorMsg :: String -> InternalError a
+ Web.Spock.Worker: WorkerDef :: WorkerConfig -> WorkHandler conn sess st err a -> ErrorHandler conn sess st err a -> WorkerDef conn sess st err a
+ Web.Spock.Worker: [wc_concurrent] :: WorkerConfig -> WorkerConcurrentStrategy
+ Web.Spock.Worker: [wc_queueLimit] :: WorkerConfig -> Int
+ Web.Spock.Worker: [wd_config] :: WorkerDef conn sess st err a -> WorkerConfig
+ Web.Spock.Worker: [wd_errorHandler] :: WorkerDef conn sess st err a -> ErrorHandler conn sess st err a
+ Web.Spock.Worker: [wd_handler] :: WorkerDef conn sess st err a -> WorkHandler conn sess st err a
+ Web.Spock.Worker: data InternalError a
+ Web.Spock.Worker: data WorkerDef conn sess st err a
+ Web.Spock.Worker: instance GHC.Classes.Eq Web.Spock.Worker.WorkResult
+ Web.Spock.Worker: instance GHC.Show.Show Web.Spock.Worker.WorkResult
+ Web.Spock.Worker.Internal.Queue: [pq_container] :: PureQueue p v -> !(Map p (Vector v))
+ Web.Spock.Worker.Internal.Queue: [pq_maxSize] :: PureQueue p v -> !Int
+ Web.Spock.Worker.Internal.Queue: instance (GHC.Classes.Eq p, GHC.Classes.Eq v) => GHC.Classes.Eq (Web.Spock.Worker.Internal.Queue.PureQueue p v)
+ Web.Spock.Worker.Internal.Queue: instance (GHC.Show.Show p, GHC.Show.Show v) => GHC.Show.Show (Web.Spock.Worker.Internal.Queue.PureQueue p v)
- Web.Spock.Worker: ErrorHandlerIO :: (InternalError -> a -> IO WorkResult) -> ErrorHandler conn sess st a
+ Web.Spock.Worker: ErrorHandlerIO :: ((InternalError err) -> a -> IO WorkResult) -> ErrorHandler conn sess st err a
- Web.Spock.Worker: ErrorHandlerSpock :: (InternalError -> a -> (WebStateM conn sess st) WorkResult) -> ErrorHandler conn sess st a
+ Web.Spock.Worker: ErrorHandlerSpock :: ((InternalError err) -> a -> (WebStateM conn sess st) WorkResult) -> ErrorHandler conn sess st err a
- Web.Spock.Worker: data ErrorHandler conn sess st a
+ Web.Spock.Worker: data ErrorHandler conn sess st err a
- Web.Spock.Worker: newWorker :: (MonadTrans t, Monad (t (WebStateM conn sess st))) => WorkerConfig -> WorkHandler conn sess st a -> ErrorHandler conn sess st a -> t (WebStateM conn sess st) (WorkQueue a)
+ Web.Spock.Worker: newWorker :: (MonadTrans t, Monad (t (WebStateM conn sess st))) => WorkerDef conn sess st err a -> t (WebStateM conn sess st) (WorkQueue a)
- Web.Spock.Worker: type WorkHandler conn sess st a = a -> ErrorT InternalError (WebStateM conn sess st) WorkResult
+ Web.Spock.Worker: type WorkHandler conn sess st err a = a -> ExceptT (InternalError err) (WebStateM conn sess st) WorkResult
Files
- Spock-worker.cabal +5/−4
- src/Web/Spock/Worker.hs +48/−35
- src/Web/Spock/Worker/Internal/Queue.hs +12/−15
- test/Tests.hs +1/−0
- test/Web/Spock/Worker/Internal/QueueTests.hs +3/−8
Spock-worker.cabal view
@@ -1,5 +1,5 @@ name: Spock-worker-version: 0.2.1.3+version: 0.3.0.0 synopsis: Background workers for Spock description: Adds a background-job queue to Spock homepage: http://github.com/agrafix/Spock-worker@@ -11,7 +11,7 @@ category: Web build-type: Simple cabal-version: >=1.10-tested-with: GHC==7.6.3, GHC==7.8.3+tested-with: GHC==7.10.3 library exposed-modules:@@ -22,10 +22,11 @@ default-language: Haskell2010 build-depends: Spock >=0.7.2,- base >=4.6 && < 5,+ base >=4.8 && < 5, containers >=0.5, lifted-base, mtl,+ errors, stm >=2.4, text >=0.11.3.1, time >=1.4,@@ -43,7 +44,7 @@ build-depends: HTF >=0.12.2.1, Spock-worker,- base >=4.6 && < 5,+ base >=4.8 && < 5, containers, stm, vector
src/Web/Spock/Worker.hs view
@@ -1,40 +1,45 @@-{-# LANGUAGE RankNTypes, ScopedTypeVariables, OverloadedStrings #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE RankNTypes #-}+{-# LANGUAGE ScopedTypeVariables #-} module Web.Spock.Worker- ( -- * Worker- WorkQueue- , WorkHandler+ ( -- * Define a Worker+ WorkHandler , WorkerConfig (..) , WorkerConcurrentStrategy (..)+ , WorkerDef (..) , newWorker- , addWork- , WorkExecution (..) , WorkResult (..)+ -- * Enqueue work+ , WorkQueue, WorkExecution (..), addWork -- * Error Handeling- , ErrorHandler(..), InternalError+ , ErrorHandler(..), InternalError (..) ) where -import Control.Concurrent-import Control.Concurrent.STM-import Control.Monad-import Control.Monad.Trans-import Control.Monad.Trans.Error-import Control.Exception.Lifted as EX-import Data.Time-import Web.Spock.Shared+import Control.Concurrent+import Control.Concurrent.STM+import Control.Error+import Control.Exception.Lifted as EX+import Control.Monad+import Control.Monad.Trans+import Data.Time+import Web.Spock.Shared import qualified Web.Spock.Worker.Internal.Queue as Q -type InternalError = String+-- | An error from a worker+data InternalError a+ = InternalErrorMsg String+ | InternalError a -- | Describe how you want to handle errors. Make sure you catch all exceptions -- that can happen inside this handler, otherwise the worker will crash!-data ErrorHandler conn sess st a- = ErrorHandlerIO (InternalError -> a -> IO WorkResult)- | ErrorHandlerSpock (InternalError -> a -> (WebStateM conn sess st) WorkResult)+data ErrorHandler conn sess st err a+ = ErrorHandlerIO ((InternalError err) -> a -> IO WorkResult)+ | ErrorHandlerSpock ((InternalError err) -> a -> (WebStateM conn sess st) WorkResult) -- | Describe how you want jobs in the queue to be performed-type WorkHandler conn sess st a- = a -> ErrorT InternalError (WebStateM conn sess st) WorkResult+type WorkHandler conn sess st err a+ = a -> ExceptT (InternalError err) (WebStateM conn sess st) WorkResult -- | The queue containing scheduled jobs newtype WorkQueue a@@ -68,21 +73,30 @@ , wc_concurrent :: WorkerConcurrentStrategy } +-- | Define a worker+data WorkerDef conn sess st err a+ = WorkerDef+ { wd_config :: WorkerConfig+ , wd_handler :: WorkHandler conn sess st err a+ , wd_errorHandler :: ErrorHandler conn sess st err a+ }+ -- | Create a new background worker and limit the size of the job queue. newWorker :: (MonadTrans t, Monad (t (WebStateM conn sess st)))- => WorkerConfig- -> WorkHandler conn sess st a- -> ErrorHandler conn sess st a+ => WorkerDef conn sess st err a -> t (WebStateM conn sess st) (WorkQueue a)-newWorker wc workHandler errorHandler =- do heart <- getSpockHeart+newWorker wdef =+ do let wc = wd_config wdef+ workHandler = wd_handler wdef+ errorHandler = wd_errorHandler wdef+ heart <- getSpockHeart q <- lift . liftIO $ Q.newQueue (wc_queueLimit wc) _ <- lift . liftIO $ forkIO (workProcessor q workHandler errorHandler heart (wc_concurrent wc)) return (WorkQueue q) workProcessor :: Q.WorkerQueue UTCTime a- -> WorkHandler conn sess st a- -> ErrorHandler conn sess st a+ -> WorkHandler conn sess st err a+ -> ErrorHandler conn sess st err a -> WebState conn sess st -> WorkerConcurrentStrategy -> IO ()@@ -92,17 +106,16 @@ where runWork work = do workRes <-- EX.catch (runSpockIO spockCore $ runErrorT $ workHandler work)- (\(e::SomeException) -> return $ Left (show e))+ EX.catch (runSpockIO spockCore $ runExceptT $ workHandler work)+ (\(e::SomeException) -> return $ Left (InternalErrorMsg $ show e)) case workRes of- Left err ->+ Left errMsg -> case errorHandler of ErrorHandlerIO h ->- h err work+ h errMsg work ErrorHandlerSpock h ->- runSpockIO spockCore $ h err work+ runSpockIO spockCore $ h errMsg work Right r -> return r- loop runningTasksV = do now <- getCurrentTime mWork <- atomically $ Q.dequeue now q@@ -126,7 +139,7 @@ loop runningTasksV launchWork runningTasksV work =- do atomically $ modifyTVar runningTasksV (\x -> x + 1)+ do atomically $ modifyTVar runningTasksV (+ 1) res <- (runWork work `EX.finally` (atomically $ modifyTVar runningTasksV (\x -> x - 1))) case res of WorkRepeatIn secs ->
src/Web/Spock/Worker/Internal/Queue.hs view
@@ -1,21 +1,20 @@ module Web.Spock.Worker.Internal.Queue where -import Control.Applicative-import Control.Concurrent.STM-import Control.Monad-import Data.Maybe-import qualified Data.Map.Strict as M-import qualified Data.Vector as V+import Control.Arrow (second)+import Control.Concurrent.STM+import Control.Monad+import qualified Data.Map.Strict as M+import Data.Maybe+import qualified Data.Vector as V data PureQueue p v = PureQueue { pq_container :: !(M.Map p (V.Vector v))- , pq_maxSize :: !Int+ , pq_maxSize :: !Int } deriving (Show, Eq) emptyPQ :: Int -> PureQueue p v-emptyPQ maxQueueSize =- PureQueue M.empty maxQueueSize+emptyPQ = PureQueue M.empty sizePQ :: PureQueue p v -> Int sizePQ (PureQueue m _) =@@ -23,11 +22,11 @@ isFullPQ :: PureQueue p v -> Bool isFullPQ pq =- sizePQ pq >= (pq_maxSize pq)+ sizePQ pq >= pq_maxSize pq toListPQ :: Ord p => PureQueue p v -> [(p, [v])] toListPQ (PureQueue m _) =- map (\(k, v) -> (k, V.toList v)) (M.toList m)+ map (second V.toList) (M.toList m) fromListPQ :: Ord p => Int -> [(p, [v])] -> Maybe (PureQueue p v) fromListPQ limit kv@@ -86,12 +85,10 @@ WorkerQueue <$> newTVarIO (emptyPQ limit) size :: WorkerQueue p v -> STM Int-size (WorkerQueue qVar) =- readTVar qVar >>= (return . sizePQ)+size (WorkerQueue qVar) = liftM sizePQ (readTVar qVar) isFull :: WorkerQueue p v -> STM Bool-isFull (WorkerQueue qVar) =- readTVar qVar >>= (return . isFullPQ)+isFull (WorkerQueue qVar) = liftM isFullPQ (readTVar qVar) enqueue :: Ord p => p -> v -> WorkerQueue p v -> STM () enqueue prio value (WorkerQueue qVar) =
test/Tests.hs view
@@ -5,4 +5,5 @@ import Test.Framework import {-@ HTF_TESTS @-} Web.Spock.Worker.Internal.QueueTests +main :: IO () main = htfMain htf_importedTests
test/Web/Spock/Worker/Internal/QueueTests.hs view
@@ -4,15 +4,10 @@ ( htf_thisModulesTests ) where -import Web.Spock.Worker.Internal.Queue+import Web.Spock.Worker.Internal.Queue -import Control.Applicative-import Control.Concurrent.STM-import Control.Monad-import Data.Maybe-import Test.Framework-import qualified Data.Map.Strict as M-import qualified Data.Vector as V+import qualified Data.Map.Strict as M+import Test.Framework tAddToMap :: Ord k => k -> a -> M.Map k [a] -> M.Map k [a] tAddToMap k val m =