diff --git a/config/bench.exs b/config/bench.exs index f5ce46da..980f571e 100644 --- a/config/bench.exs +++ b/config/bench.exs @@ -13,7 +13,8 @@ default_config = [ database: "eventstore_bench", hostname: "localhost", pool_size: 10, - serializer: EventStore.TermSerializer + serializer: EventStore.TermSerializer, + metadata_serializer: EventStore.TermSerializer ] config :eventstore, TestEventStore, default_config diff --git a/config/jsonb.exs b/config/jsonb.exs index 6b7bf9b3..0cd2afe4 100644 --- a/config/jsonb.exs +++ b/config/jsonb.exs @@ -9,12 +9,14 @@ config :ex_unit, default_config = [ column_data_type: "jsonb", + metadata_column_data_type: "jsonb", username: "postgres", password: "postgres", database: "eventstore_jsonb_test", hostname: "localhost", pool_size: 1, serializer: EventStore.JsonbSerializer, + metadata_serializer: EventStore.JsonbSerializer, subscription_retry_interval: 1_000, types: EventStore.PostgresTypes ] diff --git a/config/migration.exs b/config/migration.exs index c433d139..ef9bde02 100644 --- a/config/migration.exs +++ b/config/migration.exs @@ -14,6 +14,7 @@ default_config = [ hostname: "localhost", pool_size: 1, serializer: EventStore.JsonSerializer, + metadata_serializer: EventStore.JsonSerializer, subscription_retry_interval: 1_000 ] diff --git a/config/test.exs b/config/test.exs index c920f670..d9a059e0 100644 --- a/config/test.exs +++ b/config/test.exs @@ -15,6 +15,7 @@ default_config = [ hostname: "localhost", pool_size: 1, serializer: EventStore.JsonSerializer, + metadata_serializer: EventStore.JsonSerializer, subscription_retry_interval: 1_000 ] diff --git a/lib/event_store.ex b/lib/event_store.ex index a7b63b5f..a007d831 100644 --- a/lib/event_store.ex +++ b/lib/event_store.ex @@ -23,6 +23,7 @@ defmodule EventStore do config :my_app, MyApp.EventStore, serializer: EventStore.JsonSerializer, + metadata_serializer: EventStore.JsonSerializer, username: "postgres", password: "postgres", database: "eventstore", @@ -448,6 +449,7 @@ defmodule EventStore do conn = Keyword.fetch!(config, :conn) schema = Keyword.fetch!(config, :schema) serializer = Keyword.fetch!(config, :serializer) + metadata_serializer = Keyword.fetch!(config, :metadata_serializer) query_timeout = timeout(opts, config) @@ -467,6 +469,7 @@ defmodule EventStore do query_timeout: query_timeout, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid, subscription_name: subscription_name, start_from: start_from diff --git a/lib/event_store/config.ex b/lib/event_store/config.ex index 418323fa..a8fccc7a 100644 --- a/lib/event_store/config.ex +++ b/lib/event_store/config.ex @@ -67,6 +67,19 @@ defmodule EventStore.Config do end end + def metadata_column_data_type(event_store, config) do + case Keyword.get(config, :metadata_column_data_type, column_data_type(event_store, config)) do + valid when valid in ["bytea", "jsonb"] -> + valid + + invalid -> + raise ArgumentError, + inspect(event_store) <> + " `:column_data_type` expects either \"bytea\" or \"jsonb\" but got: " <> + inspect(invalid) + end + end + @postgrex_connection_opts [ :after_connect, :after_connect_timeout, diff --git a/lib/event_store/config/parser.ex b/lib/event_store/config/parser.ex index 70b16205..e874c570 100644 --- a/lib/event_store/config/parser.ex +++ b/lib/event_store/config/parser.ex @@ -3,6 +3,7 @@ defmodule EventStore.Config.Parser do @config_defaults [ column_data_type: "bytea", + metadata_column_data_type: "bytea", enable_hard_deletes: false, schema: "public", timeout: 15_000 diff --git a/lib/event_store/notifications/publisher.ex b/lib/event_store/notifications/publisher.ex index 4d9e65f1..b4d296e1 100644 --- a/lib/event_store/notifications/publisher.ex +++ b/lib/event_store/notifications/publisher.ex @@ -12,7 +12,15 @@ defmodule EventStore.Notifications.Publisher do alias EventStore.Notifications.Notification defmodule State do - defstruct [:conn, :event_store, :query_timeout, :schema, :serializer, :subscribe_to] + defstruct [ + :conn, + :event_store, + :query_timeout, + :schema, + :serializer, + :metadata_serializer, + :subscribe_to + ] def new(opts) do %State{ @@ -21,6 +29,7 @@ defmodule EventStore.Notifications.Publisher do query_timeout: Keyword.fetch!(opts, :query_timeout), schema: Keyword.fetch!(opts, :schema), serializer: Keyword.fetch!(opts, :serializer), + metadata_serializer: Keyword.fetch!(opts, :metadata_serializer), subscribe_to: Keyword.fetch!(opts, :subscribe_to) } end @@ -63,7 +72,13 @@ defmodule EventStore.Notifications.Publisher do to_stream_version: to_stream_version } = notification - %State{conn: conn, query_timeout: query_timeout, schema: schema, serializer: serializer} = + %State{ + conn: conn, + query_timeout: query_timeout, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer + } = state count = to_stream_version - from_stream_version + 1 @@ -74,7 +89,8 @@ defmodule EventStore.Notifications.Publisher do timeout: query_timeout ) do {:ok, events} -> - deserialized_events = deserialize_recorded_events(events, serializer) + deserialized_events = + deserialize_recorded_events(events, serializer, metadata_serializer) {stream_uuid, deserialized_events} @@ -92,8 +108,8 @@ defmodule EventStore.Notifications.Publisher do end end - defp deserialize_recorded_events(recorded_events, serializer) do - Enum.map(recorded_events, &RecordedEvent.deserialize(&1, serializer)) + defp deserialize_recorded_events(recorded_events, serializer, metadata_serializer) do + Enum.map(recorded_events, &RecordedEvent.deserialize(&1, serializer, metadata_serializer)) end defp broadcast(event_store, stream_uuid, events) do diff --git a/lib/event_store/notifications/supervisor.ex b/lib/event_store/notifications/supervisor.ex index b9b40e3b..f4474ff0 100644 --- a/lib/event_store/notifications/supervisor.ex +++ b/lib/event_store/notifications/supervisor.ex @@ -22,6 +22,7 @@ defmodule EventStore.Notifications.Supervisor do conn = Keyword.fetch!(config, :conn) schema = Keyword.fetch!(config, :schema) serializer = Keyword.fetch!(config, :serializer) + metadata_serializer = Keyword.fetch!(config, :metadata_serializer) query_timeout = Keyword.fetch!(config, :timeout) listener_name = Module.concat([event_store, Listener]) @@ -54,6 +55,7 @@ defmodule EventStore.Notifications.Supervisor do event_store: event_store, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, subscribe_to: listener_name, name: publisher_name, hibernate_after: hibernate_after} diff --git a/lib/event_store/recorded_event.ex b/lib/event_store/recorded_event.ex index f6292792..14856224 100644 --- a/lib/event_store/recorded_event.ex +++ b/lib/event_store/recorded_event.ex @@ -55,13 +55,13 @@ defmodule EventStore.RecordedEvent do :created_at ] - def deserialize(%RecordedEvent{} = recorded_event, serializer) do + def deserialize(%RecordedEvent{} = recorded_event, serializer, metadata_serializer) do %RecordedEvent{data: data, metadata: metadata, event_type: event_type} = recorded_event %RecordedEvent{ recorded_event | data: serializer.deserialize(data, type: event_type), - metadata: serializer.deserialize(metadata, []) + metadata: metadata_serializer.deserialize(metadata, []) } end diff --git a/lib/event_store/serializer.ex b/lib/event_store/serializer.ex index 3d27dea1..8f91bfd0 100644 --- a/lib/event_store/serializer.ex +++ b/lib/event_store/serializer.ex @@ -31,4 +31,19 @@ defmodule EventStore.Serializer do message: "#{inspect(event_store)} configuration expects :serializer to be configured" end end + + @doc """ + Get the metadata serializer module from the given config for the event store. + """ + def metadata_serializer(event_store, config) do + case Keyword.fetch(config, :metadata_serializer) do + {:ok, serializer} -> + serializer + + :error -> + raise ArgumentError, + message: + "#{inspect(event_store)} configuration expects :metadata_serializer to be configured" + end + end end diff --git a/lib/event_store/snapshots/snapshot_data.ex b/lib/event_store/snapshots/snapshot_data.ex index c8bbaeb3..d76f2cd3 100644 --- a/lib/event_store/snapshots/snapshot_data.ex +++ b/lib/event_store/snapshots/snapshot_data.ex @@ -16,23 +16,23 @@ defmodule EventStore.Snapshots.SnapshotData do created_at: DateTime.t() } - def serialize(%SnapshotData{} = snapshot, serializer) do + def serialize(%SnapshotData{} = snapshot, serializer, metadata_serializer) do %SnapshotData{data: data, metadata: metadata} = snapshot %SnapshotData{ snapshot | data: serializer.serialize(data), - metadata: serializer.serialize(metadata) + metadata: metadata_serializer.serialize(metadata) } end - def deserialize(%SnapshotData{} = snapshot, serializer) do + def deserialize(%SnapshotData{} = snapshot, serializer, metadata_serializer) do %SnapshotData{source_type: source_type, data: data, metadata: metadata} = snapshot %SnapshotData{ snapshot | data: serializer.deserialize(data, type: source_type), - metadata: serializer.deserialize(metadata, []) + metadata: metadata_serializer.deserialize(metadata, []) } end end diff --git a/lib/event_store/snapshots/snapshotter.ex b/lib/event_store/snapshots/snapshotter.ex index 0ada96b2..9c3c652e 100644 --- a/lib/event_store/snapshots/snapshotter.ex +++ b/lib/event_store/snapshots/snapshotter.ex @@ -9,9 +9,10 @@ defmodule EventStore.Snapshots.Snapshotter do """ def read_snapshot(conn, source_uuid, opts) do {serializer, opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, opts} = Keyword.pop(opts, :metadata_serializer) with {:ok, snapshot} <- Snapshot.read_snapshot(conn, source_uuid, opts) do - deserialized = SnapshotData.deserialize(snapshot, serializer) + deserialized = SnapshotData.deserialize(snapshot, serializer, metadata_serializer) {:ok, deserialized} end @@ -24,8 +25,9 @@ defmodule EventStore.Snapshots.Snapshotter do """ def record_snapshot(conn, %SnapshotData{} = snapshot, opts) do {serializer, opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, opts} = Keyword.pop(opts, :metadata_serializer) - serialized = SnapshotData.serialize(snapshot, serializer) + serialized = SnapshotData.serialize(snapshot, serializer, metadata_serializer) Snapshot.record_snapshot(conn, serialized, opts) end diff --git a/lib/event_store/sql/init.ex b/lib/event_store/sql/init.ex index ea529782..b39563e1 100644 --- a/lib/event_store/sql/init.ex +++ b/lib/event_store/sql/init.ex @@ -5,13 +5,14 @@ defmodule EventStore.Sql.Init do def statements(config) do column_data_type = Keyword.fetch!(config, :column_data_type) + metadata_column_data_type = Keyword.fetch!(config, :metadata_column_data_type) schema = Keyword.fetch!(config, :schema) [ ~s(SET LOCAL search_path TO "#{schema}";), create_streams_table(), create_stream_uuid_index(), - create_events_table(column_data_type), + create_events_table(column_data_type, metadata_column_data_type), create_stream_events_table(), create_stream_events_index(), create_event_store_exception_function(), @@ -26,7 +27,7 @@ defmodule EventStore.Sql.Init do create_event_notification_trigger(), create_subscriptions_table(), create_subscription_index(), - create_snapshots_table(column_data_type), + create_snapshots_table(column_data_type, metadata_column_data_type), create_schema_migrations_table(), record_event_store_schema_version() ] @@ -58,7 +59,7 @@ defmodule EventStore.Sql.Init do """ end - defp create_events_table(column_data_type) do + defp create_events_table(column_data_type, metadata_column_data_type) do """ CREATE TABLE events ( @@ -67,7 +68,7 @@ defmodule EventStore.Sql.Init do causation_id uuid NULL, correlation_id uuid NULL, data #{column_data_type} NOT NULL, - metadata #{column_data_type} NULL, + metadata #{metadata_column_data_type} NULL, created_at timestamp with time zone DEFAULT NOW() NOT NULL ); """ @@ -242,7 +243,7 @@ defmodule EventStore.Sql.Init do """ end - defp create_snapshots_table(column_data_type) do + defp create_snapshots_table(column_data_type, metadata_column_data_type) do """ CREATE TABLE snapshots ( @@ -250,7 +251,7 @@ defmodule EventStore.Sql.Init do source_version bigint NOT NULL, source_type text NOT NULL, data #{column_data_type} NOT NULL, - metadata #{column_data_type} NULL, + metadata #{metadata_column_data_type} NULL, created_at timestamp with time zone DEFAULT NOW() NOT NULL ); """ diff --git a/lib/event_store/streams/stream.ex b/lib/event_store/streams/stream.ex index f82d89b5..9c9af2e9 100644 --- a/lib/event_store/streams/stream.ex +++ b/lib/event_store/streams/stream.ex @@ -7,9 +7,19 @@ defmodule EventStore.Streams.Stream do def append_to_stream(conn, stream_uuid, expected_version, events, opts) when length(events) < 1000 do {serializer, new_opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, new_opts} = Keyword.pop(new_opts, :metadata_serializer) with {:ok, stream} <- stream_info(conn, stream_uuid, expected_version, new_opts), - :ok <- do_append_to_storage(conn, stream, events, expected_version, serializer, new_opts) do + :ok <- + do_append_to_storage( + conn, + stream, + events, + expected_version, + serializer, + metadata_serializer, + new_opts + ) do :ok end |> maybe_retry_once(conn, stream_uuid, expected_version, events, opts) @@ -17,6 +27,7 @@ defmodule EventStore.Streams.Stream do def append_to_stream(conn, stream_uuid, expected_version, events, opts) do {serializer, new_opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, new_opts} = Keyword.pop(new_opts, :metadata_serializer) transaction( conn, @@ -29,6 +40,7 @@ defmodule EventStore.Streams.Stream do events, expected_version, serializer, + metadata_serializer, new_opts ) do :ok @@ -143,18 +155,19 @@ defmodule EventStore.Streams.Stream do events, expected_version, serializer, + metadata_serializer, opts ) do - prepared_events = prepare_events(events, stream, serializer) + prepared_events = prepare_events(events, stream, serializer, metadata_serializer) write_to_stream(conn, prepared_events, stream, expected_version, opts) end - defp prepare_events(events, %StreamInfo{} = stream, serializer) do + defp prepare_events(events, %StreamInfo{} = stream, serializer, metadata_serializer) do %StreamInfo{stream_uuid: stream_uuid, stream_version: stream_version} = stream events - |> Enum.map(&map_to_recorded_event(&1, utc_now(), serializer)) + |> Enum.map(&map_to_recorded_event(&1, utc_now(), serializer, metadata_serializer)) |> Enum.with_index(1) |> Enum.map(fn {recorded_event, index} -> %RecordedEvent{ @@ -168,13 +181,19 @@ defmodule EventStore.Streams.Stream do defp map_to_recorded_event( %EventData{data: %{__struct__: event_type}, event_type: nil} = event, created_at, - serializer + serializer, + metadata_serializer ) do %{event | event_type: Atom.to_string(event_type)} - |> map_to_recorded_event(created_at, serializer) + |> map_to_recorded_event(created_at, serializer, metadata_serializer) end - defp map_to_recorded_event(%EventData{} = event_data, created_at, serializer) do + defp map_to_recorded_event( + %EventData{} = event_data, + created_at, + serializer, + metadata_serializer + ) do %EventData{ event_id: event_id, causation_id: causation_id, @@ -190,7 +209,7 @@ defmodule EventStore.Streams.Stream do correlation_id: correlation_id, event_type: event_type, data: serializer.serialize(data), - metadata: serializer.serialize(metadata), + metadata: metadata_serializer.serialize(metadata), created_at: created_at } end @@ -222,10 +241,12 @@ defmodule EventStore.Streams.Stream do %StreamInfo{stream_id: stream_id} = stream {serializer, opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, opts} = Keyword.pop(opts, :metadata_serializer) with {:ok, recorded_events} <- Storage.read_stream_forward(conn, stream_id, start_version, count, opts) do - deserialized_events = deserialize_recorded_events(recorded_events, serializer) + deserialized_events = + deserialize_recorded_events(recorded_events, serializer, metadata_serializer) {:ok, deserialized_events} end @@ -235,10 +256,12 @@ defmodule EventStore.Streams.Stream do %StreamInfo{stream_id: stream_id} = stream {serializer, opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, opts} = Keyword.pop(opts, :metadata_serializer) with {:ok, recorded_events} <- Storage.read_stream_backward(conn, stream_id, start_version, count, opts) do - deserialized_events = deserialize_recorded_events(recorded_events, serializer) + deserialized_events = + deserialize_recorded_events(recorded_events, serializer, metadata_serializer) {:ok, deserialized_events} end @@ -289,8 +312,8 @@ defmodule EventStore.Streams.Stream do ) end - defp deserialize_recorded_events(recorded_events, serializer), - do: Enum.map(recorded_events, &RecordedEvent.deserialize(&1, serializer)) + defp deserialize_recorded_events(recorded_events, serializer, metadata_serializer), + do: Enum.map(recorded_events, &RecordedEvent.deserialize(&1, serializer, metadata_serializer)) defp soft_delete_stream(conn, stream, opts) do %StreamInfo{stream_id: stream_id} = stream diff --git a/lib/event_store/subscriptions/subscription_fsm.ex b/lib/event_store/subscriptions/subscription_fsm.ex index 2db225e2..090d2a94 100644 --- a/lib/event_store/subscriptions/subscription_fsm.ex +++ b/lib/event_store/subscriptions/subscription_fsm.ex @@ -17,6 +17,7 @@ defmodule EventStore.Subscriptions.SubscriptionFsm do stream_uuid: stream_uuid, subscription_name: subscription_name, serializer: Keyword.fetch!(opts, :serializer), + metadata_serializer: Keyword.fetch!(opts, :metadata_serializer), schema: Keyword.fetch!(opts, :schema), start_from: opts[:start_from] || 0, mapper: opts[:mapper], @@ -455,6 +456,7 @@ defmodule EventStore.Subscriptions.SubscriptionFsm do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid, last_sent: last_sent, max_size: max_size, @@ -464,6 +466,7 @@ defmodule EventStore.Subscriptions.SubscriptionFsm do Stream.read_stream_forward(conn, stream_uuid, last_sent + 1, max_size, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, timeout: query_timeout ) end diff --git a/lib/event_store/subscriptions/subscription_state.ex b/lib/event_store/subscriptions/subscription_state.ex index 3bec1488..96213fa5 100644 --- a/lib/event_store/subscriptions/subscription_state.ex +++ b/lib/event_store/subscriptions/subscription_state.ex @@ -8,6 +8,7 @@ defmodule EventStore.Subscriptions.SubscriptionState do :conn, :event_store, :serializer, + :metadata_serializer, :schema, :stream_uuid, :start_from, diff --git a/lib/event_store/supervisor.ex b/lib/event_store/supervisor.ex index 5a3c95ce..96fce336 100644 --- a/lib/event_store/supervisor.ex +++ b/lib/event_store/supervisor.ex @@ -115,7 +115,9 @@ defmodule EventStore.Supervisor do defp validate_config!(event_store, name, config) do conn = postgrex_conn(name, config) column_data_type = Config.column_data_type(event_store, config) + metadata_column_data_type = Config.metadata_column_data_type(event_store, config) serializer = Serializer.serializer(event_store, config) + metadata_serializer = Serializer.metadata_serializer(event_store, config) subscription_retry_interval = Subscriptions.retry_interval(event_store, config) subscription_hibernate_after = Subscriptions.hibernate_after(event_store, config) @@ -123,6 +125,8 @@ defmodule EventStore.Supervisor do conn: conn, column_data_type: column_data_type, serializer: serializer, + metadata_serializer: metadata_serializer, + metadata_column_data_type: metadata_column_data_type, subscription_retry_interval: subscription_retry_interval, subscription_hibernate_after: subscription_hibernate_after ) diff --git a/test/config_test.exs b/test/config_test.exs index 6d5e848c..b0fa0fc1 100644 --- a/test/config_test.exs +++ b/test/config_test.exs @@ -17,6 +17,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", pool: EventStore.Config.get_pool(), timeout: 120_000, schema: "example", @@ -41,6 +42,7 @@ defmodule EventStore.ConfigTest do pool: EventStore.Config.get_pool(), enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", timeout: 120_000, password: "postgres", @@ -61,6 +63,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", timeout: 15_000, pool: EventStore.Config.get_pool(), @@ -84,6 +87,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", pool: EventStore.Config.get_pool(), timeout: 120_000, @@ -100,6 +104,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", timeout: 15_000, pool: EventStore.Config.get_pool(), @@ -116,6 +121,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", timeout: 15_000, pool: EventStore.Config.get_pool(), @@ -132,6 +138,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", timeout: 15_000, pool: DBConnection.ConnectionPool, @@ -152,6 +159,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", pool: EventStore.Config.get_pool(), username: "username", @@ -172,6 +180,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", timeout: 15_000, pool: EventStore.Config.get_pool(), @@ -204,6 +213,7 @@ defmodule EventStore.ConfigTest do assert_parsed_config(config, enable_hard_deletes: false, column_data_type: "bytea", + metadata_column_data_type: "bytea", schema: "public", pool: EventStore.Config.get_pool(), timeout: 120_000, @@ -222,6 +232,7 @@ defmodule EventStore.ConfigTest do pool: EventStore.Config.get_pool(), username: "default_username", column_data_type: "bytea", + metadata_column_data_type: "bytea", enable_hard_deletes: false, schema: "public", timeout: 15_000 diff --git a/test/runtime_config_test.exs b/test/runtime_config_test.exs index 57929eaf..019f64ca 100644 --- a/test/runtime_config_test.exs +++ b/test/runtime_config_test.exs @@ -71,7 +71,8 @@ defmodule EventStore.RuntimeConfigTest do password: "postgres", database: "eventstore_test", hostname: "localhost", - serializer: EventStore.JsonSerializer + serializer: EventStore.JsonSerializer, + metadata_serializer: EventStore.JsonSerializer ] assert {:ok, _pid} = start_supervised({RuntimeConfiguredEventStore, config}) @@ -82,6 +83,7 @@ defmodule EventStore.RuntimeConfigTest do Keyword.merge( [ column_data_type: "bytea", + metadata_column_data_type: "bytea", enable_hard_deletes: false, otp_app: :eventstore, pool: DBConnection.ConnectionPool, diff --git a/test/streams/all_stream_test.exs b/test/streams/all_stream_test.exs index 4fa74502..e93d02f0 100644 --- a/test/streams/all_stream_test.exs +++ b/test/streams/all_stream_test.exs @@ -15,12 +15,14 @@ defmodule EventStore.Streams.AllStreamTest do test "should fetch events from all streams", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do {:ok, read_events} = Stream.read_stream_forward(conn, @all_stream, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert length(read_events) == 6 @@ -34,6 +36,7 @@ defmodule EventStore.Streams.AllStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream1_uuid: stream1_uuid, stream2_uuid: stream2_uuid } do @@ -41,7 +44,8 @@ defmodule EventStore.Streams.AllStreamTest do Stream.stream_forward(conn, @all_stream, 0, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -62,13 +66,15 @@ defmodule EventStore.Streams.AllStreamTest do test "should stream events from all streams using two event batch size", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do read_events = Stream.stream_forward(conn, @all_stream, 0, read_batch_size: 2, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -80,13 +86,15 @@ defmodule EventStore.Streams.AllStreamTest do test "should stream events from all streams using large batch size", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do read_events = Stream.stream_forward(conn, @all_stream, 0, read_batch_size: 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -103,6 +111,7 @@ defmodule EventStore.Streams.AllStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream1_uuid: stream1_uuid, stream2_uuid: stream2_uuid } do @@ -110,7 +119,8 @@ defmodule EventStore.Streams.AllStreamTest do Stream.stream_backward(conn, @all_stream, -1, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -131,13 +141,15 @@ defmodule EventStore.Streams.AllStreamTest do test "should stream events from all streams using two event batch size", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do read_events = Stream.stream_backward(conn, @all_stream, -1, read_batch_size: 2, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -149,13 +161,15 @@ defmodule EventStore.Streams.AllStreamTest do test "should stream events from all streams using large batch size", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do read_events = Stream.stream_backward(conn, @all_stream, -1, read_batch_size: 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -167,13 +181,15 @@ defmodule EventStore.Streams.AllStreamTest do test "should stream events from all streams using offset", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do read_events = Stream.stream_backward(conn, @all_stream, 3, read_batch_size: 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -210,6 +226,7 @@ defmodule EventStore.Streams.AllStreamTest do test "from current should receive only new events", %{ conn: conn, serializer: serializer, + metadata_serializer: metadata_serializer, schema: schema, stream1_uuid: stream1_uuid } do @@ -230,7 +247,8 @@ defmodule EventStore.Streams.AllStreamTest do :ok = Stream.append_to_stream(conn, stream1_uuid, 3, events, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert_receive {:events, received_events} @@ -260,13 +278,26 @@ defmodule EventStore.Streams.AllStreamTest do end defp append_events_to_streams(context) do - %{conn: conn, schema: schema, serializer: serializer} = context + %{ + conn: conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer + } = context {stream1_uuid, stream1_events} = - append_events_to_stream(conn, schema: schema, serializer: serializer) + append_events_to_stream(conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer + ) {stream2_uuid, stream2_events} = - append_events_to_stream(conn, schema: schema, serializer: serializer) + append_events_to_stream(conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer + ) [ stream1_uuid: stream1_uuid, diff --git a/test/streams/single_stream_test.exs b/test/streams/single_stream_test.exs index 8789abe9..0e52b1d8 100644 --- a/test/streams/single_stream_test.exs +++ b/test/streams/single_stream_test.exs @@ -14,12 +14,14 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do {:ok, events} = Stream.read_stream_forward(conn, stream_uuid, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert length(events) == 3 @@ -29,6 +31,7 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do utc_now = DateTime.utc_now() @@ -36,7 +39,8 @@ defmodule EventStore.Streams.SingleStreamTest do {:ok, [event]} = Stream.read_stream_forward(conn, stream_uuid, 0, 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) created_at = event.created_at @@ -56,13 +60,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, events: events, stream_uuid: stream_uuid } do assert {:error, :wrong_expected_version} = Stream.append_to_stream(conn, stream_uuid, 0, events, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) end end @@ -77,13 +83,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: source_stream_uuid, other_stream_uuid: target_stream_uuid } do {:ok, source_events} = Stream.read_stream_forward(conn, source_stream_uuid, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert :ok = @@ -92,7 +100,8 @@ defmodule EventStore.Streams.SingleStreamTest do assert {:ok, events} = Stream.read_stream_forward(conn, target_stream_uuid, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert length(events) == 4 @@ -116,13 +125,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: source_stream_uuid, other_stream_uuid: target_stream_uuid } do {:ok, source_events} = Stream.read_stream_forward(conn, source_stream_uuid, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert :ok = @@ -135,13 +146,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: source_stream_uuid, other_stream_uuid: target_stream_uuid } do {:ok, source_events} = Stream.read_stream_forward(conn, source_stream_uuid, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert {:error, :wrong_expected_version} = @@ -160,13 +173,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: source_stream_uuid, other_stream_uuid: target_stream_uuid } do {:ok, source_events} = Stream.read_stream_forward(conn, source_stream_uuid, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) :ok = Stream.link_to_stream(conn, target_stream_uuid, 1, source_events, schema: schema) @@ -179,6 +194,7 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do {:ok, _subscription} = @@ -194,7 +210,8 @@ defmodule EventStore.Streams.SingleStreamTest do :ok = Stream.append_to_stream(conn, stream_uuid, :any_version, [event], schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert_receive {:events, [received_event | _]} @@ -205,21 +222,24 @@ defmodule EventStore.Streams.SingleStreamTest do test "attempt to read an unknown stream forward should error stream not found", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do unknown_stream_uuid = UUID.uuid4() assert {:error, :stream_not_found} = Stream.read_stream_forward(conn, unknown_stream_uuid, 0, 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) end test "attempt to stream an unknown stream should error stream not found", %{ conn: conn, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer } do unknown_stream_uuid = UUID.uuid4() @@ -227,7 +247,8 @@ defmodule EventStore.Streams.SingleStreamTest do Stream.stream_forward(conn, unknown_stream_uuid, 0, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) end @@ -238,12 +259,14 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do {:ok, read_events} = Stream.read_stream_forward(conn, stream_uuid, 0, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert length(read_events) == 3 @@ -258,12 +281,14 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do {:ok, read_events} = Stream.read_stream_backward(conn, stream_uuid, -1, 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert length(read_events) == 3 @@ -278,13 +303,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_forward(conn, stream_uuid, 0, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -297,13 +324,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_forward(conn, stream_uuid, 0, read_batch_size: 2, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -314,13 +343,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_forward(conn, stream_uuid, 0, read_batch_size: 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -331,13 +362,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_forward(conn, stream_uuid, 2, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -347,12 +380,19 @@ defmodule EventStore.Streams.SingleStreamTest do end test "should stream events from single stream with starting version offset outside range", - %{conn: conn, schema: schema, serializer: serializer, stream_uuid: stream_uuid} do + %{ + conn: conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer, + stream_uuid: stream_uuid + } do read_events = Stream.stream_forward(conn, stream_uuid, 4, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -367,13 +407,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_backward(conn, stream_uuid, -1, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -386,13 +428,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_backward(conn, stream_uuid, -1, read_batch_size: 2, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -405,13 +449,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_backward(conn, stream_uuid, -1, read_batch_size: 1_000, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -424,13 +470,15 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do read_events = Stream.stream_backward(conn, stream_uuid, 2, read_batch_size: 2, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -440,12 +488,19 @@ defmodule EventStore.Streams.SingleStreamTest do end test "should stream events from single stream with starting version offset outside range", - %{conn: conn, schema: schema, serializer: serializer, stream_uuid: stream_uuid} do + %{ + conn: conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer, + stream_uuid: stream_uuid + } do read_events = Stream.stream_backward(conn, stream_uuid, 0, read_batch_size: 1, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) |> Enum.to_list() @@ -474,6 +529,7 @@ defmodule EventStore.Streams.SingleStreamTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid } do {:ok, _subscription} = @@ -491,7 +547,8 @@ defmodule EventStore.Streams.SingleStreamTest do :ok = Stream.append_to_stream(conn, stream_uuid, 3, events, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert_receive {:events, received_events} @@ -507,14 +564,20 @@ defmodule EventStore.Streams.SingleStreamTest do end end - test "should return stream version", %{conn: conn, schema: schema, serializer: serializer} do + test "should return stream version", %{ + conn: conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer + } do stream_uuid = UUID.uuid4() events = EventFactory.create_events(3) :ok = Stream.append_to_stream(conn, stream_uuid, 0, events, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) # stream above needed for preventing accidental event_number/stream_version match @@ -524,14 +587,20 @@ defmodule EventStore.Streams.SingleStreamTest do :ok = Stream.append_to_stream(conn, stream_uuid, 0, events, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) assert {:ok, 3} = Stream.stream_version(conn, stream_uuid, schema: schema) end defp append_events_to_stream(context) do - %{conn: conn, schema: schema, serializer: serializer} = context + %{ + conn: conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer + } = context stream_uuid = UUID.uuid4() events = EventFactory.create_events(3) @@ -539,7 +608,8 @@ defmodule EventStore.Streams.SingleStreamTest do :ok = Stream.append_to_stream(conn, stream_uuid, 0, events, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) [ @@ -549,7 +619,12 @@ defmodule EventStore.Streams.SingleStreamTest do end defp append_event_to_another_stream(context) do - %{conn: conn, schema: schema, serializer: serializer} = context + %{ + conn: conn, + schema: schema, + serializer: serializer, + metadata_serializer: metadata_serializer + } = context stream_uuid = UUID.uuid4() events = EventFactory.create_events(1) @@ -557,7 +632,8 @@ defmodule EventStore.Streams.SingleStreamTest do :ok = Stream.append_to_stream(conn, stream_uuid, 0, events, schema: schema, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ) [ diff --git a/test/subscriptions/subscribe_to_stream_test.exs b/test/subscriptions/subscribe_to_stream_test.exs index 01d8c7b8..e9949b36 100644 --- a/test/subscriptions/subscribe_to_stream_test.exs +++ b/test/subscriptions/subscribe_to_stream_test.exs @@ -30,6 +30,7 @@ defmodule EventStore.Subscriptions.SubscribeToStreamTest do test "should receive `:subscribed` message once subscribed", %{ schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, subscription_name: subscription_name } do stream_uuid = UUID.uuid4() @@ -40,6 +41,7 @@ defmodule EventStore.Subscriptions.SubscribeToStreamTest do conn: @conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, hibernate_after: 15_000, retry_interval: 1_000, stream_uuid: stream_uuid, diff --git a/test/subscriptions/subscription_locking_test.exs b/test/subscriptions/subscription_locking_test.exs index 79992c18..5021cc3c 100644 --- a/test/subscriptions/subscription_locking_test.exs +++ b/test/subscriptions/subscription_locking_test.exs @@ -160,6 +160,7 @@ defmodule EventStore.Subscriptions.SubscriptionLockingTest do schema: schema, event_store: event_store, serializer: serializer, + metadata_serializer: metadata_serializer, subscription_name: subscription_name } = context @@ -169,6 +170,7 @@ defmodule EventStore.Subscriptions.SubscriptionLockingTest do conn: conn, schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, retry_interval: 1_000, stream_uuid: "$all", subscription_name: subscription_name, diff --git a/test/subscriptions/support/stream_subscription_test_case.ex b/test/subscriptions/support/stream_subscription_test_case.ex index 8f7fee9f..36d601ab 100644 --- a/test/subscriptions/support/stream_subscription_test_case.ex +++ b/test/subscriptions/support/stream_subscription_test_case.ex @@ -440,6 +440,7 @@ defmodule EventStore.Subscriptions.StreamSubscriptionTestCase do %{ schema: schema, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid, subscriber: subscriber } = context @@ -450,6 +451,7 @@ defmodule EventStore.Subscriptions.StreamSubscriptionTestCase do |> Keyword.put(:event_store, @event_store) |> Keyword.put(:schema, schema) |> Keyword.put(:serializer, serializer) + |> Keyword.put(:metadata_serializer, metadata_serializer) |> Keyword.put_new(:buffer_size, 3) stream_uuid diff --git a/test/support/storage_case.ex b/test/support/storage_case.ex index f4b433db..9e65b96c 100644 --- a/test/support/storage_case.ex +++ b/test/support/storage_case.ex @@ -6,6 +6,7 @@ defmodule EventStore.StorageCase do setup_all do config = Config.parsed(TestEventStore, :eventstore) serializer = Serializer.serializer(TestEventStore, config) + metadata_serializer = Serializer.metadata_serializer(TestEventStore, config) postgrex_config = Config.default_postgrex_opts(config) if Mix.env() == :migration do @@ -22,7 +23,8 @@ defmodule EventStore.StorageCase do schema: "public", event_store: TestEventStore, postgrex_config: postgrex_config, - serializer: serializer + serializer: serializer, + metadata_serializer: metadata_serializer ] end