feat: scheduler fairness [AI-159] (#14425)
# Pull Request Template ## Description Better scheduling and queueing mechanics for document auto-sync - add jitter plan wise for document sync - move auto-sync documents to purgeable queue ## Type of change Please delete options that are not relevant. - [x] New feature (non-breaking change which adds functionality) ## How Has This Been Tested? Please describe the tests that you ran to verify your changes. Provide instructions so we can reproduce. Please also list any relevant details for your test configuration. locally tested and with specs ## Checklist: - [x] My code follows the style guidelines of this project - [x] I have performed a self-review of my code - [x] I have commented on my code, particularly in hard-to-understand areas - [ ] I have made corresponding changes to the documentation - [x] My changes generate no new warnings - [x] I have added tests that prove my fix is effective or that my feature works - [x] New and existing unit tests pass locally with my changes - [x] Any dependent changes have been merged and published in downstream modules --------- Co-authored-by: Sivin Varghese <64252451+iamsivin@users.noreply.github.com> Co-authored-by: iamsivin <iamsivin@gmail.com> Co-authored-by: Muhsin Keloth <muhsinkeramam@gmail.com> Co-authored-by: Sony Mathew <sony@chatwoot.com> Co-authored-by: Vishnu Narayanan <iamwishnu@gmail.com>
This commit is contained in:
@@ -75,7 +75,6 @@ class Api::V1::Accounts::Captain::BulkActionsController < Api::V1::Accounts::Bas
|
||||
Current.account.captain_documents.where(id: params[:ids]).find_each(batch_size: 100) do |document|
|
||||
next unless document.syncable?
|
||||
next unless document.available?
|
||||
next if document.sync_in_progress?
|
||||
|
||||
document.update!(
|
||||
sync_status: :syncing,
|
||||
|
||||
@@ -37,7 +37,6 @@ class Api::V1::Accounts::Captain::DocumentsController < Api::V1::Accounts::BaseC
|
||||
def sync
|
||||
return render_could_not_create_error(I18n.t('captain.documents.sync_not_supported_for_pdf')) unless @document.syncable?
|
||||
return render_could_not_create_error(I18n.t('captain.documents.sync_only_available_documents')) unless @document.available?
|
||||
return render_could_not_create_error(I18n.t('captain.documents.sync_already_in_progress')) if @document.sync_in_progress?
|
||||
|
||||
@document.update!(
|
||||
sync_status: :syncing,
|
||||
|
||||
@@ -48,9 +48,7 @@ class Captain::Documents::PerformSyncJob < MutexApplicationJob
|
||||
return if document.pdf_document?
|
||||
|
||||
with_lock(lock_key(document), LOCK_TIMEOUT) do
|
||||
mark_sync_started(document)
|
||||
result = Captain::Documents::SyncService.new(document.reload).perform
|
||||
log_sync_outcome(document, result: result, duration_ms: duration_ms_since(start_time))
|
||||
perform_sync(document, start_time)
|
||||
end
|
||||
rescue LockAcquisitionError
|
||||
log_sync_outcome(document, result: :already_syncing)
|
||||
@@ -64,6 +62,12 @@ class Captain::Documents::PerformSyncJob < MutexApplicationJob
|
||||
|
||||
private
|
||||
|
||||
def perform_sync(document, start_time)
|
||||
mark_sync_started(document)
|
||||
result = Captain::Documents::SyncService.new(document.reload).perform
|
||||
log_sync_outcome(document, result: result, duration_ms: duration_ms_since(start_time))
|
||||
end
|
||||
|
||||
def log_sync_outcome(document, **fields)
|
||||
payload = {
|
||||
document_id: document.id,
|
||||
|
||||
@@ -1,21 +1,27 @@
|
||||
class Captain::Documents::ScheduleSyncsJob < ApplicationJob
|
||||
queue_as :scheduled_jobs
|
||||
|
||||
PER_ACCOUNT_HOURLY_CAP = 50
|
||||
GLOBAL_HOURLY_CAP = 1000
|
||||
DUE_DOCUMENT_BATCH_SIZE = PER_ACCOUNT_HOURLY_CAP * 2 # Inspite of skipping, we should at least reach the hourly cap
|
||||
DEFAULT_PER_ACCOUNT_BATCH_LIMIT = 50
|
||||
DEFAULT_GLOBAL_BATCH_LIMIT = 1000
|
||||
SYNC_STALE_TIMEOUT = Captain::Document::SYNC_STALE_TIMEOUT
|
||||
DAILY_SYNC_JITTER = 4.hours
|
||||
WEEKLY_SYNC_JITTER = 1.day
|
||||
MONTHLY_SYNC_JITTER = 4.days
|
||||
|
||||
def perform
|
||||
@remaining_global_capacity = GLOBAL_HOURLY_CAP
|
||||
def perform(plan_name = nil)
|
||||
@per_account_batch_limit = configured_sync_limit('CAPTAIN_DOCUMENT_AUTO_SYNC_PER_ACCOUNT_BATCH_LIMIT', DEFAULT_PER_ACCOUNT_BATCH_LIMIT)
|
||||
@global_batch_limit = configured_sync_limit('CAPTAIN_DOCUMENT_AUTO_SYNC_GLOBAL_BATCH_LIMIT', DEFAULT_GLOBAL_BATCH_LIMIT)
|
||||
@remaining_global_capacity = @global_batch_limit
|
||||
@plan_name = plan_name.to_s.downcase.presence
|
||||
sync_intervals = Enterprise::Account.captain_document_sync_intervals
|
||||
stats = { accounts_scanned: 0, accounts_enabled: 0, accounts_scheduled: 0, documents_enqueued: 0, documents_skipped: 0 }
|
||||
stats = { accounts_scanned: 0, accounts_enabled: 0, accounts_scheduled: 0, documents_enqueued: 0 }
|
||||
|
||||
Account.joins(:captain_documents).distinct.find_each(batch_size: 100) do |account|
|
||||
break if @remaining_global_capacity <= 0
|
||||
|
||||
stats[:accounts_scanned] += 1
|
||||
next unless account.feature_enabled?('captain_document_auto_sync')
|
||||
next unless account_in_selected_plan?(account)
|
||||
|
||||
stats[:accounts_enabled] += 1
|
||||
interval = account.captain_document_sync_interval(sync_intervals)
|
||||
@@ -24,7 +30,6 @@ class Captain::Documents::ScheduleSyncsJob < ApplicationJob
|
||||
stats[:accounts_scheduled] += 1
|
||||
result = enqueue_due_documents(account, interval)
|
||||
stats[:documents_enqueued] += result[:enqueued]
|
||||
stats[:documents_skipped] += result[:skipped]
|
||||
end
|
||||
|
||||
log_scheduler_summary(stats)
|
||||
@@ -33,92 +38,85 @@ class Captain::Documents::ScheduleSyncsJob < ApplicationJob
|
||||
private
|
||||
|
||||
def enqueue_due_documents(account, interval)
|
||||
per_account_limit = [PER_ACCOUNT_HOURLY_CAP, @remaining_global_capacity].min
|
||||
result = { enqueued: 0, skipped: 0 }
|
||||
skipped_document_ids = []
|
||||
per_account_limit = [@per_account_batch_limit, @remaining_global_capacity].min
|
||||
result = { enqueued: 0 }
|
||||
|
||||
while result[:enqueued] < per_account_limit
|
||||
|
||||
documents = due_documents(account, interval, skipped_document_ids).limit(DUE_DOCUMENT_BATCH_SIZE).to_a
|
||||
break if documents.empty?
|
||||
|
||||
documents.each do |document|
|
||||
break if result[:enqueued] >= per_account_limit
|
||||
|
||||
process_due_document(document, result, skipped_document_ids)
|
||||
end
|
||||
due_documents(account, interval).limit(per_account_limit).to_a.each do |document|
|
||||
process_due_document(document, interval, result)
|
||||
end
|
||||
|
||||
result
|
||||
end
|
||||
|
||||
def process_due_document(document, result, skipped_document_ids)
|
||||
return unless document.syncable?
|
||||
def process_due_document(document, interval, result)
|
||||
sync_execution_delay = sync_jitter(interval)
|
||||
|
||||
# Reserve the sync slot before enqueueing so later scheduler runs skip this document while the job is queued.
|
||||
unless reserve_sync_slot(document)
|
||||
result[:skipped] += 1
|
||||
skipped_document_ids << document.id
|
||||
return
|
||||
end
|
||||
|
||||
Captain::Documents::PerformSyncJob.perform_later(document)
|
||||
Captain::Documents::PerformSyncJob.set(queue: :purgable, wait: sync_execution_delay).perform_later(document)
|
||||
@remaining_global_capacity -= 1
|
||||
result[:enqueued] += 1
|
||||
end
|
||||
|
||||
def due_documents(account, interval, skipped_document_ids)
|
||||
def due_documents(account, interval)
|
||||
syncing = Captain::Document.sync_statuses[:syncing]
|
||||
synced = Captain::Document.sync_statuses[:synced]
|
||||
failed = Captain::Document.sync_statuses[:failed]
|
||||
stale_cutoff = SYNC_STALE_TIMEOUT.ago
|
||||
# The scheduler runs at predictable plan windows. Use a wider due window so
|
||||
# jittered executions do not miss the next window just because they finished later.
|
||||
sync_due_before = due_window(interval).ago
|
||||
|
||||
documents = account.captain_documents.syncable.where(status: :available).where(
|
||||
'(sync_status = ? AND last_synced_at < ?) OR (sync_status = ? AND last_sync_attempted_at < ?) OR ' \
|
||||
'(sync_status = ? AND last_sync_attempted_at < ?)',
|
||||
synced, interval.ago, failed, interval.ago, syncing, SYNC_STALE_TIMEOUT.ago
|
||||
synced, sync_due_before, failed, sync_due_before, syncing, stale_cutoff
|
||||
)
|
||||
documents = documents.where.not(id: skipped_document_ids) if skipped_document_ids.present?
|
||||
documents.order(Arel.sql('last_sync_attempted_at ASC NULLS FIRST'), :id)
|
||||
end
|
||||
|
||||
def reserve_sync_slot(document)
|
||||
mark_sync_started(document)
|
||||
true
|
||||
rescue ActiveRecord::RecordInvalid => e
|
||||
log_document_skip(document, e)
|
||||
false
|
||||
def configured_sync_limit(config_key, default)
|
||||
configured_value = InstallationConfig.find_by(name: config_key)&.value
|
||||
limit = configured_value.to_s.to_i
|
||||
limit.positive? ? limit : default
|
||||
end
|
||||
|
||||
def log_document_skip(document, error)
|
||||
payload = {
|
||||
event: 'document_skipped',
|
||||
document_id: document.id,
|
||||
account_id: document.account_id,
|
||||
assistant_id: document.assistant_id,
|
||||
error_class: error.class.name,
|
||||
error_message: error.message,
|
||||
validation_errors: document.errors.full_messages
|
||||
}
|
||||
def account_in_selected_plan?(account)
|
||||
return true if @plan_name.blank?
|
||||
|
||||
Rails.logger.warn("[Captain::Documents::ScheduleSyncsJob] #{payload.to_json}")
|
||||
account_sync_plan(account) == @plan_name
|
||||
end
|
||||
|
||||
def account_sync_plan(account)
|
||||
plan = account.custom_attributes['plan_name']
|
||||
plan = 'enterprise' if plan.blank? && ChatwootApp.self_hosted_enterprise?
|
||||
plan.to_s.downcase.presence
|
||||
end
|
||||
|
||||
def sync_jitter(interval)
|
||||
jitter_window = if interval <= 1.day
|
||||
DAILY_SYNC_JITTER
|
||||
elsif interval <= 1.week
|
||||
WEEKLY_SYNC_JITTER
|
||||
else
|
||||
MONTHLY_SYNC_JITTER
|
||||
end
|
||||
|
||||
rand(0..jitter_window.to_i).seconds
|
||||
end
|
||||
|
||||
def due_window(interval)
|
||||
(interval.to_i / 2).seconds
|
||||
end
|
||||
|
||||
def log_scheduler_summary(stats)
|
||||
payload = {
|
||||
event: 'completed',
|
||||
plan_name: @plan_name,
|
||||
global_cap_hit: @remaining_global_capacity <= 0,
|
||||
per_account_batch_limit: @per_account_batch_limit,
|
||||
global_batch_limit: @global_batch_limit,
|
||||
remaining_global_capacity: @remaining_global_capacity
|
||||
}.merge(stats)
|
||||
|
||||
Rails.logger.info("[Captain::Documents::ScheduleSyncsJob] #{payload.to_json}")
|
||||
end
|
||||
|
||||
def mark_sync_started(document)
|
||||
document.update!(
|
||||
sync_status: :syncing,
|
||||
sync_step: nil,
|
||||
last_sync_error_code: nil,
|
||||
last_sync_attempted_at: Time.current
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
module Enterprise::Internal::TriggerDailyScheduledItemsJob
|
||||
def perform
|
||||
super
|
||||
|
||||
Captain::Documents::ScheduleSyncsJob.perform_later('enterprise')
|
||||
Captain::Documents::ScheduleSyncsJob.perform_later('business') if business_auto_sync_due?
|
||||
Captain::Documents::ScheduleSyncsJob.perform_later('startups') if startup_auto_sync_due?
|
||||
end
|
||||
|
||||
private
|
||||
|
||||
def business_auto_sync_due?
|
||||
Time.current.utc.sunday?
|
||||
end
|
||||
|
||||
def startup_auto_sync_due?
|
||||
Time.current.utc.day == 1
|
||||
end
|
||||
end
|
||||
@@ -1,7 +0,0 @@
|
||||
module Enterprise::Internal::TriggerHourlyScheduledItemsJob
|
||||
def perform
|
||||
super
|
||||
|
||||
Captain::Documents::ScheduleSyncsJob.perform_later
|
||||
end
|
||||
end
|
||||
Reference in New Issue
Block a user