From 69ea9c80d2a863c5c5bd5518540f5787c3088c3d Mon Sep 17 00:00:00 2001 From: Alex Speller Date: Fri, 9 Jul 2021 01:29:58 +0100 Subject: [PATCH] Allow different serializer for metadata --- lib/event_store.ex | 21 ++++++++-- lib/event_store/config.ex | 13 ++++++ lib/event_store/notifications/reader.ex | 12 +++--- lib/event_store/notifications/supervisor.ex | 5 ++- lib/event_store/recorded_event.ex | 4 +- lib/event_store/serializer.ex | 15 +++++++ lib/event_store/snapshots/snapshot_data.ex | 8 ++-- lib/event_store/snapshots/snapshotter.ex | 14 +++++-- lib/event_store/sql/statements.ex | 13 +++--- lib/event_store/streams/stream.ex | 41 +++++++++++++------ .../subscriptions/subscription_fsm.ex | 7 +++- .../subscriptions/subscription_state.ex | 1 + lib/event_store/supervisor.ex | 9 ++-- 13 files changed, 119 insertions(+), 44 deletions(-) diff --git a/lib/event_store.ex b/lib/event_store.ex index 04e77a09..5296dff6 100644 --- a/lib/event_store.ex +++ b/lib/event_store.ex @@ -24,6 +24,7 @@ defmodule EventStore do config :my_app, MyApp.EventStore, serializer: EventStore.JsonSerializer, + metadata_serializer: EventStore.JsonSerializer, username: "postgres", password: "postgres", database: "eventstore", @@ -188,12 +189,14 @@ defmodule EventStore do {otp_app, config} = EventStore.Supervisor.compile_config(__MODULE__, opts) serializer = Serializer.serializer(__MODULE__, config) + metadata_serializer = Serializer.metadata_serializer(__MODULE__, config) registry = Registration.registry(__MODULE__, config) subscription_retry_interval = Subscriptions.retry_interval(__MODULE__, config) @otp_app otp_app @config config @serializer serializer + @metadata_serializer metadata_serializer @registry registry @subscription_retry_interval subscription_retry_interval @@ -223,7 +226,15 @@ defmodule EventStore do opts = Keyword.merge(unquote(opts), opts) name = name(opts) - EventStore.Supervisor.start_link(__MODULE__, @otp_app, @serializer, @registry, name, opts) + EventStore.Supervisor.start_link( + __MODULE__, + @otp_app, + @serializer, + @metadata_serializer, + @registry, + name, + opts + ) end def stop(supervisor, timeout \\ 5000) do @@ -321,6 +332,7 @@ defmodule EventStore do event_store: name, registry: @registry, serializer: @serializer, + metadata_serializer: @metadata_serializer, retry_interval: @subscription_retry_interval, stream_uuid: stream_uuid, subscription_name: subscription_name, @@ -366,13 +378,13 @@ defmodule EventStore do def read_snapshot(source_uuid, opts \\ []) do {conn, opts} = opts(opts) - Snapshotter.read_snapshot(conn, source_uuid, @serializer, opts) + Snapshotter.read_snapshot(conn, source_uuid, @serializer, @metadata_serializer, opts) end def record_snapshot(%SnapshotData{} = snapshot, opts \\ []) do {conn, opts} = opts(opts) - Snapshotter.record_snapshot(conn, snapshot, @serializer, opts) + Snapshotter.record_snapshot(conn, snapshot, @serializer, @metadata_serializer, opts) end def delete_snapshot(source_uuid, opts \\ []) do @@ -387,7 +399,8 @@ defmodule EventStore do timeout = timeout(opts) - {conn, [timeout: timeout, serializer: @serializer]} + {conn, + [timeout: timeout, serializer: @serializer, metadata_serializer: @metadata_serializer]} end defp name(opts) do diff --git a/lib/event_store/config.ex b/lib/event_store/config.ex index 9f87fe2d..1d5c615d 100644 --- a/lib/event_store/config.ex +++ b/lib/event_store/config.ex @@ -57,6 +57,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 [ :username, :password, diff --git a/lib/event_store/notifications/reader.ex b/lib/event_store/notifications/reader.ex index 3cc6fa8c..33a53763 100644 --- a/lib/event_store/notifications/reader.ex +++ b/lib/event_store/notifications/reader.ex @@ -9,13 +9,14 @@ defmodule EventStore.Notifications.Reader do alias EventStore.Storage defmodule State do - defstruct [:conn, :serializer, :subscribe_to] + defstruct [:conn, :serializer, :metadata_serializer, :subscribe_to] end def start_link(opts \\ []) do state = %State{ conn: Keyword.fetch!(opts, :conn), serializer: Keyword.fetch!(opts, :serializer), + metadata_serializer: Keyword.fetch!(opts, :metadata_serializer), subscribe_to: Keyword.fetch!(opts, :subscribe_to) } @@ -45,18 +46,19 @@ defmodule EventStore.Notifications.Reader do defp read_events(event, %State{} = state) do {stream_uuid, stream_id, from_stream_version, to_stream_version} = event - %State{conn: conn, serializer: serializer} = state + %State{conn: conn, serializer: serializer, metadata_serializer: metadata_serializer} = state count = to_stream_version - from_stream_version + 1 with {:ok, events} <- Storage.read_stream_forward(conn, stream_id, from_stream_version, count), - deserialized_events <- deserialize_recorded_events(events, serializer) do + deserialized_events <- + deserialize_recorded_events(events, serializer, metadata_serializer) do {stream_uuid, deserialized_events} 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 end diff --git a/lib/event_store/notifications/supervisor.ex b/lib/event_store/notifications/supervisor.ex index 68860bcb..216809d0 100644 --- a/lib/event_store/notifications/supervisor.ex +++ b/lib/event_store/notifications/supervisor.ex @@ -15,7 +15,7 @@ defmodule EventStore.Notifications.Supervisor do alias EventStore.MonitoredServer alias EventStore.Notifications.{Listener, Reader, Broadcaster} - def child_spec({name, _registry, _serializer, _config} = init_arg) do + def child_spec({name, _registry, _serializer, _metadata_serializer, _config} = init_arg) do %{id: Module.concat(name, __MODULE__), start: {__MODULE__, :start_link, [init_arg]}} end @@ -24,7 +24,7 @@ defmodule EventStore.Notifications.Supervisor do end @impl Supervisor - def init({event_store, registry, serializer, config}) do + def init({event_store, registry, serializer, metadata_serializer, config}) do schema = Keyword.fetch!(config, :schema) postgrex_config = Config.sync_connect_postgrex_opts(config) @@ -52,6 +52,7 @@ defmodule EventStore.Notifications.Supervisor do {Reader, conn: postgrex_reader_name, serializer: serializer, + metadata_serializer: metadata_serializer, subscribe_to: listener_name, name: reader_name}, {Broadcaster, diff --git a/lib/event_store/recorded_event.ex b/lib/event_store/recorded_event.ex index 7d75672e..ea72506c 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 c1f30600..b95df217 100644 --- a/lib/event_store/snapshots/snapshotter.ex +++ b/lib/event_store/snapshots/snapshotter.ex @@ -7,9 +7,9 @@ defmodule EventStore.Snapshots.Snapshotter do @doc """ Read a snapshot, if available, for a given source. """ - def read_snapshot(conn, source_uuid, serializer, opts \\ []) do + def read_snapshot(conn, source_uuid, serializer, metadata_serializer, opts \\ []) do 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 @@ -20,8 +20,14 @@ defmodule EventStore.Snapshots.Snapshotter do Returns `:ok` on success. """ - def record_snapshot(conn, %SnapshotData{} = snapshot, serializer, opts \\ []) do - serialized = SnapshotData.serialize(snapshot, serializer) + def record_snapshot( + conn, + %SnapshotData{} = snapshot, + serializer, + metadata_serializer, + opts \\ [] + ) do + serialized = SnapshotData.serialize(snapshot, serializer, metadata_serializer) Snapshot.record_snapshot(conn, serialized, opts) end diff --git a/lib/event_store/sql/statements.ex b/lib/event_store/sql/statements.ex index ca9135ff..c2b65d49 100644 --- a/lib/event_store/sql/statements.ex +++ b/lib/event_store/sql/statements.ex @@ -7,12 +7,13 @@ defmodule EventStore.Sql.Statements do def initializers(event_store, config) do column_data_type = Config.column_data_type(event_store, config) + metadata_column_data_type = Config.metadata_column_data_type(event_store, config) [ create_streams_table(), create_stream_uuid_index(), seed_all_stream(), - create_events_table(column_data_type), + create_events_table(column_data_type, metadata_column_data_type), prevent_event_update(), prevent_event_delete(), create_stream_events_table(), @@ -23,7 +24,7 @@ defmodule EventStore.Sql.Statements 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() ] @@ -80,7 +81,7 @@ defmodule EventStore.Sql.Statements 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 ( @@ -89,7 +90,7 @@ defmodule EventStore.Sql.Statements 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 ); """ @@ -197,7 +198,7 @@ defmodule EventStore.Sql.Statements 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 ( @@ -205,7 +206,7 @@ defmodule EventStore.Sql.Statements 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 d4802dcf..71b3fb17 100644 --- a/lib/event_store/streams/stream.ex +++ b/lib/event_store/streams/stream.ex @@ -9,15 +9,17 @@ defmodule EventStore.Streams.Stream do def append_to_stream(conn, stream_uuid, expected_version, events, opts \\ []) do {serializer, opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, opts} = Keyword.pop(opts, :metadata_serializer) with {:ok, stream} <- stream_info(conn, stream_uuid, opts), {:ok, stream} <- prepare_stream(conn, expected_version, stream, opts) do - do_append_to_storage(conn, events, stream, serializer, opts) + do_append_to_storage(conn, events, stream, serializer, metadata_serializer, opts) end end def link_to_stream(conn, stream_uuid, expected_version, events_or_event_ids, opts \\ []) do {_serializer, opts} = Keyword.pop(opts, :serializer) + {_metadata_serializer, opts} = Keyword.pop(opts, :metadata_serializer) with {:ok, stream} <- stream_info(conn, stream_uuid, opts), {:ok, stream} <- prepare_stream(conn, expected_version, stream, opts) do @@ -139,17 +141,24 @@ defmodule EventStore.Streams.Stream do defp prepare_stream(_conn, _expected_version, _state, _opts), do: {:error, :wrong_expected_version} - defp do_append_to_storage(conn, events, %Stream{} = stream, serializer, opts) do - prepared_events = prepare_events(events, stream, serializer) + defp do_append_to_storage( + conn, + events, + %Stream{} = stream, + serializer, + metadata_serializer, + opts + ) do + prepared_events = prepare_events(events, stream, serializer, metadata_serializer) write_to_stream(conn, prepared_events, stream, opts) end - defp prepare_events(events, %Stream{} = stream, serializer) do + defp prepare_events(events, %Stream{} = stream, serializer, metadata_serializer) do %Stream{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{ @@ -166,13 +175,19 @@ defmodule EventStore.Streams.Stream do 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{ causation_id: causation_id, correlation_id: correlation_id, @@ -187,7 +202,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 @@ -220,10 +235,12 @@ defmodule EventStore.Streams.Stream do defp read_storage_forward(conn, stream_id, start_version, count, opts) do {serializer, opts} = Keyword.pop(opts, :serializer) + {metadata_serializer, opts} = Keyword.pop(opts, :metadata_serializer) case Storage.read_stream_forward(conn, stream_id, start_version, count, opts) do {:ok, recorded_events} -> - deserialized_events = deserialize_recorded_events(recorded_events, serializer) + deserialized_events = + deserialize_recorded_events(recorded_events, serializer, metadata_serializer) {:ok, deserialized_events} @@ -254,8 +271,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 query_opts(opts), do: Keyword.take(opts, [:timeout]) end diff --git a/lib/event_store/subscriptions/subscription_fsm.ex b/lib/event_store/subscriptions/subscription_fsm.ex index 5835ae7f..c7df4a1d 100644 --- a/lib/event_store/subscriptions/subscription_fsm.ex +++ b/lib/event_store/subscriptions/subscription_fsm.ex @@ -18,6 +18,7 @@ defmodule EventStore.Subscriptions.SubscriptionFsm do subscription_name: subscription_name, registry: Keyword.fetch!(opts, :registry), serializer: Keyword.fetch!(opts, :serializer), + metadata_serializer: Keyword.fetch!(opts, :metadata_serializer), start_from: opts[:start_from] || 0, mapper: opts[:mapper], selector: opts[:selector], @@ -402,12 +403,16 @@ defmodule EventStore.Subscriptions.SubscriptionFsm do %SubscriptionState{ conn: conn, serializer: serializer, + metadata_serializer: metadata_serializer, stream_uuid: stream_uuid, last_sent: last_sent, max_size: max_size } = data - Stream.read_stream_forward(conn, stream_uuid, last_sent + 1, max_size, serializer: serializer) + Stream.read_stream_forward(conn, stream_uuid, last_sent + 1, max_size, + serializer: serializer, + metadata_serializer: metadata_serializer + ) end defp enqueue_events(%SubscriptionState{} = data, []), do: data diff --git a/lib/event_store/subscriptions/subscription_state.ex b/lib/event_store/subscriptions/subscription_state.ex index 8298c7f5..514f256c 100644 --- a/lib/event_store/subscriptions/subscription_state.ex +++ b/lib/event_store/subscriptions/subscription_state.ex @@ -6,6 +6,7 @@ defmodule EventStore.Subscriptions.SubscriptionState do :event_store, :registry, :serializer, + :metadata_serializer, :stream_uuid, :start_from, :subscription_name, diff --git a/lib/event_store/supervisor.ex b/lib/event_store/supervisor.ex index cc5fb9f5..a96da541 100644 --- a/lib/event_store/supervisor.ex +++ b/lib/event_store/supervisor.ex @@ -15,10 +15,10 @@ defmodule EventStore.Supervisor do @doc """ Starts the event store supervisor. """ - def start_link(event_store, otp_app, serializer, registry, name, opts) do + def start_link(event_store, otp_app, serializer, metadata_serializer, registry, name, opts) do Supervisor.start_link( __MODULE__, - {event_store, otp_app, serializer, registry, name, opts}, + {event_store, otp_app, serializer, metadata_serializer, registry, name, opts}, name: name ) end @@ -56,7 +56,7 @@ defmodule EventStore.Supervisor do ## Supervisor callbacks @doc false - def init({event_store, otp_app, serializer, registry, name, opts}) do + def init({event_store, otp_app, serializer, metadata_serializer, registry, name, opts}) do case runtime_config(event_store, otp_app, opts) do {:ok, config} -> advisory_locks_name = Module.concat([name, AdvisoryLocks]) @@ -77,7 +77,8 @@ defmodule EventStore.Supervisor do {Registry, keys: :unique, name: subscriptions_registry_name}, id: subscriptions_registry_name ), - {Highlander, {Notifications.Supervisor, {name, registry, serializer, config}}} + {Highlander, + {Notifications.Supervisor, {name, registry, serializer, metadata_serializer, config}}} ] ++ Registration.child_spec(name, registry) Supervisor.init(children, strategy: :one_for_all)