Merge branch 'develop' into fix/CW-6944

This commit is contained in:
Sivin Varghese
2026-05-05 15:55:39 +05:30
committed by GitHub
18 changed files with 913 additions and 18 deletions
@@ -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
@@ -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
@@ -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?
+4 -3
View File
@@ -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
+25 -2
View File
@@ -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
@@ -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
+1 -1
View File
@@ -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"
@@ -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?
@@ -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]
@@ -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
@@ -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
<account_custom_instructions>
#{@assistant.config['instructions']}
</account_custom_instructions>
<conversation_context>
#{format_conversation_context(message_history)}
</conversation_context>
<assistant_response_to_classify>
#{assistant_response}
</assistant_response_to_classify>
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
@@ -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 <account_custom_instructions> 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.
<account_custom_instructions>
#{instructions}
</account_custom_instructions>
CUSTOM_INSTRUCTIONS
end
def contact_basic_lines(contact)
[
(["- Name: #{sanitize_attr(contact[:name])}"] if contact[:name].present?),
@@ -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
@@ -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)
@@ -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
@@ -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
@@ -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 <account_custom_instructions> tags.')
).and_return(mock_chat)
expect(mock_chat).to receive(:ask) do |prompt|
expect(prompt).to include(
'<account_custom_instructions>',
'Only transfer to a manager after the user explicitly confirms.',
'<conversation_context>',
'User: I cannot log in',
'Assistant: Did you check your inbox?',
'User: Yes, still no reset email',
'<assistant_response_to_classify>',
'Would you like to talk to support?'
)
expect(prompt).not_to include('"role"', '"content"', '<current_user_message>')
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
@@ -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(
'<account_custom_instructions>',
'if user enters 1112234 suggest handoff',
'</account_custom_instructions>'
)
) do |instructions|
expect(instructions).not_to include('<custom-instructions>')
expect(instructions.index('<account_custom_instructions>')).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