Module: Pgbus::EventBus::Publisher
- Defined in:
- lib/pgbus/event_bus/publisher.rb
Class Method Summary collapse
- .build_event_data(payload, routing_key: nil) ⇒ Object
- .publish(routing_key, payload, headers: nil, delay: 0) ⇒ Object
- .publish_later(routing_key, payload, delay:, headers: nil) ⇒ Object
-
.tag_current(event_data) ⇒ Object
Current attributes publish → handler (issue #431): when config.current_attributes is set, snapshot the publisher's persisted Current classes (Pgbus::CurrentAttributes.capture — same filters, serialization and allowlist gating as jobs) into the envelope under the same
pgbus_currentkey jobs use. -
.tag_fair_share(event_data, payload, routing_key:, headers:) ⇒ Object
Fair share for consumers (issue #427): when config.event_fair_share is set, hand the callable a Pgbus::Event (routing key, the payload object as passed to publish — not its serialized form — and headers) and merge the resolved key/weight into the envelope.
Class Method Details
.build_event_data(payload, routing_key: nil) ⇒ Object
75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 |
# File 'lib/pgbus/event_bus/publisher.rb', line 75 def build_event_data(payload, routing_key: nil) event_id = SecureRandom.uuid serialized_payload = if payload.respond_to?(:to_global_id) { "_global_id" => payload.to_global_id.to_s } elsif payload.is_a?(Hash) payload else { "value" => payload } end data = { "event_id" => event_id, "payload" => serialized_payload, "published_at" => Time.now.utc.iso8601(6) } data["routing_key"] = routing_key if routing_key data end |
.publish(routing_key, payload, headers: nil, delay: 0) ⇒ Object
10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 |
# File 'lib/pgbus/event_bus/publisher.rb', line 10 def publish(routing_key, payload, headers: nil, delay: 0) event_data = build_event_data(payload, routing_key: routing_key) event_data = tag_fair_share(event_data, payload, routing_key: routing_key, headers: headers) event_data = tag_current(event_data) if defined?(Pgbus::Testing) && !Pgbus::Testing.disabled? event = Pgbus::Event.new( event_id: event_data["event_id"], payload: event_data["payload"], published_at: event_data["published_at"] ? Time.parse(event_data["published_at"]) : nil, routing_key: routing_key, headers: headers, context: event_data[Pgbus::CurrentAttributes::METADATA_KEY] ) Pgbus::Testing.store.push_event(event) if Pgbus::Testing.inline? && delay.to_i <= 0 Pgbus::EventBus::Registry.instance.handlers_for(routing_key).each do |subscriber| Pgbus::CurrentAttributes.restore(event.context) { subscriber.handler_class.new.handle(event) } end end return event_data end Pgbus.client.publish_to_topic(routing_key, event_data, headers: headers, delay: delay) end |
.publish_later(routing_key, payload, delay:, headers: nil) ⇒ Object
39 40 41 |
# File 'lib/pgbus/event_bus/publisher.rb', line 39 def publish_later(routing_key, payload, delay:, headers: nil) publish(routing_key, payload, headers: headers, delay: delay) end |
.tag_current(event_data) ⇒ Object
Current attributes publish → handler (issue #431): when
config.current_attributes is set, snapshot the publisher's persisted
Current classes (Pgbus::CurrentAttributes.capture — same filters,
serialization and allowlist gating as jobs) into the envelope under
the same pgbus_current key jobs use. Shared with Outbox.publish_event
so the capture happens where Current is set, not at relay time.
Returns event_data itself when off or nothing is assigned.
68 69 70 71 72 73 |
# File 'lib/pgbus/event_bus/publisher.rb', line 68 def tag_current(event_data) captured = Pgbus::CurrentAttributes.capture return event_data unless captured event_data.merge(Pgbus::CurrentAttributes::METADATA_KEY => captured) end |
.tag_fair_share(event_data, payload, routing_key:, headers:) ⇒ Object
Fair share for consumers (issue #427): when config.event_fair_share is set, hand the callable a Pgbus::Event (routing key, the payload object as passed to publish — not its serialized form — and headers) and merge the resolved key/weight into the envelope. Shared with Outbox.publish_event so the key is resolved where the publisher's context (Current.*) exists and rides the outbox row to the bus. Returns event_data itself when off.
49 50 51 52 53 54 55 56 57 58 59 |
# File 'lib/pgbus/event_bus/publisher.rb', line 49 def tag_fair_share(event_data, payload, routing_key:, headers:) return event_data unless FairShare.event_enabled? event = Pgbus::Event.new( event_id: event_data["event_id"], payload: payload, routing_key: routing_key, headers: headers ) FairShare.(event, event_data) end |