feat(jobs): batch job offloading

BREAKING CHANGE: Job offloading
This commit is contained in:
Gregor Kleen 2021-02-01 09:52:47 +01:00
parent cb1e715e9b
commit 09fb26f1a8
10 changed files with 584 additions and 438 deletions

View File

@ -17,6 +17,8 @@ import qualified Data.Text.Lazy.Builder as Text.Builder
import qualified Data.HashSet as HashSet import qualified Data.HashSet as HashSet
import qualified Data.HashMap.Strict as HashMap import qualified Data.HashMap.Strict as HashMap
import qualified Data.UUID as UUID
deriveJSON defaultOptions deriveJSON defaultOptions
{ constructorTagModifier = camelToPathPiece' 1 { constructorTagModifier = camelToPathPiece' 1
@ -32,9 +34,12 @@ getAdminCrontabR = do
let mCrontab = mCrontab' <&> _2 %~ filter (hasn't $ _3 . _MatchNone) let mCrontab = mCrontab' <&> _2 %~ filter (hasn't $ _3 . _MatchNone)
instanceId <- getsYesod appInstanceID
selectRep $ do selectRep $ do
provideRep $ do provideRep $ do
crontabBearer <- runMaybeT . hoist runDB $ do crontabBearer <- runMaybeT . hoist runDB $ do
guardM $ hasGlobalGetParam GetGenerateToken
uid <- MaybeT maybeAuthId uid <- MaybeT maybeAuthId
guardM . lift . existsBy $ UniqueUserGroupMember UserGroupCrontab uid guardM . lift . existsBy $ UniqueUserGroupMember UserGroupCrontab uid
@ -49,6 +54,10 @@ getAdminCrontabR = do
<section> <section>
<pre .token> <pre .token>
#{toPathPiece t} #{toPathPiece t}
<section>
<dl .deflist>
<dt .deflist__dt>_{MsgInstanceId}
<dd .deflist__dd .uuid>#{UUID.toText instanceId}
<section> <section>
$maybe (genTime, crontab) <- mCrontab $maybe (genTime, crontab) <- mCrontab
<p> <p>

View File

@ -26,6 +26,7 @@ getMetricsR = selectRep $ do
samples <- sortBy metricSort <$> collectMetrics samples <- sortBy metricSort <$> collectMetrics
metricsBearer <- runMaybeT . hoist runDB $ do metricsBearer <- runMaybeT . hoist runDB $ do
guardM $ hasGlobalGetParam GetGenerateToken
uid <- MaybeT maybeAuthId uid <- MaybeT maybeAuthId
guardM . lift . existsBy $ UniqueUserGroupMember UserGroupMetrics uid guardM . lift . existsBy $ UniqueUserGroupMember UserGroupMetrics uid

View File

