2019-07-16 03:01:45 +03:00
|
|
|
{-# OPTIONS_GHC -Wwarn #-}
|
|
|
|
|
2019-08-22 02:49:08 +03:00
|
|
|
module Vere.Pier
|
|
|
|
( booted, resumed, pier, runPersist, runCompute, generateBootSeq
|
|
|
|
) where
|
2019-05-30 23:19:26 +03:00
|
|
|
|
2019-07-16 03:01:45 +03:00
|
|
|
import UrbitPrelude
|
2019-07-24 04:34:16 +03:00
|
|
|
|
|
|
|
import Arvo
|
2019-05-30 23:19:26 +03:00
|
|
|
import Vere.Pier.Types
|
2019-08-01 08:48:08 +03:00
|
|
|
import System.Random
|
2019-07-19 03:52:53 +03:00
|
|
|
|
2019-07-21 04:29:39 +03:00
|
|
|
import System.Directory (createDirectoryIfMissing)
|
|
|
|
import System.Posix.Files (ownerModes, setFileMode)
|
2019-08-01 05:34:14 +03:00
|
|
|
import Vere.Ames (ames)
|
|
|
|
import Vere.Behn (behn)
|
2019-09-05 23:09:45 +03:00
|
|
|
import Vere.Http.Client (client)
|
2019-08-08 01:24:02 +03:00
|
|
|
import Vere.Http.Server (serv)
|
2019-07-21 04:29:39 +03:00
|
|
|
import Vere.Log (EventLog)
|
2019-09-04 01:17:20 +03:00
|
|
|
import Vere.Serf (Serf, sStderr, SerfState(..), doJob)
|
2019-08-28 23:17:01 +03:00
|
|
|
import Vere.Term
|
2019-06-19 01:38:24 +03:00
|
|
|
|
2019-07-16 03:01:45 +03:00
|
|
|
import qualified System.Entropy as Ent
|
2019-07-21 22:56:18 +03:00
|
|
|
import qualified Urbit.Time as Time
|
2019-07-16 03:01:45 +03:00
|
|
|
import qualified Vere.Log as Log
|
|
|
|
import qualified Vere.Serf as Serf
|
2019-05-30 23:19:26 +03:00
|
|
|
|
2019-06-18 02:47:20 +03:00
|
|
|
|
2019-06-25 04:10:41 +03:00
|
|
|
--------------------------------------------------------------------------------
|
2019-06-18 02:47:20 +03:00
|
|
|
|
2019-07-21 22:56:18 +03:00
|
|
|
_ioDrivers = [] :: [IODriver]
|
2019-06-18 02:47:20 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
setupPierDirectory :: FilePath -> RIO e ()
|
2019-08-15 05:42:48 +03:00
|
|
|
setupPierDirectory shipPath = do
|
2019-07-21 04:29:39 +03:00
|
|
|
for_ ["put", "get", "log", "chk"] $ \seg -> do
|
2019-08-28 14:45:49 +03:00
|
|
|
let pax = shipPath <> "/.urb/" <> seg
|
|
|
|
io $ createDirectoryIfMissing True pax
|
|
|
|
io $ setFileMode pax ownerModes
|
2019-07-21 04:29:39 +03:00
|
|
|
|
2019-06-18 02:47:20 +03:00
|
|
|
|
2019-07-21 22:56:18 +03:00
|
|
|
-- Load pill into boot sequence. -----------------------------------------------
|
2019-07-16 03:01:45 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
genEntropy :: RIO e Word512
|
|
|
|
genEntropy = fromIntegral . view (from atomBytes) <$> io (Ent.getEntropy 64)
|
2019-07-16 03:01:45 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
generateBootSeq :: Ship -> Pill -> RIO e BootSeq
|
2019-07-16 03:01:45 +03:00
|
|
|
generateBootSeq ship Pill{..} = do
|
|
|
|
ent <- genEntropy
|
|
|
|
let ovums = preKern ent <> pKernelOvums <> pUserspaceOvums
|
|
|
|
pure $ BootSeq ident pBootFormulas ovums
|
|
|
|
where
|
|
|
|
ident = LogIdentity ship True (fromIntegral $ length pBootFormulas)
|
2019-08-30 04:29:22 +03:00
|
|
|
preKern ent = [ EvBlip $ BlipEvArvo $ ArvoEvWhom () ship
|
2019-08-14 03:52:59 +03:00
|
|
|
, EvBlip $ BlipEvArvo $ ArvoEvWack () ent
|
2019-08-30 04:29:22 +03:00
|
|
|
, EvBlip $ BlipEvTerm $ TermEvBoot (1,()) (Fake (who ident))
|
2019-07-16 03:01:45 +03:00
|
|
|
]
|
|
|
|
|
|
|
|
|
2019-07-21 22:56:18 +03:00
|
|
|
-- Write a batch of jobs into the event log ------------------------------------
|
2019-07-19 03:52:53 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
writeJobs :: EventLog -> Vector Job -> RIO e ()
|
2019-07-19 03:52:53 +03:00
|
|
|
writeJobs log !jobs = do
|
|
|
|
expect <- Log.nextEv log
|
|
|
|
events <- fmap fromList $ traverse fromJob (zip [expect..] $ toList jobs)
|
|
|
|
Log.appendEvents log events
|
2019-07-17 08:32:36 +03:00
|
|
|
where
|
2019-08-28 14:45:49 +03:00
|
|
|
fromJob :: (EventId, Job) -> RIO e ByteString
|
2019-07-21 04:29:39 +03:00
|
|
|
fromJob (expectedId, job) = do
|
2019-08-28 14:45:49 +03:00
|
|
|
unless (expectedId == jobId job) $
|
|
|
|
error $ show ("bad job id!", expectedId, jobId job)
|
2019-07-21 22:56:18 +03:00
|
|
|
pure $ jamBS $ jobPayload job
|
2019-07-21 04:29:39 +03:00
|
|
|
|
|
|
|
jobPayload :: Job -> Noun
|
|
|
|
jobPayload (RunNok (LifeCyc _ m n)) = toNoun (m, n)
|
|
|
|
jobPayload (DoWork (Work _ m d o)) = toNoun (m, d, o)
|
2019-07-19 03:52:53 +03:00
|
|
|
|
2019-06-29 04:46:33 +03:00
|
|
|
|
2019-07-21 22:56:18 +03:00
|
|
|
-- Boot a new ship. ------------------------------------------------------------
|
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
booted :: HasLogFunc e
|
|
|
|
=> FilePath -> FilePath -> Serf.Flags -> Ship
|
2019-09-04 01:17:20 +03:00
|
|
|
-> RAcquire e (Serf e, EventLog, SerfState)
|
2019-08-14 03:52:59 +03:00
|
|
|
booted pillPath pierPath flags ship = do
|
2019-08-28 14:45:49 +03:00
|
|
|
rio $ logTrace "LOADING PILL"
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
pill <- io (loadFile pillPath >>= either throwIO pure)
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
rio $ logTrace "PILL LOADED"
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
seq@(BootSeq ident x y) <- rio $ generateBootSeq ship pill
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
rio $ logTrace "BootSeq Computed"
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
liftRIO (setupPierDirectory pierPath)
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
rio $ logTrace "Directory Setup"
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-14 03:52:59 +03:00
|
|
|
log <- Log.new (pierPath <> "/.urb/log") ident
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
rio $ logTrace "Event Log Initialized"
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-14 03:52:59 +03:00
|
|
|
serf <- Serf.run (Serf.Config pierPath flags)
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
rio $ logTrace "Serf Started"
|
2019-08-15 05:42:48 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
rio $ do
|
2019-07-21 22:56:18 +03:00
|
|
|
(events, serfSt) <- Serf.bootFromSeq serf seq
|
2019-08-28 14:45:49 +03:00
|
|
|
logTrace "Boot Sequence completed"
|
2019-07-21 22:56:18 +03:00
|
|
|
Serf.snapshot serf serfSt
|
2019-08-28 14:45:49 +03:00
|
|
|
logTrace "Snapshot taken"
|
2019-07-21 22:56:18 +03:00
|
|
|
writeJobs log (fromList events)
|
2019-08-28 14:45:49 +03:00
|
|
|
logTrace "Events written"
|
2019-07-21 22:56:18 +03:00
|
|
|
pure (serf, log, serfSt)
|
|
|
|
|
|
|
|
|
|
|
|
-- Resume an existing ship. ----------------------------------------------------
|
|
|
|
|
2019-08-28 15:22:56 +03:00
|
|
|
resumed :: HasLogFunc e
|
|
|
|
=> FilePath -> Serf.Flags
|
2019-09-04 01:17:20 +03:00
|
|
|
-> RAcquire e (Serf e, EventLog, SerfState)
|
2019-07-22 00:24:07 +03:00
|
|
|
resumed top flags = do
|
2019-07-21 22:56:18 +03:00
|
|
|
log <- Log.existing (top <> "/.urb/log")
|
2019-07-22 00:24:07 +03:00
|
|
|
serf <- Serf.run (Serf.Config top flags)
|
2019-08-28 15:22:56 +03:00
|
|
|
serfSt <- rio $ Serf.replay serf log
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-08-28 15:22:56 +03:00
|
|
|
rio $ Serf.snapshot serf serfSt
|
2019-07-21 22:56:18 +03:00
|
|
|
|
|
|
|
pure (serf, log, serfSt)
|
|
|
|
|
|
|
|
|
2019-06-29 04:46:33 +03:00
|
|
|
-- Run Pier --------------------------------------------------------------------
|
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
pier :: ∀e. HasLogFunc e
|
|
|
|
=> FilePath
|
2019-08-08 01:24:02 +03:00
|
|
|
-> Maybe Port
|
2019-09-04 01:17:20 +03:00
|
|
|
-> (Serf e, EventLog, SerfState)
|
2019-08-28 14:45:49 +03:00
|
|
|
-> RAcquire e ()
|
2019-08-08 01:24:02 +03:00
|
|
|
pier pierPath mPort (serf, log, ss) = do
|
2019-09-06 22:59:56 +03:00
|
|
|
computeQ <- newTQueueIO :: RAcquire e (TQueue ComputeRequest)
|
2019-08-28 14:45:49 +03:00
|
|
|
persistQ <- newTQueueIO :: RAcquire e (TQueue (Job, FX))
|
|
|
|
executeQ <- newTQueueIO :: RAcquire e (TQueue FX)
|
2019-06-01 01:55:21 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
inst <- io (KingId . UV . fromIntegral <$> randomIO @Word16)
|
2019-08-01 08:48:08 +03:00
|
|
|
|
2019-09-03 21:02:54 +03:00
|
|
|
terminalSystem <- initializeLocalTerminal
|
2019-09-04 01:17:20 +03:00
|
|
|
serf <- pure serf { sStderr = (tsStderr terminalSystem) }
|
|
|
|
|
2019-08-01 08:48:08 +03:00
|
|
|
let ship = who (Log.identity log)
|
2019-06-19 01:38:24 +03:00
|
|
|
|
2019-09-06 22:59:56 +03:00
|
|
|
let shutdownEvent = writeTQueue computeQ $ CRShutdown
|
|
|
|
let computeEvent = (\ev -> writeTQueue computeQ $ CREvent ev)
|
|
|
|
|
2019-08-01 05:34:14 +03:00
|
|
|
let (bootEvents, startDrivers) =
|
2019-09-06 22:59:56 +03:00
|
|
|
drivers pierPath inst ship mPort computeEvent shutdownEvent terminalSystem
|
2019-06-19 01:38:24 +03:00
|
|
|
|
2019-09-06 22:59:56 +03:00
|
|
|
io $ atomically $ for_ bootEvents computeEvent
|
2019-06-19 01:38:24 +03:00
|
|
|
|
2019-08-01 08:16:02 +03:00
|
|
|
tExe <- startDrivers >>= router (readTQueue executeQ)
|
2019-08-01 05:34:14 +03:00
|
|
|
tDisk <- runPersist log persistQ (writeTQueue executeQ)
|
|
|
|
tCpu <- runCompute serf ss (readTQueue computeQ) (writeTQueue persistQ)
|
2019-07-20 06:00:23 +03:00
|
|
|
|
2019-09-06 22:59:56 +03:00
|
|
|
tSaveSignal <- saveSignalThread computeQ
|
|
|
|
|
2019-08-01 08:16:02 +03:00
|
|
|
-- Wait for something to die.
|
2019-07-20 06:00:23 +03:00
|
|
|
|
2019-08-01 08:16:02 +03:00
|
|
|
let ded = asum [ death "effect thread" tExe
|
|
|
|
, death "persist thread" tDisk
|
|
|
|
, death "compute thread" tCpu
|
|
|
|
]
|
|
|
|
|
|
|
|
atomically ded >>= \case
|
2019-08-28 14:45:49 +03:00
|
|
|
Left (txt, exn) -> logError $ displayShow ("Somthing died", txt, exn)
|
|
|
|
Right tag -> logError $ displayShow ("something simply exited", tag)
|
2019-08-01 08:16:02 +03:00
|
|
|
|
|
|
|
death :: Text -> Async () -> STM (Either (Text, SomeException) Text)
|
|
|
|
death tag tid = do
|
|
|
|
waitCatchSTM tid <&> \case
|
|
|
|
Left exn -> Left (tag, exn)
|
|
|
|
Right () -> Right tag
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-09-06 22:59:56 +03:00
|
|
|
saveSignalThread :: TQueue ComputeRequest -> RAcquire e (Async ())
|
|
|
|
saveSignalThread tq = mkRAcquire start cancel
|
|
|
|
where
|
|
|
|
start = async $ forever $ do
|
|
|
|
threadDelay (120 * 1000000) -- 120 seconds
|
|
|
|
atomically $ writeTQueue tq $ CRSave
|
|
|
|
|
2019-08-01 05:34:14 +03:00
|
|
|
-- Start All Drivers -----------------------------------------------------------
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-08-29 03:26:59 +03:00
|
|
|
data Drivers e = Drivers
|
|
|
|
{ dAmes :: EffCb e AmesEf
|
|
|
|
, dBehn :: EffCb e BehnEf
|
|
|
|
, dHttpClient :: EffCb e HttpClientEf
|
|
|
|
, dHttpServer :: EffCb e HttpServerEf
|
|
|
|
, dNewt :: EffCb e NewtEf
|
|
|
|
, dSync :: EffCb e SyncEf
|
|
|
|
, dTerm :: EffCb e TermEf
|
2019-08-01 05:34:14 +03:00
|
|
|
}
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-08-29 03:26:59 +03:00
|
|
|
drivers :: HasLogFunc e
|
2019-09-06 22:59:56 +03:00
|
|
|
=> FilePath -> KingId -> Ship -> Maybe Port -> (Ev -> STM ()) -> STM()
|
2019-09-04 01:17:20 +03:00
|
|
|
-> TerminalSystem e
|
2019-08-29 03:26:59 +03:00
|
|
|
-> ([Ev], RAcquire e (Drivers e))
|
2019-09-06 22:59:56 +03:00
|
|
|
drivers pierPath inst who mPort plan shutdownSTM termSys =
|
2019-08-01 05:34:14 +03:00
|
|
|
(initialEvents, runDrivers)
|
|
|
|
where
|
|
|
|
(behnBorn, runBehn) = behn inst plan
|
|
|
|
(amesBorn, runAmes) = ames inst who mPort plan
|
2019-08-08 01:24:02 +03:00
|
|
|
(httpBorn, runHttp) = serv pierPath inst plan
|
2019-09-05 23:09:45 +03:00
|
|
|
(irisBorn, runIris) = client inst plan
|
2019-09-06 22:59:56 +03:00
|
|
|
(termBorn, runTerm) = term termSys shutdownSTM pierPath inst plan
|
2019-09-05 23:09:45 +03:00
|
|
|
initialEvents = mconcat [behnBorn, amesBorn, httpBorn, termBorn, irisBorn]
|
2019-08-01 05:34:14 +03:00
|
|
|
runDrivers = do
|
2019-08-29 03:26:59 +03:00
|
|
|
dNewt <- liftAcquire $ runAmes
|
|
|
|
dBehn <- liftAcquire $ runBehn
|
2019-08-01 08:16:02 +03:00
|
|
|
dAmes <- pure $ const $ pure ()
|
2019-09-05 23:09:45 +03:00
|
|
|
dHttpClient <- runIris
|
2019-08-08 01:24:02 +03:00
|
|
|
dHttpServer <- runHttp
|
2019-08-01 08:16:02 +03:00
|
|
|
dSync <- pure $ const $ pure ()
|
2019-09-03 21:02:54 +03:00
|
|
|
dTerm <- runTerm
|
2019-08-01 05:34:14 +03:00
|
|
|
pure (Drivers{..})
|
|
|
|
|
|
|
|
|
|
|
|
-- Route Effects to Drivers ----------------------------------------------------
|
|
|
|
|
2019-08-29 03:26:59 +03:00
|
|
|
router :: HasLogFunc e => STM FX -> Drivers e -> RAcquire e (Async ())
|
2019-08-28 14:45:49 +03:00
|
|
|
router waitFx Drivers{..} =
|
|
|
|
mkRAcquire start cancel
|
2019-08-01 05:34:14 +03:00
|
|
|
where
|
|
|
|
start = async $ forever $ do
|
|
|
|
fx <- atomically waitFx
|
2019-08-01 08:16:02 +03:00
|
|
|
for_ fx $ \ef -> do
|
2019-08-28 14:45:49 +03:00
|
|
|
logEffect ef
|
2019-08-01 08:16:02 +03:00
|
|
|
case ef of
|
|
|
|
GoodParse (EfVega _ _) -> error "TODO"
|
|
|
|
GoodParse (EfExit _ _) -> error "TODO"
|
|
|
|
GoodParse (EfVane (VEAmes ef)) -> dAmes ef
|
|
|
|
GoodParse (EfVane (VEBehn ef)) -> dBehn ef
|
|
|
|
GoodParse (EfVane (VEBoat ef)) -> dSync ef
|
|
|
|
GoodParse (EfVane (VEClay ef)) -> dSync ef
|
|
|
|
GoodParse (EfVane (VEHttpClient ef)) -> dHttpClient ef
|
|
|
|
GoodParse (EfVane (VEHttpServer ef)) -> dHttpServer ef
|
|
|
|
GoodParse (EfVane (VENewt ef)) -> dNewt ef
|
|
|
|
GoodParse (EfVane (VESync ef)) -> dSync ef
|
|
|
|
GoodParse (EfVane (VETerm ef)) -> dTerm ef
|
2019-08-28 14:45:49 +03:00
|
|
|
FailParse n -> logError
|
|
|
|
$ display
|
|
|
|
$ pack @Text (ppShow n)
|
2019-07-21 22:56:18 +03:00
|
|
|
|
2019-08-01 05:34:14 +03:00
|
|
|
|
|
|
|
-- Compute Thread --------------------------------------------------------------
|
|
|
|
|
2019-09-06 22:59:56 +03:00
|
|
|
data ComputeRequest
|
|
|
|
= CREvent Ev
|
|
|
|
| CRSave
|
|
|
|
| CRShutdown
|
|
|
|
deriving (Eq, Show)
|
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
logEvent :: HasLogFunc e => Ev -> RIO e ()
|
|
|
|
logEvent ev =
|
|
|
|
logDebug $ display $ "[EVENT]\n" <> pretty
|
|
|
|
where
|
|
|
|
pretty :: Text
|
|
|
|
pretty = pack $ unlines $ fmap ("\t" <>) $ lines $ ppShow ev
|
|
|
|
|
|
|
|
logEffect :: HasLogFunc e => Lenient Ef -> RIO e ()
|
|
|
|
logEffect ef =
|
|
|
|
logDebug $ display $ "[EFFECT]\n" <> pretty ef
|
|
|
|
where
|
|
|
|
pretty :: Lenient Ef -> Text
|
|
|
|
pretty = \case
|
|
|
|
GoodParse e -> pack $ unlines $ fmap ("\t" <>) $ lines $ ppShow e
|
|
|
|
FailParse n -> pack $ unlines $ fmap ("\t" <>) $ lines $ ppShow n
|
|
|
|
|
|
|
|
runCompute :: ∀e. HasLogFunc e
|
2019-09-06 22:59:56 +03:00
|
|
|
=> Serf e -> SerfState -> STM ComputeRequest -> ((Job, FX) -> STM ())
|
2019-08-28 14:45:49 +03:00
|
|
|
-> RAcquire e (Async ())
|
2019-08-01 05:34:14 +03:00
|
|
|
runCompute serf ss getEvent putResult =
|
2019-08-28 14:45:49 +03:00
|
|
|
mkRAcquire (async (go ss)) cancel
|
2019-08-01 05:34:14 +03:00
|
|
|
where
|
2019-08-28 14:45:49 +03:00
|
|
|
go :: SerfState -> RIO e ()
|
2019-08-01 05:34:14 +03:00
|
|
|
go ss = do
|
2019-09-06 22:59:56 +03:00
|
|
|
cr <- atomically getEvent
|
|
|
|
case cr of
|
|
|
|
CREvent ev -> do
|
|
|
|
logEvent ev
|
|
|
|
wen <- io Time.now
|
|
|
|
eId <- pure (ssNextEv ss)
|
|
|
|
mug <- pure (ssLastMug ss)
|
|
|
|
|
|
|
|
(job', ss', fx) <- doJob serf $ DoWork $ Work eId mug wen ev
|
|
|
|
atomically (putResult (job', fx))
|
|
|
|
go ss'
|
|
|
|
CRSave -> do
|
|
|
|
logDebug $ "Taking periodic snapshot"
|
|
|
|
Serf.snapshot serf ss
|
|
|
|
go ss
|
|
|
|
CRShutdown -> do
|
|
|
|
-- When shutting down, we first request a snapshot, and then we
|
|
|
|
-- just exit this recursive processing, which will cause the serf
|
|
|
|
-- to exit from its RAcquire.
|
|
|
|
logDebug $ "Shutting down compute system..."
|
|
|
|
Serf.snapshot serf ss
|
|
|
|
pure ()
|
2019-07-21 22:56:18 +03:00
|
|
|
|
|
|
|
|
2019-07-20 06:00:23 +03:00
|
|
|
-- Persist Thread --------------------------------------------------------------
|
|
|
|
|
|
|
|
data PersistExn = BadEventId EventId EventId
|
|
|
|
deriving Show
|
|
|
|
|
|
|
|
instance Exception PersistExn where
|
|
|
|
displayException (BadEventId expected got) =
|
|
|
|
unlines [ "Out-of-order event id send to persist thread."
|
|
|
|
, "\tExpected " <> show expected <> " but got " <> show got
|
|
|
|
]
|
|
|
|
|
|
|
|
runPersist :: EventLog
|
2019-08-01 05:34:14 +03:00
|
|
|
-> TQueue (Job, FX)
|
|
|
|
-> (FX -> STM ())
|
2019-08-28 14:45:49 +03:00
|
|
|
-> RAcquire e (Async ())
|
2019-08-01 05:34:14 +03:00
|
|
|
runPersist log inpQ out =
|
2019-08-28 14:45:49 +03:00
|
|
|
mkRAcquire runThread cancelWait
|
2019-07-20 06:00:23 +03:00
|
|
|
where
|
2019-08-28 14:45:49 +03:00
|
|
|
cancelWait :: Async () -> RIO e ()
|
2019-07-20 06:00:23 +03:00
|
|
|
cancelWait tid = cancel tid >> wait tid
|
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
runThread :: RIO e (Async ())
|
2019-07-20 06:00:23 +03:00
|
|
|
runThread = asyncBound $ forever $ do
|
2019-07-21 22:56:18 +03:00
|
|
|
writs <- atomically getBatchFromQueue
|
2019-08-01 05:34:14 +03:00
|
|
|
events <- validateJobsAndGetBytes (toNullable writs)
|
2019-07-20 06:00:23 +03:00
|
|
|
Log.appendEvents log events
|
2019-08-01 05:34:14 +03:00
|
|
|
atomically $ for_ writs $ \(_,fx) -> out fx
|
2019-07-20 06:00:23 +03:00
|
|
|
|
2019-08-28 14:45:49 +03:00
|
|
|
validateJobsAndGetBytes :: [(Job, FX)] -> RIO e (Vector ByteString)
|
2019-08-01 05:34:14 +03:00
|
|
|
validateJobsAndGetBytes writs = do
|
2019-07-20 06:00:23 +03:00
|
|
|
expect <- Log.nextEv log
|
|
|
|
fmap fromList
|
|
|
|
$ for (zip [expect..] writs)
|
2019-08-01 05:34:14 +03:00
|
|
|
$ \(expectedId, (j, fx)) -> do
|
|
|
|
unless (expectedId == jobId j) $
|
|
|
|
throwIO (BadEventId expectedId (jobId j))
|
|
|
|
case j of
|
|
|
|
RunNok _ ->
|
|
|
|
error "This shouldn't happen here!"
|
|
|
|
DoWork (Work eId mug wen ev) ->
|
|
|
|
pure $ jamBS $ toNoun (mug, wen, ev)
|
|
|
|
|
|
|
|
getBatchFromQueue :: STM (NonNull [(Job, FX)])
|
2019-07-20 06:00:23 +03:00
|
|
|
getBatchFromQueue =
|
|
|
|
readTQueue inpQ >>= go . singleton
|
|
|
|
where
|
|
|
|
go acc =
|
|
|
|
tryReadTQueue inpQ >>= \case
|
|
|
|
Nothing -> pure (reverse acc)
|
|
|
|
Just item -> go (item <| acc)
|