Compare commits

...
6 changed files with 71 additions and 26 deletions
@@ -313,6 +313,11 @@ const applyDelayedTrigger = () => {
? buildTriggerCondition('message_type', trigger.messageType)
: buildTriggerCondition('status', triggerStatus.value),
];
// A private note is an outgoing message, so without this an internal note would read as a reply
// and arm the customer-unresponsive wait. Incoming messages are never private.
if (trigger.messageType === 'outgoing') {
conditions.push(buildTriggerCondition('private_note', [false]));
}
if (triggerInboxes.value.length) {
conditions.push(
buildTriggerCondition(
@@ -50,7 +50,10 @@ class AutomationRules::ProcessPendingExecutionJob < ApplicationJob
).perform.present?
end
# Marked before the actions run: a row that dies here stays `executing`, which no sweep reclaims,
# so a message/email/webhook is never sent twice. Everything up to this point is still retryable.
def execute(pending_execution)
pending_execution.update!(status: :executing)
AutomationRules::ActionService.new(
pending_execution.automation_rule,
pending_execution.account,
@@ -12,7 +12,8 @@ class AutomationRules::TriggerPendingExecutionsJob < ApplicationJob
rows = AutomationRulePendingExecution.sweepable.for_enabled_accounts.order(:due_at).limit(sweep_limit).to_a
rows.each { |row| AutomationRules::ProcessPendingExecutionJob.perform_later(row) }
log_summary(enqueued: rows.size, capped: rows.size >= sweep_limit, purged: purged, started_at: started_at)
log_summary(enqueued: rows.size, capped: rows.size >= sweep_limit, purged: purged,
abandoned: AutomationRulePendingExecution.abandoned.count, started_at: started_at)
end
private
@@ -25,8 +26,9 @@ class AutomationRules::TriggerPendingExecutionsJob < ApplicationJob
(InstallationConfig.find_by(name: 'AUTOMATION_PENDING_EXECUTIONS_SWEEP_LIMIT')&.value || DEFAULT_SWEEP_LIMIT).to_i
end
def log_summary(enqueued:, capped:, purged:, started_at:)
summary = { event: 'completed', enqueued: enqueued, capped: capped, purged: purged, duration_ms: ((Time.current - started_at) * 1000).round }
def log_summary(enqueued:, capped:, purged:, abandoned:, started_at:)
summary = { event: 'completed', enqueued: enqueued, capped: capped, purged: purged, abandoned: abandoned,
duration_ms: ((Time.current - started_at) * 1000).round }
Rails.logger.info("[AutomationRules::TriggerPendingExecutionsJob] #{summary.to_json}")
end
end
+13 -3
View File
@@ -46,9 +46,8 @@ class AutomationRuleListener < BaseListener
account = conversation.account
changed_attributes = event.data[:changed_attributes]
return unless rule_present?(event_name, account)
rules = current_account_rules(event_name, account)
rules = conversation_rules(event_name, account)
return if rules.blank?
rules.each do |rule|
conditions_match = ::AutomationRules::ConditionsFilterService.new(rule, conversation, { changed_attributes: changed_attributes }).perform
@@ -56,6 +55,17 @@ class AutomationRuleListener < BaseListener
end
end
# A delayed conversation rule reads as "the conversation has been in this status for N minutes",
# so a conversation created in that status must arm it too. Creation never dispatches
# CONVERSATION_UPDATED, and both paths key the episode on the same status_changed_at, so a later
# update arming the same episode is deduped by the unique index.
def conversation_rules(event_name, account)
rules = current_account_rules(event_name, account)
return rules unless event_name == 'conversation_created'
rules + current_account_rules('conversation_updated', account).where.not(execution_delay: nil)
end
# Delayed rules record a pending execution instead of acting; the sweep re-checks and
# runs them at due time. Flag off means no arming and no immediate fallback — a delayed
# message silently becoming instant is worse than skipping.
+2 -1
View File
@@ -131,7 +131,8 @@ class AutomationRule < ApplicationRecord
end
def discard_stale_pending_executions
# armed = pending + stale processing, which the sweep would otherwise reclaim.
# armed = pending + processing, the rows the sweep would otherwise still run. Rows already
# executing are left alone: their actions are in flight and cannot be called back.
pending_executions.armed.delete_all
end
+43 -19
View File
@@ -38,13 +38,20 @@ class AutomationRulePendingExecution < ApplicationRecord
belongs_to :account
belongs_to :message, optional: true
enum status: { pending: 0, processing: 1, executed: 2, skipped: 3 }
# `processing` is claimed but not yet acting, so it is safe to reclaim and retry. `executing`
# means the actions are running: a row that dies there is never replayed, because the actions
# are customer-facing (messages, emails, webhooks) and repeating them is worse than dropping them.
enum status: { pending: 0, processing: 1, executed: 2, skipped: 3, executing: 4 }
# Rows a sweep should hand to a worker: due pending rows, plus processing rows whose lock went stale.
scope :sweepable, lambda {
pending.where(due_at: ..Time.current).or(processing.where(updated_at: ...STALE_PROCESSING_TIMEOUT.ago))
}
# Rows whose worker died mid-action. Nothing reclaims them; the sweep only counts them so a
# crash that strands customer-facing actions is visible instead of silent.
scope :abandoned, -> { executing.where(updated_at: ...STALE_PROCESSING_TIMEOUT.ago) }
# Non-terminal rows still bound to fire (a stale processing row is reclaimed by the sweep).
scope :armed, -> { where(status: [statuses[:pending], statuses[:processing]]) }
@@ -53,6 +60,11 @@ class AutomationRulePendingExecution < ApplicationRecord
scope :for_enabled_accounts, -> { joins(:account).merge(Account.feature_delayed_automations) }
def self.schedule(rule:, conversation:, message: nil)
# status_changed_at is only written from this feature onwards, so a conversation that predates it
# has no status clock. Anchoring on created_at would make every old conversation instantly
# overdue and fire on the next sweep; leave them for their next status change to arm.
return if message.nil? && conversation.status_changed_at.blank?
key = arm_episode_key_for(conversation, message)
anchor = arm_anchor_for(conversation, message)
create!(
@@ -68,22 +80,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
@@ -92,7 +109,7 @@ class AutomationRulePendingExecution < ApplicationRecord
# the timestamps the episode keys track.
def self.arm_anchor_for(conversation, message)
if message.nil?
conversation.status_changed_at.presence || conversation.created_at
conversation.status_changed_at
elsif message.incoming?
conversation.waiting_since.presence || message.created_at
else
@@ -122,7 +139,7 @@ class AutomationRulePendingExecution < ApplicationRecord
if message.nil?
# Sub-second precision so a resolve→reopen inside one second still ends the episode.
# Integer microseconds (not a float) so an in-memory arm and a DB-reloaded fire agree.
"status:#{microsecond_stamp(conversation.status_changed_at.presence || conversation.created_at)}"
"status:#{microsecond_stamp(conversation.status_changed_at)}"
elsif message.incoming?
# waiting_since is cleared on agent/bot reply, so a reply invalidates this episode. Strict
# here: at fire time a nil waiting_since means the agent replied (episode ended).
@@ -162,6 +179,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 +199,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