@ -10,6 +10,7 @@ module Jobs
import Import hiding (StateT) import Import hiding (StateT)
import Jobs.Types as Types hiding (JobCtl(JobCtlQueue)) import Jobs.Types as Types hiding (JobCtl(JobCtlQueue))
import Jobs.Queue import Jobs.Queue
import Jobs.Offload
import Jobs.Crontab import Jobs.Crontab
import qualified Data.Conduit.List as C import qualified Data.Conduit.List as C
@ -105,6 +106,7 @@ handleJobs foundation@UniWorX{..}
jobShutdown <- liftIO newEmptyTMVarIO jobShutdown <- liftIO newEmptyTMVarIO
jobCurrentCrontab <- liftIO $ newTVarIO Nothing jobCurrentCrontab <- liftIO $ newTVarIO Nothing
jobHeldLocks <- liftIO $ newTVarIO Set.empty jobHeldLocks <- liftIO $ newTVarIO Set.empty
jobOffload <- liftIO newEmptyTMVarIO
registerJobHeldLocksCount jobHeldLocks registerJobHeldLocksCount jobHeldLocks
registerJobWorkerQueueDepth appJobState registerJobWorkerQueueDepth appJobState
atomically $ putTMVar appJobState JobState atomically $ putTMVar appJobState JobState
@ -155,7 +157,9 @@ manageJobPool foundation@UniWorX{..} unmask = shutdownOnException $ \routeExc ->
atomically . asum $ atomically . asum $
[ spawnMissingWorkers [ spawnMissingWorkers
, reapDeadWorkers , reapDeadWorkers
] ++ maybe [] (\(cTime, delay) -> [return () <$ waitDelay delay, transferJobs cTime]) transferInfo ++ ] ++ maybe [] (\(cTime, delay) -> [return () <$ waitDelay delay, transferJobs cTime]) transferInfo
++ maybeToList (manageOffloadHandler <$> mkJobOffloadHandler (appDatabaseConf appSettings') (appJobMode appSettings'))
++
[ terminateGracefully terminate' [ terminateGracefully terminate'
] ]
where where
@ -286,6 +290,27 @@ manageJobPool foundation@UniWorX{..} unmask = shutdownOnException $ \routeExc ->
return $ $logWarnS "JobPoolManager" [st|Moved #{tshow (olength movePairs)} long-unadressed jobs from #{tshow (olength senders)} senders to #{tshow (olength receivers)} receivers|] return $ $logWarnS "JobPoolManager" [st|Moved #{tshow (olength movePairs)} long-unadressed jobs from #{tshow (olength senders)} senders to #{tshow (olength receivers)} receivers|]
manageOffloadHandler :: (ReaderT UniWorX m JobOffloadHandler) -> STM (ContT () m ())
manageOffloadHandler spawn = do
shouldTerminate' <- readTMVar appJobState >>= fmap not . isEmptyTMVar . jobShutdown
guard $ not shouldTerminate'
JobContext{jobOffload} <- jobContext <$> readTMVar appJobState
cOffload <- tryReadTMVar jobOffload
let respawn = do
nOffload <- lift $ runReaderT spawn foundation
atomically $ do
putTMVar jobOffload nOffload
whenIsJust cOffload $ \pOffload -> do
pOutgoing <- readTVar $ jobOffloadOutgoing pOffload
modifyTVar (jobOffloadOutgoing nOffload) (pOutgoing <>)
respawn <$ case cOffload of
Nothing -> return ()
Just JobOffloadHandler{..} -> waitSTM jobOffloadHandler
stopJobCtl :: MonadUnliftIO m => UniWorX -> m () stopJobCtl :: MonadUnliftIO m => UniWorX -> m ()
-- ^ Stop all worker threads currently running -- ^ Stop all worker threads currently running
stopJobCtl UniWorX{appJobState} = do stopJobCtl UniWorX{appJobState} = do
@ -471,7 +496,16 @@ handleJobs' wNum = C.mapM_ $ \jctl -> hoist delimitInternalState . withJobWorker
$logDebugS logIdent "JobCtlQueue..." $logDebugS logIdent "JobCtlQueue..."
lift $ queueJob' job lift $ queueJob' job
$logInfoS logIdent "JobCtlQueue" $logInfoS logIdent "JobCtlQueue"
handleCmd (JobCtlPerform jId) = handle handleQueueException . jLocked jId $ \(Entity _ j@QueuedJob{..}) -> lift $ do handleCmd (JobCtlPerform jId) = do
jMode <- getsYesod $ view _appJobMode
case jMode of
JobsLocal{} -> performLocal
JobsOffload -> performOffload
where
performOffload = hoist atomically $ do
JobOffloadHandler{..} <- lift . readTMVar =<< asks jobOffload
lift $ modifyTVar jobOffloadOutgoing (`snoc` jId)
performLocal = handle handleQueueException . jLocked jId $ \(Entity _ j@QueuedJob{..}) -> lift $ do
content <- case fromJSON queuedJobContent of content <- case fromJSON queuedJobContent of
Aeson.Success c -> return c Aeson.Success c -> return c
Aeson.Error t -> do Aeson.Error t -> do

View File

@ -27,6 +27,29 @@ determineCrontab :: DB (Crontab JobCtl)
determineCrontab = execWriterT $ do determineCrontab = execWriterT $ do
UniWorX{ appSettings' = AppSettings{..} } <- getYesod UniWorX{ appSettings' = AppSettings{..} } <- getYesod
whenIsJust appJobCronInterval $ \interval ->
tell $ HashMap.singleton
JobCtlDetermineCrontab
Cron
{ cronInitial = CronAsap
, cronRepeat = CronRepeatScheduled CronAsap
, cronRateLimit = interval
, cronNotAfter = Right CronNotScheduled
}
tell . flip foldMap universeF $ \kind ->
case appHealthCheckInterval kind of
Just int -> HashMap.singleton
(JobCtlGenerateHealthReport kind)
Cron
{ cronInitial = CronAsap
, cronRepeat = CronRepeatScheduled CronAsap
, cronRateLimit = int
, cronNotAfter = Right CronNotScheduled
}
Nothing -> mempty
when (is _JobsLocal appJobMode) $ do
case appJobFlushInterval of case appJobFlushInterval of
Just interval -> tell $ HashMap.singleton Just interval -> tell $ HashMap.singleton
JobCtlFlush JobCtlFlush
@ -38,16 +61,6 @@ determineCrontab = execWriterT $ do
} }
Nothing -> return () Nothing -> return ()
whenIsJust appJobCronInterval $ \interval ->
tell $ HashMap.singleton
JobCtlDetermineCrontab
Cron
{ cronInitial = CronAsap
, cronRepeat = CronRepeatScheduled CronAsap
, cronRateLimit = interval
, cronNotAfter = Right CronNotScheduled
}
oldestInvitationMUTC <- lift $ preview (_head . _entityVal . _invitationExpiresAt . _Just) <$> selectList [InvitationExpiresAt !=. Nothing] [Asc InvitationExpiresAt, LimitTo 1] oldestInvitationMUTC <- lift $ preview (_head . _entityVal . _invitationExpiresAt . _Just) <$> selectList [InvitationExpiresAt !=. Nothing] [Asc InvitationExpiresAt, LimitTo 1]
whenIsJust oldestInvitationMUTC $ \oldestInvUTC -> tell $ HashMap.singleton whenIsJust oldestInvitationMUTC $ \oldestInvUTC -> tell $ HashMap.singleton
(JobCtlQueue JobPruneInvitations) (JobCtlQueue JobPruneInvitations)
@ -119,18 +132,6 @@ determineCrontab = execWriterT $ do
, cronNotAfter = Right CronNotScheduled , cronNotAfter = Right CronNotScheduled
} }
tell . flip foldMap universeF $ \kind ->
case appHealthCheckInterval kind of
Just int -> HashMap.singleton
(JobCtlGenerateHealthReport kind)
Cron
{ cronInitial = CronAsap
, cronRepeat = CronRepeatScheduled CronAsap
, cronRateLimit = int
, cronNotAfter = Right CronNotScheduled
}
Nothing -> mempty
let newyear = cronCalendarAny let newyear = cronCalendarAny
{ cronDayOfYear = cronMatchOne 1 { cronDayOfYear = cronMatchOne 1
} }
@ -472,7 +473,7 @@ determineCrontab = execWriterT $ do
, cronRateLimit = appNotificationRateLimit , cronRateLimit = appNotificationRateLimit
, cronNotAfter = maybe (Right CronNotScheduled) (Right . CronTimestamp . utcToLocalTimeTZ appTZ) $ nBot =<< minimumOf (folded . _entityVal . _allocationStaffAllocationTo . to NTop . filtered (> NTop (Just staffAllocationFrom))) allocs , cronNotAfter = maybe (Right CronNotScheduled) (Right . CronTimestamp . utcToLocalTimeTZ appTZ) $ nBot =<< minimumOf (folded . _entityVal . _allocationStaffAllocationTo . to NTop . filtered (> NTop (Just staffAllocationFrom))) allocs
} }
iforM (allocationTimes AllocationRegisterTo) $ \registerTo allocs' -> do iforM_ (allocationTimes AllocationRegisterTo) $ \registerTo allocs' -> do
let allocs = flip filter allocs' $ \(Entity _ Allocation{..}) -> maybe True (> registerTo) allocationStaffAllocationTo let allocs = flip filter allocs' $ \(Entity _ Allocation{..}) -> maybe True (> registerTo) allocationStaffAllocationTo
tell $ HashMap.singleton tell $ HashMap.singleton
(JobCtlQueue $ JobQueueNotification NotificationAllocationUnratedApplications{ nAllocations = setOf (folded . _entityKey) allocs }) (JobCtlQueue $ JobQueueNotification NotificationAllocationUnratedApplications{ nAllocations = setOf (folded . _entityKey) allocs })

