diff --git a/internal/channels/bootstrap.go b/internal/channels/bootstrap.go index 60aadfb02c..7179233140 100644 --- a/internal/channels/bootstrap.go +++ b/internal/channels/bootstrap.go @@ -124,7 +124,7 @@ func (r *Runtime) Reconcile(ctx context.Context) error { break } } - if err := syncWhatsAppGateway(ctx, activeWhatsApp); err != nil && activeWhatsApp { + if err = syncWhatsAppGateway(ctx, activeWhatsApp); err != nil && activeWhatsApp { log.Printf("failed to sync WhatsApp gateway: %v", err) } @@ -138,7 +138,7 @@ func (r *Runtime) Reconcile(ctx context.Context) error { if isRunning || retryPending { continue } - if err := r.startChannel(ctx, accountID, wanted); err != nil { + if err = r.startChannel(ctx, accountID, wanted); err != nil { log.Printf("failed to start chat channel %s (%s): %v", accountID, wanted.channel, err) r.recordStartFailure(accountID, wanted.fp, now) continue diff --git a/internal/channels/dingtalk.go b/internal/channels/dingtalk.go index a49e89078d..d37f03f0eb 100644 --- a/internal/channels/dingtalk.go +++ b/internal/channels/dingtalk.go @@ -439,7 +439,7 @@ func (c *dingTalkChannel) saveSessionWebhook(chatID, sessionWebhook string) { c.sessionWebhooks[chatID] = sessionWebhook } -// postWebhook posts a markdown reply payload to a DingTalk sessionWebhook. +// postWebhook posts a Markdown reply payload to a DingTalk sessionWebhook. func (c *dingTalkChannel) postWebhook(ctx context.Context, sessionWebhook string, payload map[string]any) error { body, _ := json.Marshal(payload) req, err := http.NewRequestWithContext(ctx, http.MethodPost, sessionWebhook, bytes.NewReader(body)) diff --git a/internal/channels/discord.go b/internal/channels/discord.go index ba2dfe45d1..0216b1858f 100644 --- a/internal/channels/discord.go +++ b/internal/channels/discord.go @@ -85,7 +85,7 @@ type discordChatWorker struct { type discordGatewayPayload struct { Op int `json:"op"` // Gateway opcode, which indicates the payload type D json.RawMessage `json:"d"` // Event data - S *int64 `json:"s"` // Sequence number of event used for resuming sessions and heartbeating + S *int64 `json:"s"` // Sequence number of event used for resuming sessions and heartbeat T string `json:"t"` // Event name } diff --git a/internal/dao/api_token_test.go b/internal/dao/api_token_test.go index eb273ee62d..5a0bba25a9 100644 --- a/internal/dao/api_token_test.go +++ b/internal/dao/api_token_test.go @@ -46,7 +46,8 @@ func setupAPI4ConversationTestDB(t *testing.T) *gorm.DB { func createAPI4ConversationForDAOTest(t *testing.T, id, agentID string) { t.Helper() - if err := DB.Create(&entity.API4Conversation{ + ctx := t.Context() + if err := DB.WithContext(ctx).Create(&entity.API4Conversation{ ID: id, DialogID: agentID, UserID: "user-1", diff --git a/internal/dao/file2document_test.go b/internal/dao/file2document_test.go index 2b42ffe9b8..4dfe83595d 100644 --- a/internal/dao/file2document_test.go +++ b/internal/dao/file2document_test.go @@ -57,12 +57,13 @@ func pushDB(t *testing.T, testDB *gorm.DB) { func testFile2Document(t *testing.T, fileID, docID string) *entity.File2Document { t.Helper() + ctx := t.Context() f2d := &entity.File2Document{ ID: fileID + "_" + docID, FileID: &fileID, DocumentID: &docID, } - if err := DB.Create(f2d).Error; err != nil { + if err := DB.WithContext(ctx).Create(f2d).Error; err != nil { t.Fatalf("failed to create test record: %v", err) } return f2d diff --git a/internal/handler/agent.go b/internal/handler/agent.go index 097d1a467f..5ec14ef053 100644 --- a/internal/handler/agent.go +++ b/internal/handler/agent.go @@ -1524,7 +1524,8 @@ func (h *AgentHandler) TestDBConnection(c *gin.Context) { common.ResponseWithCodeData(c, common.CodeArgumentError, nil, "Invalid request: "+err.Error()) return } - code, err := h.agentService.TestDBConnection(user.ID, &req) + ctx := c.Request.Context() + code, err := h.agentService.TestDBConnection(ctx, user.ID, &req) if err != nil { common.ErrorWithCode(c, code, err.Error()) return diff --git a/internal/handler/tenant.go b/internal/handler/tenant.go index 21471121c3..14bccb7190 100644 --- a/internal/handler/tenant.go +++ b/internal/handler/tenant.go @@ -190,8 +190,8 @@ func (h *TenantHandler) CreateMetadataStore(c *gin.Context) { // Use user.ID as tenant ID (user IS the tenant in user mode) tenantID := user.ID - - code, err := h.tenantService.CreateMetadataStore(tenantID) + ctx := c.Request.Context() + code, err := h.tenantService.CreateMetadataStore(ctx, tenantID) if err != nil { common.ErrorWithCode(c, code, err.Error()) return @@ -219,7 +219,8 @@ func (h *TenantHandler) DeleteMetadataStore(c *gin.Context) { // Use user.ID as tenant ID (user IS the tenant in user mode) tenantID := user.ID - code, err := h.tenantService.DeleteMetadataStore(tenantID) + ctx := c.Request.Context() + code, err := h.tenantService.DeleteMetadataStore(ctx, tenantID) if err != nil { common.ErrorWithCode(c, code, err.Error()) return diff --git a/internal/ingestion/component/knowledge_compiler/structure/structure.go b/internal/ingestion/component/knowledge_compiler/structure/structure.go index 34d659eab0..c25fb49cd2 100644 --- a/internal/ingestion/component/knowledge_compiler/structure/structure.go +++ b/internal/ingestion/component/knowledge_compiler/structure/structure.go @@ -87,7 +87,6 @@ func Run(ctx context.Context, deps common.Deps, param common.Param, inputs commo perBatch := make([][]common.Product, len(batches)) jobs := make([]func() error, 0, len(batches)) for i, batch := range batches { - i, batch := i, batch jobs = append(jobs, func() error { packed, batchIDs := PackBatch(batch) if len(batchIDs) == 0 { diff --git a/internal/service/agent.go b/internal/service/agent.go index 0c57769f5d..0cadfddeb9 100644 --- a/internal/service/agent.go +++ b/internal/service/agent.go @@ -1330,9 +1330,9 @@ func (s *AgentService) RunAgent(ctx context.Context, userID, canvasID, sessionID // Cancellation of the HTTP request must reach the workflow, but it must // not stop the Redis watchers while the real Runner goroutine is still - // unwinding. Otherwise a non-cooperative external call can outlive the + // unwinding. Otherwise, a non-cooperative external call can outlive the // lease, allowing a second process to acquire the same session. The - // detached watcher context is cancelled only after inner has closed and + // detached watcher context is canceled only after inner has closed and // cleanup has taken ownership of the lease release. lifecycleDone := make(chan struct{}) go func() { @@ -2281,7 +2281,7 @@ func (s *AgentService) CancelSessionRun(ctx context.Context, userID, sessionID s if err != nil { return fmt.Errorf("agent cancel: read active session: %w: %w", err, ErrAgentStorageError) } - if err == nil && remote != nil { + if remote != nil { if remote.UserID != "" && remote.UserID != userID { return ErrAgentNotOwner } diff --git a/internal/service/agent_dbcheck.go b/internal/service/agent_dbcheck.go index 5389280c30..3f56a66677 100644 --- a/internal/service/agent_dbcheck.go +++ b/internal/service/agent_dbcheck.go @@ -226,7 +226,7 @@ func dbConnectionPort(port interface{}) string { // short timeout to keep the API responsive when targets are unreachable. // The "required argument are missing" message has a trailing semicolon // and space to stay byte-identical with the Python implementation. -func (s *AgentService) TestDBConnection(userID string, req *TestDBConnectionRequest) (common.ErrorCode, error) { +func (s *AgentService) TestDBConnection(ctx context.Context, userID string, req *TestDBConnectionRequest) (common.ErrorCode, error) { if missing := missingDBConnectionFields(req); len(missing) > 0 { return common.CodeArgumentError, fmt.Errorf("required argument are missing: %s; ", strings.Join(missing, ",")) } @@ -263,13 +263,13 @@ func (s *AgentService) TestDBConnection(userID string, req *TestDBConnectionRequ } defer db.Close() - ctx, cancel := context.WithTimeout(context.Background(), dbProbeTimeout) + newCtx, cancel := context.WithTimeout(ctx, dbProbeTimeout) defer cancel() - if err = db.PingContext(ctx); err != nil { + if err = db.PingContext(newCtx); err != nil { return common.CodeExceptionError, err } - if _, err = db.ExecContext(ctx, "SELECT 1"); err != nil { + if _, err = db.ExecContext(newCtx, "SELECT 1"); err != nil { return common.CodeExceptionError, err } default: diff --git a/internal/service/agent_test.go b/internal/service/agent_test.go index 8a5ecfc3f5..aea52e6cc5 100644 --- a/internal/service/agent_test.go +++ b/internal/service/agent_test.go @@ -1278,7 +1278,8 @@ func TestAssertHostIsSafeRejectsLocalhost(t *testing.T) { } func TestTestDBConnectionMissingFields(t *testing.T) { - code, err := NewAgentService().TestDBConnection("user-1", &TestDBConnectionRequest{DBType: "mysql"}) + ctx := t.Context() + code, err := NewAgentService().TestDBConnection(ctx, "user-1", &TestDBConnectionRequest{DBType: "mysql"}) if err == nil { t.Fatal("expected missing field error") } @@ -1292,7 +1293,8 @@ func TestTestDBConnectionMissingFields(t *testing.T) { } func TestTestDBConnectionUnsupportedDatabaseType(t *testing.T) { - code, err := NewAgentService().TestDBConnection("user-1", &TestDBConnectionRequest{ + ctx := t.Context() + code, err := NewAgentService().TestDBConnection(ctx, "user-1", &TestDBConnectionRequest{ DBType: "postgres", Database: "rag_flow", Username: "root", diff --git a/internal/service/chat_pipeline.go b/internal/service/chat_pipeline.go index b96ade8ffc..de24f100ec 100644 --- a/internal/service/chat_pipeline.go +++ b/internal/service/chat_pipeline.go @@ -403,7 +403,7 @@ func (s *ChatPipelineService) AsyncChat( } } kbinfos := map[string]interface{}{"chunks": chunks} - s.enrichChunksWithMetadata(kbinfos, chat.TenantID, metadataFields) + s.enrichChunksWithMetadata(ctx, kbinfos, chat.TenantID, metadataFields) } out <- AsyncChatResult{ @@ -826,7 +826,7 @@ func (s *ChatPipelineService) AsyncChat( // Enrich chunks with document metadata AFTER all retrieval adds. // Request values (kwargs) take precedence over config values. if includeRefMeta, metadataFields := s.resolveReferenceMetadata(promptConfig, kwargs); includeRefMeta { - s.enrichChunksWithMetadata(kbinfos, chat.TenantID, metadataFields) + s.enrichChunksWithMetadata(ctx, kbinfos, chat.TenantID, metadataFields) } timer.Exit(common.PhaseRetrieval) @@ -2358,7 +2358,7 @@ func (s *ChatPipelineService) resolveReferenceMetadata(promptConfig map[string]i // enrichChunksWithMetadata enriches chunk records in kbinfos with document-level // metadata. Mirrors Python's enrich_chunks_with_document_metadata() in // api/utils/reference_metadata_utils.py. -func (s *ChatPipelineService) enrichChunksWithMetadata(kbinfos map[string]interface{}, tenantID string, fields []string) { +func (s *ChatPipelineService) enrichChunksWithMetadata(ctx context.Context, kbinfos map[string]interface{}, tenantID string, fields []string) { chunksRaw, ok := kbinfos["chunks"].([]map[string]interface{}) if !ok || len(chunksRaw) == 0 { return @@ -2370,7 +2370,7 @@ func (s *ChatPipelineService) enrichChunksWithMetadata(kbinfos map[string]interf return } - s.MetadataSvc.EnrichChunksWithDocMetadata(chunks, tenantID, fields) + s.MetadataSvc.EnrichChunksWithDocMetadata(ctx, chunks, tenantID, fields) } // kbPrompt builds knowledge prompt blocks from retrieved chunks. diff --git a/internal/service/chat_session.go b/internal/service/chat_session.go index 8151af67a3..dd4ba6c740 100644 --- a/internal/service/chat_session.go +++ b/internal/service/chat_session.go @@ -682,9 +682,7 @@ func (s *ChatSessionService) DeleteSessionMessage(ctx context.Context, userID, c } func (s *ChatSessionService) UpdateMessageFeedback(ctx context.Context, userID, chatID, sessionID, msgID string, req map[string]interface{}) (*ChatSessionPayload, common.ErrorCode, error) { - if ctx == nil { - ctx = context.Background() - } + ownerTenantID := "" tenantIDs, err := s.userTenantDAO.GetTenantIDsByUserID(ctx, dao.DB, userID) if err != nil { @@ -840,9 +838,7 @@ func (s *ChatSessionService) applyChunkFeedback(ctx context.Context, tenantID st "disabled": true, }, nil } - if ctx == nil { - ctx = context.Background() - } + if _, ok := ctx.Deadline(); !ok { var cancel context.CancelFunc ctx, cancel = context.WithTimeout(ctx, chunkFeedbackTimeout) @@ -1452,9 +1448,6 @@ func (s *ChatSessionService) ChatCompletions( passAllHistory bool, storeHistory bool, legacy bool, stream bool, streamChan chan<- string, ) (map[string]interface{}, error) { - if ctx == nil { - ctx = context.Background() - } fail := func(err error) (map[string]interface{}, error) { if stream && streamChan != nil { diff --git a/internal/service/chunk/chunk.go b/internal/service/chunk/chunk.go index 1db2270e65..8c14cebf88 100644 --- a/internal/service/chunk/chunk.go +++ b/internal/service/chunk/chunk.go @@ -1362,7 +1362,8 @@ func (s *ChunkService) AddChunk(ctx context.Context, req *service.AddChunkReques } if req.ImageBase64 != nil { - imageBinary, err := decodeChunkImageBase64(*req.ImageBase64) + var imageBinary []byte + imageBinary, err = decodeChunkImageBase64(*req.ImageBase64) if err != nil { return nil, addChunkError{code: common.CodeDataError, message: err.Error()} } @@ -1394,7 +1395,7 @@ func (s *ChunkService) AddChunk(ctx context.Context, req *service.AddChunkReques } chunkData[fmt.Sprintf("q_%d_vec", len(mergedVec))] = mergedVec - ctx, cancel := context.WithTimeout(context.Background(), 600*time.Second) + ctx, cancel := context.WithTimeout(ctx, 600*time.Second) defer cancel() if _, err = s.docEngine.InsertChunks(ctx, []map[string]interface{}{chunkData}, indexName, req.DatasetID); err != nil { return nil, addChunkError{code: common.CodeServerError, message: fmt.Sprintf("insert chunk: %v", err)} diff --git a/internal/service/chunk_types.go b/internal/service/chunk_types.go index 625a121e00..7566187763 100644 --- a/internal/service/chunk_types.go +++ b/internal/service/chunk_types.go @@ -299,12 +299,13 @@ func (s *ChunkService) StopParsing(ctx context.Context, userID, datasetID string indexName := IndexName(kb.TenantID) - exists, err := s.docEngine.ChunkStoreExists(context.Background(), indexName, datasetID) + var exists bool + exists, err = s.docEngine.ChunkStoreExists(ctx, indexName, datasetID) if err != nil { return nil, common.CodeServerError, fmt.Errorf("failed to check chunk store %s/%s: %w", indexName, datasetID, err) } if exists { - if _, err := s.docEngine.DeleteChunks(context.Background(), map[string]interface{}{"doc_id": doc.ID}, indexName, datasetID); err != nil { + if _, err = s.docEngine.DeleteChunks(ctx, map[string]interface{}{"doc_id": doc.ID}, indexName, datasetID); err != nil { return nil, common.CodeServerError, fmt.Errorf("failed to delete chunks for document %s: %w", doc.ID, err) } } else { diff --git a/internal/service/dataset/crud.go b/internal/service/dataset/crud.go index 79d1081ac8..74145d7596 100644 --- a/internal/service/dataset/crud.go +++ b/internal/service/dataset/crud.go @@ -261,7 +261,7 @@ func (d *DatasetService) DeleteDatasets(ctx context.Context, ids []string, delet successCount := 0 errorsList := make([]string, 0) for _, kb := range kbs { - if err := d.deleteDataset(tenantID, kb); err != nil { + if err := d.deleteDataset(ctx, tenantID, kb); err != nil { errorsList = append(errorsList, err.Error()) common.Warn("deleteDataset failed", zap.String("kb_id", kb.ID), zap.Error(err)) continue @@ -275,7 +275,7 @@ func (d *DatasetService) DeleteDatasets(ctx context.Context, ids []string, delet }, common.CodeSuccess, nil } -func (d *DatasetService) deleteDataset(tenantID string, kb *entity.Knowledgebase) error { +func (d *DatasetService) deleteDataset(ctx context.Context, tenantID string, kb *entity.Knowledgebase) error { // Collect document IDs first so engine cleanup can run before the // transaction (engine ops are not transactional). var documents []entity.Document @@ -284,7 +284,7 @@ func (d *DatasetService) deleteDataset(tenantID string, kb *entity.Knowledgebase } docIDs := extractDocIDs(documents) if len(docIDs) > 0 { - d.deleteDatasetEngineData(kb, docIDs) + d.deleteDatasetEngineData(ctx, kb, docIDs) } return dao.DB.Transaction(func(tx *gorm.DB) error { @@ -552,11 +552,10 @@ func extractDocIDs(docs []entity.Document) []string { // deleteDatasetEngineData cleans up engine-level chunks and metadata for all // documents in a dataset being deleted. Called before the DB transaction // because engine operations are not transactional. -func (d *DatasetService) deleteDatasetEngineData(kb *entity.Knowledgebase, docIDs []string) { +func (d *DatasetService) deleteDatasetEngineData(ctx context.Context, kb *entity.Knowledgebase, docIDs []string) { if d.docEngine == nil || len(docIDs) == 0 { return } - ctx := context.Background() indexName := fmt.Sprintf("ragflow_%s", kb.TenantID) if _, err := d.docEngine.DeleteChunks(ctx, map[string]interface{}{"doc_id": docIDs}, indexName, kb.ID); err != nil { diff --git a/internal/service/document/document_metadata.go b/internal/service/document/document_metadata.go index 4359c49555..807525237d 100644 --- a/internal/service/document/document_metadata.go +++ b/internal/service/document/document_metadata.go @@ -25,7 +25,7 @@ func (s *DocumentService) GetMetadataSummary(ctx context.Context, kbID string, d return nil, err } - searchResult, err := s.metadataSvc.SearchMetadata(kbID, tenantID, docIDs, 1000) + searchResult, err := s.metadataSvc.SearchMetadata(ctx, kbID, tenantID, docIDs, 1000) if err != nil { return nil, err } @@ -127,7 +127,7 @@ func (s *DocumentService) GetDocumentMetadataByID(ctx context.Context, docID str return nil, err } - searchResult, err := s.metadataSvc.SearchMetadata(doc.KbID, tenantID, []string{docID}, 1) + searchResult, err := s.metadataSvc.SearchMetadata(ctx, doc.KbID, tenantID, []string{docID}, 1) if err != nil { return nil, err } diff --git a/internal/service/enrich_metadata_test.go b/internal/service/enrich_metadata_test.go index 505a7feebe..640a8a5061 100644 --- a/internal/service/enrich_metadata_test.go +++ b/internal/service/enrich_metadata_test.go @@ -254,13 +254,15 @@ func TestAttachDocMetaToChunks_EmptyMeta(t *testing.T) { // --- EnrichChunksWithDocMetadata (integration) --- func TestEnrichChunksWithDocMetadata_NoChunks(t *testing.T) { + ctx := t.Context() svc := NewMetadataService() - svc.EnrichChunksWithDocMetadata(nil, "tenant-1", nil) + svc.EnrichChunksWithDocMetadata(ctx, nil, "tenant-1", nil) // Should not panic } func TestEnrichChunksWithDocMetadata_EmptyChunks(t *testing.T) { + ctx := t.Context() svc := NewMetadataService() - svc.EnrichChunksWithDocMetadata([]map[string]interface{}{}, "tenant-1", nil) + svc.EnrichChunksWithDocMetadata(ctx, []map[string]interface{}{}, "tenant-1", nil) // Should not panic } diff --git a/internal/service/mcp.go b/internal/service/mcp.go index a28ca6e83d..b508d996f7 100644 --- a/internal/service/mcp.go +++ b/internal/service/mcp.go @@ -636,7 +636,7 @@ func (s *MCPService) ImportServers(ctx context.Context, tenantID string, servers } } - mcpCtx, cancel := context.WithTimeout(context.Background(), timeout) + mcpCtx, cancel := context.WithTimeout(ctx, timeout) tools, fetchErr := utility.FetchTools(mcpCtx, utility.FetchOptions{ URL: url, ServerType: stype, @@ -738,8 +738,9 @@ func (s *MCPService) TestServer(mcpID string, req *TestServerRequest) ([]map[str vars[k] = sv } } + ctx := context.Background() - mcpCtx, cancel := context.WithTimeout(context.Background(), timeout) + mcpCtx, cancel := context.WithTimeout(ctx, timeout) defer cancel() tools, err := utility.FetchTools(mcpCtx, utility.FetchOptions{ URL: req.URL, diff --git a/internal/service/metadata.go b/internal/service/metadata.go index 622fbfb942..f3354b88a5 100644 --- a/internal/service/metadata.go +++ b/internal/service/metadata.go @@ -105,7 +105,7 @@ type SearchMetadataResponse struct { } // SearchMetadata searches the metadata index with the given parameters -func (s *MetadataService) SearchMetadata(kbID, tenantID string, docIDs []string, size int) (*SearchMetadataResponse, error) { +func (s *MetadataService) SearchMetadata(ctx context.Context, kbID, tenantID string, docIDs []string, size int) (*SearchMetadataResponse, error) { searchReq := &types.SearchMetadataRequest{ TenantID: tenantID, Offset: 0, @@ -116,7 +116,7 @@ func (s *MetadataService) SearchMetadata(kbID, tenantID string, docIDs []string, }, } - searchResult, err := s.docEngine.SearchMetadata(context.Background(), searchReq) + searchResult, err := s.docEngine.SearchMetadata(ctx, searchReq) if err != nil { return nil, fmt.Errorf("search failed: %w", err) } @@ -306,10 +306,10 @@ func ConvertSearchResultToDocMeta(chunks []map[string]interface{}) DocMetaMap { } // FetchDocMetaByKB fetches document metadata from ES for each KB. -func (s *MetadataService) FetchDocMetaByKB(docIDsByKB KBDocIDsMap, tenantID string) DocMetaMap { +func (s *MetadataService) FetchDocMetaByKB(ctx context.Context, docIDsByKB KBDocIDsMap, tenantID string) DocMetaMap { metaByDoc := make(DocMetaMap) for kbID, docIDs := range docIDsByKB { - result, err := s.SearchMetadata(kbID, tenantID, docIDs, len(docIDs)) + result, err := s.SearchMetadata(ctx, kbID, tenantID, docIDs, len(docIDs)) if err != nil { continue } @@ -353,7 +353,7 @@ func AttachDocMetaToChunks(chunks []map[string]interface{}, metaByDoc DocMetaMap // EnrichChunksWithDocMetadata attaches document metadata to each chunk in-place. // Combines CollectDocIDsByKB, FetchDocMetaByKB, and AttachDocMetaToChunks. -func (s *MetadataService) EnrichChunksWithDocMetadata(chunks []map[string]interface{}, tenantID string, metadataFields []string) { +func (s *MetadataService) EnrichChunksWithDocMetadata(ctx context.Context, chunks []map[string]interface{}, tenantID string, metadataFields []string) { if len(chunks) == 0 || s.docEngine == nil { return } @@ -361,7 +361,7 @@ func (s *MetadataService) EnrichChunksWithDocMetadata(chunks []map[string]interf if len(docIDsByKB) == 0 { return } - metaByDoc := s.FetchDocMetaByKB(docIDsByKB, tenantID) + metaByDoc := s.FetchDocMetaByKB(ctx, docIDsByKB, tenantID) if len(metaByDoc) == 0 { return } @@ -493,7 +493,7 @@ func appendDocID(existing interface{}, docID string) []string { return result } -// ParseLengthPrefixedJSON parses Infinity's length-prefixed JSON format +// ParseLengthPrefixedJSON parses Infinity's length-prefixed JSON // Format: [4-byte length (little-endian)][JSON][4-byte length][JSON]... // Returns the FIRST valid JSON object found func ParseLengthPrefixedJSON(data []byte) map[string]interface{} { diff --git a/internal/service/openai_chat.go b/internal/service/openai_chat.go index 23fcd34c67..1c0125e11f 100644 --- a/internal/service/openai_chat.go +++ b/internal/service/openai_chat.go @@ -268,7 +268,7 @@ func (s *OpenAIChatService) OpenAIChatCompletions(c *gin.Context, userID, chatID if lfClient != nil { ctx = context.WithValue(ctx, langfuseCtxKey, lfClient) defer func() { - shutdownCtx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + shutdownCtx, cancel := context.WithTimeout(ctx, 2*time.Second) defer cancel() _ = lfClient.Shutdown(shutdownCtx) }() @@ -364,7 +364,7 @@ func (s *OpenAIChatService) OpenAIChatCompletions(c *gin.Context, userID, chatID finalReference = formatChunks(chunks) } } - s.enrichChunksWithDocumentMetadata(finalReference, dialog.TenantID, openaiReq.IncludeRefMetadata, openaiReq.MetadataFields) + s.enrichChunksWithDocumentMetadata(ctx, finalReference, dialog.TenantID, openaiReq.IncludeRefMetadata, openaiReq.MetadataFields) completionTok = tokenizer.NumTokensFromString(result.Answer) events <- OpenAIStreamEvent{ Kind: OpenAIEventFinal, @@ -405,7 +405,7 @@ func (s *OpenAIChatService) OpenAIChatCompletions(c *gin.Context, userID, chatID } } } - s.enrichChunksWithDocumentMetadata(finalReference, dialog.TenantID, openaiReq.IncludeRefMetadata, openaiReq.MetadataFields) + s.enrichChunksWithDocumentMetadata(ctx, finalReference, dialog.TenantID, openaiReq.IncludeRefMetadata, openaiReq.MetadataFields) events <- OpenAIStreamEvent{ Kind: OpenAIEventFinal, FinalAnswer: strings.TrimSpace(fullContent), @@ -449,7 +449,7 @@ func (s *OpenAIChatService) OpenAIChatCompletions(c *gin.Context, userID, chatID resp.Reference = formatChunks(chunks) } } - s.enrichChunksWithDocumentMetadata(resp.Reference, dialog.TenantID, openaiReq.IncludeRefMetadata, openaiReq.MetadataFields) + s.enrichChunksWithDocumentMetadata(ctx, resp.Reference, dialog.TenantID, openaiReq.IncludeRefMetadata, openaiReq.MetadataFields) } contextUsed := 0 @@ -669,7 +669,7 @@ func formatChunks(chunks []map[string]interface{}) []FormattedChunk { // api/utils/reference_metadata_utils.py. // When fields is a non-nil empty slice (explicitly provided as []), enrichment // is skipped — matching Python's behavior for {"fields": []}. -func (s *OpenAIChatService) enrichChunksWithDocumentMetadata(chunks []FormattedChunk, tenantID string, include bool, fields []string) { +func (s *OpenAIChatService) enrichChunksWithDocumentMetadata(ctx context.Context, chunks []FormattedChunk, tenantID string, include bool, fields []string) { if !include || len(chunks) == 0 || s == nil || s.pipeline.MetadataSvc == nil { return } @@ -684,7 +684,7 @@ func (s *OpenAIChatService) enrichChunksWithDocumentMetadata(chunks []FormattedC "document_metadata": ch.DocumentMetadata, } } - s.pipeline.MetadataSvc.EnrichChunksWithDocMetadata(maps, tenantID, fields) + s.pipeline.MetadataSvc.EnrichChunksWithDocMetadata(ctx, maps, tenantID, fields) for i, m := range maps { if md, ok := m["document_metadata"]; ok { chunks[i].DocumentMetadata = md diff --git a/internal/service/skill_space.go b/internal/service/skill_space.go index cd5bbd487e..72ea4c04d2 100644 --- a/internal/service/skill_space.go +++ b/internal/service/skill_space.go @@ -427,7 +427,7 @@ func (s *SkillSpaceService) DeleteSpace(ctx context.Context, spaceID, tenantID s // asyncDeleteSpace performs the actual deletion work in the background. // It deletes the search index, removes files via Go FileService, and soft-deletes the space record. func (s *SkillSpaceService) asyncDeleteSpace(ctx context.Context, spaceID, folderID, tenantID string, docEngine engine.DocEngine) { - bgCtx, bgCancel := context.WithTimeout(context.Background(), 120*time.Second) + bgCtx, bgCancel := context.WithTimeout(ctx, 120*time.Second) defer bgCancel() defer func() { @@ -441,7 +441,7 @@ func (s *SkillSpaceService) asyncDeleteSpace(ctx context.Context, spaceID, folde if docEngine != nil { indexName := getSkillIndexName(tenantID, spaceID) common.Info("Async deleting space index", zap.String("index", indexName), zap.String("spaceID", spaceID)) - deleteCtx, cancel := context.WithTimeout(context.Background(), 60*time.Second) + deleteCtx, cancel := context.WithTimeout(ctx, 60*time.Second) if err := docEngine.DropChunkStore(deleteCtx, indexName, "skill"); err != nil { common.Warn("Failed to delete space index during async delete", zap.String("index", indexName), zap.Error(err)) // Continue with other cleanup steps @@ -456,7 +456,7 @@ func (s *SkillSpaceService) asyncDeleteSpace(ctx context.Context, spaceID, folde // is the HTTP request context canceled when the handler returns and the // goroutine starts executing). common.Info("Async deleting space folder via Go FileService", zap.String("folderID", folderID), zap.String("spaceID", spaceID)) - ctxFS, cancelFS := context.WithTimeout(context.Background(), 60*time.Second) + ctxFS, cancelFS := context.WithTimeout(ctx, 60*time.Second) defer cancelFS() success, msg := s.fileService.DeleteFiles(ctxFS, tenantID, []string{folderID}) if !success { diff --git a/internal/service/tenant.go b/internal/service/tenant.go index a3a2e8f1f4..e6502bfa45 100644 --- a/internal/service/tenant.go +++ b/internal/service/tenant.go @@ -320,9 +320,9 @@ func (s *TenantService) GetTenantList(ctx context.Context, userID string) ([]*Te } // CreateMetadataStore creates the metadata store for a tenant -func (s *TenantService) CreateMetadataStore(tenantID string) (common.ErrorCode, error) { +func (s *TenantService) CreateMetadataStore(ctx context.Context, tenantID string) (common.ErrorCode, error) { // Call document engine to create doc meta table - err := s.docEngine.CreateMetadataStore(context.Background(), tenantID) + err := s.docEngine.CreateMetadataStore(ctx, tenantID) if err != nil { return common.CodeServerError, fmt.Errorf("failed to create metadata table: %w", err) } @@ -331,9 +331,9 @@ func (s *TenantService) CreateMetadataStore(tenantID string) (common.ErrorCode, } // DeleteMetadataStore deletes the metadata store for a tenant -func (s *TenantService) DeleteMetadataStore(tenantID string) (common.ErrorCode, error) { +func (s *TenantService) DeleteMetadataStore(ctx context.Context, tenantID string) (common.ErrorCode, error) { // Call document engine to delete doc meta table - err := s.docEngine.DropMetadataStore(context.Background(), tenantID) + err := s.docEngine.DropMetadataStore(ctx, tenantID) if err != nil { return common.CodeServerError, fmt.Errorf("failed to delete doc meta table: %w", err) } diff --git a/internal/service/think_tag_test.go b/internal/service/think_tag_test.go index ff85c4c563..d8e9374a1d 100644 --- a/internal/service/think_tag_test.go +++ b/internal/service/think_tag_test.go @@ -17,7 +17,6 @@ package service import ( - "context" "strings" "testing" ) @@ -204,6 +203,7 @@ func TestExtractVisibleAnswer_StrayTag(t *testing.T) { } func TestStreamThinkTagDelta(t *testing.T) { + ctx := t.Context() chunks := []string{"hello ", "wor", "", "think text", "", "ld", " final"} ch := make(chan string, len(chunks)) for _, c := range chunks { @@ -213,7 +213,7 @@ func TestStreamThinkTagDelta(t *testing.T) { var texts []string var markers []string - for d := range StreamThinkTagDelta(context.Background(), ch, 16) { + for d := range StreamThinkTagDelta(ctx, ch, 16) { switch d.Kind { case ThinkDeltaText: texts = append(texts, d.Value) @@ -239,6 +239,7 @@ func TestStreamThinkTagDelta(t *testing.T) { } func TestStreamThinkTagDelta_IncrementalFlush(t *testing.T) { + ctx := t.Context() chunks := []string{"1234", "5678", "90ab"} ch := make(chan string, len(chunks)) for _, c := range chunks { @@ -247,7 +248,7 @@ func TestStreamThinkTagDelta_IncrementalFlush(t *testing.T) { close(ch) var texts []string - for d := range StreamThinkTagDelta(context.Background(), ch, 1) { + for d := range StreamThinkTagDelta(ctx, ch, 1) { if d.Kind == ThinkDeltaText { texts = append(texts, d.Value) } @@ -258,6 +259,7 @@ func TestStreamThinkTagDelta_IncrementalFlush(t *testing.T) { } func TestStreamThinkTagDelta_NoThinkTags(t *testing.T) { + ctx := t.Context() chunks := []string{"just", " plain", " text"} ch := make(chan string, len(chunks)) for _, c := range chunks { @@ -266,7 +268,7 @@ func TestStreamThinkTagDelta_NoThinkTags(t *testing.T) { close(ch) var texts []string - for d := range StreamThinkTagDelta(context.Background(), ch, 4) { + for d := range StreamThinkTagDelta(ctx, ch, 4) { texts = append(texts, d.Value) } @@ -277,6 +279,7 @@ func TestStreamThinkTagDelta_NoThinkTags(t *testing.T) { } func TestStreamThinkTagDelta_DeferredClose(t *testing.T) { + ctx := t.Context() // When has no visible text after it, the marker is deferred. chunks := []string{"", "hello", "", "world"} ch := make(chan string, len(chunks)) @@ -286,7 +289,7 @@ func TestStreamThinkTagDelta_DeferredClose(t *testing.T) { close(ch) var markers []string - for d := range StreamThinkTagDelta(context.Background(), ch, 1) { + for d := range StreamThinkTagDelta(ctx, ch, 1) { if d.Kind == ThinkDeltaMarker { markers = append(markers, d.Value) } diff --git a/internal/service/web_search_provider_test.go b/internal/service/web_search_provider_test.go index d503d48b0d..0deddb2d51 100644 --- a/internal/service/web_search_provider_test.go +++ b/internal/service/web_search_provider_test.go @@ -17,7 +17,6 @@ package service import ( - "context" "encoding/json" "net/http" "net/http/httptest" @@ -134,6 +133,7 @@ func TestResolveWebSearchProviderRequiresKeyForSelectedProvider(t *testing.T) { } func TestRetrieveQueritWebSearchUsesChatDefaultsAndReturnsReferenceShape(t *testing.T) { + ctx := t.Context() var requestBody map[string]interface{} server := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { if got := request.Header.Get("Authorization"); got != "Bearer querit-test" { @@ -157,7 +157,7 @@ func TestRetrieveQueritWebSearchUsesChatDefaultsAndReturnsReferenceShape(t *test defer server.Close() result, err := retrieveQueritWebSearch( - context.Background(), + ctx, server.Client(), server.URL, "querit-test", diff --git a/internal/service/wikisearch/engine_service_test.go b/internal/service/wikisearch/engine_service_test.go index 57b4939d2f..fb71c681cc 100644 --- a/internal/service/wikisearch/engine_service_test.go +++ b/internal/service/wikisearch/engine_service_test.go @@ -72,11 +72,12 @@ func realEngineRow(id, kbID, docID, name, slug string, source []string) map[stri } func TestEngineService_QueryPages_NormalizesEngineRowShape(t *testing.T) { + ctx := t.Context() eng := &fakeDocEngine{searchRows: []map[string]interface{}{ realEngineRow("wiki/alpha", "kb1", "d1", "Alpha", "entity/alpha", []string{"c1", "c2"}), }} svc := NewEngineService(eng) - res, err := svc.QueryPages(context.Background(), "t1", []string{"kb1"}, "alpha", "", 5) + res, err := svc.QueryPages(ctx, "t1", []string{"kb1"}, "alpha", "", 5) if err != nil { t.Fatalf("QueryPages err = %v", err) } @@ -128,6 +129,7 @@ func TestEngineService_QueryPages_NormalizesEngineRowShape(t *testing.T) { // array-shaped engine keyword fields: the document engine returns slug_kwd / // docnm_kwd as arrays, and those must not be dropped. func TestEngineService_QueryPages_ArrayShapedFields(t *testing.T) { + ctx := t.Context() eng := &fakeDocEngine{searchRows: []map[string]interface{}{ { "id": "wiki/alpha", "kb_id": "kb1", "compile_kwd": "wiki_page", @@ -138,7 +140,7 @@ func TestEngineService_QueryPages_ArrayShapedFields(t *testing.T) { }, }} svc := NewEngineService(eng) - res, err := svc.QueryPages(context.Background(), "t1", []string{"kb1"}, "alpha", "", 5) + res, err := svc.QueryPages(ctx, "t1", []string{"kb1"}, "alpha", "", 5) if err != nil { t.Fatalf("QueryPages err = %v", err) } @@ -157,27 +159,29 @@ func TestEngineService_QueryPages_ArrayShapedFields(t *testing.T) { } func TestEngineService_QueryPages_DegradesEmpty(t *testing.T) { + ctx := t.Context() svc := NewEngineService(nil) - res, err := svc.QueryPages(context.Background(), "t1", []string{"kb1"}, "alpha", "", 5) + res, err := svc.QueryPages(ctx, "t1", []string{"kb1"}, "alpha", "", 5) if err != nil || len(res.Chunks) != 0 { t.Fatalf("got chunks=%d err=%v, want empty/noerr (no engine)", len(res.Chunks), err) } svc2 := NewEngineService(&fakeDocEngine{searchRows: []map[string]interface{}{ realEngineRow("x", "kb1", "d", "D", "s", nil), }}) - res2, _ := svc2.QueryPages(context.Background(), "t1", []string{"kb1"}, " ", "", 5) + res2, _ := svc2.QueryPages(ctx, "t1", []string{"kb1"}, " ", "", 5) if len(res2.Chunks) != 0 { t.Fatalf("chunks = %d, want 0 for blank query", len(res2.Chunks)) } } func TestEngineService_BackfillChunks_ByIDScoped(t *testing.T) { + ctx := t.Context() eng := &fakeDocEngine{chunks: map[string]interface{}{ "c1": map[string]interface{}{"id": "c1", "content_with_weight": "raw 1", "doc_id": "d1", "docnm_kwd": "D", "kb_id": "kb1"}, "c2": map[string]interface{}{"id": "c2", "content_with_weight": "raw 2", "doc_id": "d1", "docnm_kwd": "D", "kb_id": "kb1"}, }} svc := NewEngineService(eng) - out, err := svc.BackfillChunks(context.Background(), "t1", []string{"kb1", "kb2"}, []string{"c1", "c2", "c1", "missing"}) + out, err := svc.BackfillChunks(ctx, "t1", []string{"kb1", "kb2"}, []string{"c1", "c2", "c1", "missing"}) if err != nil { t.Fatalf("BackfillChunks err = %v", err) } @@ -205,14 +209,16 @@ func TestEngineService_BackfillChunks_ByIDScoped(t *testing.T) { } func TestEngineService_BackfillChunks_DegradesNoEngine(t *testing.T) { + ctx := t.Context() svc := NewEngineService(nil) - out, err := svc.BackfillChunks(context.Background(), "t1", []string{"kb1"}, []string{"c1"}) + out, err := svc.BackfillChunks(ctx, "t1", []string{"kb1"}, []string{"c1"}) if err != nil || len(out) != 0 { t.Fatalf("got %d err=%v, want empty/noerr (no engine)", len(out), err) } } func TestEngineService_AvailableFor_BoundedExistenceSearch(t *testing.T) { + ctx := t.Context() // A KB carrying wiki pages => AvailableFor true, and the request must be a // bounded (Limit=1) search on the tenant index, filtered to // compile_kwd="wiki_page" and the dataset KBs. @@ -220,7 +226,7 @@ func TestEngineService_AvailableFor_BoundedExistenceSearch(t *testing.T) { realEngineRow("wiki/p1", "kb1", "d1", "P", "p1", nil), }} svc := NewEngineService(eng) - if !svc.AvailableFor(context.Background(), "t1", []string{"kb1"}) { + if !svc.AvailableFor(ctx, "t1", []string{"kb1"}) { t.Errorf("AvailableFor should be true when a wiki page row exists") } if eng.searchReq == nil { @@ -239,13 +245,13 @@ func TestEngineService_AvailableFor_BoundedExistenceSearch(t *testing.T) { t.Errorf("existence KbIDs = %v, want [kb1]", eng.searchReq.KbIDs) } // No engine / empty tenant / no datasets -> false. - if NewEngineService(nil).AvailableFor(context.Background(), "t1", []string{"kb1"}) { + if NewEngineService(nil).AvailableFor(ctx, "t1", []string{"kb1"}) { t.Errorf("AvailableFor should be false with no engine") } - if NewEngineService(eng).AvailableFor(context.Background(), "", []string{"kb1"}) { + if NewEngineService(eng).AvailableFor(ctx, "", []string{"kb1"}) { t.Errorf("AvailableFor should be false with empty tenant") } - if NewEngineService(eng).AvailableFor(context.Background(), "t1", nil) { + if NewEngineService(eng).AvailableFor(ctx, "t1", nil) { t.Errorf("AvailableFor should be false with no datasets") } // No wiki page rows (only ordinary chunks) -> false, even though a chunk @@ -253,7 +259,7 @@ func TestEngineService_AvailableFor_BoundedExistenceSearch(t *testing.T) { svcNoWiki := NewEngineService(&fakeDocEngine{searchRows: []map[string]interface{}{ {"id": "c1", "kb_id": "kb1", "content_with_weight": "plain chunk"}, // no compile_kwd }}) - if svcNoWiki.AvailableFor(context.Background(), "t1", []string{"kb1"}) { + if svcNoWiki.AvailableFor(ctx, "t1", []string{"kb1"}) { t.Errorf("AvailableFor should be false for a KB with no wiki page rows") } }