packages feed

minici-0.1.9: src/Job.hs

module Job (
    Job, DeclaredJob, Job'(..),
    JobSet, DeclaredJobSet, JobSet'(..), jobsetJobs,
    JobOutput(..),
    JobName(..), stringJobName, textJobName,
    ArtifactName(..),
    JobStatus(..),
    jobStatusFinished, jobStatusFailed,
    JobManager(..), newJobManager, cancelAllJobs,
    runJobs, waitForRemainingTasks,

    prepareJob,
    getArtifactWorkPath,
    copyArtifact,

    jobStorageSubdir,

    copyRecursive,
    copyRecursiveForce,
) where

import Control.Concurrent
import Control.Concurrent.STM

import Control.Monad
import Control.Monad.Catch
import Control.Monad.Except
import Control.Monad.IO.Class

import Data.List
import Data.Map (Map)
import Data.Map qualified as M
import Data.Maybe
import Data.Set (Set)
import Data.Set qualified as S
import Data.Text (Text)
import Data.Text qualified as T
import Data.Text.IO qualified as T

import System.Directory
import System.Environment
import System.Exit
import System.FilePath
import System.FilePath.Glob
import System.IO
import System.IO.Temp
import System.Posix.Signals
import System.Process

import Destination
import Job.Types
import Output
import Repo


data JobOutput = JobOutput
    { outArtifacts :: [ArtifactOutput]
    }
    deriving (Eq)

data ArtifactOutput = ArtifactOutput
    { aoutName :: ArtifactName
    , aoutWorkPath :: FilePath
    , aoutStorePath :: FilePath
    }
    deriving (Eq)


data JobStatus a = JobQueued
                 | JobDuplicate JobId (JobStatus a)
                 | JobPreviousStatus (JobStatus a)
                 | JobWaiting [JobName]
                 | JobRunning
                 | JobSkipped
                 | JobError OutputFootnote
                 | JobFailed
                 | JobCancelled
                 | JobDone a
    deriving (Eq)

jobStatusFinished :: JobStatus a -> Bool
jobStatusFinished = \case
    JobQueued      {}   -> False
    JobDuplicate _ s    -> jobStatusFinished s
    JobPreviousStatus s -> jobStatusFinished s
    JobWaiting     {}   -> False
    JobRunning     {}   -> False
    _                   -> True

jobStatusFailed :: JobStatus a -> Bool
jobStatusFailed = \case
    JobDuplicate _ s    -> jobStatusFailed s
    JobPreviousStatus s -> jobStatusFailed s
    JobError       {}   -> True
    JobFailed      {}   -> True
    _                   -> False

jobResult :: JobStatus a -> Maybe a
jobResult = \case
    JobDone x -> Just x
    JobDuplicate _ s -> jobResult s
    JobPreviousStatus s -> jobResult s
    _ -> Nothing

textJobStatus :: JobStatus a -> Text
textJobStatus = \case
    JobQueued -> "queued"
    JobDuplicate {} -> "duplicate"
    JobPreviousStatus s -> textJobStatus s
    JobWaiting _ -> "waiting"
    JobRunning -> "running"
    JobSkipped -> "skipped"
    JobError _ -> "error"
    JobFailed -> "failed"
    JobCancelled -> "cancelled"
    JobDone _ -> "done"

readJobStatus :: (MonadIO m) => Output -> Text -> m a -> m (Maybe (JobStatus a))
readJobStatus tout text readResult = case T.lines text of
    "queued" : _ -> return (Just JobQueued)
    "running" : _ -> return (Just JobRunning)
    "skipped" : _ -> return (Just JobSkipped)
    "error" : note : _ -> Just . JobError <$> liftIO (outputFootnote tout note)
    "failed" : _ -> return (Just JobFailed)
    "cancelled" : _ -> return (Just JobCancelled)
    "done" : _ -> Just . JobDone <$> readResult
    _ -> return Nothing

textJobStatusDetails :: JobStatus a -> Text
textJobStatusDetails = \case
    JobError err -> footnoteText err <> "\n"
    JobPreviousStatus s -> textJobStatusDetails s
    _ -> ""