70
src/Jobs/Offload.hs Normal file
View File

@ -0,0 +1,70 @@
module Jobs.Offload
( mkJobOffloadHandler
) where
import Import hiding (bracket, js)
import Jobs.Types
import Jobs.Queue
import qualified Database.PostgreSQL.Simple as PG
import qualified Database.PostgreSQL.Simple.Types as PG
import qualified Database.PostgreSQL.Simple.Notification as PG
import Database.Persist.Postgresql (PostgresConf, pgConnStr)
import Data.Text.Encoding (decodeUtf8')
import UnliftIO.Exception (bracket)
jobOffloadChannel :: Text
jobOffloadChannel = "job-offload"
mkJobOffloadHandler :: forall m.
( MonadResource m
, MonadUnliftIO m
, MonadThrow m, MonadReader UniWorX m
, MonadLogger m
)
=> PostgresConf -> JobMode
-> Maybe (m JobOffloadHandler)
mkJobOffloadHandler dbConf jMode
| is _JobsLocal jMode, hasn't (_jobsAcceptOffload . only True) jMode = Nothing
| otherwise = Just $ do
jobOffloadOutgoing <- newTVarIO mempty
jobOffloadHandler <- allocateAsync . bracket (liftIO . PG.connectPostgreSQL $ pgConnStr dbConf) (liftIO . PG.close) $ \pgConn -> do
myPid <- liftIO $ PG.getBackendPID pgConn
let shouldListen = has (_jobsAcceptOffload . only True) jMode
when shouldListen $
void . liftIO $ PG.execute pgConn "LISTEN ?" (PG.Only $ PG.Identifier jobOffloadChannel)
foreverBreak $ \(($ ()) -> terminate) -> do
UniWorX{appJobState} <- ask
shouldTerminate <- atomically $ readTMVar appJobState >>= fmap not . isEmptyTMVar . jobShutdown
when shouldTerminate terminate
let
getInput = do
n@PG.Notification{..} <- liftIO $ PG.getNotification pgConn
if | notificationPid == myPid || notificationChannel /= (encodeUtf8 jobOffloadChannel) -> getInput
| otherwise -> return n
getOutput = atomically $ do
jQueue <- readTVar jobOffloadOutgoing
case jQueue of
j :< js -> j <$ writeTVar jobOffloadOutgoing js
_other -> mzero
io <- lift $ if
| shouldListen -> getInput `race` getOutput
| otherwise -> Right <$> getOutput
case io of
Left PG.Notification{..}
| Just jId <- fromPathPiece =<< either (const Nothing) Just (decodeUtf8' notificationData)
-> writeJobCtl $ JobCtlPerform jId
| otherwise
-> $logErrorS "JobOffloadHandler" $ "Could not parse incoming notification data: " <> tshow notificationData
Right jId -> void . liftIO $ PG.execute pgConn "NOTIFY ?, ?" (PG.Identifier jobOffloadChannel, encodeUtf8 $ toPathPiece jId)
return JobOffloadHandler{..}

View File

@ -9,7 +9,7 @@ module Jobs.Types
, classifyJobCtl , classifyJobCtl
, YesodJobDB , YesodJobDB
, JobHandler(..), _JobHandlerAtomic, _JobHandlerException , JobHandler(..), _JobHandlerAtomic, _JobHandlerException
, JobContext(..) , JobOffloadHandler(..), JobContext(..)
, JobState(..), _jobWorkers, _jobWorkerName, _jobContext, _jobPoolManager, _jobCron, _jobShutdown, _jobCurrentCrontab , JobState(..), _jobWorkers, _jobWorkerName, _jobContext, _jobPoolManager, _jobCron, _jobShutdown, _jobCurrentCrontab
, jobWorkerNames , jobWorkerNames
, JobWorkerState(..), _jobWorkerJobCtl, _jobWorkerJob , JobWorkerState(..), _jobWorkerJobCtl, _jobWorkerJob
@ -238,10 +238,16 @@ showWorkerId = tshow . hashUnique . jobWorkerUnique
newWorkerId :: MonadIO m => m JobWorkerId newWorkerId :: MonadIO m => m JobWorkerId
newWorkerId = JobWorkerId <$> liftIO newUnique newWorkerId = JobWorkerId <$> liftIO newUnique
data JobOffloadHandler = JobOffloadHandler
{ jobOffloadHandler :: Async ()
, jobOffloadOutgoing :: TVar (Seq QueuedJobId)
}
data JobContext = JobContext data JobContext = JobContext
{ jobCrontab :: TVar (Crontab JobCtl) { jobCrontab :: TVar (Crontab JobCtl)
, jobConfirm :: TVar (HashMap JobCtl (NonEmpty (TMVar (Maybe SomeException)))) , jobConfirm :: TVar (HashMap JobCtl (NonEmpty (TMVar (Maybe SomeException))))
, jobHeldLocks :: TVar (Set QueuedJobId) , jobHeldLocks :: TVar (Set QueuedJobId)
, jobOffload :: TMVar JobOffloadHandler
} }

View File

@ -206,9 +206,16 @@ data AppSettings = AppSettings
, appInitialInstanceID :: Maybe (Either FilePath UUID) , appInitialInstanceID :: Maybe (Either FilePath UUID)
, appRibbon :: Maybe Text , appRibbon :: Maybe Text
, appJobMode :: JobMode
, appMemcacheAuth :: Bool , appMemcacheAuth :: Bool
} deriving Show } deriving Show
data JobMode = JobsLocal { jobsAcceptOffload :: Bool }
| JobsOffload
deriving (Eq, Ord, Read, Show, Generic, Typeable)
deriving anyclass (Hashable)
data ApprootScope = ApprootUserGenerated | ApprootDefault data ApprootScope = ApprootUserGenerated | ApprootDefault
deriving (Eq, Ord, Read, Show, Enum, Bounded, Generic, Typeable) deriving (Eq, Ord, Read, Show, Enum, Bounded, Generic, Typeable)
deriving anyclass (Universe, Finite, Hashable) deriving anyclass (Universe, Finite, Hashable)
@ -342,6 +349,11 @@ deriveFromJSON defaultOptions
{ fieldLabelModifier = camelToPathPiece' 2 { fieldLabelModifier = camelToPathPiece' 2
} ''UserDefaultConf } ''UserDefaultConf
deriveJSON defaultOptions
{ fieldLabelModifier = camelToPathPiece' 1
, constructorTagModifier = camelToPathPiece' 1
} ''JobMode
instance FromJSON LdapConf where instance FromJSON LdapConf where
parseJSON = withObject "LdapConf" $ \o -> do parseJSON = withObject "LdapConf" $ \o -> do
ldapTls <- o .:? "tls" ldapTls <- o .:? "tls"
@ -596,6 +608,8 @@ instance FromJSON AppSettings where
appMemcacheAuth <- o .:? "memcache-auth" .!= False appMemcacheAuth <- o .:? "memcache-auth" .!= False
appJobMode <- o .:? "job-mode" .!= JobsLocal True
return AppSettings{..} return AppSettings{..}
makeClassy_ ''AppSettings makeClassy_ ''AppSettings

