diff --git a/bindings/haskell/event-sorcery.cabal b/bindings/haskell/event-sorcery.cabal index 3db2b1c..2db0973 100644 --- a/bindings/haskell/event-sorcery.cabal +++ b/bindings/haskell/event-sorcery.cabal @@ -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 @@ -105,6 +106,7 @@ test-suite spec JobExecutionSpec JobSpec JobWorkerSpec + SchemaSpec StoreSpec WireSpec diff --git a/bindings/haskell/src/Event/Sorcery/Aggregate.hs b/bindings/haskell/src/Event/Sorcery/Aggregate.hs index ff3952a..eff0061 100644 --- a/bindings/haskell/src/Event/Sorcery/Aggregate.hs +++ b/bindings/haskell/src/Event/Sorcery/Aggregate.hs @@ -1,6 +1,7 @@ {-# LANGUAGE FunctionalDependencies #-} module Event.Sorcery.Aggregate ( + CompactionPolicy (..), DecodeCause (..), DispatchIntent, Dispatches (..), @@ -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) @@ -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 diff --git a/bindings/haskell/src/Event/Sorcery/Engine/Internal/FFI.hs b/bindings/haskell/src/Event/Sorcery/Engine/Internal/FFI.hs index d03bdde..e314c39 100644 --- a/bindings/haskell/src/Event/Sorcery/Engine/Internal/FFI.hs +++ b/bindings/haskell/src/Event/Sorcery/Engine/Internal/FFI.hs @@ -17,6 +17,8 @@ module Event.Sorcery.Engine.Internal.FFI ( esJobRetry, esLoadStream, esOpen, + esSchemaReconcile, + esSchemaRecord, esSnapshotDiscard, esSnapshotLoad, esSnapshotStore, @@ -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 diff --git a/bindings/haskell/src/Event/Sorcery/Schema.hs b/bindings/haskell/src/Event/Sorcery/Schema.hs new file mode 100644 index 0000000..19c8afd --- /dev/null +++ b/bindings/haskell/src/Event/Sorcery/Schema.hs @@ -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)) diff --git a/bindings/haskell/test/SchemaSpec.hs b/bindings/haskell/test/SchemaSpec.hs new file mode 100644 index 0000000..abd73a8 --- /dev/null +++ b/bindings/haskell/test/SchemaSpec.hs @@ -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 ()