Class: PgEventstore::Client

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

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(config) ⇒ Client

@param config

Parameters:



14
15
16
# File 'lib/pg_eventstore/client.rb', line 14

def initialize(config)
  @config = config
end

Instance Attribute Details

#configConfig

Returns the value of attribute config.

Returns:



10
11
12
# File 'lib/pg_eventstore/client.rb', line 10

def config
  @config
end

Instance Method Details

#append_to_stream(stream, events_or_event, options: {}, middlewares: nil) ⇒ PgEventstore::Event+

Append the event or multiple events to the stream. This operation is atomic, meaning that no other event can be appended by parallel process between the given events.

Parameters:

Options Hash (options:):

  • :expected_revision (Integer)

    expected stream revision

  • :expected_revision (Symbol)

    provide one of next values: :any, :no_stream or :stream_exists

  • :expected_revision (Hash<String, Integer>, Hash<String, Symbol>)

    -to- map. Allows to define expected stream revision at the given event type. Useful when implementing Dynamic Consistency Boundaries

  • :expected_revision (Hash<String, Hash>)

    -to- map. Allows to define expected stream revision at the given event type, scoped to the specific marker(s). Useful when implementing Dynamic Consistency Boundaries

Returns:

Raises:



34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
# File 'lib/pg_eventstore/client.rb', line 34

def append_to_stream(stream, events_or_event, options: {}, middlewares: nil)
  Utils.assert_node_role!(config, Config::NodeRole.writable)
  middlewares = self.middlewares(middlewares)
  event_modifier = Commands::EventModifiers::PrepareRegularEvent.new(EventSerializer.new(middlewares))
  queries = Queries.new(
    partitions: partition_queries,
    events: event_queries,
    transactions: transaction_queries,
    events_global_index: events_global_index_queries,
    streams_global_index: streams_global_index_queries,
    event_subscription_positions: event_subscription_position_queries,
    event_markers: event_marker_queries,
    index_filtering: index_filtering_queries
  )
  result = Commands::Append.new(queries).call(
    stream, *events_or_event, event_modifier:, deserializer: event_deserializer(middlewares), options:
  )
  events_or_event.is_a?(Array) ? result : result.first
end

#connectionPgEventstore::Connection



267
268
269
# File 'lib/pg_eventstore/client.rb', line 267

def connection
  PgEventstore.connection(config.name)
end

#event_deserializer(middlewares) ⇒ PgEventstore::EventDeserializer

Parameters:

Returns:

  • (PgEventstore::EventDeserializer)


288
289
290
# File 'lib/pg_eventstore/client.rb', line 288

def event_deserializer(middlewares)
  EventDeserializer.new(middlewares, config.event_class_resolver)
end

#event_marker_queriesPgEventstore::EventMarkerQueries

Returns:

  • (PgEventstore::EventMarkerQueries)


308
309
310
# File 'lib/pg_eventstore/client.rb', line 308

def event_marker_queries
  EventMarkerQueries.new(connection, QueryStrategy::Foreground.new(connection))
end

#event_queriesPgEventstore::EventQueries

Returns:

  • (PgEventstore::EventQueries)


282
283
284
# File 'lib/pg_eventstore/client.rb', line 282

def event_queries
  EventQueries.new(connection)
end

#event_subscription_position_queriesPgEventstore::EventSubscriptionPositionQueries

Returns:

  • (PgEventstore::EventSubscriptionPositionQueries)


303
304
305
# File 'lib/pg_eventstore/client.rb', line 303

def event_subscription_position_queries
  EventSubscriptionPositionQueries.new(connection)
end

#events_global_index_queriesPgEventstore::EventsGlobalIndexQueries

Returns:

  • (PgEventstore::EventsGlobalIndexQueries)


293
294
295
# File 'lib/pg_eventstore/client.rb', line 293

def events_global_index_queries
  EventsGlobalIndexQueries.new(connection, QueryStrategy::Foreground.new(connection))
end

#index_filtering_queriesPgEventstore::IndexFilteringQueries



313
314
315
# File 'lib/pg_eventstore/client.rb', line 313

