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
-
.boot_from_schema!(schema_path:, registry:) ⇒ Object
-
.configuration ⇒ Object
-
.configure {|configuration| ... } ⇒ Object
-
.definition_port ⇒ Object
-
.discovered_schema_paths(port = definition_port) ⇒ Object
-
.dispatch(event) ⇒ Object
-
.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
-
.enriched_metadata(call_site_metadata) ⇒ Object
-
.evaluated_metadata_defaults ⇒ Object
-
.file_schema_registry(schema_path: configuration.schema_path) ⇒ Object
-
.handler_registry ⇒ Object
-
.process(event) ⇒ Object
-
.processor_registry ⇒ Object
-
.register_definition_publisher!(port = definition_port) ⇒ Object
-
.register_handler(handler, process_types:) ⇒ Object
-
.register_processor(name, processor) ⇒ Object
-
.register_slice!(schema_path:) ⇒ Object
-
.reset_handlers! ⇒ Object
-
.reset_processors! ⇒ Object
-
.schema_sources(port = definition_port) ⇒ Object
-
.unrouted_events ⇒ Object
-
.validate_rules! ⇒ Object
Class Attribute Details
.processing_rules ⇒ Object
68
69
70
|
# File 'lib/event_engine.rb', line 68
def processing_rules
@processing_rules ||= ProcessingRules.load(configuration.rules_path)
end
|
.schema_registry ⇒ Object
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
|
.configuration ⇒ Object
25
26
27
|
# File 'lib/event_engine.rb', line 25
def configuration
@configuration ||= Configuration.new
end
|
29
30
31
|
# File 'lib/event_engine.rb', line 29
def configure
yield(configuration)
end
|
.definition_port ⇒ Object
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] = enriched_metadata(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
|
114
115
116
117
118
119
|
# File 'lib/event_engine.rb', line 114
def enriched_metadata(call_site_metadata)
defaults = evaluated_metadata_defaults
return call_site_metadata if defaults.nil?
defaults.merge(call_site_metadata || {})
end
|
121
122
123
124
125
126
127
128
129
|
# File 'lib/event_engine.rb', line 121
def evaluated_metadata_defaults
callable = configuration.metadata_defaults
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_registry ⇒ Object
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_registry ⇒ Object
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_events ⇒ Object
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
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
|