fix AI agent reply, turn-cap, and stats bugs

Stop dropping customer follow-ups sent mid-response, reset the turn cap only on reassignment, send max_tokens for non-reasoning models, and make FAQ approval atomic. Also exclude CSAT surveys from assistant stats and remove unused scoped-search code.
This commit is contained in:
Abhinav Raut
2026-07-12 12:15:32 +05:30
parent 76d16285ab
commit 44925e69a7
8 changed files with 108 additions and 77 deletions
+3 -29
View File
@@ -58,15 +58,8 @@ func (ix *embeddingIndex) removeSource(sourceType string, sourceID int) {
ix.replaceSource(sourceType, sourceID, nil)
}
// sourceKey identifies one embedding source in the in-memory index.
type sourceKey struct {
sourceType string
sourceID int
}
// search returns the top-k matches and the count of chunks skipped for mismatched vector dimensions.
// A non-nil allowed set restricts the search to those sources; nil searches the whole index.
func (ix *embeddingIndex) search(query []float32, k int, allowed map[sourceKey]bool) ([]models.SearchResult, int) {
func (ix *embeddingIndex) search(query []float32, k int) ([]models.SearchResult, int) {
ix.mu.RLock()
defer ix.mu.RUnlock()
@@ -78,9 +71,6 @@ func (ix *embeddingIndex) search(query []float32, k int, allowed map[sourceKey]b
dimMismatch := 0
results := make([]models.SearchResult, 0, len(ix.chunks))
for _, c := range ix.chunks {
if allowed != nil && !allowed[sourceKey{c.sourceType, c.sourceID}] {
continue
}
if len(c.vec) != len(query) {
dimMismatch++
continue
@@ -106,31 +96,15 @@ func (ix *embeddingIndex) search(query []float32, k int, allowed map[sourceKey]b
// Search embeds the query and returns the top-k most similar chunks across the whole index.
func (m *Manager) Search(ctx context.Context, query string, k int) ([]models.SearchResult, error) {
return m.search(query, k, nil)
}
// SearchScoped restricts retrieval to the given sources; an empty refs slice returns no results.
func (m *Manager) SearchScoped(ctx context.Context, query string, k int, refs []models.SourceRef) ([]models.SearchResult, error) {
if len(refs) == 0 {
return nil, nil
}
allowed := make(map[sourceKey]bool, len(refs))
for _, r := range refs {
allowed[sourceKey{r.SourceType, r.SourceID}] = true
}
return m.search(query, k, allowed)
}
func (m *Manager) search(query string, k int, allowed map[sourceKey]bool) ([]models.SearchResult, error) {
qvec, err := m.GetEmbeddings(query)
if err != nil {
return nil, err
}
results, dimMismatch := m.index.search(qvec, k, allowed)
results, dimMismatch := m.index.search(qvec, k)
if dimMismatch > 0 {
m.lo.Warn("skipped stale embeddings with mismatched dimensions; reindex the knowledge base after changing the embedding model", "count", dimMismatch, "query_dimensions", len(qvec))
}
m.lo.Debug("rag search", "query", query, "scoped", allowed != nil, "hits", len(results))
m.lo.Debug("rag search", "query", query, "hits", len(results))
for i, r := range results {
preview := r.ChunkText
if len(preview) > 120 {
-6
View File
@@ -113,12 +113,6 @@ type SearchResult struct {
Score float64 `json:"score"`
}
// SourceRef identifies one embedding source, used to scope a search to a subset of the index.
type SourceRef struct {
SourceType string `json:"source_type"`
SourceID int `json:"source_id"`
}
// ChatMessage is one OpenAI-compatible chat message.
type ChatMessage struct {
Role string `json:"role"`
+9 -3
View File
@@ -84,9 +84,15 @@ func (o *OpenAIClient) SendChatCompletion(payload models.ChatCompletionPayload)
o.lo.Debug("chat completion request", "model", model, "messages", len(messages), "images", sentImages, "vision", o.cfg.Vision, "tools", len(payload.Tools))
body := map[string]any{
"model": model,
"messages": messages,
"max_completion_tokens": maxTokens,
"model": model,
"messages": messages,
}
// Reasoning models require max_completion_tokens and reject the older max_tokens; other models
// (including most OpenAI-compatible third-party servers) only understand max_tokens.
if o.cfg.ReasoningEffort != "" {
body["max_completion_tokens"] = maxTokens
} else {
body["max_tokens"] = maxTokens
}
// Only send optional params the admin set. Reasoning models (e.g. GPT-5.x) reject a non-default
// temperature and require reasoning_effort "none" to use tools on /chat/completions; leave
+54 -27
View File
@@ -32,29 +32,31 @@ const (
var efs embed.FS
type queries struct {
GetAssistants *sqlx.Stmt `query:"get-assistants"`
GetAssistant *sqlx.Stmt `query:"get-assistant"`
GetAssistantByUserID *sqlx.Stmt `query:"get-assistant-by-user-id"`
GetAssistantUserIDs *sqlx.Stmt `query:"get-assistant-user-ids"`
InsertAssistantUser *sqlx.Stmt `query:"insert-assistant-user"`
InsertAssistant *sqlx.Stmt `query:"insert-assistant"`
UpdateAssistant *sqlx.Stmt `query:"update-assistant"`
UpdateAssistantUser *sqlx.Stmt `query:"update-assistant-user"`
SoftDeleteAssistantUser *sqlx.Stmt `query:"soft-delete-assistant-user"`
DeleteAssistant *sqlx.Stmt `query:"delete-assistant"`
GetAssistantTools *sqlx.Stmt `query:"get-assistant-tools"`
DeleteAssistantTools *sqlx.Stmt `query:"delete-assistant-tools"`
InsertAssistantTool *sqlx.Stmt `query:"insert-assistant-tool"`
CountAITurns *sqlx.Stmt `query:"count-ai-turns-since-assignment"`
GetRecentContactConvos *sqlx.Stmt `query:"get-recent-contact-conversations"`
GetAssistantWindowStats *sqlx.Stmt `query:"get-assistant-window-stats"`
InsertAIAgentEvent *sqlx.Stmt `query:"insert-ai-agent-event"`
GetAssistants *sqlx.Stmt `query:"get-assistants"`
GetAssistant *sqlx.Stmt `query:"get-assistant"`
GetAssistantByUserID *sqlx.Stmt `query:"get-assistant-by-user-id"`
GetAssistantUserIDs *sqlx.Stmt `query:"get-assistant-user-ids"`
InsertAssistantUser *sqlx.Stmt `query:"insert-assistant-user"`
InsertAssistant *sqlx.Stmt `query:"insert-assistant"`
UpdateAssistant *sqlx.Stmt `query:"update-assistant"`
UpdateAssistantUser *sqlx.Stmt `query:"update-assistant-user"`
SoftDeleteAssistantUser *sqlx.Stmt `query:"soft-delete-assistant-user"`
DeleteAssistant *sqlx.Stmt `query:"delete-assistant"`
GetAssistantTools *sqlx.Stmt `query:"get-assistant-tools"`
GetAllAssistantTools *sqlx.Stmt `query:"get-all-assistant-tools"`
DeleteAssistantTools *sqlx.Stmt `query:"delete-assistant-tools"`
InsertAssistantTool *sqlx.Stmt `query:"insert-assistant-tool"`
CountAITurns *sqlx.Stmt `query:"count-ai-turns-since-assignment"`
GetRecentContactConvos *sqlx.Stmt `query:"get-recent-contact-conversations"`
GetAssistantWindowStats *sqlx.Stmt `query:"get-assistant-window-stats"`
InsertAIAgentEvent *sqlx.Stmt `query:"insert-ai-agent-event"`
InsertFAQSuggestion *sqlx.Stmt `query:"insert-faq-suggestion"`
CountFAQByConversation *sqlx.Stmt `query:"count-faq-suggestions-by-conversation"`
GetFAQSuggestions *sqlx.Stmt `query:"get-faq-suggestions"`
GetFAQSuggestion *sqlx.Stmt `query:"get-faq-suggestion"`
UpdateFAQSuggestionStatus *sqlx.Stmt `query:"update-faq-suggestion-status"`
InsertFAQSuggestion *sqlx.Stmt `query:"insert-faq-suggestion"`
CountFAQByConversation *sqlx.Stmt `query:"count-faq-suggestions-by-conversation"`
GetFAQSuggestions *sqlx.Stmt `query:"get-faq-suggestions"`
GetFAQSuggestion *sqlx.Stmt `query:"get-faq-suggestion"`
UpdateFAQSuggestionStatus *sqlx.Stmt `query:"update-faq-suggestion-status"`
ApproveFAQSuggestionIfPending *sqlx.Stmt `query:"approve-faq-suggestion-if-pending"`
}
type Manager struct {
@@ -69,7 +71,10 @@ type Manager struct {
queue chan int
inflight map[int]bool
mu sync.Mutex
// pending marks a conversation that received a fresh event while its response was in flight, so
// markDone re-enqueues it instead of dropping the follow-up.
pending map[int]bool
mu sync.Mutex
miningQueue chan int
miningInflight map[int]bool
@@ -103,6 +108,7 @@ func New(opts Opts, aiManager *ai.Manager, convo *conversation.Manager, mediaMan
setting: settingManager,
queue: make(chan int, opts.QueueSize),
inflight: map[int]bool{},
pending: map[int]bool{},
miningQueue: make(chan int, opts.QueueSize),
miningInflight: map[int]bool{},
assistantUserIDs: map[int]bool{},
@@ -118,16 +124,37 @@ func (m *Manager) GetAssistants() ([]models.Assistant, error) {
m.lo.Error("error fetching assistants", "error", err)
return nil, envelope.NewError(envelope.GeneralError, m.i18n.T("globals.messages.somethingWentWrong"), nil)
}
toolsByAssistant, err := m.allAssistantTools()
if err != nil {
return nil, err
}
for i := range assistants {
toolIDs, err := m.getTools(assistants[i].ID)
if err != nil {
return nil, err
ids := toolsByAssistant[assistants[i].ID]
if ids == nil {
ids = []int{}
}
assistants[i].ToolIDs = toolIDs
assistants[i].ToolIDs = ids
}
return assistants, nil
}
// allAssistantTools returns tool ids grouped by assistant id in a single query.
func (m *Manager) allAssistantTools() (map[int][]int, error) {
var rows []struct {
AssistantID int `db:"assistant_id"`
ToolID int `db:"tool_id"`
}
if err := m.q.GetAllAssistantTools.Select(&rows); err != nil {
m.lo.Error("error fetching assistant tools", "error", err)
return nil, envelope.NewError(envelope.GeneralError, m.i18n.T("globals.messages.somethingWentWrong"), nil)
}
out := make(map[int][]int)
for _, r := range rows {
out[r.AssistantID] = append(out[r.AssistantID], r.ToolID)
}
return out, nil
}
// GetAssistant returns one assistant with its knowledge links.
func (m *Manager) GetAssistant(id int) (models.Assistant, error) {
var a models.Assistant
+11 -4
View File
@@ -80,13 +80,20 @@ func (m *Manager) ApproveFAQSuggestion(id int, question, answer string, reviewer
if answer == "" {
answer = s.Answer
}
if _, err := m.ai.CreateKnowledgeBaseItem(question, answer, aimodels.KnowledgeSourceConversation, true); err != nil {
return err
}
if _, err := m.q.UpdateFAQSuggestionStatus.Exec(id, models.FAQStatusApproved, reviewerID); err != nil {
// Claim the suggestion before creating the snippet: only the caller that flips pending->approved
// creates one, so a concurrent or retried approval can't produce duplicate snippets.
res, err := m.q.ApproveFAQSuggestionIfPending.Exec(id, reviewerID)
if err != nil {
m.lo.Error("error approving faq suggestion", "id", id, "error", err)
return envelope.NewError(envelope.GeneralError, m.i18n.T("globals.messages.somethingWentWrong"), nil)
}
if n, _ := res.RowsAffected(); n == 0 {
return envelope.NewError(envelope.ConflictError, m.i18n.T("ai.faqAlreadyReviewed"), nil)
}
if _, err := m.ai.CreateKnowledgeBaseItem(question, answer, aimodels.KnowledgeSourceConversation, true); err != nil {
m.lo.Error("faq suggestion approved but snippet creation failed", "id", id, "error", err)
return err
}
return nil
}
+18 -8
View File
@@ -53,6 +53,9 @@ DELETE FROM ai_assistants WHERE id = $1;
-- name: get-assistant-tools
SELECT tool_id FROM ai_assistant_tools WHERE assistant_id = $1 ORDER BY tool_id;
-- name: get-all-assistant-tools
SELECT assistant_id, tool_id FROM ai_assistant_tools ORDER BY assistant_id, tool_id;
-- name: delete-assistant-tools
DELETE FROM ai_assistant_tools WHERE assistant_id = $1;
@@ -66,29 +69,32 @@ INSERT INTO ai_agent_events (assistant_id, conversation_id, type) VALUES ($1, $2
-- name: get-assistant-window-stats
-- $1 = assistant user id (message sender), $2 = assistant id (events), $3 = window start, $4 = window end.
-- CSAT survey messages are sent under the assistant's identity but are not genuine replies, so they
-- are excluded from the reply/conversation counts.
SELECT
(SELECT count(DISTINCT conversation_id) FROM conversation_messages WHERE sender_id = $1 AND type = 'outgoing' AND private = false AND created_at >= $3 AND created_at < $4) AS conversations,
(SELECT count(*) FROM conversation_messages WHERE sender_id = $1 AND type = 'outgoing' AND private = false AND created_at >= $3 AND created_at < $4) AS replies,
(SELECT count(DISTINCT conversation_id) FROM conversation_messages WHERE sender_id = $1 AND type = 'outgoing' AND private = false AND created_at >= $3 AND created_at < $4 AND NOT COALESCE((meta->>'is_csat')::boolean, false)) AS conversations,
(SELECT count(*) FROM conversation_messages WHERE sender_id = $1 AND type = 'outgoing' AND private = false AND created_at >= $3 AND created_at < $4 AND NOT COALESCE((meta->>'is_csat')::boolean, false)) AS replies,
(SELECT count(DISTINCT conversation_id) FROM ai_agent_events WHERE assistant_id = $2 AND type = 'handoff' AND created_at >= $3 AND created_at < $4) AS handoffs,
(SELECT count(DISTINCT conversation_id) FROM ai_agent_events WHERE assistant_id = $2 AND type = 'resolve' AND created_at >= $3 AND created_at < $4) AS resolves,
(SELECT count(DISTINCT e.conversation_id) FROM ai_agent_events e JOIN conversations c ON c.id = e.conversation_id JOIN conversation_statuses s ON s.id = c.status_id
WHERE e.assistant_id = $2 AND e.type = 'resolve' AND e.created_at >= $3 AND e.created_at < $4 AND s.name <> 'Resolved') AS reopened,
(SELECT count(*) FROM csat_responses cr WHERE cr.rating > 0 AND cr.created_at >= $3 AND cr.created_at < $4 AND EXISTS (
SELECT 1 FROM conversation_messages m WHERE m.conversation_id = cr.conversation_id AND m.sender_id = $1 AND m.type = 'outgoing' AND m.private = false)) AS csat_count,
SELECT 1 FROM conversation_messages m WHERE m.conversation_id = cr.conversation_id AND m.sender_id = $1 AND m.type = 'outgoing' AND m.private = false AND NOT COALESCE((m.meta->>'is_csat')::boolean, false))) AS csat_count,
COALESCE((SELECT round(avg(cr.rating)::numeric, 2) FROM csat_responses cr WHERE cr.rating > 0 AND cr.created_at >= $3 AND cr.created_at < $4 AND EXISTS (
SELECT 1 FROM conversation_messages m WHERE m.conversation_id = cr.conversation_id AND m.sender_id = $1 AND m.type = 'outgoing' AND m.private = false)), 0)::float8 AS csat_avg,
SELECT 1 FROM conversation_messages m WHERE m.conversation_id = cr.conversation_id AND m.sender_id = $1 AND m.type = 'outgoing' AND m.private = false AND NOT COALESCE((m.meta->>'is_csat')::boolean, false))), 0)::float8 AS csat_avg,
COALESCE((SELECT round((count(*) FILTER (WHERE cr.rating >= 4))::numeric / NULLIF(count(*), 0) * 100, 1) FROM csat_responses cr WHERE cr.rating > 0 AND cr.created_at >= $3 AND cr.created_at < $4 AND EXISTS (
SELECT 1 FROM conversation_messages m WHERE m.conversation_id = cr.conversation_id AND m.sender_id = $1 AND m.type = 'outgoing' AND m.private = false)), 0)::float8 AS csat_positive;
SELECT 1 FROM conversation_messages m WHERE m.conversation_id = cr.conversation_id AND m.sender_id = $1 AND m.type = 'outgoing' AND m.private = false AND NOT COALESCE((m.meta->>'is_csat')::boolean, false))), 0)::float8 AS csat_positive;
-- name: count-ai-turns-since-assignment
-- Counts the assistant's public replies in the current engagement, i.e. since it was last
-- (re)assigned. Any assignment/status change writes an activity row, so the last activity marks
-- the start of this engagement and a fresh assignment resets the turn budget.
-- Counts the assistant's public replies since it was last (re)assigned, so a fresh assignment resets
-- the turn budget. Keyed off the last assignment activity only; unrelated activity (priority, status,
-- tag changes) must not reset the cap.
SELECT count(*) FROM conversation_messages
WHERE conversation_id = $1 AND sender_id = $2 AND type = 'outgoing' AND private = false
AND created_at > COALESCE((
SELECT max(created_at) FROM conversation_messages
WHERE conversation_id = $1 AND type = 'activity'
AND meta->>'activity_type' IN ('assigned_user_change', 'self_assign')
), to_timestamp(0));
-- name: get-recent-contact-conversations
@@ -119,3 +125,7 @@ FROM ai_faq_suggestions WHERE id = $1;
-- name: update-faq-suggestion-status
UPDATE ai_faq_suggestions SET status = $2, reviewed_by_id = $3, reviewed_at = now(), updated_at = now() WHERE id = $1;
-- name: approve-faq-suggestion-if-pending
UPDATE ai_faq_suggestions SET status = 'approved', reviewed_by_id = $2, reviewed_at = now(), updated_at = now()
WHERE id = $1 AND status = 'pending';
+9
View File
@@ -75,6 +75,9 @@ func (m *Manager) HandleConversationEvent(conversationID, assigneeUserID int) {
func (m *Manager) enqueue(convID int) {
m.mu.Lock()
if m.inflight[convID] {
// A response is already running; remember to run again so a message that arrives mid-response
// still gets answered.
m.pending[convID] = true
m.mu.Unlock()
return
}
@@ -86,6 +89,7 @@ func (m *Manager) enqueue(convID int) {
default:
m.mu.Lock()
delete(m.inflight, convID)
delete(m.pending, convID)
m.mu.Unlock()
m.lo.Warn("ai agent queue full, dropping response job", "conversation_id", convID)
}
@@ -93,8 +97,13 @@ func (m *Manager) enqueue(convID int) {
func (m *Manager) markDone(convID int) {
m.mu.Lock()
requeue := m.pending[convID]
delete(m.pending, convID)
delete(m.inflight, convID)
m.mu.Unlock()
if requeue {
m.enqueue(convID)
}
}
func (m *Manager) handle(ctx context.Context, convID int) {
+4
View File
@@ -701,6 +701,9 @@ func (m *Manager) InsertConversationActivity(activityType, conversationUUID, new
return envelope.NewError(envelope.GeneralError, m.i18n.T("globals.messages.somethingWentWrong"), nil)
}
// Store the activity type structurally so callers can filter activities without parsing i18n content.
meta, _ := json.Marshal(map[string]string{"activity_type": activityType})
message := models.Message{
Type: models.MessageActivity,
Status: models.MessageStatusSent,
@@ -710,6 +713,7 @@ func (m *Manager) InsertConversationActivity(activityType, conversationUUID, new
Private: true,
SenderID: actor.ID,
SenderType: models.SenderTypeAgent,
Meta: meta,
}
if err := m.InsertMessage(&message); err != nil {