def index_filtering_queries
  IndexFilteringQueries.new(connection, QueryStrategy::Foreground.new(connection))
end

Links event from one stream into another stream. You can later access it by providing :resolve_link_tos option when reading from a stream. Only existing events can be linked.

Parameters:

  • stream (PgEventstore::Stream)
  • events_or_event (PgEventstore::Event, Array<PgEventstore::Event>)
  • options (Hash) (defaults to: {})
  • middlewares (Array<Symbol>) (defaults to: [])

    provide a list of middleware names to use. Defaults to empty array, meaning no middlewares will be applied to the "link" event

Options Hash (options:):

  • :expected_revision (Integer)

    provide your own revision number

  • :expected_revision (Symbol)

    provide one of next values: :any, :no_stream or :stream_exists

Returns:

Raises:



227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
# File 'lib/pg_eventstore/client.rb', line 227

def link_to(stream, events_or_event, options: {}, middlewares: [])
  Utils.assert_node_role!(config, Config::NodeRole.writable)
  middlewares = self.middlewares(middlewares)
  event_modifier = Commands::EventModifiers::PrepareLinkEvent.new(
    partition_queries, EventSerializer.new(middlewares)
  )
  queries = Queries.new(
    partitions: partition_queries,
    events: event_queries,
    transactions: transaction_queries,
    events_global_index: events_global_index_queries,
    streams_global_index: streams_global_index_queries,
    event_subscription_positions: event_subscription_position_queries,
    event_markers: event_marker_queries,
    index_filtering: index_filtering_queries
  )
  result = Commands::LinkTo.new(queries).call(
    stream, *events_or_event, event_modifier:, deserializer: event_deserializer(middlewares), options:
  )
  events_or_event.is_a?(Array) ? result : result.first
end

#middlewares(middlewares = nil) ⇒ Array<PgEventstore::Middleware>

Parameters:

  • middlewares (Array<Symbol>, nil) (defaults to: nil)

Returns:



260
261
262
263
264
# File 'lib/pg_eventstore/client.rb', line 260

def middlewares(middlewares = nil)
  return config.middlewares.values unless middlewares

  config.middlewares.slice(*middlewares).values
end

#multiple(read_only: false) { ... } ⇒ Object

Allows you to make several different commands atomic by wrapping then into a block. Order of events, produced by multiple commands, belonging to different streams - is unbreakable. So, if you append event1 to stream1 and event2 to stream2 using this method, then thet appear in the same order in the "all" stream. Example:

PgEventstore.client.multiple do
PgEventstore.client.read(...)
PgEventstore.client.append_to_stream(...)
PgEventstore.client.append_to_stream(...)
end

Parameters:

  • read_only (Boolean) (defaults to: false)

    whether transaction is read-only. Running mutation queries within read-only transaction will result in exception

  • read_only: (Boolean) (defaults to: false)

Yields:

Yield Returns:

  • (Object)

Returns:

  • the result of the given block



67
68
69
# File 'lib/pg_eventstore/client.rb', line 67

def multiple(read_only: false, &)
  Commands::Multiple.new(Queries.new(transactions: transaction_queries)).call(read_only:, &)
end

#partition_queriesPgEventstore::PartitionQueries

Returns:

  • (PgEventstore::PartitionQueries)


272
273
274
# File 'lib/pg_eventstore/client.rb', line 272

def partition_queries
  PartitionQueries.new(connection)
end

#read(stream, options: {}, middlewares: nil) ⇒ Array<PgEventstore::Event>

Read events from the specific stream or from "all" stream.

Parameters:

  • stream (PgEventstore::Stream)
  • options (Hash) (defaults to: {})

    request options

  • middlewares (Array<Symbol>, nil) (defaults to: nil)

    provide a list of middleware names to override a config's middlewares

  • options: (::Hash[untyped, untyped]) (defaults to: {})
  • middlewares: (::Array[::Symbol], nil) (defaults to: nil)

