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 dead for 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

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