Module: EventEngine

Defined in:
lib/event_engine.rb,
lib/event_engine/event.rb,
lib/event_engine/railtie.rb,
lib/event_engine/version.rb,
lib/event_engine/rules_file.rb,
lib/event_engine/event_schema.rb,
lib/event_engine/catalog_entry.rb,
lib/event_engine/configuration.rb,
lib/event_engine/event_builder.rb,
lib/event_engine/schema_registry.rb,
lib/event_engine/handler_registry.rb,
lib/event_engine/processing_rules.rb,
lib/event_engine/processor_registry.rb,
lib/event_engine/processor_resolver.rb,
lib/event_engine/invalid_rules_error.rb,
lib/event_engine/definition_publisher.rb,
lib/event_engine/unrouted_events_error.rb,
lib/event_engine/schema_catalog_builder.rb,
lib/event_engine/unroutable_event_error.rb,
lib/event_engine/event_schema_json_loader.rb,
lib/event_engine/unregistered_processor_error.rb

Defined Under Namespace

Classes: CatalogEntry, Configuration, DefinitionPublisher, Event, EventBuilder, EventSchema, EventSchemaJsonLoader, HandlerRegistry, InvalidRulesError, ProcessingRules, ProcessorRegistry, ProcessorResolver, Railtie, RulesFile, SchemaCatalogBuilder, SchemaRegistry, UnregisteredProcessorError, UnroutableEventError, UnroutedEventsError

Constant Summary collapse

VERSION =
"0.2.0"

Class Attribute Summary collapse

Class Method Summary collapse

Class Attribute Details

.processing_rulesObject



68
69
70
# File 'lib/event_engine.rb', line 68

def processing_rules
  @processing_rules ||= ProcessingRules.load(configuration.rules_path)
end

.schema_registryObject



41
42
43
# File 'lib/event_engine.rb', line 41

def schema_registry
  @schema_registry ||= SchemaRegistry.new
end

Class Method Details

.boot_from_schema!(schema_path:, registry:) ⇒ Object



162
163
164
165
166
167
168
169
170
171
# File 'lib/event_engine.rb', line 162

def boot_from_schema!(schema_path:, registry:)
  event_schema = EventSchemaJsonLoader.load(schema_path)

  registry.reset!
  registry.load_from_schema!(event_schema)

  self.schema_registry = registry

  event_schema
end

.configurationObject



25
26
27
# File 'lib/event_engine.rb', line 25

def configuration
  @configuration ||= Configuration.new
end

.configure {|configuration| ... } ⇒ Object

Yields:



29
30
31
# File 'lib/event_engine.rb', line 29

def configure
  yield(configuration)
end

.definition_portObject



110
111
112
# File 'lib/event_engine.rb', line 110

def definition_port
  ::EventEngine::Definition if defined?(::EventEngine::Definition)
end

.discovered_schema_paths(port = definition_port) ⇒ Object



98
99
100
101
102
# File 'lib/event_engine.rb', line 98

def discovered_schema_paths(port = definition_port)
  return [] unless port.respond_to?(:pack_schema_paths)

  port.pack_schema_paths
end

.dispatch(event) ⇒ Object



139
140
141
# File 'lib/event_engine.rb', line 139

def dispatch(event)
  handler_registry.dispatch(event)
end

.emit(event_name, inputs:, domain: nil, event_version: nil, occurred_at: nil, metadata: nil, idempotency_key: nil, aggregate_type: nil, aggregate_id: nil, aggregate_version: nil) ⇒ Object



47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
# File 'lib/event_engine.rb', line 47

def emit(event_name, inputs:, domain: nil, event_version: nil, occurred_at: nil,
         metadata: nil, idempotency_key: nil, aggregate_type: nil,
         aggregate_id: nil, aggregate_version: nil)
  schema = schema_registry.schema(event_name, version: event_version, domain: domain)

  attrs = EventBuilder.build(schema: schema, data: inputs)
  attrs[:occurred_at] = occurred_at || Time.current
  attrs[:metadata] = ()
  attrs[:idempotency_key] = idempotency_key || SecureRandom.uuid
  attrs[:aggregate_type] = aggregate_type
  attrs[:aggregate_id] = aggregate_id
  attrs[:aggregate_version] = aggregate_version
  attrs[:process_type] = processing_rules.for(event_name: schema.event_name, pack: schema.domain)
  attrs[:subject] = schema.subject
  attrs[:domain] = schema.domain

  event = Event.new(**attrs)
  process(event)
  dispatch(event)
