fix(automation): serialize pending execution re-arm and leave live workers alone
This commit is contained in:
@@ -68,22 +68,27 @@ class AutomationRulePendingExecution < ApplicationRecord
|
||||
def self.rearm_or_advance_episode(rule, conversation, key, message, anchor)
|
||||
return unless message
|
||||
|
||||
row = find_by!(automation_rule_id: rule.id, conversation_id: conversation.id, episode_key: key)
|
||||
# Jobs can arrive out of order; only a strictly newer message advances or re-arms, so a late
|
||||
# older message can't pull due_at backwards and fire before the delay elapses.
|
||||
return unless message.id > row.message_id
|
||||
|
||||
due_at = rule.execution_delay.minutes.since(anchor)
|
||||
if row.condition_skipped?
|
||||
# A message episode key can recur (no new incoming reply) while conditions swing back into
|
||||
# match, so a later qualifying message re-arms the condition-only skip instead of dropping.
|
||||
row.update!(status: :pending, skip_reason: nil, due_at: due_at, message_id: message.id)
|
||||
elsif !row.terminal?
|
||||
# Track the newest qualifying message. Reply-chase advances due_at with each agent reply;
|
||||
# awaiting-agent keeps its first clock (its anchor is the stable waiting_since, so due_at is
|
||||
# unchanged). Re-anchoring a row still processing (its worker died mid-run) back to pending
|
||||
# also keeps a stale reclaim from firing the old clock instead of the latest one.
|
||||
row.update!(status: :pending, due_at: due_at, message_id: message.id)
|
||||
row = find_by!(automation_rule_id: rule.id, conversation_id: conversation.id, episode_key: key)
|
||||
# The lock (and the reload it does) makes the compare-and-write atomic. Two listeners racing on
|
||||
# the same episode would otherwise both read the old message_id and let whichever wrote last
|
||||
# win, so an older message could overwrite a newer one and pull due_at backwards.
|
||||
row.with_lock do
|
||||
# Jobs can arrive out of order; only a strictly newer message advances or re-arms, so a late
|
||||
# older message can't pull due_at backwards and fire before the delay elapses.
|
||||
next unless message.id > row.message_id
|
||||
|
||||
if row.condition_skipped?
|
||||
# A message episode key can recur (no new incoming reply) while conditions swing back into
|
||||
# match, so a later qualifying message re-arms the condition-only skip instead of dropping.
|
||||
row.update!(status: :pending, skip_reason: nil, due_at: due_at, message_id: message.id)
|
||||
elsif row.pending? || row.stale_processing?
|
||||
# Track the newest qualifying message. Reply-chase advances due_at with each agent reply;
|
||||
# awaiting-agent keeps its first clock (its anchor is the stable waiting_since, so due_at is
|
||||
# unchanged). Re-anchoring a row whose worker died mid-run back to pending also keeps a
|
||||
# stale reclaim from firing the old clock instead of the latest one.
|
||||
row.update!(status: :pending, due_at: due_at, message_id: message.id)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -162,6 +167,13 @@ class AutomationRulePendingExecution < ApplicationRecord
|
||||
self.class.episode_key_for(conversation, message) == episode_key
|
||||
end
|
||||
|
||||
# The claim renews updated_at, so a processing row past the timeout means its worker died. Only
|
||||
# then may a re-arm take the row back: pulling a live worker's row to pending would let the sweep
|
||||
# claim it and run the same actions alongside the worker still executing them.
|
||||
def stale_processing?
|
||||
processing? && updated_at < STALE_PROCESSING_TIMEOUT.ago
|
||||
end
|
||||
|
||||
def condition_skipped?
|
||||
skipped? && skip_reason == CONDITIONS_CHANGED_SKIP
|
||||
end
|
||||
@@ -175,6 +187,6 @@ class AutomationRulePendingExecution < ApplicationRecord
|
||||
def claimable?
|
||||
# due_at guard: a reply-chase reschedule can push due_at forward after this row was enqueued;
|
||||
# such a row must wait for a later sweep instead of firing early.
|
||||
(pending? && due_at <= Time.current) || (processing? && updated_at < STALE_PROCESSING_TIMEOUT.ago)
|
||||
(pending? && due_at <= Time.current) || stale_processing?
|
||||
end
|
||||
end
|
||||
|
||||
Reference in New Issue
Block a user