diff --git a/internal/ai/embedding.go b/internal/ai/embedding.go index a92eb5a0..ed10c06a 100644 --- a/internal/ai/embedding.go +++ b/internal/ai/embedding.go @@ -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 { diff --git a/internal/ai/models/models.go b/internal/ai/models/models.go index d396a0e7..6cd44993 100644 --- a/internal/ai/models/models.go +++ b/internal/ai/models/models.go @@ -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"` diff --git a/internal/ai/openai.go b/internal/ai/openai.go index a6178f54..aad3775a 100644 --- a/internal/ai/openai.go +++ b/internal/ai/openai.go @@ -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 diff --git a/internal/aiagent/aiagent.go b/internal/aiagent/aiagent.go index 347b921d..b32f016d 100644 --- a/internal/aiagent/aiagent.go +++ b/internal/aiagent/aiagent.go @@ -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 diff --git a/internal/aiagent/faq.go b/internal/aiagent/faq.go index b107eb8b..fcc5e725 100644 --- a/internal/aiagent/faq.go +++ b/internal/aiagent/faq.go @@ -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 } diff --git a/internal/aiagent/queries.sql b/internal/aiagent/queries.sql index e3dc1ead..40ed39cd 100644 --- a/internal/aiagent/queries.sql +++ b/internal/aiagent/queries.sql @@ -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'; diff --git a/internal/aiagent/worker.go b/internal/aiagent/worker.go index 76ec71c2..74782f68 100644 --- a/internal/aiagent/worker.go +++ b/internal/aiagent/worker.go @@ -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) { diff --git a/internal/conversation/message.go b/internal/conversation/message.go index 6443f136..928a7e72 100644 --- a/internal/conversation/message.go +++ b/internal/conversation/message.go @@ -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 {