data JobManager = JobManager
    { jmSemaphore :: TVar Int
    , jmDataDir :: FilePath
    , jmJobs :: TVar (Map JobId (TVar (JobStatus JobOutput)))
    , jmNextTaskId :: TVar TaskId
    , jmNextTask :: TVar (Maybe TaskId)
    , jmReadyTasks :: TVar (Set TaskId)
    , jmRunningTasks :: TVar (Map TaskId ThreadId)
    , jmCancelled :: TVar Bool
    , jmOpenStatusUpdates :: TVar Int
    }

newtype TaskId = TaskId Int
    deriving (Eq, Ord)

data JobCancelledException = JobCancelledException
    deriving (Show)

instance Exception JobCancelledException


newJobManager :: FilePath -> Int -> IO JobManager
newJobManager jmDataDir queueLen = do
    jmSemaphore <- newTVarIO queueLen
    jmJobs <- newTVarIO M.empty
    jmNextTaskId <- newTVarIO (TaskId 0)
    jmNextTask <- newTVarIO Nothing
    jmReadyTasks <- newTVarIO S.empty
    jmRunningTasks <- newTVarIO M.empty
    jmCancelled <- newTVarIO False
    jmOpenStatusUpdates <- newTVarIO 0
    return JobManager {..}

cancelAllJobs :: JobManager -> IO ()
cancelAllJobs JobManager {..} = do
    threads <- atomically $ do
        writeTVar jmCancelled True
        M.elems <$> readTVar jmRunningTasks

    mapM_ (`throwTo` JobCancelledException) threads

reserveTaskId :: JobManager -> STM TaskId
reserveTaskId JobManager {..} = do
    tid@(TaskId n) <- readTVar jmNextTaskId
    writeTVar jmNextTaskId (TaskId (n + 1))
    return tid

