Class: PgEventstore::QueryBuilders::IndexBasedEventsFiltering

Inherits:
Object
  • Object
show all
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 =

Returns:

  • (Integer)
1_000
ABSTRACT_TABLE_NAME =

Returns:

  • (String)
'events_index'
ReadCursor =

Returns:

  • (:StreamCursor cursor)

Constants included from BasicFiltering

BasicFiltering::SQL_DIRECTIONS

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initializeIndexBasedEventsFiltering

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

Parameters:

  • filter_rows (Array<PgEventstore::QueryBuilders::Filters::FilterRow, PgEventstore::QueryBuilders::Filters::MarkerFilterRow>)

Returns:

  • (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.from_options({})
  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

Parameters:

Returns:

  • (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.from_stream_and_options(
    stream, { direction: :desc, max_count: 1 }
  )
  filters_collection = Filters::Collection.from_stream_and_options(
    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

Parameters:

Returns:

  • (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.from_stream_and_options(
    stream, { direction: :desc, max_count: 1 }
  )
  filters_collection = Filters::Collection.from_stream_and_options(
    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

Parameters:

Returns:

  • (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.from_stream_and_options(
    stream, { direction: :desc, max_count: 1 }
  )
  filters_collection = Filters::Collection.from_stream_and_options(
    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

Parameters:

  • filter_rows (Array<PgEventstore::QueryBuilders::Filters::FilterRow, PgEventstore::QueryBuilders::Filters::MarkerFilterRow>)
  • cursor (PgEventstore::QueryBuilders::ReadCursor::StreamCursor)

Returns:

  • (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

Parameters:

  • filter_rows (Array<PgEventstore::QueryBuilders::Filters::FilterRow, PgEventstore::QueryBuilders::Filters::MarkerFilterRow>)
  • cursor (PgEventstore::QueryBuilders::ReadCursor::StreamCursor)

Returns:

  • (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

Parameters:

  • filter_rows (Array<PgEventstore::QueryBuilders::Filters::FilterRow, PgEventstore::QueryBuilders::Filters::MarkerFilterRow>)

Returns:

  • (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.

Parameters:

  • direction (String, Symbol, nil)


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.

Parameters:

  • limit (Integer, nil)


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_builderObject



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_nameObject



239
240
241
# File 'lib/pg_eventstore/query_builders/index_based_events_filtering.rb', line 239

def to_table_name
  ''
end