Class: PgEventstore::IndexFilteringQueries
- Inherits:
-
Object
- Object
- PgEventstore::IndexFilteringQueries
- Defined in:
- lib/pg_eventstore/queries/index_filtering_queries.rb,
sig/pg_eventstore/queries/index_filtering_queries.rbs
Instance Attribute Summary collapse
-
#connection ⇒ Connection
Returns the value of attribute connection.
Instance Method Summary collapse
- #compute_read_api_chunks_repo(indexes, resolve_link_tos) ⇒ PgEventstore::Chunks::Repository
- #deserialize(pg_result, repr: nil) ⇒ Array<PgEventstore::EventGlobalIndex>, ...
- #expand_event_types(filters_collection) ⇒ Array<PgEventstore::QueryBuilders::Filters::FilterRow>
- #fetch_grouped_indexes_for_read_api(filters_collection, cursor) ⇒ Array<PgEventstore::EventGlobalIndex::ReadApiRepr>
- #fetch_indexes_for_read_api(filters_collection, cursor) ⇒ Array<PgEventstore::EventGlobalIndex::ReadApiRepr>
- #fetch_indexes_for_revision_validation(stream, expected_revisions) ⇒ Array<PgEventstore::EventGlobalIndex::RevisionCheckRepr>
-
#initialize(connection, query_strategy) ⇒ IndexFilteringQueries
constructor
A new instance of IndexFilteringQueries.
- #partition_queries ⇒ PgEventstore::PartitionQueries
Constructor Details
#initialize(connection, query_strategy) ⇒ IndexFilteringQueries
Returns a new instance of IndexFilteringQueries.
12 13 14 15 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 12 def initialize(connection, query_strategy) @connection = connection @query_strategy = query_strategy end |
Instance Attribute Details
#connection ⇒ Connection
Returns the value of attribute connection.
7 8 9 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 7 def connection @connection end |
Instance Method Details
#compute_read_api_chunks_repo(indexes, resolve_link_tos) ⇒ PgEventstore::Chunks::Repository
78 79 80 81 82 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 78 def compute_read_api_chunks_repo(indexes, resolve_link_tos) repo = Chunks::Repository.new repo.add_chunk(Chunks::ReadApiEventsIndexChunk.new(indexes, connection, @query_strategy, resolve_link_tos)) repo end |
#deserialize(pg_result, repr: nil) ⇒ Array<PgEventstore::EventGlobalIndex>, ...
173 174 175 176 177 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 173 def deserialize(pg_result, repr: nil) pg_result.map do |attrs| EventGlobalIndex.create_representation(attrs.transform_keys(&:to_sym), repr:) end end |
#expand_event_types(filters_collection) ⇒ Array<PgEventstore::QueryBuilders::Filters::FilterRow>
93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 93 def (filters_collection) filter_rows = filters_collection.collection = filter_rows.select(&:ambiguous_event_type?) rest_filter_rows = filter_rows - partition_builders = .map do |filter_row| partitions_filtering = QueryBuilders::PartitionsFiltering.new partitions_filtering.with_event_types filter_row = filter_row.to_filter_row if filter_row.is_a?(QueryBuilders::Filters::MarkerFilterRow) partitions_filtering.add_filter_row(filter_row) partitions_filtering.to_sql_builder.unselect.select('context, stream_name, event_type') end final_partition_builder = SQLBuilder.union_builders(partition_builders) = @query_strategy.exec_params(*final_partition_builder.to_exec_params) = .group_by { _1['event_type'] } = .flat_map do |filter_row| case filter_row when QueryBuilders::Filters::FilterRow filter_row.event_type_filters.flat_map do |event_type_filter| = if event_type_filter.prefix? .select do |event_type, _| event_type.start_with?(event_type_filter.event_type) end.values.flatten else [event_type_filter.event_type] end next [] unless if filter_row.stream_filter = .select do |attrs| if filter_row.stream_filter.context? attrs['context'] == filter_row.stream_filter.context else attrs['context'] == filter_row.stream_filter.context && attrs['stream_name'] == filter_row.stream_filter.stream_name end end end .map do |partition_attrs| stream_filter = QueryBuilders::Filters::StreamFilter.new( context: partition_attrs['context'], stream_name: partition_attrs['stream_name'] ) event_type_filter = QueryBuilders::Filters::EventTypeFilter.new( event_type: partition_attrs['event_type'], prefix: false ) QueryBuilders::Filters::FilterRow.new(stream_filter:, event_type_filters: [event_type_filter]) end end when QueryBuilders::Filters::MarkerFilterRow = [filter_row.marker_filter.event_type] next [] unless if filter_row.stream_filter = .select do |attrs| attrs['context'] == filter_row.stream_filter.context end end .map do |partition_attrs| stream_filter = QueryBuilders::Filters::StreamFilter.new( context: partition_attrs['context'], stream_name: partition_attrs['stream_name'] ) QueryBuilders::Filters::MarkerFilterRow.new(stream_filter:, marker_filter: filter_row.marker_filter.dup) end else Utils.missing_implementation!(filter_row) end end adjusted_rows = + rest_filter_rows return [QueryBuilders::Filters::FilterRow.null_filter_row] if adjusted_rows.empty? adjusted_rows end |
#fetch_grouped_indexes_for_read_api(filters_collection, cursor) ⇒ Array<PgEventstore::EventGlobalIndex::ReadApiRepr>
63 64 65 66 67 68 69 70 71 72 73 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 63 def fetch_grouped_indexes_for_read_api(filters_collection, cursor) filter_rows = if filters_collection.collection.any?(&:ambiguous_event_type?) (filters_collection) else filters_collection.collection end sql_builder = QueryBuilders::IndexBasedEventsFiltering.sql_builder_for_read_grouped(filter_rows, cursor) deserialize(@query_strategy.exec_params(*sql_builder.to_exec_params), repr: EventGlobalIndex::ReprType::READ_API) end |
#fetch_indexes_for_read_api(filters_collection, cursor) ⇒ Array<PgEventstore::EventGlobalIndex::ReadApiRepr>
48 49 50 51 52 53 54 55 56 57 58 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 48 def fetch_indexes_for_read_api(filters_collection, cursor) filter_rows = if filters_collection.collection.any?(&:ambiguous_event_type?) (filters_collection) else filters_collection.collection end sql_builder = QueryBuilders::IndexBasedEventsFiltering.sql_builder_for_read_common(filter_rows, cursor) deserialize(@query_strategy.exec_params(*sql_builder.to_exec_params), repr: EventGlobalIndex::ReprType::READ_API) end |
#fetch_indexes_for_revision_validation(stream, expected_revisions) ⇒ Array<PgEventstore::EventGlobalIndex::RevisionCheckRepr>
22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 22 def fetch_indexes_for_revision_validation(stream, expected_revisions) builders = expected_revisions.map do |expected_revision| case expected_revision when Commands::RevisionCheck::ExpectedRevision::EventTypeRevision QueryBuilders::IndexBasedEventsFiltering.sql_builder_for_event_type_revision_validation( stream, expected_revision ) when Commands::RevisionCheck::ExpectedRevision::EventTypeRevisionWithMarkers QueryBuilders::IndexBasedEventsFiltering.sql_builder_for_event_type_with_markers_revision_validation( stream, expected_revision ) when Commands::RevisionCheck::ExpectedRevision::MarkersRevision QueryBuilders::IndexBasedEventsFiltering.sql_builder_for_markers_revision_validation( stream, expected_revision ) else Utils.missing_implementation!(expected_revision) end end final = SQLBuilder.union_builders(builders) deserialize(@query_strategy.exec_params(*final.to_exec_params), repr: EventGlobalIndex::ReprType::REVISION_CHECK) end |
#partition_queries ⇒ PgEventstore::PartitionQueries
87 88 89 |
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 87 def partition_queries PartitionQueries.new(connection) end |