baikai-0.8.0.0: src/Baikai/Provider/Cli/Process/Internal.hs
{-# LANGUAGE CPP #-}
{-# LANGUAGE ForeignFunctionInterface #-}
-- | Internal batch-process ownership shared by the vendor adapters. Not
-- covered by the public PVP interface. POSIX groups contain inherited children,
-- not processes deliberately escaping into another group or session.
module Baikai.Provider.Cli.Process.Internal
( withOwnedProcess,
withOwnedWorker,
)
where
import Control.Concurrent (forkIO, forkIOWithUnmask, killThread, threadDelay)
import Control.Concurrent.MVar (newEmptyMVar, putMVar, readMVar, tryReadMVar)
import Control.Exception
( SomeAsyncException,
SomeException,
catch,
displayException,
finally,
fromException,
mask,
onException,
throwIO,
try,
)
import Control.Monad (forM_, void)
import Data.ByteString qualified as BS
import Data.Text qualified as Text
import Data.Text.Encoding qualified as Text
import System.IO (Handle, hClose, stderr)
import System.Process qualified as P
#ifndef mingw32_HOST_OS
import Data.ByteString.Char8 qualified as BS8
import Data.Maybe (mapMaybe)
import Foreign.C.Error (throwErrnoIfMinus1)
import Foreign.C.Types (CInt (..))
import GHC.Clock (getMonotonicTimeNSec)
import System.Exit (ExitCode (..))
import System.Posix.Process (getProcessStatus)
import System.Posix.Signals (Signal, sigINT, sigTERM, sigKILL, signalProcess, signalProcessGroup)
import System.IO.Error (isDoesNotExistError)
import System.Posix.Types (ProcessID, ProcessGroupID)
import Text.Read (readMaybe)
#endif
-- | Run a callback with ordinary interruptibility, then finish release in a
-- private masked worker. Repeated cancellation interrupts only the join, never
-- release. Every worker exit publishes completion. Preserve the first async
-- exception, including one delivered after a successful callback.
ownedScope :: IO resource -> (resource -> IO ()) -> (resource -> IO a) -> IO a
ownedScope acquire release use = mask $ \restore -> do
done <- newEmptyMVar
resource <- acquire
outcome <- try (restore (use resource))
_ <- forkIO $ try (release resource) >>= putMVar done
let initial = case outcome of
Left e | isAsync e -> Just e
_ -> Nothing
join first =
((,first) <$> readMVar done) `catch` \e ->
if isAsync e
then join (case first of Nothing -> Just e; Just _ -> first)
else throwIO e
(cleaned, cancellation) <- join initial
case cancellation of
Just e -> do
-- Cancellation remains asynchronous even if cleanup diagnosed a failure.
-- Make that exceptional failure visible rather than silently losing it.
case cleaned of
Left failure -> void (try (BS.hPut stderr (Text.encodeUtf8 (Text.pack ("baikai CLI cleanup failed: " <> displayException failure <> "\n")))) :: IO (Either SomeException ()))
Right () -> pure ()
throwIO e
Nothing -> case cleaned of
Left e -> throwIO (e :: SomeException)
Right () -> either throwIO pure outcome
where
isAsync e = case fromException e :: Maybe SomeAsyncException of
Just _ -> True
Nothing -> False
-- | The supplied join action returns the reader's value or exception. On every
-- scope exit an unfinished reader is cancelled and joined before pipe closure.
withOwnedWorker :: forall a b. IO a -> (IO a -> IO b) -> IO b
withOwnedWorker action use = ownedScope acquire release $ \(_, done) ->
use (readMVar done >>= either throwIO pure)
where
acquire = do
done <- newEmptyMVar
tid <- forkIOWithUnmask $ \unmask -> do
result <- try (unmask action)
putMVar done (result :: Either SomeException a)
pure (tid, done)
release (tid, done) = do
completed <- tryReadMVar done
case completed of
Nothing -> killThread tid
Just _ -> pure ()
void (readMVar done)
-- | Own handles, group identity, and synchronous direct-child reaping. Reader
-- scopes must be nested inside this callback, so no reader retains a handle lock
-- when release closes the pipes. Windows has direct-child-only cleanup.
withOwnedProcess ::
P.CreateProcess ->
(Maybe Handle -> Maybe Handle -> Maybe Handle -> P.ProcessHandle -> IO a) ->
IO a
withOwnedProcess spec use = ownedScope (acquireProcess spec) releaseProcess $ \(handles, _) ->
let (input, output, err, ph) = handles in use input output err ph
type ProcessHandles = (Maybe Handle, Maybe Handle, Maybe Handle, P.ProcessHandle)
#ifdef mingw32_HOST_OS
type GroupIdentity = ()
acquireProcess :: P.CreateProcess -> IO (ProcessHandles, GroupIdentity)
acquireProcess spec = do
handles <- P.createProcess spec {P.create_group = True}
pure (handles, ())
releaseProcess :: (ProcessHandles, GroupIdentity) -> IO ()
releaseProcess (handles@(_, _, _, ph), _) =
(P.terminateProcess ph >> void (P.waitForProcess ph)) `finally` closeHandles handles
#else
type GroupIdentity = (ProcessGroupID, ProcessID)
acquireProcess :: P.CreateProcess -> IO (ProcessHandles, GroupIdentity)
acquireProcess spec = do
handles@(_, _, _, ph) <- P.createProcess spec
{P.create_group = True, P.new_session = False, P.delegate_ctlc = False}
-- No callback can reap the leader before the anchor has joined. Even a
-- quickly exiting leader is still an unreaped group member at this point.
anchored <- try $ do
group <- P.getPid ph >>= maybe (ioError (userError "CLI leader PID unavailable")) pure
anchor <- throwErrnoIfMinus1 "CLI group anchor" (c_anchor (fromIntegral group))
pure (group, fromIntegral anchor)
case anchored of
Right identity -> pure (handles, identity)
Left (e :: SomeException) ->
-- The unreaped leader reserves the group during acquisition failure.
-- This nested scope also protects failed-acquisition cleanup from repeats.
ownedScope (pure handles) cleanupFailure (\_ -> throwIO e)
where
cleanupFailure handles@(_, _, _, ph) =
(do
pid <- P.getPid ph
forM_ pid $ \group -> terminateGroup group (-1) `onException` signalGroup sigKILL group
) `finally` (void (P.waitForProcess ph) `finally` closeHandles handles)
releaseProcess :: (ProcessHandles, GroupIdentity) -> IO ()
releaseProcess (handles@(_, _, _, ph), (group, anchor)) =
((terminateGroup group anchor `onException` signalGroup sigKILL group)
`finally` void (P.waitForProcess ph))
`finally` (finishAnchor `finally` closeHandles handles)
where
finishAnchor = do
-- Darwin reports EPERM when signalling a zombie-only group/PID.
-- Signalling is already complete, so we may now reap the anchor.
status <- getProcessStatus False False anchor
case status of
Just _ -> pure ()
Nothing -> signalPid sigKILL anchor >> void (getProcessStatus True False anchor)
#endif
closeHandles :: ProcessHandles -> IO ()
closeHandles (input, output, err, _) =
-- finally ensures one failed close cannot skip the other owned descriptors.
forM_ input hClose `finally` (forM_ output hClose `finally` forM_ err hClose)
#ifndef mingw32_HOST_OS
foreign import ccall safe "baikai_cli_group_anchor" c_anchor :: CInt -> IO CInt
signalGroup :: Signal -> ProcessGroupID -> IO ()
signalGroup signal group = absentOnly (signalProcessGroup signal group)
signalPid :: Signal -> ProcessID -> IO ()
signalPid signal pid = absentOnly (signalProcess signal pid)
-- Only ESRCH is benign. EPERM and all other observation/signal errors propagate.
absentOnly :: IO () -> IO ()
absentOnly action = action `catch` \e -> do
if isDoesNotExistError e then pure () else throwIO (e :: IOError)
terminateGroup :: ProcessGroupID -> ProcessID -> IO ()
terminateGroup group anchor = do
live <- runningMembers group anchor
if null live then pure () else do
signalGroup sigINT group
interrupted <- settle 100000
if interrupted then pure () else do
signalGroup sigTERM group
terminated <- settle 500000
if terminated then pure () else do
-- The anchor is killed too, but stays unreaped until all observations
-- and group signals finish. Its zombie membership reserves the PGID.
signalGroup sigKILL group
killed <- settle 1000000
if killed then pure () else ioError (userError "CLI group has live survivors after SIGKILL")
where
settle :: Int -> IO Bool
settle micros = do
start <- getMonotonicTimeNSec
let deadline = start + fromIntegral micros * 1000
poll = do
members <- runningMembers group anchor
if null members then pure True else do
now <- getMonotonicTimeNSec
if now >= deadline then pure False else threadDelay 10000 >> poll
poll
-- Darwin and Linux both expose PID, PGID, and process state via POSIX ps.
-- A zombie is no longer running; adopted grandchildren are reaped elsewhere.
-- Never infer successful cleanup from kill(0), which includes zombies.
runningMembers :: ProcessGroupID -> ProcessID -> IO [ProcessID]
runningMembers group anchor = do
(code, bytes) <- P.withCreateProcess
(P.proc "/bin/ps" ["-axo", "pid=,pgid=,stat="])
{ P.std_in = P.NoStream, P.std_out = P.CreatePipe, P.std_err = P.CreatePipe }
$ \_ output err ph -> case (output, err) of
(Just out, Just errors) -> withOwnedWorker (BS.hGetContents errors) $ \joinErr -> do
captured <- BS.hGetContents out
void joinErr
status <- P.waitForProcess ph
pure (status, captured)
_ -> ioError (userError "ps capture handles unavailable")
case code of
ExitFailure n -> ioError (userError ("CLI group observation failed: ps exit " <> show n))
ExitSuccess -> do
rows <- traverse parseRow (BS8.lines bytes)
pure (mapMaybe running rows)
where
parseRow row = case BS8.words row of
[pid, pgid, state] -> case (readMaybe (BS8.unpack pid), readMaybe (BS8.unpack pgid)) of
(Just p, Just g) -> pure (p, g, state)
_ -> invalid
_ -> invalid
where invalid = ioError (userError "CLI group observation: malformed ps row")
running (pid, pgid, state)
| pgid == group && pid /= anchor && not (BS8.elem 'Z' state) = Just pid
| otherwise = Nothing
#endif