runManagedJob :: (MonadIO m, MonadMask m) => JobManager -> TaskId -> m a -> m a -> m a
runManagedJob JobManager {..} tid cancel job = bracket acquire release $ \case
    True -> cancel
    False -> job
  where
    acquire = liftIO $ do
        atomically $ do
            writeTVar jmReadyTasks . S.insert tid =<< readTVar jmReadyTasks
            trySelectNext
        threadId <- myThreadId
        atomically $ do
            readTVar jmCancelled >>= \case
                True -> return True
                False -> readTVar jmNextTask >>= \case
                    Just tid' | tid' == tid -> do
                        writeTVar jmNextTask Nothing
                        writeTVar jmRunningTasks . M.insert tid threadId =<< readTVar jmRunningTasks
                        return False
                    _ -> retry

    release False = liftIO $ atomically $ do
        free <- readTVar jmSemaphore
        writeTVar jmSemaphore $ free + 1
        trySelectNext
    release True = return ()

    trySelectNext = do
        readTVar jmNextTask >>= \case
            Just _ -> return ()
            Nothing -> do
                readTVar jmSemaphore >>= \case
                    0   -> return ()
                    sem -> (S.minView <$> readTVar jmReadyTasks) >>= \case
                        Nothing -> return ()
                        Just ( tid', ready ) -> do
                            writeTVar jmReadyTasks ready
                            writeTVar jmSemaphore (sem - 1)
                            writeTVar jmNextTask (Just tid')
                            writeTVar jmRunningTasks . M.delete tid =<< readTVar jmRunningTasks


runJobs :: JobManager -> Output -> [ Job ]
        -> (JobId -> JobStatus JobOutput -> Bool) -- ^ Rerun condition
        -> IO [ ( Job, TVar (JobStatus JobOutput) ) ]
runJobs mngr@JobManager {..} tout jobs rerun = do
    results <- atomically $ do
        forM jobs $ \job -> do
            tid <- reserveTaskId mngr
            managed <- readTVar jmJobs
            ( job, tid, ) <$> case M.lookup (jobId job) managed of
                Just origVar -> do
                    newTVar . JobDuplicate (jobId job) =<< readTVar origVar

                Nothing -> do
                    statusVar <- newTVar JobQueued
                    writeTVar jmJobs $ M.insert (jobId job) statusVar managed
                    return statusVar

    forM_ results $ \( job, tid, outVar ) -> void $ forkIO $ do
        let handler e = do
                status <- if
                    | Just JobCancelledException <- fromException e -> do
                        return JobCancelled
                    | otherwise -> do
                        JobError <$> outputFootnote tout (T.pack $ displayException e)
                atomically $ writeTVar outVar status
                outputJobFinishedEvent tout job status
        handle handler $ do
            res <- runExceptT $ do
                duplicate <- liftIO $ atomically $ do
                    readTVar outVar >>= \case
                        JobDuplicate jid _ -> do
                            fmap ( jid, ) . M.lookup jid <$> readTVar jmJobs
                        _ -> do
                            return Nothing

                case duplicate of
                    Nothing -> do
                        let jdir = jmDataDir </> jobStorageSubdir (jobId job)
                        readStatusFile tout job jdir >>= \case
                            Just status | status /= JobCancelled && not (rerun (jobId job) status) -> do
                                let status' = JobPreviousStatus status
                                liftIO $ atomically $ writeTVar outVar status'
                                return status'
                            mbStatus -> do
                                when (isJust mbStatus) $ do
                                    liftIO $ removeDirectoryRecursive jdir
                                uses <- waitForUsedArtifacts tout job results outVar
                                runManagedJob mngr tid (return JobCancelled) $ do
                                    liftIO $ atomically $ writeTVar outVar JobRunning
                                    liftIO $ outputEvent tout $ JobStarted (jobId job)
                                    prepareJob jmDataDir job $ \checkoutPath -> do
                                        updateStatusFile mngr jdir outVar
                                        JobDone <$> runJob job uses checkoutPath jdir

                    Just ( jid, origVar ) -> do
                        let wait = do
                                status <- atomically $ do
                                    status <- readTVar origVar
                                    out <- readTVar outVar
                                    if status == out
                                      then retry
                                      else do
                                        writeTVar outVar $ JobDuplicate jid status
                                        return status
                                if jobStatusFinished status
                                  then return $ JobDuplicate jid status
                                  else wait
                        liftIO wait

            atomically $ writeTVar outVar $ either id id res
            outputJobFinishedEvent tout job $ either id id res
    return $ map (\( job, _, var ) -> ( job, var )) results

waitForRemainingTasks :: JobManager -> IO ()
waitForRemainingTasks JobManager {..} = do
    atomically $ do
        remainingStatusUpdates <- readTVar jmOpenStatusUpdates
        when (remainingStatusUpdates > 0) retry

waitForUsedArtifacts
    :: (MonadIO m, MonadError (JobStatus JobOutput) m)
    => Output -> Job
    -> [ ( Job, TaskId, TVar (JobStatus JobOutput) ) ]
    -> TVar (JobStatus JobOutput)
    -> m [ ( ArtifactSpec Evaluated, ArtifactOutput ) ]
waitForUsedArtifacts tout job results outVar = do
    origState <- liftIO $ atomically $ readTVar outVar
    let ( selfSpecs, artSpecs ) = partition ((jobId job ==) . fst) $ jobRequiredArtifacts job

    forM_ selfSpecs $ \( _, artName@(ArtifactName tname) ) -> do
        when (not (artName `elem` map fst (jobArtifacts job))) $ do
            throwError . JobError =<< liftIO (outputFootnote tout $ "Artifact ‘" <> tname <> "’ not produced by the job")

    ujobs <- forM artSpecs $ \( ujobId, uartName ) -> do
        case find (\( j, _, _ ) -> jobId j == ujobId) results of
            Just ( _, _, var ) -> return ( var, ( ujobId, uartName ))
            Nothing -> throwError . JobError =<< liftIO (outputFootnote tout $ "Job ‘" <> textJobId ujobId <> "’ not found")

    let loop prev = do
            ustatuses <- atomically $ do
                ustatuses <- forM ujobs $ \( uoutVar, uartSpec ) -> do
                    (, uartSpec) <$> readTVar uoutVar
                when (Just (map fst ustatuses) == prev) retry
                let remains = map (fromMaybe (JobName "?") . lastJobNameId . fst . snd) $
                        filter (not . jobStatusFinished . fst) ustatuses
                writeTVar outVar $ if null remains then origState else JobWaiting remains
                return ustatuses
            if all (jobStatusFinished . fst) ustatuses
               then return ustatuses
               else loop $ Just $ map fst ustatuses
    ustatuses <- liftIO $ loop Nothing

    forM ustatuses $ \( ustatus, spec@( tjobId, uartName@(ArtifactName tartName)) ) -> do
        case jobResult ustatus of
            Just out -> case find ((==uartName) . aoutName) $ outArtifacts out of
                Just art -> return ( spec, art )
                Nothing -> throwError . JobError =<< liftIO (outputFootnote tout $ "Artifact ‘" <> textJobId tjobId <> "." <> tartName <> "’ not found")
            _ -> throwError JobSkipped

outputJobFinishedEvent :: Output -> Job -> JobStatus a -> IO ()
outputJobFinishedEvent tout job = \case
    JobDuplicate _ s    -> outputEvent tout $ JobIsDuplicate (jobId job) (textJobStatus s)
    JobPreviousStatus s -> outputEvent tout $ JobPreviouslyFinished (jobId job) (textJobStatus s)
    JobSkipped          -> outputEvent tout $ JobWasSkipped (jobId job)
    s                   -> outputEvent tout $ JobFinished (jobId job) (textJobStatus s)

readStatusFile :: (MonadIO m, MonadCatch m) => Output -> Job -> FilePath -> m (Maybe (JobStatus JobOutput))
readStatusFile tout job jdir = do
    handleIOError (\_ -> return Nothing) $ do
        text <- liftIO $ T.readFile (jdir </> "status")
        readJobStatus tout text $ do
            artifacts <- forM (jobArtifacts job) $ \( aoutName@(ArtifactName tname), _ ) -> do
                let adir = jdir </> "artifacts" </> T.unpack tname
                    aoutStorePath = adir </> "data"
                aoutWorkPath <- fmap T.unpack $ liftIO $ T.readFile (adir </> "path")
                return ArtifactOutput {..}

            return JobOutput
                { outArtifacts = artifacts
                }

updateStatusFile :: MonadIO m => JobManager -> FilePath -> TVar (JobStatus JobOutput) -> m ()
updateStatusFile JobManager {..} jdir outVar = liftIO $ do
    atomically $ writeTVar jmOpenStatusUpdates . (+ 1) =<< readTVar jmOpenStatusUpdates
    void $ forkIO $ loop Nothing
  where
    loop prev = do
        status <- atomically $ do
            status <- readTVar outVar
            when (Just status == prev) retry
            return status
        T.writeFile (jdir </> "status") $ textJobStatus status <> "\n" <> textJobStatusDetails status
        if (not (jobStatusFinished status))
          then loop $ Just status
          else atomically $ writeTVar jmOpenStatusUpdates . (subtract 1) =<< readTVar jmOpenStatusUpdates

jobStorageSubdir :: JobId -> FilePath
jobStorageSubdir (JobId jidParts) = "jobs" </> joinPath (map (T.unpack . textJobIdPart) (jidParts))


prepareJob :: (MonadIO m, MonadMask m, MonadFail m) => FilePath -> Job -> (FilePath -> m a) -> m a
prepareJob dir job inner = do
    withSystemTempDirectory "minici" $ \checkoutPath -> do
        forM_ (jobCheckout job) $ \(JobCheckout tree mbsub dest) -> do
            subtree <- maybe return (getSubtree Nothing . makeRelative (treeSubdir tree)) mbsub $ tree
            checkoutAt subtree $ checkoutPath </> fromMaybe "" dest

        liftIO $ forM_ (jobUses job) $ \( jid, aname ) -> do
            modifyError (userError . T.unpack) $ do
                wpath <- getArtifactWorkPath dir jid aname
                let target = checkoutPath </> wpath
                liftIO $ createDirectoryIfMissing True $ takeDirectory target
                copyArtifact dir jid aname target

        let jdir = dir </> jobStorageSubdir (jobId job)
        liftIO $ createDirectoryIfMissing True jdir
        inner checkoutPath

getArtifactStoredPath :: (MonadIO m, MonadError Text m) => FilePath -> JobId -> ArtifactName -> m FilePath
getArtifactStoredPath storageDir jid@(JobId ids) (ArtifactName aname) = do
    let jdir = joinPath $ (storageDir :) $ ("jobs" :) $ map (T.unpack . textJobIdPart) ids
        adir = jdir </> "artifacts" </> T.unpack aname

    liftIO (doesDirectoryExist jdir) >>= \case
        True -> return ()
        False -> throwError $ "job ‘" <> textJobId jid <> "’ not yet executed"

    liftIO (doesDirectoryExist adir) >>= \case
        True -> return ()
        False -> throwError $ "artifact ‘" <> aname <> "’ of job ‘" <> textJobId jid <> "’ not found"

    return adir

getArtifactWorkPath :: (MonadIO m, MonadError Text m) => FilePath -> JobId -> ArtifactName -> m FilePath
getArtifactWorkPath storageDir jid aname = do
    adir <- getArtifactStoredPath storageDir jid aname
    liftIO $ readFile (adir </> "path")

copyArtifact :: (MonadIO m, MonadError Text m) => FilePath -> JobId -> ArtifactName -> FilePath -> m ()
copyArtifact storageDir jid aname tpath = do
    adir <- getArtifactStoredPath storageDir jid aname
    liftIO $ copyRecursive (adir </> "data") tpath


runJob :: Job -> [ ( ArtifactSpec Evaluated, ArtifactOutput) ] -> FilePath -> FilePath -> ExceptT (JobStatus JobOutput) IO JobOutput
runJob job uses checkoutPath jdir = do
    bracket (liftIO $ openFile (jdir </> "log") WriteMode) (liftIO . hClose) $ \logs -> do
        forM_ (fromMaybe [] $ jobRecipe job) $ \ep -> do
            ( p, input ) <- case ep of
                Left p -> return ( p, "" )
                Right script -> do
                    sh <- fromMaybe "/bin/sh" <$> liftIO (lookupEnv "SHELL")
                    return ( proc sh [], script )
            (Just hin, _, _, hp) <- liftIO $ createProcess_ "" p
                { cwd = Just checkoutPath
                , std_in = CreatePipe
                , std_out = UseHandle logs
                , std_err = UseHandle logs
                }
            liftIO $ void $ forkIO $ do
                T.hPutStr hin input
                hClose hin
            liftIO (waitForProcess hp) >>= \case
                ExitSuccess -> return ()
                ExitFailure n
                    | fromIntegral n == -sigINT -> throwError JobCancelled
                    | otherwise -> throwError JobFailed

        artifacts <- forM (jobArtifacts job) $ \( name@(ArtifactName tname), pathPattern ) -> do
            let adir = jdir </> "artifacts" </> T.unpack tname
            path <- liftIO (globDir1 pathPattern checkoutPath) >>= \case
                [ path ] -> return path
                found -> do
                    liftIO $ hPutStrLn logs $
                        (if null found then "no file" else "multiple files") <> " found matching pattern ‘" <>
                        decompile pathPattern <> "’ for artifact ‘" <> T.unpack tname <> "’"
                    throwError JobFailed
            let target = adir </> "data"
                workPath = makeRelative checkoutPath path
            liftIO $ do
                createDirectoryIfMissing True $ takeDirectory target
                copyRecursiveForce path target
                T.writeFile (adir </> "path") $ T.pack workPath
            return $ ArtifactOutput
                { aoutName = name
                , aoutWorkPath = workPath
                , aoutStorePath = target
                }

        forM_ (jobPublish job) $ \pub -> do
            Just aout <- return $ lookup (jpArtifact pub) $ map (\aout -> ( ( jobId job, aoutName aout ), aout )) artifacts ++ uses
            let ppath = case jpPath pub of
                    Just path
                        | hasTrailingPathSeparator path -> path </> takeFileName (aoutWorkPath aout)
                        | otherwise                     -> path
                    Nothing                             -> aoutWorkPath aout
            copyToDestination (aoutStorePath aout) (jpDestination pub) ppath

        return JobOutput
            { outArtifacts = artifacts
            }