Options Hash (options:):

  • :direction (String)

    read direction. Allowed values are "Forwards", "Backwards", "asc", "desc", :asc, :desc

  • :from_revision (Integer)

    a starting revision number. Use this option when stream name is a normal stream name

  • :to_revision (Integer)

    ending revision number. Use this option when stream name is a normal stream name

  • :from_position (Integer)

    a starting global position number. Use this option when reading from "all" stream

  • :to_position (Integer)

    ending global position number. Use this option when reading from "all" stream

  • :max_count (Integer)

    max number of events to return in one response. Defaults to config.max_count

  • :resolve_link_tos (Boolean)

    When using projections to create new events you can set whether the generated events are pointers to existing events. Setting this option to true tells PgEventstore to return the original event instead a link event.

  • :filter (Hash)

    provide it to filter events. You can filter by: stream, event type and event markers. Filtering by stream is only available when reading from "all" stream. Basically, you can mix almost all combinations of stream attributes, event types and markers. That are some limitations, though. Please refer to docs/reading_events.md for details. Here are some examples:

    # Filtering by stream's context. This will return all events which #context is 'User'
    PgEventstore.client.read(
    PgEventstore::Stream.all_stream,
    options: { filter: { streams: [{ context: 'User' }] } }
    )
    
    # Filtering by several stream's contexts. This will return all events which #context is either 'User' or
    # 'Profile'
    PgEventstore.client.read(
    PgEventstore::Stream.all_stream,
    options: { filter: { streams: [{ context: 'User' }, { context: 'Profile' }] } }
    )
    
    # Filtering by a mix of specific stream and a context. This will return all events which #context is 'User' or
    # events belonging to the stream with { context: 'Profile', stream_name: 'ProfileFields', stream_id: '123' }
    PgEventstore.client.read(
    PgEventstore::Stream.all_stream,
    options: {
      filter: {
        streams: [
          { context: 'User' },
          { context: 'Profile', stream_name: 'ProfileFields', stream_id: '123' }
        ]
      }
    }
    )
    
    # Filtering a mix of context and event type
    PgEventstore.client.read(
    PgEventstore::Stream.all_stream,
    options: { filter: { streams: [{ context: 'User' }], event_types: ['MyAwesomeEvent'] } }
    )
    
    # Filtering by specific event when reading from the specific stream
    PgEventstore.client.read(stream, options: { filter: { event_types: ['MyAwesomeEvent'] } })
    
    # Filtering a specific stream by event markers
    PgEventstore.client.read(stream, options: { filter: { event_types: [{ markers: ['foo'] }] } })
    
    # Filtering a specific stream by event type with markers and event types
    PgEventstore.client.read(
    stream, options: { filter: { event_types: [{ type: 'Foo', markers: ['foo'] }, 'Bar'] } }
    )
    
    # Filtering all events by markers
    PgEventstore.client.read(
    PgEventstore::Stream.all_stream, options: { filter: { event_types: [{ markers: ['foo'] }] } }
    )

Returns:

Raises:



144
145
146
147
148
149
150
151
152
153
154
# File 'lib/pg_eventstore/client.rb', line 144

def read(stream, options: {}, middlewares: nil)
  queries = Queries.new(
    index_filtering: index_filtering_queries,
    streams_global_index: streams_global_index_queries
  )
  Commands::Read.new(queries).call(
    stream,
    deserializer: event_deserializer(middlewares(middlewares)),
    options: { max_count: config.max_count }.merge(options)
  )
end

#read_grouped(stream, options: {}, middlewares: nil) ⇒ Array<PgEventstore::Event>

Takes a stream, event types filter and returns most recent(or very first - depending on :direction option) events, one of each given type. The result size is almost always less than or equal to event types list size, so passing :max_count option does not take any effect. In case if event of same type appears in different context/stream name - it will be counted as a different event, thus, may appear several times in the result, scoped to each context and stream name in the result. Useful when implementing Dynamic Consistency Boundaries.

Parameters:

  • stream (PgEventstore::Stream)
  • options (Hash) (defaults to: {})

    request options

  • middlewares (Array<Symbol>, nil) (defaults to: nil)
  • options: (Hash[untyped, untyped]) (defaults to: {})
  • middlewares: (::Array[::Symbol], nil) (defaults to: nil)