end

.enriched_metadata(call_site_metadata) ⇒ Object



114
115
116
117
118
119
# File 'lib/event_engine.rb', line 114

def ()
  defaults = 
  return  if defaults.nil?

  defaults.merge( || {})
end

.evaluated_metadata_defaultsObject



121
122
123
124
125
126
127
128
129
# File 'lib/event_engine.rb', line 121

def 
  callable = configuration.
  return nil unless callable

  callable.call
rescue => error
  configuration.logger&.error("EventEngine metadata_defaults raised: #{error.message}")
  nil
end

.file_schema_registry(schema_path: configuration.schema_path) ⇒ Object



184
185
186
187
188
189
# File 'lib/event_engine.rb', line 184

def file_schema_registry(schema_path: configuration.schema_path)
  loaded = EventSchemaJsonLoader.load(schema_path)
  registry = SchemaRegistry.new
  registry.load_from_schema!(loaded)
  registry
end

.handler_registryObject



33
34
35
# File 'lib/event_engine.rb', line 33

def handler_registry
  @handler_registry ||= HandlerRegistry.new
end

.process(event) ⇒ Object



143
144
145
146
147
148
149
150
151
152
# File 'lib/event_engine.rb', line 143

def process(event)
  resolver = ProcessorResolver.new(processing_rules)
  return event unless resolver.routes?

  name = resolver.resolve(event)
  processor = processor_registry.fetch(name) || raise(UnregisteredProcessorError.new(name, event))

  processor.call(event)
  event
end

.processor_registryObject



37
38
39
# File 'lib/event_engine.rb', line 37

def processor_registry
  @processor_registry ||= ProcessorRegistry.new
end

.register_definition_publisher!(port = definition_port) ⇒ Object



104
105
106
107
108
# File 'lib/event_engine.rb', line 104

def register_definition_publisher!(port = definition_port)
  return nil unless port.respond_to?(:publisher=)

  port.publisher = DefinitionPublisher.new
end

.register_handler(handler, process_types:) ⇒ Object



131
132
133
# File 'lib/event_engine.rb', line 131

def register_handler(handler, process_types:)
  handler_registry.register(handler, process_types: process_types)
end

.register_processor(name, processor) ⇒ Object



135
136
137
# File 'lib/event_engine.rb', line 135

def register_processor(name, processor)
  processor_registry.register(name, processor)
end

.register_slice!(schema_path:) ⇒ Object



173
174
175
176
177
178
179
180
181
182
# File 'lib/event_engine.rb', line 173

def register_slice!(schema_path:)
  slice = EventSchemaJsonLoader.load(schema_path)

  schema_registry.load_from_schema!(EventSchema.new) unless schema_registry.loaded?
  slice.event_schema.schemas_by_event.each_value do |versions|
    versions.each_value { |schema| schema_registry.register(schema) }
  end

  schema_registry
end

.reset_handlers!Object



154
155
156
# File 'lib/event_engine.rb', line 154

def reset_handlers!
  handler_registry.clear!
end

.reset_processors!Object



158
159
160
# File 'lib/event_engine.rb', line 158

def reset_processors!
  processor_registry.clear!
end

.schema_sources(port = definition_port) ⇒ Object



91
92
93
94
95
96
# File 'lib/event_engine.rb', line 91

def schema_sources(port = definition_port)
  configured = configuration.publisher_schema_paths
  return configured if configured.any?

  discovered_schema_paths(port)
end

.unrouted_eventsObject



83
84
85
86
87
88
89
# File 'lib/event_engine.rb', line 83

def unrouted_events
  return [] unless processing_rules.any?

  schema_registry.events.reject do |event_name|
    processing_rules.for(event_name: event_name, pack: schema_registry.latest_for(event_name).domain)
  end
end

.validate_rules!Object

Raises:



74
75
76
77
78
79
80
81
# File 'lib/event_engine.rb', line 74

def validate_rules!
  missing = processing_rules.processor_names.reject { |name| processor_registry.fetch(name) }
  raise InvalidRulesError, missing if missing.any?

  raise UnroutedEventsError, unrouted_events if unrouted_events.any?

  true
end