diff --git a/app/models/automation_rule_pending_execution.rb b/app/models/automation_rule_pending_execution.rb index 3f287b201..9c574fedb 100644 --- a/app/models/automation_rule_pending_execution.rb +++ b/app/models/automation_rule_pending_execution.rb @@ -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