Class: PgEventstore::QueryBuilders::IndexBasedEventsFiltering
- Inherits:
-
Object
- Object
- PgEventstore::QueryBuilders::IndexBasedEventsFiltering
- Defined in:
- lib/pg_eventstore/query_builders/index_based_events_filtering.rb,
sig/pg_eventstore/query_builders/index_based_events_filtering.rbs
Constant Summary collapse
- DEFAULT_LIMIT =
1_000- ABSTRACT_TABLE_NAME =
'events_index'- ReadCursor =
Constants included from BasicFiltering
BasicFiltering::SQL_DIRECTIONS
Class Method Summary collapse
- .sql_builder_for_estimate_count(filter_rows) ⇒ PgEventstore::SQLBuilder
- .sql_builder_for_event_type_revision_validation(stream, expected_revision) ⇒ PgEventstore::SQLBuilder
- .sql_builder_for_event_type_with_markers_revision_validation(stream, expected_revision) ⇒ PgEventstore::SQLBuilder
- .sql_builder_for_markers_revision_validation(stream, expected_revision) ⇒ PgEventstore::SQLBuilder
- .sql_builder_for_read_common(filter_rows, cursor) ⇒ PgEventstore::SQLBuilder
- .sql_builder_for_read_grouped(filter_rows, cursor) ⇒ PgEventstore::SQLBuilder
- .sql_builder_for_subscriptions(filter_rows) ⇒ PgEventstore::SQLBuilder
Instance Method Summary collapse
- #add_global_position_direction(direction) ⇒ void
- #add_limit(limit) ⇒ void
-
#initialize ⇒ IndexBasedEventsFiltering
constructor
A new instance of IndexBasedEventsFiltering.
- #to_sql_builder ⇒ Object
- #to_table_name ⇒ Object
Constructor Details
#initialize ⇒ IndexBasedEventsFiltering
Returns a new instance of IndexBasedEventsFiltering.
234 235 236 237 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 234 def initialize @sql_builder = SQLBuilder.new @sql_builder.select('global_position, event_type_partition_id') end |
Class Method Details
.sql_builder_for_estimate_count(filter_rows) ⇒ PgEventstore::SQLBuilder
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 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 140 def sql_builder_for_estimate_count(filter_rows) cursor = ReadCursor::StreamCursor.({}) if filter_rows.empty? return EventsGlobalIndexFiltering.default_filtering(cursor).to_sql_builder.remove_limit.remove_order end has_markers = false builders = filter_rows.flat_map do |filter_row| case filter_row when Filters::FilterRow filter_row.flatten.map do |flattened| EventsGlobalIndexFiltering.for_read_common(flattened, cursor) end when Filters::MarkerFilterRow has_markers ||= true filter_row.flatten.map do |flattened| EventMarkersIndexFiltering.for_read_common(flattened, cursor) end else Utils.missing_implementation!(filter_row) end end builders.each { _1.remove_limit.remove_order } top_builder = SQLBuilder.new.select('*') top_builder.from(SQLBuilder.union_builders(builders, mode: has_markers ? :union_distinct : :union_all)) top_builder end |
.sql_builder_for_event_type_revision_validation(stream, expected_revision) ⇒ PgEventstore::SQLBuilder
17 18 19 20 21 22 23 24 25 26 27 28 29 30 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 17 def sql_builder_for_event_type_revision_validation(stream, expected_revision) cursor = QueryBuilders::ReadCursor::StreamCursor.( stream, { direction: :desc, max_count: 1 } ) filters_collection = Filters::Collection.( stream, { filter: { event_types: [expected_revision.event_type] } } ) rows = filters_collection.collection Utils.assert!(rows.size == 1, filters_collection.inspect) EventsGlobalIndexFiltering.for_revision_validation_per_type( rows.first, cursor, expected_revision.sequence_number ) end |
.sql_builder_for_event_type_with_markers_revision_validation(stream, expected_revision) ⇒ PgEventstore::SQLBuilder
35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 35 def sql_builder_for_event_type_with_markers_revision_validation(stream, expected_revision) cursor = QueryBuilders::ReadCursor::StreamCursor.( stream, { direction: :desc, max_count: 1 } ) filters_collection = Filters::Collection.( stream, { filter: { event_types: [{ type: expected_revision.event_type, markers: expected_revision.markers }], }, } ) rows = filters_collection.collection Utils.assert!(rows.size == 1, filters_collection.inspect) builders = rows.first.flatten.map do |flattened| EventMarkersIndexFiltering.for_revision_validation_per_type( flattened, cursor, expected_revision.sequence_number ) end final = SQLBuilder.new.select('*') final.from(SQLBuilder.union_builders(builders, mode: :union_distinct)) final.order('stream_revision desc') final.limit(1) end |
.sql_builder_for_markers_revision_validation(stream, expected_revision) ⇒ PgEventstore::SQLBuilder
64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 64 def sql_builder_for_markers_revision_validation(stream, expected_revision) cursor = QueryBuilders::ReadCursor::StreamCursor.( stream, { direction: :desc, max_count: 1 } ) filters_collection = Filters::Collection.( stream, { filter: { event_types: [{ markers: expected_revision.markers }] } } ) rows = filters_collection.collection Utils.assert!(rows.size == 1, filters_collection.inspect) builders = rows.first.flatten.map do |flattened| EventMarkersIndexFiltering.for_revision_validation_per_type( flattened, cursor, expected_revision.sequence_number ) end final = SQLBuilder.new.select('*') final.from(SQLBuilder.union_builders(builders, mode: :union_distinct)) final.order('stream_revision desc') final.limit(1) end |
.sql_builder_for_read_common(filter_rows, cursor) ⇒ PgEventstore::SQLBuilder
115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 115 def sql_builder_for_read_common(filter_rows, cursor) return EventsGlobalIndexFiltering.default_filtering(cursor).to_sql_builder if filter_rows.empty? has_markers = false builders = filter_rows.flat_map do |filter_row| case filter_row when Filters::FilterRow filter_row.flatten.map do |flattened| EventsGlobalIndexFiltering.for_read_common(flattened, cursor) end when Filters::MarkerFilterRow has_markers ||= true filter_row.flatten.map do |flattened| EventMarkersIndexFiltering.for_read_common(flattened, cursor) end else Utils.missing_implementation!(filter_row) end end union_builders(builders, cursor, mode: has_markers ? :union_distinct : :union_all) end |
.sql_builder_for_read_grouped(filter_rows, cursor) ⇒ PgEventstore::SQLBuilder
89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 89 def sql_builder_for_read_grouped(filter_rows, cursor) cursor = cursor.dup cursor.max_count = 1 has_markers = false builders = filter_rows.flat_map do |filter_row| case filter_row when Filters::FilterRow filter_row.flatten.map do |flattened| EventsGlobalIndexFiltering.for_read_common(flattened, cursor) end when Filters::MarkerFilterRow has_markers ||= true filter_row.flatten.map do |flattened| EventMarkersIndexFiltering.for_read_common(flattened, cursor) end else Utils.missing_implementation!(filter_row) end end union_builders(builders, cursor, mode: has_markers ? :union_distinct : :union_all, requires_limit: false) end |
.sql_builder_for_subscriptions(filter_rows) ⇒ PgEventstore::SQLBuilder
171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 171 def sql_builder_for_subscriptions(filter_rows) if filter_rows.empty? sql_builder = EventsGlobalIndexFiltering.new.tap(&:for_subscription).to_sql_builder return finalize_subscription_builders([sql_builder]) end builders = [] grouped = filter_rows.group_by(&:class) grouped.each_value do |rows| case rows.first when Filters::FilterRow filtering = EventsGlobalIndexFiltering.new filtering.for_subscription rows.each(&filtering.method(:add_filter_row)) builders.push(filtering.to_sql_builder) when Filters::MarkerFilterRow filtering = EventMarkersIndexFiltering.new filtering.for_subscription rows.each(&filtering.method(:add_marker_filter_row)) builders.push(filtering.to_sql_builder) else Utils.missing_implementation!(rows.first) end end finalize_subscription_builders(builders) end |
Instance Method Details
#add_global_position_direction(direction) ⇒ void
This method returns an undefined value.
249 250 251 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 249 def add_global_position_direction(direction) @sql_builder.order("global_position #{SQL_DIRECTIONS[direction]}") end |
#add_limit(limit) ⇒ void
This method returns an undefined value.
255 256 257 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 255 def add_limit(limit) @sql_builder.limit(limit || DEFAULT_LIMIT) end |
#to_sql_builder ⇒ Object
243 244 245 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 243 def to_sql_builder @sql_builder end |
#to_table_name ⇒ Object
239 240 241 |
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 239 def to_table_name '' end |