View File

@ -70,6 +70,7 @@ import Control.Monad.Writer.Class (MonadWriter(..))
import Control.Monad.Catch import Control.Monad.Catch
import Control.Monad.Morph (hoist) import Control.Monad.Morph (hoist)
import Control.Monad.Fail import Control.Monad.Fail
import Control.Monad.Trans.Cont (ContT, evalContT, callCC)
import Language.Haskell.TH import Language.Haskell.TH
import Language.Haskell.TH.Instances () import Language.Haskell.TH.Instances ()
@ -943,6 +944,11 @@ forever' :: Monad m
-> m b -> m b
forever' start cont = cont start >>= flip forever' cont forever' start cont = cont start >>= flip forever' cont
foreverBreak :: Monad m
=> ((r -> ContT r m b) -> ContT r m a)
-> m r
foreverBreak cont = evalContT . callCC $ forever . cont
-------------- --------------
-- Foldable -- -- Foldable --

View File

@ -5,6 +5,7 @@
module Utils.Lens ( module Utils.Lens ) where module Utils.Lens ( module Utils.Lens ) where
import Import.NoModel import Import.NoModel
import Settings
import Model import Model
import Model.Rating import Model.Rating
import qualified ClassyPrelude.Yesod as Yesod (HasHttpManager(..)) import qualified ClassyPrelude.Yesod as Yesod (HasHttpManager(..))
@ -272,6 +273,9 @@ makePrisms ''AllocationPriority
makePrisms ''RoomReference makePrisms ''RoomReference
makeLenses_ ''RoomReference makeLenses_ ''RoomReference
makePrisms ''JobMode
makeLenses_ ''JobMode
-- makeClassy_ ''Load -- makeClassy_ ''Load
-------------------------- --------------------------

View File

@ -32,6 +32,7 @@ data GlobalGetParam = GetLang
| GetDownload | GetDownload
| GetError | GetError
| GetSelectTable | GetSelectTable
| GetGenerateToken
deriving (Eq, Ord, Enum, Bounded, Read, Show, Generic) deriving (Eq, Ord, Enum, Bounded, Read, Show, Generic)
deriving anyclass (Universe, Finite) deriving anyclass (Universe, Finite)