Class: PgEventstore::IndexFilteringQueries

Inherits:
Object
  • Object
show all
Defined in:
lib/pg_eventstore/queries/index_filtering_queries.rb,
sig/pg_eventstore/queries/index_filtering_queries.rbs

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(connection, query_strategy) ⇒ IndexFilteringQueries

Returns a new instance of IndexFilteringQueries.

Parameters:



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

#connectionConnection

Returns the value of attribute connection.

Returns:



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

Parameters:

Returns:

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

Parameters:

  • repr (Symbol, nil) (defaults to: nil)
  • pg_result (PG::Result)
  • repr: (Symbol, nil) (defaults to: nil)

Returns:



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>

Parameters:

  • filters_collection (PgEventstore::QueryBuilders::Filters::Collection)

Returns:

  • (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 expand_event_types(filters_collection)
  filter_rows = filters_collection.collection
  to_expand = filter_rows.select(&:ambiguous_event_type?)
  rest_filter_rows = filter_rows - to_expand
  partition_builders = to_expand.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)
  expanded_partitions_list = @query_strategy.exec_params(*final_partition_builder.to_exec_params)
  expanded_partitions_list = expanded_partitions_list.group_by { _1['event_type'] }
  expanded_filter_rows =
    to_expand.flat_map do |filter_row|
      case filter_row
      when QueryBuilders::Filters::FilterRow
        filter_row.event_type_filters.flat_map do |event_type_filter|
          related_partitions =
            if event_type_filter.prefix?
              expanded_partitions_list.select do |event_type, _|
                event_type.start_with?(event_type_filter.event_type)
              end.values.flatten
            else
              expanded_partitions_list[event_type_filter.event_type]
            end
          next [] unless related_partitions

          if filter_row.stream_filter
            related_partitions = related_partitions.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
          related_partitions.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
        related_partitions = expanded_partitions_list[filter_row.marker_filter.event_type]
        next [] unless related_partitions

        if filter_row.stream_filter
          related_partitions = related_partitions.select do |attrs|
            attrs['context'] == filter_row.stream_filter.context
          end
        end
        related_partitions.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 = expanded_filter_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>

Parameters:

  • filters_collection (PgEventstore::QueryBuilders::Filters::Collection)
  • cursor (PgEventstore::QueryBuilders::ReadCursor::StreamCursor)

Returns:



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?)
      expand_event_types(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>

Parameters:

  • filters_collection (PgEventstore::QueryBuilders::Filters::Collection)
  • cursor (PgEventstore::QueryBuilders::ReadCursor::StreamCursor)

Returns:



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?)
      expand_event_types(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_queriesPgEventstore::PartitionQueries

Returns:

  • (PgEventstore::PartitionQueries)


87
88
89
# File 'lib/pg_eventstore/queries/index_filtering_queries.rb', line 87

def partition_queries
  PartitionQueries.new(connection)
end