Module: OpenLoam::DurableEvents
- Defined in:
- lib/open_loam/durable_events.rb
Overview
Durable (persistent) event subscribers — the twin of OpenLoam::Events.subscribe, and the formal contract L-706 pins down.
THE CONTRACT
- Ephemeral — OpenLoam::Events.subscribe(&block): runs INLINE in the publisher's thread, synchronously, best-effort. An exception in the block PROPAGATES into whatever published the event. No persistence, no retry. Right for cheap in-process fan-out (the webhook dispatcher enqueues its own jobs this way).
- Persistent — OpenLoam::DurableEvents.register(...): publishing commits a
OpenLoam::EventDelivery row in the event's tenant, then hands the handler to a
background job. The handler's exception is CONTAINED in the job; delivery
is retried with backoff and, past MAX_ATTEMPTS, parked as
deadfor an operator (Admin::EventDeliveriesController).
GUARANTEE: at-least-once, UNORDERED. Handlers MUST be idempotent — a retry or the sweep may deliver the same event twice. Durability is of DELIVERY, not CAPTURE: an event whose process dies between the after_commit and the publish leaves no row and is lost, exactly as today. This feature makes what WAS published arrive; it does not resurrect what was never published.
SECURITY: a handler is resolved from THIS in-memory registry by key, populated
at boot from trusted code — never constantized from the stored row. If the key
is unknown at delivery time (handler removed since enqueue) the row is parked
dead; an arbitrary class is never executed off a DB value. (Same posture as
the scheduler's job_class allowlist.)
Constant Summary collapse
- MAX_ATTEMPTS =
5- BACKOFF =
Backoff (seconds) indexed by attempt number; the last value is the cap.
[ 0, 60, 300, 1800, 7200 ].freeze
- SWEEP_KEY =
"open_loam_event_redelivery_sweep".freeze
Class Method Summary collapse
-
.capture(event_name, payload) ⇒ Object
Persist a delivery row per matching durable subscriber, in the event's tenant, then nudge a job per row.
-
.deliver(delivery, now: Time.current) ⇒ Object
Run one delivery and advance the row's state.
- .handler_for(key) ⇒ Object
-
.redeliver_stuck(now: Time.current, limit: 500) ⇒ Object
THE durability guarantee: re-enqueue due-but-undelivered rows whose job was lost (worker crash, dropped message, async adapter racing the txn).
-
.register(key:, to:, call:) ⇒ Object
key: a stable identifier stored on every delivery row it produces.
- .registered ⇒ Object
- .reset_registry! ⇒ Object
- .run_delivery(delivery, now) ⇒ Object
-
.subscribe! ⇒ Object
Wired once from OpenLoam::Engine.
- .subscribers_for(event_name) ⇒ Object
Class Method Details
.capture(event_name, payload) ⇒ Object
Persist a delivery row per matching durable subscriber, in the event's tenant, then nudge a job per row. The ROW is the durable record; the job is only the accelerator (the sweep redelivers rows whose job was lost).
73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 |
# File 'lib/open_loam/durable_events.rb', line 73 def capture(event_name, payload) tenant = OpenLoam::Tenant.find_by(id: payload[:tenant_id]) return if tenant.nil? # nil-tenant events are not durably delivered (as Webhooks.dispatch) matches = subscribers_for(event_name) return if matches.empty? # Only JSON-safe primitives cross into the row/job — payloads are ids and # scalars by convention, never records (same rule as the webhook path). deliverable = payload.transform_keys(&:to_s) OpenLoam.as_tenant(tenant) do matches.each do |sub| delivery = OpenLoam::EventDelivery.create!( subscriber_key: sub[:key], event_name: event_name.to_s, payload: deliverable, status: "pending", attempts: 0 ) OpenLoam::EventDeliveryJob.perform_later(tenant.id, delivery.id) end end end |
.deliver(delivery, now: Time.current) ⇒ Object
Run one delivery and advance the row's state. Manual, ROW-STATE retries — deliberately NOT ActiveJob retry_on, which keeps retry state in the queue (the very thing this feature stops trusting). Safe to call twice on one row (at-least-once): a delivered/dead/not-yet-due row is a no-op.
101 102 103 104 105 106 107 108 109 |
# File 'lib/open_loam/durable_events.rb', line 101 def deliver(delivery, now: Time.current) return unless delivery.status == "pending" return if delivery.next_attempt_at && delivery.next_attempt_at > now # duplicate nudge, backoff not elapsed OpenLoam::Telemetry.span("durable_event_delivery", subscriber_key: delivery.subscriber_key, event_name: delivery.event_name) do run_delivery(delivery, now) end end |
.handler_for(key) ⇒ Object
57 58 59 60 |
# File 'lib/open_loam/durable_events.rb', line 57 def handler_for(key) entry = registry[key.to_s] entry && resolve(entry[:handler]) end |
.redeliver_stuck(now: Time.current, limit: 500) ⇒ Object
THE durability guarantee: re-enqueue due-but-undelivered rows whose job was lost (worker crash, dropped message, async adapter racing the txn). Tenant-scoped — called from the per-tenant sweep the engine registers, so no cross-tenant scan is needed.
137 138 139 140 141 142 143 144 |
# File 'lib/open_loam/durable_events.rb', line 137 def redeliver_stuck(now: Time.current, limit: 500) count = 0 OpenLoam::EventDelivery.due(now).limit(limit).find_each do |delivery| OpenLoam::EventDeliveryJob.perform_later(delivery.tenant_id, delivery.id) count += 1 end count end |
.register(key:, to:, call:) ⇒ Object
key: a stable identifier stored on every delivery row it produces. to: an event name ("billing.invoice.paid") or a domain prefix ("billing.") — same matching rule as Events/webhooks. call: a class (or anything) responding to .call(event_name, payload), or its name as a String. Resolved from this registry, never the row.
42 43 44 45 |
# File 'lib/open_loam/durable_events.rb', line 42 def register(key:, to:, call:) registry[key.to_s] = { key: key.to_s, pattern: to.to_s, handler: call } key.to_s end |
.registered ⇒ Object
47 |
# File 'lib/open_loam/durable_events.rb', line 47 def registered = registry.values |
.reset_registry! ⇒ Object
49 50 51 |
# File 'lib/open_loam/durable_events.rb', line 49 def reset_registry! @registry = {} end |
.run_delivery(delivery, now) ⇒ Object
111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 |
# File 'lib/open_loam/durable_events.rb', line 111 def run_delivery(delivery, now) handler = handler_for(delivery.subscriber_key) if handler.nil? delivery.update!(status: "dead", last_error: "no registered subscriber #{delivery.subscriber_key.inspect}") return end begin handler.call(delivery.event_name, delivery.payload_hash) delivery.update!(status: "delivered", delivered_at: now, last_error: nil) rescue StandardError => error attempts = delivery.attempts + 1 if attempts >= MAX_ATTEMPTS delivery.update!(status: "dead", attempts: attempts, last_error: format_error(error)) else delivery.update!(attempts: attempts, next_attempt_at: now + backoff_for(attempts), last_error: format_error(error)) end end end |
.subscribe! ⇒ Object
Wired once from OpenLoam::Engine. Idempotent: subscribing twice would persist every event twice.
66 67 68 |
# File 'lib/open_loam/durable_events.rb', line 66 def subscribe! @subscription ||= OpenLoam::Events.subscribe_all { |event_name, payload| capture(event_name, payload) } end |
.subscribers_for(event_name) ⇒ Object
53 54 55 |
# File 'lib/open_loam/durable_events.rb', line 53 def subscribers_for(event_name) registered.select { |entry| OpenLoam::Events.pattern_matches?(entry[:pattern], event_name) } end |