Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ef5d2704eb | ||
|
|
c59dac2fcf | ||
|
|
6f26f0e34b | ||
|
|
51e154790c | ||
|
|
7a2578e814 | ||
|
|
c0f47f1319 |
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user