Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
269 changes: 269 additions & 0 deletions bindings/haskell/benchmark/Main.hs
Original file line number Diff line number Diff line change
@@ -0,0 +1,269 @@
module Main (main) where

import Control.DeepSeq (NFData (rnf))
import Criterion.Main (
Benchmark,
bench,
bgroup,
defaultMain,
envWithCleanup,
whnfIO,
)
import Data.ByteString qualified as ByteString
import Data.IORef (IORef, atomicModifyIORef', newIORef)
import Data.List (foldl')
import Data.List.NonEmpty (NonEmpty ((:|)))
import Data.Maybe (Maybe (..))
import Data.Text qualified as Text
import Data.Unrestricted.Linear (Ur)
import Data.Word (Word64)
import EventSorcery.Engine (
EngineError,
OpenOptions (OpenOptions),
Store,
closeStore,
openStore,
)
import EventSorcery.Job (
ClaimBudget (ClaimBudget),
ClaimedJob,
JobClaimDetails,
JobClaimResult (JobClaimed),
JobId,
JobInstant (JobInstant),
JobKind (JobKind),
JobSeed (JobSeed),
JobSettlement (SettlementApplied),
JobSettlementToken,
LeaseDuration (LeaseDuration),
WorkerId (WorkerId),
acknowledgeJob,
claimJob,
enqueueJob,
mkJobId,
settlementToken,
)
import EventSorcery.Stream (
AggregateId (AggregateId),
AggregateType (AggregateType),
EventType (EventType),
EventVersion (EventVersion),
ProposedEvent (ProposedEvent),
StoredEvent (StoredEvent),
StreamIdentity (StreamIdentity),
commit,
loadStream,
)
import Prelude (
Either (..),
IO,
Int,
String,
error,
fromIntegral,
pure,
replicate,
seq,
show,
($),
(+),
(<>),
)


data BenchmarkEnvironment = BenchmarkEnvironment
{ store :: Store
, nextStream :: IORef Int
, nextJob :: IORef Int
}


instance NFData BenchmarkEnvironment where
rnf (BenchmarkEnvironment store nextStream nextJob) =
store `seq` nextStream `seq` nextJob `seq` ()


main :: IO ()
main =
defaultMain
[ envWithCleanup
openBenchmarkEnvironment
closeBenchmarkEnvironment
benchmarks
]


benchmarks :: BenchmarkEnvironment -> Benchmark
benchmarks environment =
bgroup
"engine-backed Haskell binding"
[ bench
"load and force 1,000 stored events"
(whnfIO (loadAndForceEvents environment))
, bench
"commit one event to a fresh stream"
(whnfIO (commitFreshStream environment))
, bench
"enqueue, claim, and acknowledge one job"
(whnfIO (runFreshJob environment))
]


openBenchmarkEnvironment :: IO BenchmarkEnvironment
openBenchmarkEnvironment = do
opened <- openStore (OpenOptions "sqlite::memory:" 5_000 1 1)

store <- expectRight "failed to open benchmark store" opened
preloadReplayStream store
nextStream <- newIORef 0
nextJob <- newIORef 0

pure BenchmarkEnvironment {store, nextStream, nextJob}


closeBenchmarkEnvironment :: BenchmarkEnvironment -> IO ()
closeBenchmarkEnvironment environment = do
closed <- closeStore environment.store
expectRight "failed to close benchmark store" closed


preloadReplayStream :: Store -> IO ()
preloadReplayStream store = do
committed <-
commit
store
replayStream
0
(benchmarkEvent :| replicate 999 benchmarkEvent)
expectRight "failed to preload replay benchmark" committed


loadAndForceEvents :: BenchmarkEnvironment -> IO Word64
loadAndForceEvents environment = do
loaded <- loadStream environment.store replayStream Nothing
events <- expectRight "failed to load replay benchmark stream" loaded

pure (foldl' forceStoredEvent 0 events)


forceStoredEvent :: Word64 -> StoredEvent -> Word64
forceStoredEvent
total
( StoredEvent
sequence
(EventType eventTypeText)
(EventVersion eventVersionText)
payload
) =
total
+ sequence
+ fromIntegral (Text.length eventTypeText)
+ fromIntegral (Text.length eventVersionText)
+ fromIntegral (ByteString.length payload)


commitFreshStream :: BenchmarkEnvironment -> IO Int
commitFreshStream environment = do
identifier <- nextIdentifier environment.nextStream
let stream =
StreamIdentity
(AggregateType "benchmark-commit")
(AggregateId ("stream-" <> Text.pack (show identifier)))

committed <- commit environment.store stream 0 (benchmarkEvent :| [])
expectRight "failed to commit benchmark event" committed

pure identifier


runFreshJob :: BenchmarkEnvironment -> IO Int
runFreshJob environment = do
identifier <- nextIdentifier environment.nextJob
jobId <-
expectJust
"generated an invalid benchmark job identifier"
(mkJobId (Text.justifyRight 26 '0' (Text.pack (show identifier))))
enqueued <-
enqueueJob
environment.store
(JobSeed jobId benchmarkJobKind benchmarkPayload benchmarkNow)
expectRight "failed to enqueue benchmark job" enqueued

token <- claimBenchmarkJob environment.store jobId
settled <- acknowledgeJob environment.store token

case settled of
Right SettlementApplied -> pure identifier
_ -> error "failed to acknowledge benchmark job"


claimBenchmarkJob :: Store -> JobId -> IO JobSettlementToken
claimBenchmarkJob store jobId = do
claimed <-
claimJob
store
jobId
benchmarkWorker
benchmarkNow
(LeaseDuration 30_000)
(ClaimBudget 50)
releaseClaim

case claimed of
Right (JobClaimed token) -> pure token
_ -> error "failed to claim benchmark job"


releaseClaim
:: JobClaimDetails
-> ClaimedJob
%1 -> Ur JobSettlementToken
releaseClaim _ = settlementToken


nextIdentifier :: IORef Int -> IO Int
nextIdentifier counter =
atomicModifyIORef' counter $ \current ->
let next = current + 1
in (next, next)


expectRight :: String -> Either EngineError value -> IO value
expectRight _ (Right value) = pure value
expectRight message (Left _) = error message


expectJust :: String -> Maybe value -> IO value
expectJust _ (Just value) = pure value
expectJust message Nothing = error message


replayStream :: StreamIdentity
replayStream =
StreamIdentity
(AggregateType "benchmark-replay")
(AggregateId "fixed-stream")


benchmarkEvent :: ProposedEvent
benchmarkEvent =
ProposedEvent
(EventType "balance-adjusted")
(EventVersion "1")
(ByteString.replicate 256 42)


benchmarkJobKind :: JobKind
benchmarkJobKind = JobKind "benchmark-job"


benchmarkWorker :: WorkerId
benchmarkWorker = WorkerId "benchmark-worker"


benchmarkNow :: JobInstant
benchmarkNow = JobInstant 1_000


benchmarkPayload :: ByteString.ByteString
benchmarkPayload = ByteString.replicate 256 42
17 changes: 17 additions & 0 deletions bindings/haskell/event-sorcery.cabal
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ library
, linear-base >=0.5 && <0.6
, text >=2.0 && <2.2
, transformers >=0.6 && <0.7
, ulid >=0.3 && <0.4

test-suite wire-spec
import: strict
Expand Down Expand Up @@ -140,6 +141,7 @@ test-suite domain-spec
, base >=4.18 && <5
, bytestring >=0.11 && <0.13
, event-sorcery
, text >=2.0 && <2.2

test-suite store-spec
import: strict
Expand Down Expand Up @@ -196,3 +198,18 @@ test-suite dispatch-worker-spec
, bytestring >=0.11 && <0.13
, event-sorcery
, text >=2.0 && <2.2

benchmark event-sorcery-benchmarks
import: strict
type: exitcode-stdio-1.0
hs-source-dirs: benchmark
main-is: Main.hs
ghc-options: -O2 -threaded -rtsopts
build-depends:
, base >=4.18 && <5
, bytestring >=0.11 && <0.13
, criterion >=1.6 && <1.7
, deepseq >=1.4 && <1.6
, event-sorcery
, linear-base >=0.5 && <0.6
, text >=2.0 && <2.2
15 changes: 11 additions & 4 deletions bindings/haskell/src/EventSorcery/Job/Definition.hs
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,10 @@ import Data.Maybe (Maybe (..))
import Data.Proxy (Proxy (Proxy))
import Data.Text (Text)
import Data.Text qualified as Text
import Data.ULID (ulidFromInteger)
import Data.ULID.Base32 qualified as Base32
import GHC.TypeLits (KnownSymbol, Symbol, symbolVal)
import Prelude (Either, Eq, Ord, Show, otherwise, (==))
import Prelude (Either (..), Eq, Integer, Ord, Show)


newtype JobId = JobId Text
Expand Down Expand Up @@ -49,9 +51,14 @@ class KnownSymbol (JobType job) => Job job where


mkJobId :: Text -> Maybe JobId
mkJobId value
| value == "" = Nothing
| otherwise = Just (JobId value)
mkJobId value =
case Base32.decode 26 value :: [(Integer, Text)] of
[(decoded, remaining)]
| Text.null remaining ->
case ulidFromInteger decoded of
Left _ -> Nothing
Right _ -> Just (JobId value)
_ -> Nothing


jobIdText :: JobId -> Text
Expand Down
4 changes: 2 additions & 2 deletions bindings/haskell/test/DispatchSpec.hs
Original file line number Diff line number Diff line change
Expand Up @@ -69,8 +69,8 @@ instance Job ChargeCard where

main :: IO ()
main = do
let first = requireJobId "job-one"
second = requireJobId "job-two"
let first = requireJobId "01ARZ3NDEKTSV4RRFFQ69G5FAV"
second = requireJobId "01ARZ3NDEKTSV4RRFFQ69G5FAW"
dispatched = Dispatched first ChargeCard
confirmed = confirmedOutcome @ChargeCard first Receipt 2
failed =
Expand Down
13 changes: 9 additions & 4 deletions bindings/haskell/test/DomainSpec.hs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ module Main (main) where

import Data.ByteString qualified as ByteString
import Data.List.NonEmpty (NonEmpty ((:|)))
import Data.Text (Text)
import EventSorcery.Aggregate (
DecodeCause (DecodeCause),
Dispatches (..),
Expand Down Expand Up @@ -56,17 +57,17 @@ import Prelude (
main :: IO ()
main = do
identifier <-
case mkJobId "job-1" of
case mkJobId testJobIdText of
Just value -> pure value
Nothing -> error "non-empty job identifier was rejected"
Nothing -> error "valid job identifier was rejected"

case mkJobId "" of
Nothing -> pure ()
Just _ -> error "empty job identifier was accepted"

let intent = dispatchIntent identifier SendWelcomeEmail

if dispatchJobId intent == identifier && jobIdText identifier == "job-1"
if dispatchJobId intent == identifier && jobIdText identifier == testJobIdText
then pure ()
else error "dispatch intent did not preserve the job identity"

Expand Down Expand Up @@ -225,6 +226,10 @@ instance EventSourced Account where


testJobId :: JobId
testJobId = case mkJobId "job-1" of
testJobId = case mkJobId testJobIdText of
Just identifier -> identifier
Nothing -> error "valid test job identifier was rejected"


testJobIdText :: Text
testJobIdText = "01ARZ3NDEKTSV4RRFFQ69G5FAV"
Loading