diff --git a/app/controllers/api/v1/accounts/assignment_policies_controller.rb b/app/controllers/api/v1/accounts/assignment_policies_controller.rb new file mode 100644 index 000000000..acc7c1186 --- /dev/null +++ b/app/controllers/api/v1/accounts/assignment_policies_controller.rb @@ -0,0 +1,79 @@ +# frozen_string_literal: true + +class Api::V1::Accounts::AssignmentPoliciesController < Api::V1::Accounts::BaseController + before_action :fetch_assignment_policy, only: [:show, :update, :destroy] + before_action :check_authorization + + def index + @assignment_policies = Current.account.assignment_policies.includes(:inboxes) + render json: { assignment_policies: serialize_assignment_policies(@assignment_policies) } + end + + def show + render json: { assignment_policy: serialize_assignment_policy(@assignment_policy) } + end + + def create + @assignment_policy = Current.account.assignment_policies.build(assignment_policy_params) + + if @assignment_policy.save + render json: { assignment_policy: serialize_assignment_policy(@assignment_policy) }, status: :created + else + render json: { errors: @assignment_policy.errors.full_messages }, status: :unprocessable_entity + end + end + + def update + if @assignment_policy.update(assignment_policy_params) + render json: { assignment_policy: serialize_assignment_policy(@assignment_policy) } + else + render json: { errors: @assignment_policy.errors.full_messages }, status: :unprocessable_entity + end + end + + def destroy + if @assignment_policy.destroy + head :ok + else + render json: { errors: @assignment_policy.errors.full_messages }, status: :unprocessable_entity + end + end + + private + + def fetch_assignment_policy + @assignment_policy = Current.account.assignment_policies.find(params[:id]) + end + + def assignment_policy_params + params.require(:assignment_policy).permit( + :name, :description, :assignment_order, :conversation_priority, + :fair_distribution_limit, :fair_distribution_window, :enabled + ) + end + + def serialize_assignment_policy(policy) + { + id: policy.id, + name: policy.name, + description: policy.description, + assignment_order: policy.assignment_order, + conversation_priority: policy.conversation_priority, + fair_distribution_limit: policy.fair_distribution_limit, + fair_distribution_window: policy.fair_distribution_window, + enabled: policy.enabled, + inbox_count: policy.inboxes.count, + inboxes: policy.inboxes.map { |inbox| { id: inbox.id, name: inbox.name } }, + created_at: policy.created_at, + updated_at: policy.updated_at + } + end + + def serialize_assignment_policies(policies) + policies.map { |policy| serialize_assignment_policy(policy) } + end + + def check_authorization + authorize(AssignmentPolicy) + end +end diff --git a/app/controllers/api/v1/accounts/inbox_assignment_policies_controller.rb b/app/controllers/api/v1/accounts/inbox_assignment_policies_controller.rb new file mode 100644 index 000000000..395cd9622 --- /dev/null +++ b/app/controllers/api/v1/accounts/inbox_assignment_policies_controller.rb @@ -0,0 +1,81 @@ +# frozen_string_literal: true + +class Api::V1::Accounts::InboxAssignmentPoliciesController < Api::V1::Accounts::BaseController + before_action :fetch_inbox + before_action :check_authorization + + def show + @inbox_assignment_policy = @inbox.inbox_assignment_policy + + if @inbox_assignment_policy + render json: { + inbox_assignment_policy: serialize_inbox_assignment_policy(@inbox_assignment_policy) + } + else + render json: { + inbox_assignment_policy: nil, + message: 'No assignment policy assigned to this inbox' + } + end + end + + def create + # Remove existing assignment if any + @inbox.inbox_assignment_policy&.destroy + + @assignment_policy = Current.account.assignment_policies.find(params[:assignment_policy_id]) + @inbox_assignment_policy = @inbox.build_inbox_assignment_policy(assignment_policy: @assignment_policy) + + if @inbox_assignment_policy.save + render json: { + inbox_assignment_policy: serialize_inbox_assignment_policy(@inbox_assignment_policy) + }, status: :created + else + render json: { errors: @inbox_assignment_policy.errors.full_messages }, status: :unprocessable_entity + end + end + + def destroy + @inbox_assignment_policy = @inbox.inbox_assignment_policy + + if @inbox_assignment_policy + if @inbox_assignment_policy.destroy + head :ok + else + render json: { errors: @inbox_assignment_policy.errors.full_messages }, status: :unprocessable_entity + end + else + render json: { error: 'No assignment policy found for this inbox' }, status: :not_found + end + end + + private + + def fetch_inbox + @inbox = Current.account.inboxes.find(params[:inbox_id]) + end + + def check_authorization + authorize(@inbox, :update?) + end + + def serialize_inbox_assignment_policy(inbox_assignment_policy) + { + id: inbox_assignment_policy.id, + inbox_id: inbox_assignment_policy.inbox_id, + assignment_policy_id: inbox_assignment_policy.assignment_policy_id, + assignment_policy: { + id: inbox_assignment_policy.assignment_policy.id, + name: inbox_assignment_policy.assignment_policy.name, + description: inbox_assignment_policy.assignment_policy.description, + assignment_order: inbox_assignment_policy.assignment_policy.assignment_order, + conversation_priority: inbox_assignment_policy.assignment_policy.conversation_priority, + fair_distribution_limit: inbox_assignment_policy.assignment_policy.fair_distribution_limit, + fair_distribution_window: inbox_assignment_policy.assignment_policy.fair_distribution_window, + enabled: inbox_assignment_policy.assignment_policy.enabled + }, + created_at: inbox_assignment_policy.created_at, + updated_at: inbox_assignment_policy.updated_at + } + end +end diff --git a/app/jobs/assignment_v2/assignment_job.rb b/app/jobs/assignment_v2/assignment_job.rb new file mode 100644 index 000000000..647021e47 --- /dev/null +++ b/app/jobs/assignment_v2/assignment_job.rb @@ -0,0 +1,36 @@ +# frozen_string_literal: true + +class AssignmentV2::AssignmentJob < ApplicationJob + queue_as :low + + def perform(inbox_id: nil, conversation_id: nil) + if conversation_id + assign_single_conversation(conversation_id) + elsif inbox_id + assign_inbox_conversations(inbox_id) + else + Rails.logger.error 'AssignmentV2::AssignmentJob: No inbox_id or conversation_id provided' + end + end + + private + + def assign_single_conversation(conversation_id) + conversation = Conversation.find_by(id: conversation_id) + return unless conversation + + service = AssignmentV2::AssignmentService.new(inbox: conversation.inbox) + service.perform_for_conversation(conversation) + end + + def assign_inbox_conversations(inbox_id) + inbox = Inbox.find_by(id: inbox_id) + return unless inbox + return unless inbox.assignment_v2_enabled? + + service = AssignmentV2::AssignmentService.new(inbox: inbox) + assigned_count = service.perform_bulk_assignment + + Rails.logger.info "AssignmentV2::AssignmentJob: Assigned #{assigned_count} conversations for inbox #{inbox_id}" + end +end diff --git a/app/models/assignment_policy.rb b/app/models/assignment_policy.rb new file mode 100644 index 000000000..b7a898569 --- /dev/null +++ b/app/models/assignment_policy.rb @@ -0,0 +1,103 @@ +# frozen_string_literal: true + +# == Schema Information +# +# Table name: assignment_policies +# +# id :bigint not null, primary key +# assignment_order :integer default("round_robin"), not null +# conversation_priority :integer default("earliest_created"), not null +# description :text +# enabled :boolean default(TRUE), not null +# fair_distribution_limit :integer default(10), not null +# fair_distribution_window :integer default(3600), not null +# name :string(255) not null +# created_at :datetime not null +# updated_at :datetime not null +# account_id :bigint not null +# +# Indexes +# +# index_assignment_policies_on_account_id (account_id) +# index_assignment_policies_on_enabled (enabled) +# unique_assignment_policy_name_per_account (account_id,name) UNIQUE +# + +class AssignmentPolicy < ApplicationRecord + include AccountCacheRevalidator + + # Enums + enum assignment_order: { round_robin: 0, balanced: 1 } + enum conversation_priority: { earliest_created: 0, longest_waiting: 1 } + + # Associations + belongs_to :account + has_many :inbox_assignment_policies, dependent: :destroy + has_many :inboxes, through: :inbox_assignment_policies + + # Validations + validates :name, presence: true, uniqueness: { scope: :account_id } + validates :name, length: { maximum: 255 } + validates :description, length: { maximum: 1000 } + validates :fair_distribution_limit, presence: true, numericality: { greater_than: 0, less_than_or_equal_to: 100 } + validates :fair_distribution_window, presence: true, numericality: { greater_than: 60, less_than_or_equal_to: 86_400 } + validates :assignment_order, inclusion: { in: assignment_orders.keys } + validates :conversation_priority, inclusion: { in: conversation_priorities.keys } + + # Validate balanced assignment is only available for enterprise + validate :validate_balanced_assignment_enterprise_only + + # Server-side validation to prevent bypass + before_save :enforce_enterprise_features + + # Scopes + scope :enabled, -> { where(enabled: true) } + scope :disabled, -> { where(enabled: false) } + + # Callbacks + after_update_commit :clear_assignment_caches + after_destroy :clear_assignment_caches + + def can_use_balanced_assignment? + account.feature_enabled?(:enterprise_agent_capacity) if account.respond_to?(:feature_enabled?) + end + + def webhook_data + { + id: id, + name: name, + description: description, + assignment_order: assignment_order, + conversation_priority: conversation_priority, + fair_distribution_limit: fair_distribution_limit, + fair_distribution_window: fair_distribution_window, + enabled: enabled + } + end + + private + + def validate_balanced_assignment_enterprise_only + return unless balanced? + return if can_use_balanced_assignment? + + errors.add(:assignment_order, 'Balanced assignment is only available for enterprise accounts') + end + + def enforce_enterprise_features + # Server-side enforcement to prevent API bypass + return unless balanced? && !can_use_balanced_assignment? + + # Force to round_robin if enterprise not available + self.assignment_order = 'round_robin' + Rails.logger.warn("Assignment V2: Forced assignment_order to round_robin for non-enterprise account #{account_id}") + end + + def clear_assignment_caches + # Clear Redis caches when policy is updated + Rails.cache.delete("assignment_v2:policy:#{id}") + inboxes.find_each do |inbox| + Rails.cache.delete("assignment_v2:inbox_policy:#{inbox.id}") + end + end +end diff --git a/app/models/concerns/assignment_v2_feature_flag.rb b/app/models/concerns/assignment_v2_feature_flag.rb new file mode 100644 index 000000000..d2263f290 --- /dev/null +++ b/app/models/concerns/assignment_v2_feature_flag.rb @@ -0,0 +1,37 @@ +# frozen_string_literal: true + +module AssignmentV2FeatureFlag + extend ActiveSupport::Concern + + def assignment_v2_enabled? + # Check account-level feature flag + return false unless feature_enabled_for_account? + + # Check for any inbox-level overrides + return false if inbox_level_override_disabled? + + # Check system-wide killswitch + !GlobalConfig.get('assignment_v2_disabled', false) + end + + private + + def feature_enabled_for_account? + config = GlobalConfig.get('assignment_v2') + return false unless config&.dig('enabled') + + # If no account allowlist, enable for all + allowed_accounts = config['accounts'] + return true if allowed_accounts.blank? + + # Check if account is in allowlist + allowed_accounts.include?(id) + end + + def inbox_level_override_disabled? + return false unless respond_to?(:id) + + # Allow per-inbox disabling during migration + GlobalConfig.get('assignment_v2_disabled_inboxes', []).include?(id) + end +end diff --git a/app/models/concerns/auto_assignment_handler.rb b/app/models/concerns/auto_assignment_handler.rb index de5fdae7d..1cd2c2c37 100644 --- a/app/models/concerns/auto_assignment_handler.rb +++ b/app/models/concerns/auto_assignment_handler.rb @@ -14,11 +14,18 @@ module AutoAssignmentHandler return unless conversation_status_changed_to_open? return unless should_run_auto_assignment? - ::AutoAssignment::AgentAssignmentService.new(conversation: self, allowed_agent_ids: inbox.member_ids_with_assignment_capacity).perform + if inbox.assignment_v2_enabled? + # Use Assignment V2 system + AssignmentV2::AssignmentJob.perform_later(conversation_id: id) + else + # Use legacy assignment system + ::AutoAssignment::AgentAssignmentService.new(conversation: self, allowed_agent_ids: inbox.member_ids_with_assignment_capacity).perform + end end def should_run_auto_assignment? - return false unless inbox.enable_auto_assignment? + # Check auto assignment is enabled (either legacy or v2) + return false unless inbox.auto_assignment_enabled? # run only if assignee is blank or doesn't have access to inbox assignee.blank? || inbox.members.exclude?(assignee) diff --git a/app/models/inbox.rb b/app/models/inbox.rb index 1c898ba7f..e5348867e 100644 --- a/app/models/inbox.rb +++ b/app/models/inbox.rb @@ -44,6 +44,11 @@ class Inbox < ApplicationRecord include Avatarable include OutOfOffisable include AccountCacheRevalidator + include AssignmentV2FeatureFlag + include InboxAgentAvailability + include InboxChannelTypes + include InboxWebhooks + include InboxNameSanitization # Not allowing characters: validates :name, presence: true @@ -72,6 +77,10 @@ class Inbox < ApplicationRecord has_many :webhooks, dependent: :destroy_async has_many :hooks, dependent: :destroy_async, class_name: 'Integrations::Hook' + # Assignment V2 associations + has_one :inbox_assignment_policy, dependent: :destroy + has_one :assignment_policy, through: :inbox_assignment_policy + enum sender_name_type: { friendly: 0, professional: 1 } after_destroy :delete_round_robin_agents @@ -97,109 +106,167 @@ class Inbox < ApplicationRecord update_account_cache end - # Sanitizes inbox name for balanced email provider compatibility - # ALLOWS: /'._- and Unicode letters/numbers/emojis - # REMOVES: Forbidden chars (\<>@") + spam-trigger symbols (!#$%&*+=?^`{|}~) - def sanitized_name - return default_name_for_blank_name if name.blank? - - sanitized = apply_sanitization_rules(name) - sanitized.blank? && email? ? display_name_from_email : sanitized - end - - def sms? - channel_type == 'Channel::Sms' - end - - def facebook? - channel_type == 'Channel::FacebookPage' - end - - def instagram? - (facebook? || instagram_direct?) && channel.instagram_id.present? - end - - def instagram_direct? - channel_type == 'Channel::Instagram' - end - - def web_widget? - channel_type == 'Channel::WebWidget' - end - - def api? - channel_type == 'Channel::Api' - end - - def email? - channel_type == 'Channel::Email' - end - - def twilio? - channel_type == 'Channel::TwilioSms' - end - - def twitter? - channel_type == 'Channel::TwitterProfile' - end - - def whatsapp? - channel_type == 'Channel::Whatsapp' - end - def assignable_agents (account.users.where(id: members.select(:user_id)) + account.administrators).uniq end - def active_bot? - agent_bot_inbox&.active? || hooks.where(app_id: %w[dialogflow], - status: 'enabled').count.positive? + # Assignment V2 methods + def assignment_v2_enabled? + account.assignment_v2_enabled? && assignment_policy.present? && assignment_policy.enabled? end - def inbox_type - channel.name - end - - def webhook_data - { - id: id, - name: name - } - end - - def callback_webhook_url - case channel_type - when 'Channel::TwilioSms' - "#{ENV.fetch('FRONTEND_URL', nil)}/twilio/callback" - when 'Channel::Sms' - "#{ENV.fetch('FRONTEND_URL', nil)}/webhooks/sms/#{channel.phone_number.delete_prefix('+')}" - when 'Channel::Line' - "#{ENV.fetch('FRONTEND_URL', nil)}/webhooks/line/#{channel.line_channel_id}" - when 'Channel::Whatsapp' - "#{ENV.fetch('FRONTEND_URL', nil)}/webhooks/whatsapp/#{channel.phone_number}" + def auto_assignment_enabled? + if assignment_v2_enabled? + assignment_policy.present? && assignment_policy.enabled? + else + enable_auto_assignment? end end - def member_ids_with_assignment_capacity - members.ids + # Returns inbox members who are available for assignment + # This method performs all filtering upfront at the database level for optimal performance + # + # Filters applied: + # 1. Online status - Only agents marked as 'online' in OnlineStatusTracker + # 2. Capacity limits (Enterprise) - Agents who haven't reached their conversation limit + # 3. Rate limiting - Agents who haven't exceeded rate limits (when implemented) + # 4. User exclusions - Specific users can be excluded (e.g., for reassignment) + # + # @param options [Hash] Additional filter options + # @option options [Boolean] :check_capacity (true) Whether to check capacity limits + # @option options [Boolean] :check_rate_limits (false) Whether to check rate limits + # @option options [Array] :exclude_user_ids Users to exclude from results + # + # @return [ActiveRecord::Relation] Available inbox members with preloaded users + # + # @example Get all available agents + # inbox.available_agents + # + # @example Get available agents excluding specific users + # inbox.available_agents(exclude_user_ids: [1, 2, 3]) + # + # @example Get available agents without capacity check (faster but less accurate) + # inbox.available_agents(check_capacity: false) + def available_agents(options = {}) + options = { check_capacity: true }.merge(options) + + # Get online agent IDs + online_agent_ids = fetch_online_agent_ids + return inbox_members.none if online_agent_ids.empty? + + # Base query - only online agents + scope = build_online_agents_scope(online_agent_ids) + + # Apply filters + apply_agent_filters(scope, options) end private - def default_name_for_blank_name - email? ? display_name_from_email : '' + def build_online_agents_scope(online_agent_ids) + inbox_members + .joins(:user) + .where(users: { id: online_agent_ids }) + .includes(:user) end - def apply_sanitization_rules(name) - name.gsub(/[\\<>@"!#$%&*+=?^`{|}~:;]/, '') # Remove forbidden chars - .gsub(/[\x00-\x1F\x7F]/, ' ') # Replace control chars with spaces - .gsub(/\A[[:punct:]]+|[[:punct:]]+\z/, '') # Remove leading/trailing punctuation - .gsub(/\s+/, ' ') # Normalize spaces - .strip + def apply_agent_filters(scope, options) + # Exclude specific users if requested + scope = scope.where.not(users: { id: options[:exclude_user_ids] }) if options[:exclude_user_ids].present? + + # Apply capacity filtering for enterprise accounts + scope = filter_by_capacity(scope) if options[:check_capacity] && enterprise_capacity_enabled? + + # Apply rate limiting if implemented + scope = filter_by_rate_limits(scope) if options[:check_rate_limits] && defined?(AssignmentV2::RateLimiter) + + # Exclude agents who are on leave + scope = filter_agents_on_leave(scope) if options[:exclude_on_leave] != false + + scope end - def display_name_from_email - channel.email.split('@').first.parameterize.titleize + def fetch_online_agent_ids + OnlineStatusTracker.get_available_users(account_id) + .select { |_key, value| value.eql?('online') } + .keys + .map(&:to_i) + end + + def enterprise_capacity_enabled? + defined?(Enterprise) && + account.custom_attributes&.dig('enterprise_features', 'capacity_management').present? + end + + def filter_by_capacity(inbox_members_scope) + return inbox_members_scope unless capacity_check_required? + + assignment_counts = fetch_assignment_counts + + inbox_members_scope.select do |inbox_member| + agent_has_capacity?(inbox_member, assignment_counts) + end + end + + def capacity_check_required? + defined?(Enterprise::InboxCapacityLimit) && + account.account_users.joins(:agent_capacity_policy).exists? + end + + def fetch_assignment_counts + conversations + .where(status: :open) + .where.not(assignee_id: nil) + .group(:assignee_id) + .count + end + + def agent_has_capacity?(inbox_member, assignment_counts) + user = inbox_member.user + account_user = account.account_users.find_by(user: user) + + return true unless account_user&.agent_capacity_policy_id + + capacity_limit = fetch_capacity_limit(account_user.agent_capacity_policy_id) + return true unless capacity_limit&.conversation_limit + + current_count = assignment_counts[user.id] || 0 + current_count < capacity_limit.conversation_limit + end + + def fetch_capacity_limit(policy_id) + Enterprise::InboxCapacityLimit + .where(agent_capacity_policy_id: policy_id) + .find_by(inbox_id: id) + end + + def filter_by_rate_limits(inbox_members_scope) + # Filter out agents who have exceeded rate limits + return inbox_members_scope unless assignment_policy&.enabled? + + # Get IDs of inbox members within rate limits + valid_inbox_member_ids = inbox_members_scope.find_each.filter_map do |inbox_member| + rate_limiter = AssignmentV2::RateLimiter.new(inbox: self, user: inbox_member.user) + inbox_member.id if rate_limiter.within_limits? + end + + # Return filtered scope maintaining ActiveRecord relation + inbox_members_scope.where(id: valid_inbox_member_ids) + end + + def filter_agents_on_leave(inbox_members_scope) + # Get account users on active leave + account_user_ids_on_leave = account.account_users + .joins(:leaves) + .where(leaves: { status: 'approved' }) + .where('leaves.start_date <= ? AND leaves.end_date >= ?', Date.current, Date.current) + .pluck(:id) + + return inbox_members_scope if account_user_ids_on_leave.empty? + + # Exclude inbox members whose account_users are on leave + user_ids_on_leave = account.account_users.where(id: account_user_ids_on_leave).pluck(:user_id) + inbox_members_scope.where.not(user_id: user_ids_on_leave) end def dispatch_create_event diff --git a/app/models/inbox_assignment_policy.rb b/app/models/inbox_assignment_policy.rb new file mode 100644 index 000000000..2540fadf1 --- /dev/null +++ b/app/models/inbox_assignment_policy.rb @@ -0,0 +1,67 @@ +# frozen_string_literal: true + +# == Schema Information +# +# Table name: inbox_assignment_policies +# +# id :bigint not null, primary key +# created_at :datetime not null +# updated_at :datetime not null +# assignment_policy_id :bigint not null +# inbox_id :bigint not null +# +# Indexes +# +# index_inbox_assignment_policies_on_assignment_policy_id (assignment_policy_id) +# index_inbox_assignment_policies_on_inbox_id (inbox_id) +# unique_inbox_assignment_policy (inbox_id) UNIQUE +# + +class InboxAssignmentPolicy < ApplicationRecord + include AccountCacheRevalidator + + # Associations + belongs_to :inbox + belongs_to :assignment_policy + + # Validations + validates :inbox_id, uniqueness: true + validate :inbox_belongs_to_same_account + + # Delegations + delegate :account, to: :inbox + delegate :name, :description, :assignment_order, :conversation_priority, + :fair_distribution_limit, :fair_distribution_window, :enabled?, + to: :assignment_policy, prefix: :policy + + # Callbacks + after_commit :clear_inbox_cache + + # Scopes + scope :enabled, -> { joins(:assignment_policy).where(assignment_policies: { enabled: true }) } + scope :disabled, -> { joins(:assignment_policy).where(assignment_policies: { enabled: false }) } + + def webhook_data + { + id: id, + inbox_id: inbox_id, + assignment_policy_id: assignment_policy_id, + policy: assignment_policy.webhook_data + } + end + + private + + def inbox_belongs_to_same_account + return unless inbox && assignment_policy + + return if inbox.account_id == assignment_policy.account_id + + errors.add(:inbox, 'must belong to the same account as the assignment policy') + end + + def clear_inbox_cache + Rails.cache.delete("assignment_v2:inbox_policy:#{inbox_id}") + update_account_cache + end +end diff --git a/app/policies/assignment_policy_policy.rb b/app/policies/assignment_policy_policy.rb new file mode 100644 index 000000000..7a1f32ff8 --- /dev/null +++ b/app/policies/assignment_policy_policy.rb @@ -0,0 +1,23 @@ +# frozen_string_literal: true + +class AssignmentPolicyPolicy < ApplicationPolicy + def index? + @account_user.administrator? + end + + def show? + @account_user.administrator? + end + + def create? + @account_user.administrator? + end + + def update? + @account_user.administrator? + end + + def destroy? + @account_user.administrator? + end +end diff --git a/app/services/assignment_v2/assignment_service.rb b/app/services/assignment_v2/assignment_service.rb new file mode 100644 index 000000000..bb7260593 --- /dev/null +++ b/app/services/assignment_v2/assignment_service.rb @@ -0,0 +1,112 @@ +# frozen_string_literal: true + +class AssignmentV2::AssignmentService + pattr_initialize [:inbox!] + + def perform_for_conversation(conversation) + return false unless can_assign?(conversation) + + agent = find_agent_for_conversation(conversation) + return false unless agent + + assign_conversation_to_agent(conversation, agent) + end + + def perform_bulk_assignment(limit: 50) + return 0 unless assignment_enabled? + + conversations = unassigned_conversations(limit) + assigned_count = 0 + + conversations.find_each do |conversation| + assigned_count += 1 if perform_for_conversation(conversation) + end + + assigned_count + end + + private + + def policy + @policy ||= inbox.assignment_policy + end + + def assignment_enabled? + policy&.enabled? + end + + def can_assign?(conversation) + assignment_enabled? && + conversation.status == 'open' && + conversation.assignee_id.nil? + end + + def find_agent_for_conversation(_conversation) + available_agents = inbox.available_agents(check_rate_limits: true) + + if available_agents.empty? + log_no_agents_available + return nil + end + + selector_service.select_agent(available_agents) + end + + def selector_service + @selector_service ||= if policy.assignment_order == 'balanced' && enterprise_enabled? && policy.can_use_balanced_assignment? + ::Enterprise::AssignmentV2::BalancedSelector.new(inbox: inbox) + else + AssignmentV2::RoundRobinSelector.new(inbox: inbox) + end + end + + def unassigned_conversations(limit) + scope = inbox.conversations + .unassigned + .open + + # Apply conversation priority ordering + scope = case policy.conversation_priority + when 'longest_waiting' + scope.order(last_activity_at: :asc, created_at: :asc) + else + scope.order(created_at: :asc) + end + + scope.limit(limit) + end + + def assign_conversation_to_agent(conversation, agent) + conversation.update!(assignee: agent) + create_assignment_activity(conversation, agent) + record_assignment_in_rate_limiter(conversation, agent) + true + rescue ActiveRecord::RecordInvalid => e + Rails.logger.error "AssignmentV2: Failed to assign conversation #{conversation.id}: #{e.message}" + false + end + + def create_assignment_activity(conversation, agent) + Rails.configuration.dispatcher.dispatch( + Events::Types::ASSIGNEE_CHANGED, + Time.zone.now, + conversation: conversation, + user: agent + ) + end + + def enterprise_enabled? + @enterprise_enabled ||= defined?(Enterprise) + end + + def log_no_agents_available + Rails.logger.warn("AssignmentV2: No agents available for inbox #{inbox.id}") + end + + def record_assignment_in_rate_limiter(conversation, agent) + rate_limiter = AssignmentV2::RateLimiter.new(inbox: inbox, user: agent) + rate_limiter.record_assignment(conversation) + rescue StandardError => e + Rails.logger.error "AssignmentV2: Failed to record assignment in rate limiter: #{e.message}" + end +end diff --git a/app/services/assignment_v2/rate_limiter.rb b/app/services/assignment_v2/rate_limiter.rb new file mode 100644 index 000000000..da1f956ab --- /dev/null +++ b/app/services/assignment_v2/rate_limiter.rb @@ -0,0 +1,85 @@ +# frozen_string_literal: true + +# Rate limiter for assignment operations +# Uses Redis to track assignment counts per agent per time window +# based on assignment policy's fair_distribution_limit and fair_distribution_window +class AssignmentV2::RateLimiter + pattr_initialize [:inbox!, :user!] + + # Check if the user has exceeded rate limits + # @return [Boolean] true if within limits, false if exceeded + def within_limits? + return true unless policy_exists? + + current_count < rate_limit + end + + # Record an assignment for rate limiting purposes + # @param conversation [Conversation] The conversation being assigned + def record_assignment(_conversation) + return unless policy_exists? + + key = rate_limit_key + redis = Redis.new(Redis::Config.app) + redis.multi do |multi| + multi.incr(key) + multi.expire(key, time_window) + end + end + + # Get current rate limit status for the user + # @return [Hash] Rate limit status information + def status + if policy_exists? + { + within_limits: within_limits?, + current_count: current_count, + limit: rate_limit, + reset_at: Time.zone.at(next_window_start) + } + else + { + within_limits: true, + current_count: 0, + limit: Float::INFINITY, + reset_at: nil + } + end + end + + private + + def policy + @policy ||= inbox.assignment_policy + end + + def policy_exists? + policy.present? && policy.enabled? + end + + def current_count + key = rate_limit_key + redis = Redis.new(Redis::Config.app) + redis.get(key).to_i + end + + def rate_limit + policy&.fair_distribution_limit || 10 + end + + def time_window + policy&.fair_distribution_window || 3600 + end + + def rate_limit_key + "assignment_v2:rate_limit:#{user.id}:#{current_window}" + end + + def current_window + (Time.current.to_i / time_window) * time_window + end + + def next_window_start + current_window + time_window + end +end diff --git a/app/services/assignment_v2/round_robin_selector.rb b/app/services/assignment_v2/round_robin_selector.rb new file mode 100644 index 000000000..939593dcd --- /dev/null +++ b/app/services/assignment_v2/round_robin_selector.rb @@ -0,0 +1,37 @@ +# frozen_string_literal: true + +class AssignmentV2::RoundRobinSelector + pattr_initialize [:inbox!] + + def select_agent(available_agents) + return nil if available_agents.empty? + + # Extract user IDs from inbox members + agent_user_ids = available_agents.map(&:user_id).map(&:to_s) + + # Use Redis queue for round robin + selected_user_id = round_robin_service.available_agent(allowed_agent_ids: agent_user_ids) + return nil unless selected_user_id + + # Return the user object + available_agents.find { |inbox_member| inbox_member.user_id.to_s == selected_user_id }&.user + end + + def add_agent_to_queue(user_id) + round_robin_service.add_agent_to_queue(user_id) + end + + def remove_agent_from_queue(user_id) + round_robin_service.remove_agent_from_queue(user_id) + end + + def reset_queue + round_robin_service.reset_queue + end + + private + + def round_robin_service + @round_robin_service ||= AutoAssignment::InboxRoundRobinService.new(inbox: inbox) + end +end diff --git a/config/routes.rb b/config/routes.rb index 7fb348084..81715922c 100644 --- a/config/routes.rb +++ b/config/routes.rb @@ -217,6 +217,33 @@ Rails.application.routes.draw do end end + resources :leaves do + member do + post :approve + post :reject + end + end + + # Assignment V2 Routes + resources :assignment_policies + + resources :inboxes, only: [] do + resource :assignment_policy, only: [:show, :create, :destroy], controller: 'inbox_assignment_policies' + end + + # Agent Capacity Management (Enterprise) + resources :agent_capacity_policies do + member do + post 'inbox_limits/:inbox_id', to: 'agent_capacity_policies#set_inbox_limit' + delete 'inbox_limits/:inbox_id', to: 'agent_capacity_policies#remove_inbox_limit' + post 'users', to: 'agent_capacity_policies#assign_user' + delete 'users/:user_id', to: 'agent_capacity_policies#remove_user' + end + end + + # Agent capacity status + get 'agents/:agent_id/capacity', to: 'agent_capacity_policies#agent_capacity' + namespace :twitter do resource :authorization, only: [:create] end @@ -290,8 +317,6 @@ Rails.application.routes.draw do member do patch :archive delete :logo - post :send_instructions - get :ssl_status end resources :categories resources :articles do diff --git a/lib/tasks/assignment_v2.rake b/lib/tasks/assignment_v2.rake new file mode 100644 index 000000000..bba23477e --- /dev/null +++ b/lib/tasks/assignment_v2.rake @@ -0,0 +1,65 @@ +# frozen_string_literal: true + +# rubocop:disable Metrics/BlockLength +namespace :assignment_v2 do + desc 'Enable Assignment V2 for an account with default policy' + task :enable, [:account_id] => :environment do |_task, args| + account = Account.find(args[:account_id]) + + # Create default assignment policy if not exists + policy = account.assignment_policies.find_or_create_by!(name: 'Default Policy') do |p| + p.description = 'Default round-robin assignment policy' + p.assignment_order = 'round_robin' + p.conversation_priority = 'earliest_created' + p.enabled = true + end + + puts "Assignment V2 enabled for account #{account.name} with policy #{policy.name}" + end + + desc 'Disable Assignment V2 for an account' + task :disable, [:account_id] => :environment do |_task, args| + account = Account.find(args[:account_id]) + + # Disable all assignment policies + account.assignment_policies.find_each do |policy| + policy.update!(enabled: false) + end + + puts "Assignment V2 disabled for account #{account.name}" + end + + desc 'Run assignment for all enabled inboxes' + task run_all: :environment do + count = 0 + + Inbox.joins(:assignment_policy) + .where(assignment_policies: { enabled: true }) + .find_each do |inbox| + AssignmentV2::AssignmentJob.perform_later(inbox_id: inbox.id) + count += 1 + end + + puts "Queued assignment jobs for #{count} inboxes" + end + + desc 'Show Assignment V2 status' + task status: :environment do + puts 'Assignment V2 Status' + puts '===================' + + total_policies = AssignmentPolicy.count + enabled_policies = AssignmentPolicy.enabled.count + + puts "Total Assignment Policies: #{total_policies}" + puts "Enabled Policies: #{enabled_policies}" + puts + + inboxes_with_v2 = Inbox.joins(:assignment_policy).count + total_inboxes = Inbox.count + + puts "Total Inboxes: #{total_inboxes}" + puts "Inboxes with Assignment V2: #{inboxes_with_v2}" + end +end +# rubocop:enable Metrics/BlockLength diff --git a/spec/factories/assignment_policies.rb b/spec/factories/assignment_policies.rb new file mode 100644 index 000000000..84937d09f --- /dev/null +++ b/spec/factories/assignment_policies.rb @@ -0,0 +1,34 @@ +# frozen_string_literal: true + +FactoryBot.define do + factory :assignment_policy do + account + sequence(:name) { |n| "Assignment Policy #{n}" } + description { 'Test assignment policy' } + assignment_order { :round_robin } + conversation_priority { :earliest_created } + fair_distribution_limit { 10 } + fair_distribution_window { 3600 } + enabled { true } + + trait :balanced do + assignment_order { :balanced } + end + + trait :disabled do + enabled { false } + end + + trait :longest_waiting do + conversation_priority { :longest_waiting } + end + + trait :with_high_limit do + fair_distribution_limit { 50 } + end + + trait :with_short_window do + fair_distribution_window { 300 } # 5 minutes + end + end +end diff --git a/spec/factories/inbox_assignment_policies.rb b/spec/factories/inbox_assignment_policies.rb new file mode 100644 index 000000000..759852585 --- /dev/null +++ b/spec/factories/inbox_assignment_policies.rb @@ -0,0 +1,13 @@ +# frozen_string_literal: true + +FactoryBot.define do + factory :inbox_assignment_policy do + inbox + assignment_policy + + # Ensure inbox and policy belong to same account + after(:build) do |inbox_policy| + inbox_policy.assignment_policy.account = inbox_policy.inbox.account if inbox_policy.inbox && inbox_policy.assignment_policy + end + end +end diff --git a/spec/jobs/assignment_v2/assignment_job_spec.rb b/spec/jobs/assignment_v2/assignment_job_spec.rb new file mode 100644 index 000000000..62b4e476f --- /dev/null +++ b/spec/jobs/assignment_v2/assignment_job_spec.rb @@ -0,0 +1,194 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe AssignmentV2::AssignmentJob, type: :job do + before do + # Mock GlobalConfig to avoid InstallationConfig issues + allow(GlobalConfig).to receive(:get).and_return({}) + end + + let(:account) { create(:account) } + let(:inbox) { create(:inbox, account: account) } + let(:conversation) { create(:conversation, inbox: inbox, assignee: nil) } + let(:assignment_policy) { create(:assignment_policy, account: account, enabled: true) } + let!(:inbox_assignment_policy) { create(:inbox_assignment_policy, inbox: inbox, assignment_policy: assignment_policy) } + + describe '#perform' do + context 'with conversation_id' do + it 'assigns a single conversation' do + service = instance_double(AssignmentV2::AssignmentService) + expect(AssignmentV2::AssignmentService).to receive(:new).with(inbox: inbox).and_return(service) + expect(service).to receive(:perform_for_conversation).with(conversation) + + described_class.new.perform(conversation_id: conversation.id) + end + + it 'handles non-existent conversation gracefully' do + expect(AssignmentV2::AssignmentService).not_to receive(:new) + + # Should not raise error + expect do + described_class.new.perform(conversation_id: 999_999) + end.not_to raise_error + end + end + + context 'with inbox_id' do + let!(:agent) { create(:user, account: account, role: :agent, availability: :online) } + + before do + create_list(:conversation, 3, inbox: inbox, assignee: nil) + create(:inbox_member, inbox: inbox, user: agent) + end + + it 'assigns multiple conversations for inbox' do + allow(Inbox).to receive(:find_by).with(id: inbox.id).and_return(inbox) + allow(account).to receive(:assignment_v2_enabled?).and_return(true) + + service = instance_double(AssignmentV2::AssignmentService) + expect(AssignmentV2::AssignmentService).to receive(:new).with(inbox: inbox).and_return(service) + expect(service).to receive(:perform_bulk_assignment).and_return(3) + + described_class.new.perform(inbox_id: inbox.id) + end + + it 'logs the number of assigned conversations' do + allow(Inbox).to receive(:find_by).with(id: inbox.id).and_return(inbox) + allow(account).to receive(:assignment_v2_enabled?).and_return(true) + + service = instance_double(AssignmentV2::AssignmentService) + allow(AssignmentV2::AssignmentService).to receive(:new).with(inbox: inbox).and_return(service) + allow(service).to receive(:perform_bulk_assignment).and_return(2) + + expect(Rails.logger).to receive(:info).with("AssignmentV2::AssignmentJob: Assigned 2 conversations for inbox #{inbox.id}") + + described_class.new.perform(inbox_id: inbox.id) + end + + it 'skips assignment when inbox has no policy' do + inbox_assignment_policy.destroy! + + expect(AssignmentV2::AssignmentService).not_to receive(:new) + + described_class.new.perform(inbox_id: inbox.id) + end + + it 'skips assignment when policy is disabled' do + assignment_policy.update!(enabled: false) + + expect(AssignmentV2::AssignmentService).not_to receive(:new) + + described_class.new.perform(inbox_id: inbox.id) + end + + it 'handles non-existent inbox gracefully' do + expect(AssignmentV2::AssignmentService).not_to receive(:new) + + # Should not raise error + expect do + described_class.new.perform(inbox_id: 999_999) + end.not_to raise_error + end + end + + context 'without parameters' do + it 'logs error when no parameters provided' do + expect(Rails.logger).to receive(:error).with('AssignmentV2::AssignmentJob: No inbox_id or conversation_id provided') + + described_class.new.perform + end + + it 'does not attempt assignment' do + expect(AssignmentV2::AssignmentService).not_to receive(:new) + + described_class.new.perform + end + end + + context 'with both parameters' do + it 'prioritizes conversation_id over inbox_id' do + service = instance_double(AssignmentV2::AssignmentService) + expect(AssignmentV2::AssignmentService).to receive(:new).with(inbox: inbox).and_return(service) + expect(service).to receive(:perform_for_conversation).with(conversation) + expect(service).not_to receive(:perform_bulk_assignment) + + described_class.new.perform(conversation_id: conversation.id, inbox_id: inbox.id) + end + end + end + + describe 'job configuration' do + it 'uses the low queue' do + expect(described_class.new.queue_name).to eq('low') + end + end + + describe 'error handling' do + context 'when assignment service raises error' do + it 'propagates the error for retry' do + service = instance_double(AssignmentV2::AssignmentService) + allow(AssignmentV2::AssignmentService).to receive(:new).and_return(service) + allow(service).to receive(:perform_for_conversation).and_raise(StandardError, 'Assignment failed') + + expect do + described_class.new.perform(conversation_id: conversation.id) + end.to raise_error(StandardError, 'Assignment failed') + end + end + + context 'when database connection fails' do + it 'raises error for retry' do + allow(Conversation).to receive(:find_by).and_raise(ActiveRecord::ConnectionNotEstablished) + + expect do + described_class.new.perform(conversation_id: conversation.id) + end.to raise_error(ActiveRecord::ConnectionNotEstablished) + end + end + end + + describe 'concurrency and idempotency' do + it 'handles concurrent job execution safely' do + # Create multiple jobs for same inbox + jobs = [] + 3.times { jobs << described_class.new } + + # All should execute without issues + expect do + jobs.each { |job| job.perform(inbox_id: inbox.id) } + end.not_to raise_error + end + + it 'is idempotent for conversation assignment' do + service = instance_double(AssignmentV2::AssignmentService) + allow(AssignmentV2::AssignmentService).to receive(:new).and_return(service) + + # First call assigns + expect(service).to receive(:perform_for_conversation).and_return(true) + described_class.new.perform(conversation_id: conversation.id) + + # Second call should handle already assigned conversation + expect(service).to receive(:perform_for_conversation).and_return(false) + expect { described_class.new.perform(conversation_id: conversation.id) }.not_to raise_error + end + end + + describe 'performance considerations' do + it 'processes large inbox assignments in batches' do + # Create many unassigned conversations + create_list(:conversation, 100, inbox: inbox, assignee: nil) + + allow(Inbox).to receive(:find_by).with(id: inbox.id).and_return(inbox) + allow(account).to receive(:assignment_v2_enabled?).and_return(true) + + service = instance_double(AssignmentV2::AssignmentService) + allow(AssignmentV2::AssignmentService).to receive(:new).and_return(service) + + # Service should be called with default limit + expect(service).to receive(:perform_bulk_assignment).with(no_args).and_return(50) + + described_class.new.perform(inbox_id: inbox.id) + end + end +end diff --git a/spec/models/assignment_policy_spec.rb b/spec/models/assignment_policy_spec.rb new file mode 100644 index 000000000..4495297f3 --- /dev/null +++ b/spec/models/assignment_policy_spec.rb @@ -0,0 +1,127 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe AssignmentPolicy, type: :model do + let(:account) { create(:account) } + let(:assignment_policy) { create(:assignment_policy, account: account) } + + describe 'associations' do + it { is_expected.to belong_to(:account) } + it { is_expected.to have_many(:inbox_assignment_policies).dependent(:destroy) } + it { is_expected.to have_many(:inboxes).through(:inbox_assignment_policies) } + end + + describe 'validations' do + subject { assignment_policy } + + it { is_expected.to validate_presence_of(:name) } + it { is_expected.to validate_uniqueness_of(:name).scoped_to(:account_id) } + it { is_expected.to validate_length_of(:name).is_at_most(255) } + it { is_expected.to validate_length_of(:description).is_at_most(1000) } + + it { is_expected.to validate_presence_of(:fair_distribution_limit) } + it { is_expected.to validate_numericality_of(:fair_distribution_limit).is_greater_than(0).is_less_than_or_equal_to(100) } + + it { is_expected.to validate_presence_of(:fair_distribution_window) } + it { is_expected.to validate_numericality_of(:fair_distribution_window).is_greater_than(60).is_less_than_or_equal_to(86_400) } + + context 'with balanced assignment validation' do + let(:enterprise_account) { create(:account) } + + before do + allow(enterprise_account).to receive(:feature_enabled?).with(:enterprise_agent_capacity).and_return(true) + end + + it 'allows balanced assignment for enterprise accounts' do + policy = build(:assignment_policy, account: enterprise_account, assignment_order: :balanced) + expect(policy).to be_valid + end + + it 'rejects balanced assignment for non-enterprise accounts' do + policy = build(:assignment_policy, account: account, assignment_order: :balanced) + expect(policy).not_to be_valid + expect(policy.errors[:assignment_order]).to include('Balanced assignment is only available for enterprise accounts') + end + end + end + + describe 'enums' do + it { is_expected.to define_enum_for(:assignment_order).with_values(round_robin: 0, balanced: 1) } + it { is_expected.to define_enum_for(:conversation_priority).with_values(earliest_created: 0, longest_waiting: 1) } + end + + describe 'scopes' do + let!(:enabled_policy) { create(:assignment_policy, account: account, enabled: true) } + let!(:disabled_policy) { create(:assignment_policy, account: account, enabled: false) } + + it 'filters enabled policies' do + expect(described_class.enabled).to include(enabled_policy) + expect(described_class.enabled).not_to include(disabled_policy) + end + + it 'filters disabled policies' do + expect(described_class.disabled).to include(disabled_policy) + expect(described_class.disabled).not_to include(enabled_policy) + end + end + + describe '#can_use_balanced_assignment?' do + context 'when account has enterprise agent capacity feature' do + before do + allow(account).to receive(:feature_enabled?).with(:enterprise_agent_capacity).and_return(true) + end + + it 'returns true' do + expect(assignment_policy.can_use_balanced_assignment?).to be true + end + end + + context 'when account does not have enterprise features' do + before do + allow(account).to receive(:feature_enabled?).with(:enterprise_agent_capacity).and_return(false) + end + + it 'returns false' do + expect(assignment_policy.can_use_balanced_assignment?).to be false + end + end + end + + describe '#webhook_data' do + it 'returns correct data structure' do + data = assignment_policy.webhook_data + + expect(data).to include( + id: assignment_policy.id, + name: assignment_policy.name, + description: assignment_policy.description, + assignment_order: assignment_policy.assignment_order, + conversation_priority: assignment_policy.conversation_priority, + fair_distribution_limit: assignment_policy.fair_distribution_limit, + fair_distribution_window: assignment_policy.fair_distribution_window, + enabled: assignment_policy.enabled + ) + end + end + + describe 'cache invalidation' do + let(:inbox) { create(:inbox, account: account) } + + before { create(:inbox_assignment_policy, inbox: inbox, assignment_policy: assignment_policy) } + + it 'clears assignment caches on update' do + expect(Rails.cache).to receive(:delete).with("assignment_v2:policy:#{assignment_policy.id}") + expect(Rails.cache).to receive(:delete).with("assignment_v2:inbox_policy:#{inbox.id}") + + assignment_policy.update!(name: 'Updated Policy') + end + + it 'clears assignment caches on destroy' do + expect(Rails.cache).to receive(:delete).with("assignment_v2:policy:#{assignment_policy.id}") + expect(Rails.cache).to receive(:delete).with("assignment_v2:inbox_policy:#{inbox.id}") + + assignment_policy.destroy! + end + end +end diff --git a/spec/models/inbox_assignment_policy_spec.rb b/spec/models/inbox_assignment_policy_spec.rb new file mode 100644 index 000000000..b46da4b72 --- /dev/null +++ b/spec/models/inbox_assignment_policy_spec.rb @@ -0,0 +1,165 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe InboxAssignmentPolicy, type: :model do + let(:account) { create(:account) } + let(:inbox) { create(:inbox, account: account) } + let(:assignment_policy) { create(:assignment_policy, account: account) } + let(:inbox_assignment_policy) { create(:inbox_assignment_policy, inbox: inbox, assignment_policy: assignment_policy) } + + describe 'associations' do + it { is_expected.to belong_to(:inbox) } + it { is_expected.to belong_to(:assignment_policy) } + end + + describe 'validations' do + subject { inbox_assignment_policy } + + it { is_expected.to validate_uniqueness_of(:inbox_id) } + + context 'with inbox and policy from different accounts' do + let(:other_account) { create(:account) } + let(:other_policy) { create(:assignment_policy, account: other_account) } + + it 'validates inbox belongs to same account as policy' do + # Build without the factory callback that sets the accounts to be the same + invalid_policy = described_class.new(inbox: inbox, assignment_policy: other_policy) + + expect(invalid_policy).not_to be_valid + expect(invalid_policy.errors[:inbox]).to include('must belong to the same account as the assignment policy') + end + end + end + + describe 'delegations' do + it 'delegates account to inbox' do + expect(inbox_assignment_policy.account).to eq(account) + end + + it 'delegates policy attributes' do + expect(inbox_assignment_policy.policy_name).to eq(assignment_policy.name) + expect(inbox_assignment_policy.policy_description).to eq(assignment_policy.description) + expect(inbox_assignment_policy.policy_assignment_order).to eq(assignment_policy.assignment_order) + expect(inbox_assignment_policy.policy_conversation_priority).to eq(assignment_policy.conversation_priority) + expect(inbox_assignment_policy.policy_fair_distribution_limit).to eq(assignment_policy.fair_distribution_limit) + expect(inbox_assignment_policy.policy_fair_distribution_window).to eq(assignment_policy.fair_distribution_window) + expect(inbox_assignment_policy.policy_enabled?).to eq(assignment_policy.enabled?) + end + end + + describe 'scopes' do + let!(:enabled_policy) { create(:assignment_policy, account: account, enabled: true) } + let!(:disabled_policy) { create(:assignment_policy, account: account, enabled: false) } + let(:inbox2) { create(:inbox, account: account) } + let!(:enabled_inbox_policy) { create(:inbox_assignment_policy, inbox: inbox, assignment_policy: enabled_policy) } + let!(:disabled_inbox_policy) { create(:inbox_assignment_policy, inbox: inbox2, assignment_policy: disabled_policy) } + + describe '.enabled' do + it 'returns only inbox policies with enabled assignment policies' do + expect(described_class.enabled).to include(enabled_inbox_policy) + expect(described_class.enabled).not_to include(disabled_inbox_policy) + end + end + + describe '.disabled' do + it 'returns only inbox policies with disabled assignment policies' do + expect(described_class.disabled).to include(disabled_inbox_policy) + expect(described_class.disabled).not_to include(enabled_inbox_policy) + end + end + end + + describe '#webhook_data' do + it 'returns correct data structure' do + data = inbox_assignment_policy.webhook_data + + expect(data).to include( + id: inbox_assignment_policy.id, + inbox_id: inbox.id, + assignment_policy_id: assignment_policy.id + ) + + expect(data[:policy]).to eq(assignment_policy.webhook_data) + end + end + + describe 'cache management' do + it 'clears inbox cache on create' do + expect(Rails.cache).to receive(:delete).with("assignment_v2:inbox_policy:#{inbox.id}") + + create(:inbox_assignment_policy, inbox: inbox, assignment_policy: assignment_policy) + end + + it 'clears inbox cache on update' do + expect(Rails.cache).to receive(:delete).with("assignment_v2:inbox_policy:#{inbox.id}").at_least(:once) + + inbox_assignment_policy.update!(updated_at: Time.current) + end + + it 'clears inbox cache on destroy' do + expect(Rails.cache).to receive(:delete).with("assignment_v2:inbox_policy:#{inbox.id}").at_least(:once) + + inbox_assignment_policy.destroy! + end + + it 'updates account cache' do + # AccountCacheRevalidator concern should trigger cache update + expect(inbox_assignment_policy).to receive(:update_account_cache) + + inbox_assignment_policy.send(:clear_inbox_cache) + end + end + + describe 'business logic constraints' do + it 'prevents multiple policies per inbox' do + # Ensure first policy exists + inbox_assignment_policy + + policy2 = create(:assignment_policy, account: account) + + # Try to create a second policy for the same inbox + duplicate_policy = described_class.new(inbox: inbox, assignment_policy: policy2) + expect(duplicate_policy).not_to be_valid + expect(duplicate_policy.errors[:inbox_id]).to include('has already been taken') + end + + it 'allows reassigning to different policy' do + policy2 = create(:assignment_policy, account: account) + + expect do + inbox_assignment_policy.update!(assignment_policy: policy2) + end.not_to raise_error + + expect(inbox_assignment_policy.reload.assignment_policy).to eq(policy2) + end + end + + describe 'edge cases' do + it 'handles nil associations gracefully' do + # Build without saving to test nil handling + policy = build(:inbox_assignment_policy, inbox: nil, assignment_policy: nil) + + expect { policy.valid? }.not_to raise_error + expect(policy).not_to be_valid + end + + it 'handles policy deletion cascade' do + inbox_policy_id = inbox_assignment_policy.id + + # Deleting policy should delete inbox assignment + assignment_policy.destroy! + + expect(described_class.find_by(id: inbox_policy_id)).to be_nil + end + + it 'handles inbox deletion cascade' do + inbox_policy_id = inbox_assignment_policy.id + + # Deleting inbox should delete inbox assignment + inbox.destroy! + + expect(described_class.find_by(id: inbox_policy_id)).to be_nil + end + end +end diff --git a/spec/services/assignment_v2/assignment_service_spec.rb b/spec/services/assignment_v2/assignment_service_spec.rb new file mode 100644 index 000000000..4b03c8972 --- /dev/null +++ b/spec/services/assignment_v2/assignment_service_spec.rb @@ -0,0 +1,329 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe AssignmentV2::AssignmentService do + before do + # Mock the GlobalConfig to avoid InstallationConfig issues + allow(GlobalConfig).to receive(:get).and_return({}) + + # Define the constant if not already defined + stub_const('ASSIGNEE_CHANGED', 'assignee.changed') unless defined?(ASSIGNEE_CHANGED) + create(:inbox_member, inbox: inbox, user: agent1) + create(:inbox_member, inbox: inbox, user: agent2) + create(:inbox_member, inbox: inbox, user: agent3) + + # Mock available agents to return inbox members + online_members = InboxMember.joins(:user).where(inbox: inbox, user: [agent1, agent2]) + allow(inbox).to receive(:available_agents).and_return(online_members) + end + + let(:account) { create(:account) } + let(:inbox) { create(:inbox, account: account) } + let(:assignment_policy) { create(:assignment_policy, account: account, enabled: true) } + let!(:inbox_assignment_policy) { create(:inbox_assignment_policy, inbox: inbox, assignment_policy: assignment_policy) } + let(:service) { described_class.new(inbox: inbox) } + + # Create agents + let!(:agent1) { create(:user, account: account, role: :agent, availability: :online) } + let!(:agent2) { create(:user, account: account, role: :agent, availability: :online) } + let!(:agent3) { create(:user, account: account, role: :agent, availability: :offline) } + + # Make agents members of inbox + + describe '#perform_for_conversation' do + let(:conversation) { create(:conversation, inbox: inbox, assignee: nil) } + + context 'when policy is enabled' do + before do + # Mock the selector to return an agent + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + allow(selector).to receive(:select_agent).and_return(agent1) + end + + it 'assigns conversation to an available agent' do + expect(service.perform_for_conversation(conversation)).to be true + expect(conversation.reload.assignee).to eq(agent1) + end + + it 'dispatches assignment event' do + # The dispatcher is called from the assignment service and also from conversation model + allow(Rails.configuration.dispatcher).to receive(:dispatch) + + service.perform_for_conversation(conversation) + + expect(Rails.configuration.dispatcher).to have_received(:dispatch).with( + 'assignee.changed', + anything, + hash_including(conversation: conversation, user: agent1) + ).at_least(:once) + end + + it 'returns false when no agents are available' do + allow(inbox).to receive(:available_agents).and_return(InboxMember.none) + allow(Rails.logger).to receive(:warn) + + expect(service.perform_for_conversation(conversation)).to be false + expect(conversation.reload.assignee).to be_nil + end + end + + context 'when policy is disabled' do + before { assignment_policy.update!(enabled: false) } + + it 'does not assign conversation' do + expect(service.perform_for_conversation(conversation)).to be false + expect(conversation.reload.assignee).to be_nil + end + end + + context 'when conversation is already assigned' do + before { conversation.update!(assignee: agent1) } + + it 'does not reassign conversation' do + expect(service.perform_for_conversation(conversation)).to be false + expect(conversation.reload.assignee).to eq(agent1) + end + end + + context 'with round robin assignment' do + before do + assignment_policy.update!(assignment_order: :round_robin) + + # Mock round robin selector to return agents in rotation + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + agent_index = 0 + allow(selector).to receive(:select_agent) do + agent = [agent1, agent2][agent_index % 2] + agent_index += 1 + agent + end + end + + it 'assigns agents in rotation' do + conversations = create_list(:conversation, 4, inbox: inbox, assignee: nil) + + assignments = conversations.map do |conv| + service.perform_for_conversation(conv) + conv.reload.assignee + end + + # Should rotate between available agents + expect(assignments[0]).to eq(agent1) + expect(assignments[1]).to eq(agent2) + expect(assignments[2]).to eq(agent1) # Back to first agent + expect(assignments[3]).to eq(agent2) # Back to second agent + end + end + + context 'with balanced assignment' do + before do + # For now, just use round robin since balanced is enterprise only + # The test is verifying the service works, not the specific algorithm + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + allow(selector).to receive(:select_agent).and_return(agent2) + end + + it 'assigns conversations successfully' do + # Create existing assignments + create_list(:conversation, 3, inbox: inbox, assignee: agent1, status: :open) + create(:conversation, inbox: inbox, assignee: agent2, status: :open) + + new_conversation = create(:conversation, inbox: inbox, assignee: nil) + + expect(service.perform_for_conversation(new_conversation)).to be true + expect(new_conversation.reload.assignee).to eq(agent2) + end + + it 'handles different conversation statuses' do + # Create resolved conversations (should not count) + create_list(:conversation, 5, inbox: inbox, assignee: agent1, status: :resolved) + + # Create open conversation + create(:conversation, inbox: inbox, assignee: agent2, status: :open) + + new_conversation = create(:conversation, inbox: inbox, assignee: nil) + + expect(service.perform_for_conversation(new_conversation)).to be true + expect(new_conversation.reload.assignee).to eq(agent2) # Selected by mock + end + end + + context 'when error occurs' do + before do + # Mock the selector to return an agent + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + allow(selector).to receive(:select_agent).and_return(agent1) + end + + it 'returns false and logs error on assignment failure' do + allow(conversation).to receive(:update!).and_raise(ActiveRecord::RecordInvalid.new(conversation)) + expect(Rails.logger).to receive(:error).with(/Failed to assign conversation/) + + expect(service.perform_for_conversation(conversation)).to be false + end + end + end + + describe '#perform_bulk_assignment' do + before do + create_list(:conversation, 5, inbox: inbox, assignee: nil, status: :open) + + # Mock the selector to return agents + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + call_count = 0 + allow(selector).to receive(:select_agent) do + call_count += 1 + call_count.odd? ? agent1 : agent2 + end + end + + context 'when policy is enabled' do + it 'assigns multiple conversations' do + assigned_count = service.perform_bulk_assignment(limit: 3) + + expect(assigned_count).to eq(3) + expect(inbox.conversations.unassigned.count).to eq(2) + end + + it 'respects conversation priority order' do + # Clear existing conversations first + Conversation.destroy_all + + # Create conversations with different timestamps + old_conversation = create(:conversation, inbox: inbox, assignee: nil, status: :open, created_at: 1.hour.ago) + new_conversation = create(:conversation, inbox: inbox, assignee: nil, status: :open, created_at: 1.minute.ago) + + assignment_policy.update!(conversation_priority: :earliest_created) + + # Re-create service after policy change + service_with_priority = described_class.new(inbox: inbox) + + service_with_priority.perform_bulk_assignment(limit: 1) + + expect(old_conversation.reload.assignee).not_to be_nil + expect(new_conversation.reload.assignee).to be_nil + end + + it 'handles longest_waiting priority' do + # Clear existing conversations first + Conversation.destroy_all + + # Create conversations with different last activity + inactive_conversation = create(:conversation, inbox: inbox, assignee: nil, status: :open, last_activity_at: 2.hours.ago) + active_conversation = create(:conversation, inbox: inbox, assignee: nil, status: :open, last_activity_at: 5.minutes.ago) + + assignment_policy.update!(conversation_priority: :longest_waiting) + + # Re-create service after policy change + service_with_priority = described_class.new(inbox: inbox) + + service_with_priority.perform_bulk_assignment(limit: 1) + + expect(inactive_conversation.reload.assignee).not_to be_nil + expect(active_conversation.reload.assignee).to be_nil + end + + it 'returns 0 when no conversations to assign' do + Conversation.find_each { |c| c.update!(assignee_id: agent1.id) } + + expect(service.perform_bulk_assignment).to eq(0) + end + end + + context 'when policy is disabled' do + before { assignment_policy.update!(enabled: false) } + + it 'does not assign any conversations' do + expect(service.perform_bulk_assignment).to eq(0) + expect(inbox.conversations.unassigned.count).to eq(5) + end + end + end + + describe 'enterprise capacity features' do + let(:conversation) { create(:conversation, inbox: inbox, assignee: nil) } + + before do + # Mock enterprise availability + stub_const('Enterprise', Module.new) + stub_const('Enterprise::AssignmentV2::CapacityManager', Class.new) + + # Mock the selector to return agent1 + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + allow(selector).to receive(:select_agent).and_return(agent1) + end + + it 'uses round robin when enterprise features are available' do + expect(service.perform_for_conversation(conversation)).to be true + expect(conversation.reload.assignee).to eq(agent1) + end + + it 'handles absence of enterprise features gracefully' do + # Remove enterprise constant + hide_const('Enterprise') + + # Service should still work with round robin + expect(service.perform_for_conversation(conversation)).to be true + expect(conversation.reload.assignee).to eq(agent1) + end + end + + describe 'cache management' do + it 'uses cache for round robin state' do + assignment_policy.update!(assignment_order: :round_robin) + + # Mock the selector and round robin service + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + allow(selector).to receive(:select_agent).and_return(agent1) + + # Create and assign conversations + conversation1 = create(:conversation, inbox: inbox, assignee: nil) + service.perform_for_conversation(conversation1) + + conversation2 = create(:conversation, inbox: inbox, assignee: nil) + service.perform_for_conversation(conversation2) + + # Just verify assignments worked + expect(conversation1.reload.assignee).to eq(agent1) + expect(conversation2.reload.assignee).to eq(agent1) + end + end + + describe 'edge cases' do + it 'handles inbox without policy gracefully' do + inbox_assignment_policy.destroy! + conversation = create(:conversation, inbox: inbox, assignee: nil) + + expect(service.perform_for_conversation(conversation)).to be false + end + + it 'handles empty agent list' do + allow(inbox).to receive(:available_agents).and_return(InboxMember.none) + conversation = create(:conversation, inbox: inbox, assignee: nil) + + expect(service.perform_for_conversation(conversation)).to be false + end + + it 'filters out agents without inbox membership' do + non_member_agent = create(:user, account: account, role: :agent, availability: :online) + conversation = create(:conversation, inbox: inbox, assignee: nil) + + # Mock selector to return agent1 (who is a member) + selector = instance_double(AssignmentV2::RoundRobinSelector) + allow(AssignmentV2::RoundRobinSelector).to receive(:new).and_return(selector) + allow(selector).to receive(:select_agent).and_return(agent1) + + expect(service.perform_for_conversation(conversation)).to be true + expect(conversation.reload.assignee).not_to eq(non_member_agent) + expect(conversation.reload.assignee).to eq(agent1) + end + end +end diff --git a/spec/services/assignment_v2/rate_limiter_spec.rb b/spec/services/assignment_v2/rate_limiter_spec.rb new file mode 100644 index 000000000..67ddf6639 --- /dev/null +++ b/spec/services/assignment_v2/rate_limiter_spec.rb @@ -0,0 +1,293 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe AssignmentV2::RateLimiter, type: :service do + before do + # Mock GlobalConfig to avoid InstallationConfig issues + allow(GlobalConfig).to receive(:get).and_return({}) + redis = Redis.new(Redis::Config.app) + redis.flushdb if Rails.env.test? + + # Ensure inbox_assignment_policy exists so the inbox has a policy + inbox_assignment_policy + end + + let(:account) { create(:account) } + let(:policy) { create(:assignment_policy, account: account, fair_distribution_limit: 5, fair_distribution_window: 3600) } + let(:agent) { create(:user, account: account) } + let(:inbox) { create(:inbox, account: account) } + let(:inbox_assignment_policy) { create(:inbox_assignment_policy, inbox: inbox, assignment_policy: policy) } + let(:rate_limiter) { described_class.new(inbox: inbox, user: agent) } + + describe '#initialize' do + it 'sets up rate limiter with inbox and user' do + expect(rate_limiter.instance_variable_get(:@inbox)).to eq(inbox) + expect(rate_limiter.instance_variable_get(:@user)).to eq(agent) + end + end + + describe '#within_limits?' do + context 'when agent has no assignments in current window' do + it 'returns true' do + expect(rate_limiter.within_limits?).to be true + end + end + + context 'when agent is below limit' do + before do + # Simulate 3 assignments in current window + 3.times { rate_limiter.record_assignment(create(:conversation)) } + end + + it 'returns true' do + expect(rate_limiter.within_limits?).to be true + end + end + + context 'when agent reaches limit' do + before do + # Simulate reaching the limit (5 assignments) + 5.times { rate_limiter.record_assignment(create(:conversation)) } + end + + it 'returns false' do + expect(rate_limiter.within_limits?).to be false + end + end + + context 'when agent exceeds limit' do + before do + # Simulate exceeding the limit + 6.times { rate_limiter.record_assignment(create(:conversation)) } + end + + it 'returns false' do + expect(rate_limiter.within_limits?).to be false + end + end + end + + describe '#record_assignment' do + let(:conversation) { create(:conversation, inbox: inbox) } + + it 'increments assignment count for agent' do + initial_status = rate_limiter.status + rate_limiter.record_assignment(conversation) + new_status = rate_limiter.status + + expect(new_status[:current_count]).to eq(initial_status[:current_count] + 1) + end + + it 'sets expiration on the key' do + rate_limiter.record_assignment(conversation) + + current_window = (Time.current.to_i / 3600) * 3600 + key = "assignment_v2:rate_limit:#{agent.id}:#{current_window}" + + redis = Redis.new(Redis::Config.app) + ttl = redis.ttl(key) + expect(ttl).to be > 0 + expect(ttl).to be <= 3600 + end + + context 'when Redis fails' do + before do + redis_double = instance_double(Redis) + allow(Redis).to receive(:new).and_return(redis_double) + allow(redis_double).to receive(:multi).and_raise(Redis::ConnectionError) + allow(Rails.logger).to receive(:error) + end + + it 'raises error' do + expect { rate_limiter.record_assignment(conversation) }.to raise_error(Redis::ConnectionError) + end + end + end + + describe '#status' do + it 'returns correct status for agent with no assignments' do + status = rate_limiter.status + expect(status[:current_count]).to eq(0) + expect(status[:within_limits]).to be true + expect(status[:limit]).to eq(5) + end + + it 'returns correct count after assignments' do + 3.times { rate_limiter.record_assignment(create(:conversation)) } + status = rate_limiter.status + expect(status[:current_count]).to eq(3) + expect(status[:within_limits]).to be true + end + + context 'when Redis fails' do + before do + redis_double = instance_double(Redis) + allow(Redis).to receive(:new).and_return(redis_double) + allow(redis_double).to receive(:get).and_raise(Redis::ConnectionError) + allow(Rails.logger).to receive(:error) + end + + it 'raises error' do + expect { rate_limiter.status }.to raise_error(Redis::ConnectionError) + end + end + end + + describe 'remaining assignments' do + it 'returns full limit when no assignments made' do + status = rate_limiter.status + expect(status[:limit] - status[:current_count]).to eq(5) + end + + it 'returns correct remaining count' do + 2.times { rate_limiter.record_assignment(create(:conversation)) } + status = rate_limiter.status + expect(status[:limit] - status[:current_count]).to eq(3) + end + + it 'returns 0 when limit reached' do + 5.times { rate_limiter.record_assignment(create(:conversation)) } + status = rate_limiter.status + expect(status[:limit] - status[:current_count]).to eq(0) + end + + it 'returns negative when limit exceeded' do + 6.times { rate_limiter.record_assignment(create(:conversation)) } + status = rate_limiter.status + expect(status[:limit] - status[:current_count]).to eq(-1) + end + end + + describe 'assignment capacity checks' do + it 'returns true when agent has remaining capacity' do + 2.times { rate_limiter.record_assignment(create(:conversation)) } + expect(rate_limiter.within_limits?).to be true + end + + it 'returns false when agent has no capacity' do + 5.times { rate_limiter.record_assignment(create(:conversation)) } + expect(rate_limiter.within_limits?).to be false + end + + it 'correctly tracks multiple assignments' do + 3.times { rate_limiter.record_assignment(create(:conversation)) } + + status = rate_limiter.status + expect(status[:current_count]).to eq(3) + expect(status[:within_limits]).to be true + expect(status[:limit] - status[:current_count]).to eq(2) + end + end + + describe 'multiple agents' do + let(:agent2) { create(:user, account: account) } + let(:rate_limiter2) { described_class.new(inbox: inbox, user: agent2) } + + before do + 2.times { rate_limiter.record_assignment(create(:conversation)) } + 4.times { rate_limiter2.record_assignment(create(:conversation)) } + end + + it 'tracks status independently for each agent' do + status1 = rate_limiter.status + status2 = rate_limiter2.status + + expect(status1[:current_count]).to eq(2) + expect(status1[:within_limits]).to be true + expect(status1[:limit] - status1[:current_count]).to eq(3) + + expect(status2[:current_count]).to eq(4) + expect(status2[:within_limits]).to be true + expect(status2[:limit] - status2[:current_count]).to eq(1) + end + end + + describe 'reset functionality' do + before do + 3.times { rate_limiter.record_assignment(create(:conversation)) } + end + + it 'can be reset by clearing Redis key' do + status_before = rate_limiter.status + expect(status_before[:current_count]).to eq(3) + + # Manually clear the key + redis = Redis.new(Redis::Config.app) + current_window = (Time.current.to_i / 3600) * 3600 + key = "assignment_v2:rate_limit:#{agent.id}:#{current_window}" + redis.del(key) + + status_after = rate_limiter.status + expect(status_after[:current_count]).to eq(0) + end + + context 'when Redis fails' do + before do + redis_double = instance_double(Redis) + allow(Redis).to receive(:new).and_return(redis_double) + allow(redis_double).to receive(:del).and_raise(Redis::ConnectionError) + allow(Rails.logger).to receive(:error) + end + + it 'raises error' do + redis = Redis.new(Redis::Config.app) + current_window = (Time.current.to_i / 3600) * 3600 + key = "assignment_v2:rate_limit:#{agent.id}:#{current_window}" + expect { redis.del(key) }.to raise_error(Redis::ConnectionError) + end + end + end + + describe 'window timing' do + it 'calculates reset time correctly' do + # Mock current time to make test predictable + travel_to(Time.zone.parse('2024-01-01 10:30:00')) do + status = rate_limiter.status + reset_time = status[:reset_at] + + expect(reset_time).to be_a(Time) + expect(reset_time).to be > Time.current + expect(reset_time - Time.current).to be <= 3600 + end + end + end + + describe 'window boundaries' do + it 'resets count in new window' do + # Set up assignment in current window + 2.times { rate_limiter.record_assignment(create(:conversation)) } + expect(rate_limiter.status[:current_count]).to eq(2) + + # Travel to next window (advance by window size) + travel(3601.seconds) do + expect(rate_limiter.status[:current_count]).to eq(0) + expect(rate_limiter.within_limits?).to be true + end + end + end + + describe 'concurrent access' do + it 'handles concurrent increments correctly' do + threads = [] + results = [] + mutex = Mutex.new + + # Simulate concurrent assignment requests + 5.times do + threads << Thread.new do + within_limits = rate_limiter.within_limits? + mutex.synchronize { results << within_limits } + rate_limiter.record_assignment(create(:conversation)) if within_limits + end + end + + threads.each(&:join) + + # Final count should not exceed the limit + final_count = rate_limiter.status[:current_count] + expect(final_count).to be <= 5 + expect(results.count(true)).to eq(final_count) + end + end +end diff --git a/spec/services/assignment_v2/round_robin_selector_spec.rb b/spec/services/assignment_v2/round_robin_selector_spec.rb new file mode 100644 index 000000000..2311b83b1 --- /dev/null +++ b/spec/services/assignment_v2/round_robin_selector_spec.rb @@ -0,0 +1,145 @@ +# frozen_string_literal: true + +require 'rails_helper' + +RSpec.describe AssignmentV2::RoundRobinSelector, type: :service do + before do + # Mock GlobalConfig to avoid InstallationConfig issues + allow(GlobalConfig).to receive(:get).and_return({}) + create(:inbox_member, inbox: inbox, user: user1) + create(:inbox_member, inbox: inbox, user: user2) + create(:inbox_member, inbox: inbox, user: user3) + end + + let(:account) { create(:account) } + let(:inbox) { create(:inbox, account: account) } + let(:policy) { create(:assignment_policy, account: account) } + let(:user1) { create(:user, account: account, availability: :online) } + let(:user2) { create(:user, account: account, availability: :online) } + let(:user3) { create(:user, account: account, availability: :offline) } + + describe '#select_agent' do + let(:selector) { described_class.new(inbox: inbox) } + let(:round_robin_service) { instance_double(AutoAssignment::InboxRoundRobinService) } + let(:available_agents) { InboxMember.where(inbox: inbox, user: [user1, user2]) } + + before do + allow(AutoAssignment::InboxRoundRobinService).to receive(:new).with(inbox: inbox).and_return(round_robin_service) + end + + context 'when Redis is available' do + before do + allow(round_robin_service).to receive(:available_agent).with(allowed_agent_ids: [user1.id.to_s, user2.id.to_s]).and_return(user1.id.to_s) + end + + it 'returns an online agent' do + result = selector.select_agent(available_agents) + expect(result).to eq(user1) + end + + it 'excludes offline agents' do + result = selector.select_agent(available_agents) + expect(result).not_to eq(user3) + end + + it 'handles no available agent gracefully' do + allow(round_robin_service).to receive(:available_agent).with(allowed_agent_ids: [user1.id.to_s, user2.id.to_s]).and_return(nil) + result = selector.select_agent(available_agents) + expect(result).to be_nil + end + end + + context 'when Redis fails' do + before do + allow(round_robin_service).to receive(:available_agent).and_raise(Redis::CannotConnectError) + end + + it 'raises the error' do + expect { selector.select_agent(available_agents) }.to raise_error(Redis::CannotConnectError) + end + end + + context 'with empty available agents' do + it 'returns nil when no agents are available' do + result = selector.select_agent(InboxMember.none) + expect(result).to be_nil + end + end + + context 'with different user IDs' do + it 'correctly finds the inbox member by user_id' do + allow(round_robin_service).to receive(:available_agent).with(allowed_agent_ids: [user1.id.to_s, user2.id.to_s]).and_return(user2.id.to_s) + + result = selector.select_agent(available_agents) + expect(result).to eq(user2) + end + end + end + + describe '#add_agent_to_queue' do + let(:selector) { described_class.new(inbox: inbox) } + let(:round_robin_service) { instance_double(AutoAssignment::InboxRoundRobinService) } + + before do + allow(AutoAssignment::InboxRoundRobinService).to receive(:new).with(inbox: inbox).and_return(round_robin_service) + end + + it 'delegates to round robin service' do + expect(round_robin_service).to receive(:add_agent_to_queue).with(user1.id) + selector.add_agent_to_queue(user1.id) + end + end + + describe '#remove_agent_from_queue' do + let(:selector) { described_class.new(inbox: inbox) } + let(:round_robin_service) { instance_double(AutoAssignment::InboxRoundRobinService) } + + before do + allow(AutoAssignment::InboxRoundRobinService).to receive(:new).with(inbox: inbox).and_return(round_robin_service) + end + + it 'delegates to round robin service' do + expect(round_robin_service).to receive(:remove_agent_from_queue).with(user1.id) + selector.remove_agent_from_queue(user1.id) + end + end + + describe '#reset_queue' do + let(:selector) { described_class.new(inbox: inbox) } + let(:round_robin_service) { instance_double(AutoAssignment::InboxRoundRobinService) } + + before do + allow(AutoAssignment::InboxRoundRobinService).to receive(:new).with(inbox: inbox).and_return(round_robin_service) + end + + it 'delegates to round robin service' do + expect(round_robin_service).to receive(:reset_queue) + selector.reset_queue + end + end + + describe 'edge cases' do + let(:selector) { described_class.new(inbox: inbox) } + let(:round_robin_service) { instance_double(AutoAssignment::InboxRoundRobinService) } + let(:available_agents) { InboxMember.where(inbox: inbox, user: [user1, user2]) } + + before do + allow(AutoAssignment::InboxRoundRobinService).to receive(:new).with(inbox: inbox).and_return(round_robin_service) + end + + it 'handles invalid user_id from round robin service' do + allow(round_robin_service).to receive(:available_agent).and_return('invalid_id') + + result = selector.select_agent(available_agents) + expect(result).to be_nil + end + + it 'handles user_id not in available agents' do + other_user = create(:user, account: account) + allow(round_robin_service).to receive(:available_agent).and_return(other_user.id.to_s) + + result = selector.select_agent(available_agents) + expect(result).to be_nil + end + end +end