Merge pull request #1136 from trheyi/main
Refactor document processing functions to improve job handling
This commit is contained in:
commit
783ab7256e
4 changed files with 6 additions and 111 deletions
|
|
@ -352,7 +352,7 @@ func AddFileAsync(c *gin.Context) {
|
||||||
|
|
||||||
err = j.Add(&job.ExecutionOptions{
|
err = j.Add(&job.ExecutionOptions{
|
||||||
Priority: 1,
|
Priority: 1,
|
||||||
}, "kb.documents.processfile", jobData)
|
}, "kb.documents.addfile", jobData)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("AddFileAsync: Failed to add job execution: %v", err)
|
log.Error("AddFileAsync: Failed to add job execution: %v", err)
|
||||||
// Rollback: remove document record
|
// Rollback: remove document record
|
||||||
|
|
@ -452,40 +452,6 @@ func ProcessAddFile(process *process.Process) interface{} {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ProcessProcessFile documents.processfile Knowledge Base process file content processor (async version)
|
|
||||||
// Args[0] map: Request parameters {"collection_id": "collection", "file_id": "file123", "uploader": "local", ...}
|
|
||||||
// Return: map: Response data {"doc_id": "document_id"}
|
|
||||||
func ProcessProcessFile(process *process.Process) interface{} {
|
|
||||||
process.ValidateArgNums(1)
|
|
||||||
|
|
||||||
// Get parameters
|
|
||||||
reqMap := process.ArgsMap(0)
|
|
||||||
|
|
||||||
// Check knowledge base instance
|
|
||||||
if kb.Instance == nil {
|
|
||||||
exception.New("knowledge base not initialized", 500).Throw()
|
|
||||||
}
|
|
||||||
|
|
||||||
// Convert parameters to AddFileRequest structure
|
|
||||||
req := parseAddFileRequest(reqMap)
|
|
||||||
|
|
||||||
// This is async version - document record should already exist
|
|
||||||
// Just process the file content
|
|
||||||
ctx := process.Context
|
|
||||||
if ctx == nil {
|
|
||||||
ctx = context.Background()
|
|
||||||
}
|
|
||||||
err := HandleFileContent(ctx, req)
|
|
||||||
if err != nil {
|
|
||||||
exception.New("failed to process file content: %s", 500, err.Error()).Throw()
|
|
||||||
}
|
|
||||||
|
|
||||||
// Return result
|
|
||||||
return maps.MapStrAny{
|
|
||||||
"doc_id": req.DocID,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// parseAddFileRequest parses request map into AddFileRequest structure
|
// parseAddFileRequest parses request map into AddFileRequest structure
|
||||||
func parseAddFileRequest(reqMap map[string]interface{}) *AddFileRequest {
|
func parseAddFileRequest(reqMap map[string]interface{}) *AddFileRequest {
|
||||||
req := &AddFileRequest{}
|
req := &AddFileRequest{}
|
||||||
|
|
|
||||||
|
|
@ -308,7 +308,7 @@ func AddTextAsync(c *gin.Context) {
|
||||||
|
|
||||||
err = j.Add(&job.ExecutionOptions{
|
err = j.Add(&job.ExecutionOptions{
|
||||||
Priority: 1,
|
Priority: 1,
|
||||||
}, "kb.documents.processtext", jobData)
|
}, "kb.documents.addtext", jobData)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("AddTextAsync: Failed to add job execution: %v", err)
|
log.Error("AddTextAsync: Failed to add job execution: %v", err)
|
||||||
// Rollback: remove document record
|
// Rollback: remove document record
|
||||||
|
|
@ -408,40 +408,6 @@ func ProcessAddText(process *process.Process) interface{} {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ProcessProcessText documents.processtext Knowledge Base process text content processor (async version)
|
|
||||||
// Args[0] map: Request parameters {"collection_id": "collection", "text": "content", ...}
|
|
||||||
// Return: map: Response data {"doc_id": "document_id"}
|
|
||||||
func ProcessProcessText(process *process.Process) interface{} {
|
|
||||||
process.ValidateArgNums(1)
|
|
||||||
|
|
||||||
// Get parameters
|
|
||||||
reqMap := process.ArgsMap(0)
|
|
||||||
|
|
||||||
// Check knowledge base instance
|
|
||||||
if kb.Instance == nil {
|
|
||||||
exception.New("knowledge base not initialized", 500).Throw()
|
|
||||||
}
|
|
||||||
|
|
||||||
// Convert parameters to AddTextRequest structure
|
|
||||||
req := parseAddTextRequest(reqMap)
|
|
||||||
|
|
||||||
// This is async version - document record should already exist
|
|
||||||
// Just process the text content
|
|
||||||
ctx := process.Context
|
|
||||||
if ctx == nil {
|
|
||||||
ctx = context.Background()
|
|
||||||
}
|
|
||||||
err := HandleTextContent(ctx, req)
|
|
||||||
if err != nil {
|
|
||||||
exception.New("failed to process text content: %s", 500, err.Error()).Throw()
|
|
||||||
}
|
|
||||||
|
|
||||||
// Return result
|
|
||||||
return maps.MapStrAny{
|
|
||||||
"doc_id": req.DocID,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// parseAddTextRequest parses request map into AddTextRequest structure
|
// parseAddTextRequest parses request map into AddTextRequest structure
|
||||||
func parseAddTextRequest(reqMap map[string]interface{}) *AddTextRequest {
|
func parseAddTextRequest(reqMap map[string]interface{}) *AddTextRequest {
|
||||||
req := &AddTextRequest{}
|
req := &AddTextRequest{}
|
||||||
|
|
|
||||||
|
|
@ -308,7 +308,7 @@ func AddURLAsync(c *gin.Context) {
|
||||||
|
|
||||||
err = j.Add(&job.ExecutionOptions{
|
err = j.Add(&job.ExecutionOptions{
|
||||||
Priority: 1,
|
Priority: 1,
|
||||||
}, "kb.documents.processurl", jobData)
|
}, "kb.documents.addurl", jobData)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("AddURLAsync: Failed to add job execution: %v", err)
|
log.Error("AddURLAsync: Failed to add job execution: %v", err)
|
||||||
// Rollback: remove document record
|
// Rollback: remove document record
|
||||||
|
|
@ -408,40 +408,6 @@ func ProcessAddURL(process *process.Process) interface{} {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// ProcessProcessURL documents.processurl Knowledge Base process URL content processor (async version)
|
|
||||||
// Args[0] map: Request parameters {"collection_id": "collection", "url": "https://example.com", ...}
|
|
||||||
// Return: map: Response data {"doc_id": "document_id"}
|
|
||||||
func ProcessProcessURL(process *process.Process) interface{} {
|
|
||||||
process.ValidateArgNums(1)
|
|
||||||
|
|
||||||
// Get parameters
|
|
||||||
reqMap := process.ArgsMap(0)
|
|
||||||
|
|
||||||
// Check knowledge base instance
|
|
||||||
if kb.Instance == nil {
|
|
||||||
exception.New("knowledge base not initialized", 500).Throw()
|
|
||||||
}
|
|
||||||
|
|
||||||
// Convert parameters to AddURLRequest structure
|
|
||||||
req := parseAddURLRequest(reqMap)
|
|
||||||
|
|
||||||
// This is async version - document record should already exist
|
|
||||||
// Just process the URL content
|
|
||||||
ctx := process.Context
|
|
||||||
if ctx == nil {
|
|
||||||
ctx = context.Background()
|
|
||||||
}
|
|
||||||
err := HandleURLContent(ctx, req)
|
|
||||||
if err != nil {
|
|
||||||
exception.New("failed to process URL content: %s", 500, err.Error()).Throw()
|
|
||||||
}
|
|
||||||
|
|
||||||
// Return result
|
|
||||||
return maps.MapStrAny{
|
|
||||||
"doc_id": req.DocID,
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
// parseAddURLRequest parses request map into AddURLRequest structure
|
// parseAddURLRequest parses request map into AddURLRequest structure
|
||||||
func parseAddURLRequest(reqMap map[string]interface{}) *AddURLRequest {
|
func parseAddURLRequest(reqMap map[string]interface{}) *AddURLRequest {
|
||||||
req := &AddURLRequest{}
|
req := &AddURLRequest{}
|
||||||
|
|
|
||||||
|
|
@ -11,12 +11,9 @@ import (
|
||||||
func init() {
|
func init() {
|
||||||
// Register kb process handlers
|
// Register kb process handlers
|
||||||
process.RegisterGroup("kb", map[string]process.Handler{
|
process.RegisterGroup("kb", map[string]process.Handler{
|
||||||
"documents.addfile": ProcessAddFile,
|
"documents.addfile": ProcessAddFile,
|
||||||
"documents.addtext": ProcessAddText,
|
"documents.addtext": ProcessAddText,
|
||||||
"documents.addurl": ProcessAddURL,
|
"documents.addurl": ProcessAddURL,
|
||||||
"documents.processfile": ProcessProcessFile,
|
|
||||||
"documents.processtext": ProcessProcessText,
|
|
||||||
"documents.processurl": ProcessProcessURL,
|
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue