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
2 changes: 2 additions & 0 deletions bindings/haskell/event-sorcery.cabal
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ library
Event.Sorcery.Job
Event.Sorcery.Job.Execution
Event.Sorcery.Job.Worker
Event.Sorcery.Schema
Event.Sorcery.Snapshot
Event.Sorcery.Store
Event.Sorcery.Stream
Expand Down Expand Up @@ -105,6 +106,7 @@ test-suite spec
JobExecutionSpec
JobSpec
JobWorkerSpec
SchemaSpec
StoreSpec
WireSpec

Expand Down
9 changes: 9 additions & 0 deletions bindings/haskell/src/Event/Sorcery/Aggregate.hs
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
{-# LANGUAGE FunctionalDependencies #-}

module Event.Sorcery.Aggregate (
CompactionPolicy (..),
DecodeCause (..),
DispatchIntent,
Dispatches (..),
Expand Down Expand Up @@ -37,6 +38,12 @@ newtype SchemaVersion = SchemaVersion Word16
deriving stock (Eq, Ord, Show)


data CompactionPolicy
= Retain
| CompactAfterSnapshot
deriving stock (Eq, Show)


newtype DecodeCause = DecodeCause Text
deriving stock (Eq, Show)

Expand Down Expand Up @@ -81,6 +88,8 @@ class EventSourced entity where
eventType :: Event entity -> Text
eventVersion :: Event entity -> EventVersion
schemaVersion :: Proxy entity -> SchemaVersion
compactionPolicy :: Proxy entity -> CompactionPolicy
compactionPolicy _ = Retain
encodeEvent :: Event entity -> ByteString
decodeEvent :: ByteString -> Either DecodeCause (Event entity)
encodeSnapshot :: entity -> ByteString
Expand Down
11 changes: 11 additions & 0 deletions bindings/haskell/src/Event/Sorcery/Engine/Internal/FFI.hs
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ module Event.Sorcery.Engine.Internal.FFI (
esJobRetry,
esLoadStream,
esOpen,
esSchemaReconcile,
esSchemaRecord,
esSnapshotDiscard,
esSnapshotLoad,
esSnapshotStore,
Expand Down Expand Up @@ -96,6 +98,15 @@ foreign import capi safe "event_sorcery.h es_commit_with_job"
esCommitWithJob :: Ptr (Ptr EsStore) -> Ptr EsBuf -> Ptr EsBuf -> IO CInt


foreign import capi safe "event_sorcery.h es_schema_reconcile"
esSchemaReconcile
:: Ptr (Ptr EsStore) -> Ptr EsBuf -> Ptr Word8 -> Ptr EsBuf -> IO CInt


foreign import capi safe "event_sorcery.h es_schema_record"
esSchemaRecord :: Ptr (Ptr EsStore) -> Ptr EsBuf -> Ptr EsBuf -> IO CInt


foreign import capi safe "event_sorcery.h es_snapshot_load"
esSnapshotLoad
:: Ptr (Ptr EsStore) -> Ptr EsBuf -> Ptr EsBuf -> Ptr EsBuf -> IO CInt
Expand Down
113 changes: 113 additions & 0 deletions bindings/haskell/src/Event/Sorcery/Schema.hs
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
module Event.Sorcery.Schema (
SchemaReconciliation (..),
reconcileSchema,
recordSchema,
) where

import Codec.CBOR.Encoding (
encodeListLen,
encodeString,
encodeWord,
encodeWord64,
)
import Codec.CBOR.Write (toStrictByteString)
import Data.ByteString (ByteString)
import Data.Proxy (Proxy)
import Data.Word (Word64, Word8)
import Event.Sorcery.Aggregate (
CompactionPolicy (CompactAfterSnapshot, Retain),
EventSourced (aggregateType, compactionPolicy, schemaVersion),
SchemaVersion (SchemaVersion),
)
import Event.Sorcery.Engine (
BindingFault (UnknownResultTag),
EngineError (BindingProtocolError),
Store,
)
import Event.Sorcery.Engine.Internal (
callWithoutOutput,
withInputBuffer,
withOpenStore,
)
import Event.Sorcery.Engine.Internal.FFI (
esSchemaReconcile,
esSchemaRecord,
)
import Foreign.Marshal.Alloc (alloca)
import Foreign.Storable (peek, poke)
import Prelude (
Either (Left, Right),
Eq,
IO,
Show,
Word,
fromIntegral,
pure,
($),
(.),
(<$>),
(<>),
)


data SchemaReconciliation
= Changed
| Unchanged
deriving stock (Eq, Show)


reconcileSchema
:: EventSourced entity
=> Store
-> Proxy entity
-> IO (Either EngineError SchemaReconciliation)
reconcileSchema store entity =
withOpenStore store $ \handle ->
withInputBuffer (encodeSchemaTarget entity) $ \request ->
alloca $ \outReconciliation -> do
poke outReconciliation 0
reconciled <-
callWithoutOutput
(esSchemaReconcile handle request outReconciliation)

case reconciled of
Left failure -> pure (Left failure)
Right () -> decodeReconciliation <$> peek outReconciliation


recordSchema
:: EventSourced entity
=> Store
-> Proxy entity
-> IO (Either EngineError ())
recordSchema store entity =
withOpenStore store $ \handle ->
withInputBuffer
(encodeSchemaTarget entity)
(callWithoutOutput . esSchemaRecord handle)


encodeSchemaTarget :: EventSourced entity => Proxy entity -> ByteString
encodeSchemaTarget entity =
toStrictByteString $
encodeListLen 4
<> encodeWord 1
<> encodeString (aggregateType entity)
<> encodeWord64 (schemaVersionWord64 (schemaVersion entity))
<> encodeWord (compactionTag (compactionPolicy entity))


schemaVersionWord64 :: SchemaVersion -> Word64
schemaVersionWord64 (SchemaVersion version) = fromIntegral version


compactionTag :: CompactionPolicy -> Word
compactionTag Retain = 0
compactionTag CompactAfterSnapshot = 1


decodeReconciliation :: Word8 -> Either EngineError SchemaReconciliation
decodeReconciliation 0 = Right Changed
decodeReconciliation 1 = Right Unchanged
decodeReconciliation tag =
Left (BindingProtocolError (UnknownResultTag tag))
121 changes: 121 additions & 0 deletions bindings/haskell/test/SchemaSpec.hs
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
module SchemaSpec (spec) where

import Data.Proxy (Proxy (Proxy))
import Event.Sorcery.Aggregate (
CompactionPolicy (CompactAfterSnapshot),
EventSourced (..),
SchemaVersion (SchemaVersion),
)
import Event.Sorcery.Engine (
EngineError (InvalidState),
OpenOptions (OpenOptions),
Store,
closeStore,
openStore,
)
import Event.Sorcery.Schema (
SchemaReconciliation (Changed, Unchanged),
reconcileSchema,
recordSchema,
)
import Test.Hspec (Spec, it, shouldReturn)
import Prelude (Either (Left, Right), IO, error, ($))


data Account


data AccountV2


data AccountId


data AccountV2Id


data AccountCommand


data AccountEvent


data AccountV2Event


data AccountCommandError


data AccountApplyError


instance EventSourced Account where
type EntityId Account = AccountId
type Command Account = AccountCommand
type Event Account = AccountEvent
type CommandError Account = AccountCommandError
type ApplyError Account = AccountApplyError
type Jobs Account = '[]


aggregateType _ = "account"
encodeEntityId _ = "account-id"
eventType _ = "event"
eventVersion _ = error "unused event version"
schemaVersion _ = SchemaVersion 1
encodeEvent _ = error "unused event encoder"
decodeEvent _ = error "unused event decoder"
encodeSnapshot _ = error "unused snapshot encoder"
decodeSnapshot _ = error "unused snapshot decoder"
originate _ = error "unused origin fold"
evolve _ _ = error "unused evolution fold"
initialize _ = error "unused initialization handler"
transition _ _ = error "unused transition handler"


instance EventSourced AccountV2 where
type EntityId AccountV2 = AccountV2Id
type Command AccountV2 = AccountCommand
type Event AccountV2 = AccountV2Event
type CommandError AccountV2 = AccountCommandError
type ApplyError AccountV2 = AccountApplyError
type Jobs AccountV2 = '[]


aggregateType _ = "account"
encodeEntityId _ = "account-id"
eventType _ = "event"
eventVersion _ = error "unused event version"
schemaVersion _ = SchemaVersion 2
compactionPolicy _ = CompactAfterSnapshot
encodeEvent _ = error "unused event encoder"
decodeEvent _ = error "unused event decoder"
encodeSnapshot _ = error "unused snapshot encoder"
decodeSnapshot _ = error "unused snapshot decoder"
originate _ = error "unused origin fold"
evolve _ _ = error "unused evolution fold"
initialize _ = error "unused initialization handler"
transition _ _ = error "unused transition handler"


spec :: Spec
spec = it "reconciles aggregate schemas through the shared engine" do
withStore $ \store -> do
reconcileSchema store (Proxy @Account) `shouldReturn` Right Changed
recordSchema store (Proxy @Account) `shouldReturn` Right ()
reconcileSchema store (Proxy @Account) `shouldReturn` Right Unchanged

reconcileSchema store (Proxy @AccountV2)
`shouldReturn` Left
(InvalidState "compacted snapshot retention cannot be cleared")


withStore :: (Store -> IO ()) -> IO ()
withStore action = do
opened <- openStore (OpenOptions "sqlite::memory:" 5000 1 1)

case opened of
Left _ -> error "failed to open the shared engine"
Right store -> do
action store
closeStore store `shouldReturn` Right ()
Loading