diff --git a/app/controllers/concerns/meta_token_verify_concern.rb b/app/controllers/concerns/meta_token_verify_concern.rb index b3f920644..42fe918cc 100644 --- a/app/controllers/concerns/meta_token_verify_concern.rb +++ b/app/controllers/concerns/meta_token_verify_concern.rb @@ -2,6 +2,10 @@ # This concern handles the token verification step. module MetaTokenVerifyConcern + CHANNEL_APP_SECRET_KEYS = %w[app_secret app_secret_key client_secret api_secret].freeze + META_SIGNATURE_HEADER = 'X-Hub-Signature-256'.freeze + META_SIGNATURE_PREFIX = 'sha256='.freeze + def verify service = is_a?(Webhooks::WhatsappController) ? 'whatsapp' : 'instagram' if valid_token?(params['hub.verify_token']) @@ -14,6 +18,53 @@ module MetaTokenVerifyConcern private + def verify_meta_signature! + return unless meta_signature_verification_required? + return if valid_meta_signature? + + head :unauthorized + end + + def valid_meta_signature? + signature = request.headers[META_SIGNATURE_HEADER] + return false unless signature&.start_with?(META_SIGNATURE_PREFIX) + + meta_app_secrets.any? do |secret| + next false if secret.blank? + + expected_signature = "#{META_SIGNATURE_PREFIX}#{OpenSSL::HMAC.hexdigest('SHA256', secret, meta_request_body)}" + ActiveSupport::SecurityUtils.secure_compare(expected_signature, signature) + end + end + + def meta_request_body + @meta_request_body ||= request.raw_post + end + + def meta_app_secrets + raise 'Overwrite this method in your controller' + end + + def meta_signature_verification_required? + true + end + + def channel_meta_app_secrets(channel) + return [] if channel.blank? + + secrets = [] + secrets << channel.app_secret if channel.respond_to?(:app_secret) + secrets.concat(provider_config_meta_app_secrets(channel)) + secrets.compact_blank.uniq + end + + def provider_config_meta_app_secrets(channel) + return [] unless channel.respond_to?(:provider_config) + + provider_config = channel.provider_config.to_h.with_indifferent_access + CHANNEL_APP_SECRET_KEYS.filter_map { |key| provider_config[key].presence } + end + def valid_token?(_token) raise 'Overwrite this method your controller' end diff --git a/app/controllers/webhooks/instagram_controller.rb b/app/controllers/webhooks/instagram_controller.rb index 569c2524b..6e5168634 100644 --- a/app/controllers/webhooks/instagram_controller.rb +++ b/app/controllers/webhooks/instagram_controller.rb @@ -1,6 +1,8 @@ class Webhooks::InstagramController < ActionController::API include MetaTokenVerifyConcern + before_action :verify_meta_signature!, only: :events + def events Rails.logger.info('Instagram webhook received events') if params['object'].casecmp('instagram').zero? @@ -39,4 +41,38 @@ class Webhooks::InstagramController < ActionController::API token == GlobalConfigService.load('IG_VERIFY_TOKEN', '') || token == GlobalConfigService.load('INSTAGRAM_VERIFY_TOKEN', '') end + + def meta_app_secrets + [ + *instagram_channel_meta_app_secrets, + GlobalConfigService.load('INSTAGRAM_APP_SECRET', nil), + GlobalConfigService.load('FB_APP_SECRET', nil) + ] + end + + def instagram_channel_meta_app_secrets + instagram_channels_from_payload.flat_map { |channel| channel_meta_app_secrets(channel) } + end + + def instagram_channels_from_payload + Array(params.to_unsafe_hash[:entry]).flat_map do |entry| + instagram_ids_from_entry(entry.with_indifferent_access).flat_map do |instagram_id| + [ + Channel::Instagram.find_by(instagram_id: instagram_id), + Channel::FacebookPage.find_by(instagram_id: instagram_id) + ] + end + end.compact.uniq + end + + def instagram_ids_from_entry(entry) + messages = entry[:messaging].presence || entry[:standby] || [] + messages.filter_map { |messaging| instagram_id_from_messaging(messaging.with_indifferent_access) } + end + + def instagram_id_from_messaging(messaging) + return messaging.dig(:sender, :id) if messaging.dig(:message, :is_echo).present? + + messaging.dig(:recipient, :id) + end end diff --git a/app/controllers/webhooks/whatsapp_controller.rb b/app/controllers/webhooks/whatsapp_controller.rb index c4c376e5c..ee71f3c92 100644 --- a/app/controllers/webhooks/whatsapp_controller.rb +++ b/app/controllers/webhooks/whatsapp_controller.rb @@ -1,6 +1,8 @@ class Webhooks::WhatsappController < ActionController::API include MetaTokenVerifyConcern + before_action :verify_meta_signature!, only: :process_payload + def process_payload if inactive_whatsapp_number? Rails.logger.warn("Rejected webhook for inactive WhatsApp number: #{params[:phone_number]}") @@ -20,6 +22,45 @@ class Webhooks::WhatsappController < ActionController::API token == whatsapp_webhook_verify_token if whatsapp_webhook_verify_token.present? end + def meta_app_secrets + [ + *channel_meta_app_secrets(whatsapp_channel), + GlobalConfigService.load('WHATSAPP_APP_SECRET', nil) + ] + end + + def whatsapp_channel + @whatsapp_channel ||= whatsapp_business_payload_channel || Channel::Whatsapp.find_by(phone_number: params[:phone_number]) + end + + def meta_signature_verification_required? + return true if whatsapp_channel.blank? + return false unless whatsapp_channel.provider == 'whatsapp_cloud' + return true if channel_meta_app_secrets(whatsapp_channel).present? + + whatsapp_channel.provider_config['source'] == 'embedded_signup' + end + + def whatsapp_business_payload_channel + return unless params[:object] == 'whatsapp_business_account' + + metadata = params.dig(:entry, 0, :changes, 0, :value, :metadata) + return if metadata.blank? + + phone_number = normalized_phone_number(metadata[:display_phone_number]) + phone_number_id = metadata[:phone_number_id] + channel = Channel::Whatsapp.find_by(phone_number: phone_number) + + return channel if channel && channel.provider_config['phone_number_id'] == phone_number_id + end + + def normalized_phone_number(phone_number) + return if phone_number.blank? + + phone_number = phone_number.to_s + phone_number.start_with?('+') ? phone_number : "+#{phone_number}" + end + def inactive_whatsapp_number? phone_number = params[:phone_number] return false if phone_number.blank? diff --git a/config/features.yml b/config/features.yml index f16199932..92d004f81 100644 --- a/config/features.yml +++ b/config/features.yml @@ -144,10 +144,11 @@ - name: chatwoot_v4 display_name: Chatwoot V4 enabled: true -- name: report_v4 - display_name: Report V4 +- name: captain_v1_action_classifier + display_name: Captain V1 Action Classifier enabled: false - deprecated: true + premium: true + chatwoot_internal: true - name: contact_chatwoot_support_team display_name: Contact Chatwoot Support Team enabled: true diff --git a/config/initializers/facebook_messenger.rb b/config/initializers/facebook_messenger.rb index f93829e65..047715cfb 100644 --- a/config/initializers/facebook_messenger.rb +++ b/config/initializers/facebook_messenger.rb @@ -1,11 +1,13 @@ # ref: https://github.com/jgorset/facebook-messenger#make-a-configuration-provider class ChatwootFbProvider < Facebook::Messenger::Configuration::Providers::Base + CHANNEL_APP_SECRET_KEYS = %w[app_secret app_secret_key client_secret api_secret].freeze + def valid_verify_token?(_verify_token) GlobalConfigService.load('FB_VERIFY_TOKEN', '') end - def app_secret_for(_page_id) - GlobalConfigService.load('FB_APP_SECRET', '') + def app_secret_for(page_id) + channel_app_secret_for(page_id).presence || GlobalConfigService.load('FB_APP_SECRET', '') end def access_token_for(page_id) @@ -14,6 +16,27 @@ class ChatwootFbProvider < Facebook::Messenger::Configuration::Providers::Base private + def channel_app_secret_for(page_id) + channel = Channel::FacebookPage.where(page_id: page_id).last + return if channel.blank? + + channel_app_secret_candidates(channel).first + end + + def channel_app_secret_candidates(channel) + secrets = [] + secrets << channel.app_secret if channel.respond_to?(:app_secret) + secrets.concat(provider_config_app_secrets(channel)) + secrets.compact_blank.uniq + end + + def provider_config_app_secrets(channel) + return [] unless channel.respond_to?(:provider_config) + + provider_config = channel.provider_config.to_h.with_indifferent_access + CHANNEL_APP_SECRET_KEYS.filter_map { |key| provider_config[key].presence } + end + def bot Chatwoot::Bot end diff --git a/db/migrate/20260430114500_repurpose_report_v4_flag_for_captain_v1_action_classifier.rb b/db/migrate/20260430114500_repurpose_report_v4_flag_for_captain_v1_action_classifier.rb new file mode 100644 index 000000000..8f7056ba7 --- /dev/null +++ b/db/migrate/20260430114500_repurpose_report_v4_flag_for_captain_v1_action_classifier.rb @@ -0,0 +1,15 @@ +class RepurposeReportV4FlagForCaptainV1ActionClassifier < ActiveRecord::Migration[7.1] + def up + Account.feature_captain_v1_action_classifier.find_each(batch_size: 100) do |account| + account.disable_features(:captain_v1_action_classifier) + account.save!(validate: false) + end + + config = InstallationConfig.find_by(name: 'ACCOUNT_LEVEL_FEATURE_DEFAULTS') + return if config&.value.blank? + + config.value = config.value.reject { |feature| feature['name'] == 'report_v4' } + config.save! + GlobalConfig.clear_cache + end +end diff --git a/db/schema.rb b/db/schema.rb index 4ff250ada..9ce734ba7 100644 --- a/db/schema.rb +++ b/db/schema.rb @@ -10,7 +10,7 @@ # # It's strongly recommended that you check this file into your version control system. -ActiveRecord::Schema[7.1].define(version: 2026_04_28_120000) do +ActiveRecord::Schema[7.1].define(version: 2026_04_30_114500) do # These extensions should be enabled to support this database enable_extension "pg_stat_statements" enable_extension "pg_trgm" diff --git a/enterprise/app/helpers/captain/chat_response_helper.rb b/enterprise/app/helpers/captain/chat_response_helper.rb index bfb8adc11..e8ac1044c 100644 --- a/enterprise/app/helpers/captain/chat_response_helper.rb +++ b/enterprise/app/helpers/captain/chat_response_helper.rb @@ -35,7 +35,10 @@ module Captain::ChatResponseHelper def credit_used_for_response?(parsed_response) response = parsed_response['response'] - response.present? && response != 'conversation_handoff' + + # The classifier can still decide to hand off after this trace is written. + # Actual response usage is charged later in ResponseBuilderJob, so billing stays correct. + response.present? && response != 'conversation_handoff' && parsed_response['action'] != 'handoff' end def captain_v1_assistant? diff --git a/enterprise/app/jobs/captain/conversation/response_builder_job.rb b/enterprise/app/jobs/captain/conversation/response_builder_job.rb index 5e7c5b3c2..5050f11b2 100644 --- a/enterprise/app/jobs/captain/conversation/response_builder_job.rb +++ b/enterprise/app/jobs/captain/conversation/response_builder_job.rb @@ -1,4 +1,6 @@ class Captain::Conversation::ResponseBuilderJob < ApplicationJob + include Captain::Conversation::V1ActionClassifier + MAX_MESSAGE_LENGTH = 10_000 retry_on ActiveStorage::FileNotFoundError, attempts: 3, wait: 2.seconds retry_on Faraday::BadRequestError, attempts: 3, wait: 2.seconds @@ -31,9 +33,11 @@ class Captain::Conversation::ResponseBuilderJob < ApplicationJob delegate :account, :inbox, to: :@conversation def generate_and_process_response + message_history = collect_previous_messages @response = Captain::Llm::AssistantChatService.new(assistant: @assistant, conversation: @conversation).generate_response( - message_history: collect_previous_messages + message_history: message_history ) + classify_v1_response_action(message_history) if conversation_pending? process_response end @@ -102,6 +106,14 @@ class Captain::Conversation::ResponseBuilderJob < ApplicationJob end def v1_handoff_requested? + legacy_v1_handoff_token? || classifier_v1_handoff_requested? + end + + def classifier_v1_handoff_requested? + @response['action'] == 'handoff' + end + + def legacy_v1_handoff_token? @response['response'] == 'conversation_handoff' end @@ -111,8 +123,13 @@ class Captain::Conversation::ResponseBuilderJob < ApplicationJob def process_v1_handoff I18n.with_locale(@assistant.account.locale) do + Rails.logger.info( + "[CAPTAIN][ResponseBuilderJob] V1 handoff requested for account=#{account.id} conversation=#{@conversation.display_id} " \ + "source=#{@response&.dig('action_source') || 'legacy'} reason=#{@response&.dig('action_reason')}" + ) create_handoff_message @conversation.bot_handoff! + report_v1_handoff_not_executed if conversation_pending? send_out_of_office_message_if_applicable end end @@ -166,6 +183,9 @@ class Captain::Conversation::ResponseBuilderJob < ApplicationJob def handle_error(error) log_error(error) + @response ||= {} + @response['action_source'] ||= 'error' + @response['action_reason'] ||= error_action_reason(error) process_v1_handoff if conversation_pending? true end @@ -174,10 +194,23 @@ class Captain::Conversation::ResponseBuilderJob < ApplicationJob ChatwootExceptionTracker.new(error, account: account).capture_exception end + def error_action_reason(error) + error.class.name.underscore.tr('/', '_') + end + def captain_v2_enabled? account.feature_enabled?('captain_integration_v2') end + def report_v1_handoff_not_executed + error = StandardError.new("Captain V1 handoff requested but conversation #{@conversation.display_id} is still pending") + ChatwootExceptionTracker.new(error, account: account).capture_exception + Rails.logger.error( + "[CAPTAIN][ResponseBuilderJob] V1 handoff requested but not executed for account=#{account.id} " \ + "conversation=#{@conversation.display_id}" + ) + end + def conversation_pending? status = Conversation.uncached { Conversation.where(id: @conversation.id).pick(:status) } status == 'pending' || status == Conversation.statuses[:pending] diff --git a/enterprise/app/jobs/captain/conversation/v1_action_classifier.rb b/enterprise/app/jobs/captain/conversation/v1_action_classifier.rb new file mode 100644 index 000000000..9d2010b79 --- /dev/null +++ b/enterprise/app/jobs/captain/conversation/v1_action_classifier.rb @@ -0,0 +1,57 @@ +module Captain::Conversation::V1ActionClassifier + private + + def v1_action_classifier_enabled? + account.feature_enabled?('captain_v1_action_classifier') + end + + def classify_v1_response_action(message_history) + return unless v1_action_classifier_enabled? + return if legacy_v1_handoff_token? + + classification = Captain::Llm::AssistantActionClassifierService.new( + assistant: @assistant, + conversation: @conversation + ).classify(message_history: message_history, assistant_response: @response['response']) + + apply_v1_action_classification(classification) + rescue StandardError => e + ChatwootExceptionTracker.new(e, account: account).capture_exception + Rails.logger.warn( + "[CAPTAIN][ResponseBuilderJob] V1 action classifier failed for account=#{account.id} " \ + "conversation=#{@conversation.display_id}: #{e.class.name}: #{e.message}" + ) + end + + def apply_v1_action_classification(classification) + action = classification['action'] + return log_invalid_v1_action_classification(classification) unless valid_v1_action_classification?(action) + + @response.merge!( + 'action' => action, + 'action_reason' => classification['action_reason'], + 'action_source' => 'classifier', + 'action_classifier_model' => classification['model'] + ) + + log_v1_action_classification(action, classification) + end + + def log_v1_action_classification(action, classification) + Rails.logger.info( + "[CAPTAIN][ResponseBuilderJob] V1 action classifier account=#{account.id} conversation=#{@conversation.display_id} " \ + "action=#{action} reason=#{classification['action_reason']} model=#{classification['model']}" + ) + end + + def valid_v1_action_classification?(action) + Captain::AssistantActionSchema::ACTIONS.include?(action) + end + + def log_invalid_v1_action_classification(classification) + Rails.logger.warn( + '[CAPTAIN][ResponseBuilderJob] V1 action classifier returned invalid action; falling back to assistant response ' \ + "for account=#{account.id} conversation=#{@conversation.display_id}: #{classification['error'] || classification['raw_response']}" + ) + end +end diff --git a/enterprise/app/services/captain/llm/assistant_action_classifier_service.rb b/enterprise/app/services/captain/llm/assistant_action_classifier_service.rb new file mode 100644 index 000000000..7c0f1e91e --- /dev/null +++ b/enterprise/app/services/captain/llm/assistant_action_classifier_service.rb @@ -0,0 +1,148 @@ +class Captain::Llm::AssistantActionClassifierService < Llm::BaseAiService + include Integrations::LlmInstrumentation + + MAX_CONTEXT_MESSAGES = 10 + + def initialize(assistant:, conversation:) + super() + @assistant = assistant + @conversation = conversation + @temperature = 0.0 + end + + def classify(message_history:, assistant_response:) + user_prompt = classification_user_prompt( + message_history: message_history, + assistant_response: assistant_response + ) + + response = instrument_llm_call(instrumentation_params(user_prompt)) do + chat(model: @model, temperature: @temperature) + .with_schema(Captain::AssistantActionSchema) + .with_instructions(system_prompt) + .ask(user_prompt) + end + + parsed = parse_response(response.content) + normalize_response(parsed, response.content) + rescue StandardError => e + ChatwootExceptionTracker.new(e, account: @conversation.account).capture_exception + Rails.logger.warn( + "[CAPTAIN][AssistantActionClassifier] Failed for conversation #{@conversation.display_id}: #{e.class.name}: #{e.message}" + ) + { 'action' => nil, 'action_reason' => nil, 'error' => e.message, 'model' => @model } + end + + private + + def classification_user_prompt(message_history:, assistant_response:) + <<~PROMPT + + #{@assistant.config['instructions']} + + + + #{format_conversation_context(message_history)} + + + + #{assistant_response} + + PROMPT + end + + def normalize_messages(message_history) + message_history.filter_map do |message| + role = message[:role] || message['role'] + next if role.blank? + + { role: role.to_s, content: normalize_content(message[:content] || message['content']) } + end + end + + def normalize_content(content) + return content if content.is_a?(String) + return content.filter_map { |part| part[:text] || part['text'] if text_part?(part) }.join("\n") if content.is_a?(Array) + + content.to_s + end + + def text_part?(part) + return false unless part.is_a?(Hash) + + (part[:type] || part['type']).to_s == 'text' + end + + def format_conversation_context(messages) + normalize_messages(messages).last(MAX_CONTEXT_MESSAGES).filter_map do |message| + content = message[:content].to_s.strip + next if content.blank? + + "#{role_label(message[:role])}: #{content}" + end.join("\n") + end + + def role_label(role) + return 'User' if role == 'user' + return 'Assistant' if role == 'assistant' + + role.to_s.titleize + end + + def parse_response(content) + return content if content.is_a?(Hash) + + JSON.parse(sanitize_json_response(content)) + rescue JSON::ParserError, TypeError + {} + end + + def normalize_response(parsed, raw_content) + action = parsed['action'].to_s + reason = parsed['action_reason'].to_s + return invalid_response(raw_content) unless Captain::AssistantActionSchema::ACTIONS.include?(action) + + { + 'action' => action, + 'action_reason' => reason.presence, + 'raw_response' => raw_content, + 'model' => @model + } + end + + def invalid_response(raw_content) + { + 'action' => nil, + 'action_reason' => nil, + 'raw_response' => raw_content, + 'error' => 'invalid_classifier_response', + 'model' => @model + } + end + + def instrumentation_params(user_prompt) + { + span_name: 'llm.captain.assistant_action_classifier', + model: @model, + temperature: @temperature, + account_id: @conversation.account_id, + conversation_id: @conversation.display_id, + feature_name: 'assistant_action_classifier', + messages: [ + { role: 'system', content: system_prompt }, + { role: 'user', content: user_prompt } + ], + metadata: { + assistant_id: @assistant.id, + channel_type: @conversation.inbox&.channel_type, + source: 'v1_response_builder' + } + } + end + + def system_prompt + Captain::Llm::SystemPromptsService.assistant_action_classifier( + has_custom_instructions: @assistant.config['instructions'].present? + ) + end +end diff --git a/enterprise/app/services/captain/llm/system_prompts_service.rb b/enterprise/app/services/captain/llm/system_prompts_service.rb index eb8c334f4..9520330f6 100644 --- a/enterprise/app/services/captain/llm/system_prompts_service.rb +++ b/enterprise/app/services/captain/llm/system_prompts_service.rb @@ -93,6 +93,50 @@ class Captain::Llm::SystemPromptsService SYSTEM_PROMPT_MESSAGE end + def assistant_action_classifier(has_custom_instructions: false) + <<~PROMPT + You are a routing classifier for a customer-support assistant. + + Decide whether the current conversation should stay with the assistant or be transferred to a human agent now. + + The action field MUST be one of: + - "continue": keep the current conversation with the assistant. + - "handoff": transfer the current conversation to a human agent now. + + The action_reason field MUST be one of: + - "general_product_question" + - "missing_docs_bounded_answer" + - "clarifying_question_needed" + - "collect_required_identifier" + - "external_contact_or_lead_routing" + - "out_of_scope_bounded_answer" + - "explicit_human_request" + - "human_offer_accepted" + - "account_or_transaction_verification" + - "operational_issue_needs_inspection" + - "repeated_frustration_or_loop" + - "custom_instruction_transfer" + + Use "continue" when: + - The user has a general product, pricing, capability, setup, pre-sales, or how-to question. + - The assistant can give a bounded answer, ask one useful clarifying question, collect a missing identifier, or share an approved external contact path. + - The assistant says someone will contact the user outside this conversation, but the current conversation itself does not need to be transferred now. + - The user has not explicitly asked for a human and the assistant is still collecting required details. + + Use "handoff" when: + - The user explicitly asks for a human, agent, representative, phone call, callback, or escalation. + - The user accepts an offer to speak with a human. + - The user has provided enough detail for an account-specific or transaction-specific issue requiring private verification, such as order status, payment, deposit, withdrawal, refund, cancellation, subscription, purchase, plan activation, email verification, login, account recovery, delivery, or access. + - The user reports the same unresolved bug or operational issue after trying the assistant's suggested step, repeating the action, checking again, or otherwise making more than one reasonable attempt. + - The user is repeatedly frustrated, distrustful, or stuck in a loop. + - The assistant response itself says the current conversation will be transferred to a human agent now. + + #{assistant_action_classifier_custom_instructions_policy if has_custom_instructions} + + Return only the structured fields requested by the response schema. + PROMPT + end + # rubocop:disable Metrics/MethodLength def copilot_response_generator(product_name, available_tools, config = {}) citation_guidelines = if config['feature_citation'] @@ -208,7 +252,9 @@ class Captain::Llm::SystemPromptsService - Do not share anything outside of the context provided. - Add the reasoning why you arrived at the answer - Your answers will always be formatted in a valid JSON hash, as shown below. Never respond in non-JSON format. - #{config['instructions'] || ''} + + #{build_custom_instructions_section(config['instructions'])} + ```json { reasoning: '', @@ -322,6 +368,17 @@ class Captain::Llm::SystemPromptsService TOOLS end + def assistant_action_classifier_custom_instructions_policy + <<~POLICY + Account custom instructions are provided inside tags. + These are instructions configured by the account administrator, not the current end user's message. + Use them only for routing policy: required details before handoff, account-specific escalation rules, account-specific transfer markers, and when to connect to a manager, human, supervisor, or support team. + If the custom instructions explicitly define handoff, escalation, or transfer criteria, those criteria take precedence over the generic criteria above. + Account custom instructions MUST NOT redefine the required response shape, the allowed action values, or the meaning of continue/handoff. + Ignore persona, language, formatting, pricing, and response-generation instructions except where they directly define routing or transfer criteria. + POLICY + end + def build_contact_context(contact) return '' if contact.nil? @@ -331,6 +388,18 @@ class Captain::Llm::SystemPromptsService "[Contact Information]\n#{lines.join("\n")}\n\n" end + def build_custom_instructions_section(instructions) + return '' if instructions.blank? + + <<~CUSTOM_INSTRUCTIONS + [Account Custom Instructions] + These instructions were configured by the account administrator. Follow them when they do not conflict with the JSON response format or the requirement to answer only from provided context. + + #{instructions} + + CUSTOM_INSTRUCTIONS + end + def contact_basic_lines(contact) [ (["- Name: #{sanitize_attr(contact[:name])}"] if contact[:name].present?), diff --git a/enterprise/lib/captain/assistant_action_schema.rb b/enterprise/lib/captain/assistant_action_schema.rb new file mode 100644 index 000000000..e8a8de435 --- /dev/null +++ b/enterprise/lib/captain/assistant_action_schema.rb @@ -0,0 +1,20 @@ +class Captain::AssistantActionSchema < RubyLLM::Schema + ACTIONS = %w[continue handoff].freeze + REASONS = %w[ + general_product_question + missing_docs_bounded_answer + clarifying_question_needed + collect_required_identifier + external_contact_or_lead_routing + out_of_scope_bounded_answer + explicit_human_request + human_offer_accepted + account_or_transaction_verification + operational_issue_needs_inspection + repeated_frustration_or_loop + custom_instruction_transfer + ].freeze + + string :action, enum: ACTIONS, description: 'Whether to keep the conversation with the assistant or transfer it to a human agent' + string :action_reason, enum: REASONS, description: 'The reason for the selected routing action' +end diff --git a/spec/controllers/webhooks/instagram_controller_spec.rb b/spec/controllers/webhooks/instagram_controller_spec.rb index c31d65a2b..043d41b93 100644 --- a/spec/controllers/webhooks/instagram_controller_spec.rb +++ b/spec/controllers/webhooks/instagram_controller_spec.rb @@ -1,6 +1,25 @@ require 'rails_helper' RSpec.describe 'Webhooks::InstagramController', type: :request do + let(:client_secret) { 'test-instagram-secret' } + + def signature_for(body, secret = client_secret) + "sha256=#{OpenSSL::HMAC.hexdigest('SHA256', secret, body)}" + end + + def post_instagram_webhook(body, signature: signature_for(body), env: { INSTAGRAM_APP_SECRET: client_secret }) + with_modified_env env do + post '/webhooks/instagram', + params: body, + headers: { 'CONTENT_TYPE' => 'application/json', 'X-Hub-Signature-256' => signature } + end + end + + before do + InstallationConfig.where(name: %w[FB_APP_SECRET IG_VERIFY_TOKEN INSTAGRAM_APP_SECRET INSTAGRAM_VERIFY_TOKEN]).delete_all + GlobalConfig.clear_cache + end + describe 'GET /webhooks/verify' do it 'returns 401 when valid params are not present' do get '/webhooks/instagram/verify' @@ -24,26 +43,62 @@ RSpec.describe 'Webhooks::InstagramController', type: :request do describe 'POST /webhooks/instagram' do let!(:dm_params) { build(:instagram_message_create_event).with_indifferent_access } + let(:body) { dm_params.merge(object: 'instagram').to_json } - it 'call the instagram events job with the params' do + it 'calls the instagram events job with the params for a valid signature' do allow(Webhooks::InstagramEventsJob).to receive(:perform_later) expect(Webhooks::InstagramEventsJob).to receive(:perform_later) - instagram_params = dm_params.merge(object: 'instagram') - post '/webhooks/instagram', params: instagram_params + post_instagram_webhook(body) expect(response).to have_http_status(:success) end + it 'accepts webhook payloads signed with the Facebook app secret' do + allow(Webhooks::InstagramEventsJob).to receive(:perform_later) + expect(Webhooks::InstagramEventsJob).to receive(:perform_later) + + facebook_secret = 'test-facebook-secret' + post_instagram_webhook( + body, + signature: signature_for(body, facebook_secret), + env: { FB_APP_SECRET: facebook_secret } + ) + + expect(response).to have_http_status(:success) + end + + it 'returns unauthorized when signature is missing' do + allow(Webhooks::InstagramEventsJob).to receive(:perform_later) + + with_modified_env INSTAGRAM_APP_SECRET: client_secret do + post '/webhooks/instagram', + params: body, + headers: { 'CONTENT_TYPE' => 'application/json' } + end + + expect(response).to have_http_status(:unauthorized) + expect(Webhooks::InstagramEventsJob).not_to have_received(:perform_later) + end + + it 'returns unauthorized when signature is invalid' do + allow(Webhooks::InstagramEventsJob).to receive(:perform_later) + + post_instagram_webhook(body, signature: 'sha256=invalid-signature') + + expect(response).to have_http_status(:unauthorized) + expect(Webhooks::InstagramEventsJob).not_to have_received(:perform_later) + end + context 'when processing echo events' do let!(:echo_params) { build(:instagram_story_mention_event_with_echo).with_indifferent_access } + let(:echo_body) { echo_params.merge(object: 'instagram').to_json } it 'delays processing for echo events by 2 seconds' do job_double = class_double(Webhooks::InstagramEventsJob) allow(Webhooks::InstagramEventsJob).to receive(:set).with(wait: 2.seconds).and_return(job_double) allow(job_double).to receive(:perform_later) - instagram_params = echo_params.merge(object: 'instagram') - post '/webhooks/instagram', params: instagram_params + post_instagram_webhook(echo_body) expect(response).to have_http_status(:success) expect(Webhooks::InstagramEventsJob).to have_received(:set).with(wait: 2.seconds) expect(job_double).to have_received(:perform_later) diff --git a/spec/controllers/webhooks/whatsapp_controller_spec.rb b/spec/controllers/webhooks/whatsapp_controller_spec.rb index e5fa392b0..05816094b 100644 --- a/spec/controllers/webhooks/whatsapp_controller_spec.rb +++ b/spec/controllers/webhooks/whatsapp_controller_spec.rb @@ -2,6 +2,33 @@ require 'rails_helper' RSpec.describe 'Webhooks::WhatsappController', type: :request do let(:channel) { create(:channel_whatsapp, provider: 'whatsapp_cloud', sync_templates: false, validate_provider_config: false) } + let(:client_secret) { 'test-whatsapp-secret' } + let(:body) { { content: 'hello' }.to_json } + + def signature_for(body, secret = client_secret) + "sha256=#{OpenSSL::HMAC.hexdigest('SHA256', secret, body)}" + end + + def post_whatsapp_webhook(path, body, signature: signature_for(body), env: { WHATSAPP_APP_SECRET: client_secret }) + with_modified_env env do + post path, + params: body, + headers: { 'CONTENT_TYPE' => 'application/json', 'X-Hub-Signature-256' => signature } + end + end + + def post_unsigned_whatsapp_webhook(path, body, env: { WHATSAPP_APP_SECRET: client_secret }) + with_modified_env env do + post path, + params: body, + headers: { 'CONTENT_TYPE' => 'application/json' } + end + end + + before do + InstallationConfig.where(name: 'WHATSAPP_APP_SECRET').delete_all + GlobalConfig.clear_cache + end describe 'GET /webhooks/verify' do it 'returns 401 when valid params are not present' do @@ -23,13 +50,103 @@ RSpec.describe 'Webhooks::WhatsappController', type: :request do end describe 'POST /webhooks/whatsapp/{:phone_number}' do - it 'call the whatsapp events job with the params' do + it 'calls the whatsapp events job with the params for a valid signature' do allow(Webhooks::WhatsappEventsJob).to receive(:perform_later) expect(Webhooks::WhatsappEventsJob).to receive(:perform_later) - post '/webhooks/whatsapp/123221321', params: { content: 'hello' } + post_whatsapp_webhook('/webhooks/whatsapp/123221321', body) expect(response).to have_http_status(:success) end + it 'accepts webhook payloads signed with the channel app secret' do + channel_secret = 'channel-whatsapp-secret' + channel.provider_config = channel.provider_config.merge('app_secret' => channel_secret) + channel.save! + + allow(Webhooks::WhatsappEventsJob).to receive(:perform_later) + expect(Webhooks::WhatsappEventsJob).to receive(:perform_later) + + channel_body = { + object: 'whatsapp_business_account', + entry: [{ + changes: [{ + value: { + metadata: { + display_phone_number: channel.phone_number.delete_prefix('+'), + phone_number_id: channel.provider_config['phone_number_id'] + } + } + }] + }] + }.to_json + + post_whatsapp_webhook( + "/webhooks/whatsapp/#{channel.phone_number}", + channel_body, + signature: signature_for(channel_body, channel_secret), + env: {} + ) + + expect(response).to have_http_status(:success) + end + + it 'skips signature validation for 360dialog channels' do + dialog_channel = create(:channel_whatsapp, provider: 'default', sync_templates: false, validate_provider_config: false) + allow(Webhooks::WhatsappEventsJob).to receive(:perform_later) + expect(Webhooks::WhatsappEventsJob).to receive(:perform_later) + + post_unsigned_whatsapp_webhook("/webhooks/whatsapp/#{dialog_channel.phone_number}", body) + + expect(response).to have_http_status(:success) + end + + it 'skips signature validation for manual whatsapp cloud channels without an app secret' do + channel.update!( + provider_config: channel.provider_config.except('app_secret', 'app_secret_key', 'api_secret', 'client_secret', 'source') + ) + allow(Webhooks::WhatsappEventsJob).to receive(:perform_later) + expect(Webhooks::WhatsappEventsJob).to receive(:perform_later) + + channel_body = { + object: 'whatsapp_business_account', + entry: [{ + changes: [{ + value: { + metadata: { + display_phone_number: channel.phone_number.delete_prefix('+'), + phone_number_id: channel.provider_config['phone_number_id'] + } + } + }] + }] + }.to_json + + post_unsigned_whatsapp_webhook("/webhooks/whatsapp/#{channel.phone_number}", channel_body) + + expect(response).to have_http_status(:success) + end + + it 'returns unauthorized when signature is missing' do + allow(Webhooks::WhatsappEventsJob).to receive(:perform_later) + + with_modified_env WHATSAPP_APP_SECRET: client_secret do + post '/webhooks/whatsapp/123221321', + params: body, + headers: { 'CONTENT_TYPE' => 'application/json' } + end + + expect(response).to have_http_status(:unauthorized) + expect(Webhooks::WhatsappEventsJob).not_to have_received(:perform_later) + end + + it 'returns unauthorized when signature is invalid' do + allow(Webhooks::WhatsappEventsJob).to receive(:perform_later) + + post_whatsapp_webhook('/webhooks/whatsapp/123221321', body, signature: 'sha256=invalid-signature') + + expect(response).to have_http_status(:unauthorized) + expect(Webhooks::WhatsappEventsJob).not_to have_received(:perform_later) + end + context 'when phone number is in inactive list' do before do allow(GlobalConfig).to receive(:get_value).with('INACTIVE_WHATSAPP_NUMBERS').and_return('+1234567890,+9876543210') @@ -39,7 +156,7 @@ RSpec.describe 'Webhooks::WhatsappController', type: :request do allow(Rails.logger).to receive(:warn) expect(Rails.logger).to receive(:warn).with('Rejected webhook for inactive WhatsApp number: +1234567890') - post '/webhooks/whatsapp/+1234567890', params: { content: 'hello' } + post_whatsapp_webhook('/webhooks/whatsapp/+1234567890', body) expect(response).to have_http_status(:unprocessable_entity) expect(response.parsed_body['error']).to eq('Inactive WhatsApp number') end @@ -54,7 +171,7 @@ RSpec.describe 'Webhooks::WhatsappController', type: :request do allow(Webhooks::WhatsappEventsJob).to receive(:perform_later) expect(Webhooks::WhatsappEventsJob).to receive(:perform_later) - post '/webhooks/whatsapp/+1234567890', params: { content: 'hello' } + post_whatsapp_webhook('/webhooks/whatsapp/+1234567890', body) expect(response).to have_http_status(:success) end end diff --git a/spec/enterprise/jobs/captain/conversation/response_builder_job_spec.rb b/spec/enterprise/jobs/captain/conversation/response_builder_job_spec.rb index 548b84992..8fac81d60 100644 --- a/spec/enterprise/jobs/captain/conversation/response_builder_job_spec.rb +++ b/spec/enterprise/jobs/captain/conversation/response_builder_job_spec.rb @@ -10,6 +10,7 @@ RSpec.describe Captain::Conversation::ResponseBuilderJob, type: :job do let(:conversation) { create(:conversation, inbox: inbox, account: account, status: :pending) } let(:mock_llm_chat_service) { instance_double(Captain::Llm::AssistantChatService) } let(:mock_agent_runner_service) { instance_double(Captain::Assistant::AgentRunnerService) } + let(:mock_action_classifier_service) { instance_double(Captain::Llm::AssistantActionClassifierService) } before do create(:message, conversation: conversation, content: 'Hello', message_type: :incoming) @@ -19,6 +20,8 @@ RSpec.describe Captain::Conversation::ResponseBuilderJob, type: :job do allow(mock_llm_chat_service).to receive(:generate_response).and_return({ 'response' => 'Hey, welcome to Captain Specs' }) allow(Captain::Assistant::AgentRunnerService).to receive(:new).and_return(mock_agent_runner_service) allow(mock_agent_runner_service).to receive(:generate_response).and_return({ 'response' => 'Hey, welcome to Captain V2' }) + allow(Captain::Llm::AssistantActionClassifierService).to receive(:new).and_return(mock_action_classifier_service) + allow(mock_action_classifier_service).to receive(:classify).and_return({ 'action' => 'continue' }) end context 'when captain_v2 is disabled' do @@ -48,6 +51,107 @@ RSpec.describe Captain::Conversation::ResponseBuilderJob, type: :job do expect(account.usage_limits[:captain][:responses][:consumed]).to eq(1) end + it 'does not run the action classifier when the classifier feature is disabled' do + expect(Captain::Llm::AssistantActionClassifierService).not_to receive(:new) + + described_class.perform_now(conversation, assistant) + + expect(conversation.messages.last.content).to eq('Hey, welcome to Captain Specs') + end + + context 'when V1 action classifier is enabled' do + before do + allow(account).to receive(:feature_enabled?).and_return(false) + allow(account).to receive(:feature_enabled?).with('captain_integration_v2').and_return(false) + allow(account).to receive(:feature_enabled?).with('captain_v1_action_classifier').and_return(true) + end + + it 'keeps the conversation pending when the classifier returns continue' do + expect(Captain::Llm::AssistantActionClassifierService).to receive(:new).with( + assistant: assistant, + conversation: conversation + ).and_return(mock_action_classifier_service) + expect(mock_action_classifier_service).to receive(:classify).with( + message_history: [{ content: 'Hello', role: 'user' }], + assistant_response: 'Hey, welcome to Captain Specs' + ).and_return({ + 'action' => 'continue', + 'action_reason' => 'general_product_question', + 'model' => 'gpt-4.1' + }) + + described_class.perform_now(conversation, assistant) + + expect(conversation.reload.status).to eq('pending') + expect(conversation.messages.outgoing.last.content).to eq('Hey, welcome to Captain Specs') + expect(account.reload.usage_limits[:captain][:responses][:consumed]).to eq(1) + end + + it 'hands off without incrementing response usage when the classifier returns handoff' do + allow(mock_action_classifier_service).to receive(:classify).and_return({ + 'action' => 'handoff', + 'action_reason' => 'explicit_human_request', + 'model' => 'gpt-4.1' + }) + + described_class.perform_now(conversation, assistant) + + expect(conversation.reload.status).to eq('open') + expect(conversation.messages.outgoing.last.content).to eq(I18n.t('conversations.captain.handoff')) + expect(account.reload.usage_limits[:captain][:responses][:consumed]).to eq(0) + end + + it 'skips the classifier when the legacy handoff token is returned' do + allow(mock_llm_chat_service).to receive(:generate_response).and_return({ 'response' => 'conversation_handoff' }) + expect(Captain::Llm::AssistantActionClassifierService).not_to receive(:new) + + described_class.perform_now(conversation, assistant) + + expect(conversation.reload.status).to eq('open') + expect(conversation.messages.outgoing.last.content).to eq(I18n.t('conversations.captain.handoff')) + end + + it 'falls back to the assistant response when the classifier fails' do + error = StandardError.new('classifier unavailable') + allow(mock_action_classifier_service).to receive(:classify).and_raise(error) + allow(ChatwootExceptionTracker).to receive(:new).and_call_original + + described_class.perform_now(conversation, assistant) + + expect(ChatwootExceptionTracker).to have_received(:new).with(error, account: account) + expect(conversation.reload.status).to eq('pending') + expect(conversation.messages.outgoing.last.content).to eq('Hey, welcome to Captain Specs') + expect(account.reload.usage_limits[:captain][:responses][:consumed]).to eq(1) + end + + it 'falls back to the assistant response when the classifier returns an invalid action' do + allow(mock_action_classifier_service).to receive(:classify).and_return({ + 'action' => nil, + 'error' => 'invalid_classifier_response' + }) + + described_class.perform_now(conversation, assistant) + + expect(conversation.reload.status).to eq('pending') + expect(conversation.messages.outgoing.last.content).to eq('Hey, welcome to Captain Specs') + expect(account.reload.usage_limits[:captain][:responses][:consumed]).to eq(1) + end + + it 'skips the classifier when the conversation is no longer pending after response generation' do + allow(mock_llm_chat_service).to receive(:generate_response) do + conversation.open! + { 'response' => 'Hey, welcome to Captain Specs' } + end + + expect(Captain::Llm::AssistantActionClassifierService).not_to receive(:new) + + described_class.perform_now(conversation, assistant) + + expect(conversation.messages.outgoing.count).to eq(0) + expect(account.reload.usage_limits[:captain][:responses][:consumed]).to eq(0) + end + end + it 'does not send a response when the conversation is no longer pending' do conversation.open! @@ -292,9 +396,11 @@ RSpec.describe Captain::Conversation::ResponseBuilderJob, type: :job do it 'handles API errors and triggers handoff' do allow(mock_llm_chat_service).to receive(:generate_response) .and_raise(Faraday::BadRequestError, 'Bad request to image service') + allow(Rails.logger).to receive(:info).and_call_original described_class.perform_now(conversation, assistant) expect(conversation.reload.status).to eq('open') + expect(Rails.logger).to have_received(:info).with(include('source=error reason=faraday_bad_request_error')) end it 'succeeds when no error occurs' do diff --git a/spec/enterprise/services/captain/llm/assistant_action_classifier_service_spec.rb b/spec/enterprise/services/captain/llm/assistant_action_classifier_service_spec.rb new file mode 100644 index 000000000..260e3f4f7 --- /dev/null +++ b/spec/enterprise/services/captain/llm/assistant_action_classifier_service_spec.rb @@ -0,0 +1,95 @@ +require 'rails_helper' + +RSpec.describe Captain::Llm::AssistantActionClassifierService do + let(:account) { create(:account) } + let(:assistant) do + create( + :captain_assistant, + account: account, + config: { + 'instructions' => 'Only transfer to a manager after the user explicitly confirms.' + } + ) + end + let(:conversation) { create(:conversation, account: account) } + let(:service) { described_class.new(assistant: assistant, conversation: conversation) } + let(:mock_chat) { instance_double(RubyLLM::Chat) } + let(:mock_response) do + instance_double( + RubyLLM::Message, + content: { 'action' => 'handoff', 'action_reason' => 'human_offer_accepted' } + ) + end + + before do + allow(RubyLLM).to receive(:chat).and_return(mock_chat) + allow(mock_chat).to receive(:with_temperature).and_return(mock_chat) + allow(mock_chat).to receive(:with_schema).and_return(mock_chat) + allow(mock_chat).to receive(:with_instructions).and_return(mock_chat) + end + + describe '#classify' do + let(:message_history) do + [ + { role: 'user', content: 'I cannot log in' }, + { role: 'assistant', content: 'Did you check your inbox?' }, + { role: 'user', content: 'Yes, still no reset email' } + ] + end + + it 'passes delimited custom instructions and classifier context to the LLM' do + expect(mock_chat).to receive(:with_schema).with(Captain::AssistantActionSchema).and_return(mock_chat) + expect(mock_chat).to receive(:with_instructions).with( + a_string_including('Account custom instructions are provided inside tags.') + ).and_return(mock_chat) + expect(mock_chat).to receive(:ask) do |prompt| + expect(prompt).to include( + '', + 'Only transfer to a manager after the user explicitly confirms.', + '', + 'User: I cannot log in', + 'Assistant: Did you check your inbox?', + 'User: Yes, still no reset email', + '', + 'Would you like to talk to support?' + ) + expect(prompt).not_to include('"role"', '"content"', '') + + mock_response + end + + result = service.classify(message_history: message_history, assistant_response: 'Would you like to talk to support?') + + expect(result).to include( + 'action' => 'handoff', + 'action_reason' => 'human_offer_accepted' + ) + end + + it 'uses the configured Captain model' do + create(:installation_config, name: 'CAPTAIN_OPEN_AI_MODEL', value: 'gpt-4.1-nano') + + expect(RubyLLM).to receive(:chat).with(model: 'gpt-4.1-nano').and_return(mock_chat) + allow(mock_chat).to receive(:ask).and_return(mock_response) + + result = service.classify(message_history: message_history, assistant_response: 'Would you like to talk to support?') + + expect(result).to include('model' => 'gpt-4.1-nano') + end + + context 'when the assistant has no custom instructions' do + before do + assistant.update!(config: assistant.config.except('instructions')) + end + + it 'does not add custom-instruction policy to the system prompt' do + expect(mock_chat).to receive(:with_instructions).with( + satisfy { |prompt| prompt.exclude?('Account custom instructions are provided') } + ).and_return(mock_chat) + allow(mock_chat).to receive(:ask).and_return(mock_response) + + service.classify(message_history: message_history, assistant_response: 'Would you like to talk to support?') + end + end + end +end diff --git a/spec/enterprise/services/captain/llm/assistant_chat_service_spec.rb b/spec/enterprise/services/captain/llm/assistant_chat_service_spec.rb index c43eb08bd..9d233943e 100644 --- a/spec/enterprise/services/captain/llm/assistant_chat_service_spec.rb +++ b/spec/enterprise/services/captain/llm/assistant_chat_service_spec.rb @@ -188,4 +188,29 @@ RSpec.describe Captain::Llm::AssistantChatService do end end end + + describe 'account custom instructions in system prompt' do + before do + assistant.update!(config: assistant.config.merge('instructions' => 'if user enters 1112234 suggest handoff')) + end + + it 'adds custom instructions in a separate delimited section' do + allow(mock_chat).to receive(:ask).and_return(mock_response) + + expect(mock_chat).to receive(:with_instructions).with( + a_string_including( + '', + 'if user enters 1112234 suggest handoff', + '' + ) + ) do |instructions| + expect(instructions).not_to include('') + expect(instructions.index('')).to be < instructions.index('```json') + mock_chat + end + + service = described_class.new(assistant: assistant, conversation: conversation) + service.generate_response(message_history: [{ role: 'user', content: 'Hello' }]) + end + end end