assignment_v2: import assignment policy controllers/models/services/job/policy/specs

This commit is contained in:
Tanmay Sharma
2025-08-11 09:06:30 +05:30
parent 304c938260
commit b2ada112d7
22 changed files with 2214 additions and 90 deletions
@@ -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
@@ -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
@@ -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