Returns:

See Also:

  • for the detailed docs


184
185
186
187
188
189
190
191
192
# File 'lib/pg_eventstore/client.rb', line 184

def read_grouped(stream, options: {}, middlewares: nil)
  queries = Queries.new(
    index_filtering: index_filtering_queries,
    streams_global_index: streams_global_index_queries
  )
  Commands::ReadGrouped.new(queries).call(
    stream, deserializer: event_deserializer(middlewares(middlewares)), options:
  )
end

#read_paginated(stream, options: {}, middlewares: nil) ⇒ Enumerator

Returns enumerator will yield ArrayPgEventstore::Event.

Parameters:

  • stream (PgEventstore::Stream)
  • options (Hash) (defaults to: {})

    request options

  • middlewares (Array<Symbol>, nil) (defaults to: nil)

Returns:

See Also:

  • for the detailed docs


161
162
163
164
165
166
167
168
169
170
171
172
# File 'lib/pg_eventstore/client.rb', line 161

def read_paginated(stream, options: {}, middlewares: nil)
  cmd_class = stream.system? ? Commands::SystemStreamReadPaginated : Commands::RegularStreamReadPaginated
  queries = Queries.new(
    index_filtering: index_filtering_queries,
    streams_global_index: streams_global_index_queries
  )
  cmd_class.new(queries).call(
    stream,
    deserializer: event_deserializer(middlewares(middlewares)),
    options: { max_count: config.max_count }.merge(options)
  )
end

#read_streams(options: {}) ⇒ Array<PgEventstore::Stream>

Parameters:

  • options (Hash) (defaults to: {})

    request options

  • options: (::Hash[untyped, untyped]) (defaults to: {})

Options Hash (options:):

  • :direction (String)

    read direction. Allowed values are "Forwards", "Backwards", "asc", "desc", :asc, :desc

  • :from_position (Integer)

    a starting global position number

  • :max_count (Integer)

    max number of streams to return in one response. Defaults to config.max_count

Returns:



200
201
202
203
# File 'lib/pg_eventstore/client.rb', line 200

def read_streams(options: {})
  queries = Queries.new(streams_global_index: streams_global_index_queries)
  Commands::ReadStreams.new(queries).call(options: { max_count: config.max_count }.merge(options))
end

#read_streams_paginated(options: {}) ⇒ Enumerator

Returns yields ArrayPgEventstore::Stream.

Parameters:

  • options (Hash) (defaults to: {})

    request options

  • options: (::Hash[untyped, untyped]) (defaults to: {})

Options Hash (options:):

  • :direction (String)

    read direction. Allowed values are "Forwards", "Backwards", "asc", "desc", :asc, :desc

  • :from_position (Integer)

    a starting global position number

  • :max_count (Integer)

    max number of streams to yield. Defaults to config.max_count

Returns:



211
212
213
214
# File 'lib/pg_eventstore/client.rb', line 211

def read_streams_paginated(options: {})
  queries = Queries.new(streams_global_index: streams_global_index_queries)
  Commands::ReadStreamsPaginated.new(queries).call(options: { max_count: config.max_count }.merge(options))
end

#stream_revision(stream) ⇒ Integer

Parameters:

Returns:

  • (Integer)


251
252
253
254
# File 'lib/pg_eventstore/client.rb', line 251

def stream_revision(stream)
  queries = Queries.new(streams_global_index: streams_global_index_queries)
  Commands::StreamRevision.new(queries).call(stream)
end

#streams_global_index_queriesPgEventstore::StreamsGlobalIndexQueries

Returns:

  • (PgEventstore::StreamsGlobalIndexQueries)


298
299
300
# File 'lib/pg_eventstore/client.rb', line 298

def streams_global_index_queries
  StreamsGlobalIndexQueries.new(connection, QueryStrategy::Foreground.new(connection))
end

#transaction_queriesPgEventstore::TransactionQueries

Returns:

  • (PgEventstore::TransactionQueries)


277
278
279
# File 'lib/pg_eventstore/client.rb', line 277

def transaction_queries
  TransactionQueries.new(connection)
end