From e1f2899723216cec3e007b79c60de5e13590118d Mon Sep 17 00:00:00 2001 From: Max Date: Sun, 1 Jun 2025 16:26:39 +0800 Subject: [PATCH] Add attachment status management and filtering in neo package - Introduced new status field in the attachment table to track processing states: uploading, uploaded, indexing, indexed, upload_failed, and index_failed. - Added progress and error fields to provide detailed information during the attachment processing workflow. - Implemented AttachmentFilter struct to support filtering attachments by status. - Updated SaveAttachment and GetAttachments methods to handle new status and progress information. - Enhanced tests to validate attachment status management and filtering functionality. --- neo/store/README.md | 69 ++++++++++++++++++ neo/store/types.go | 1 + neo/store/xun.go | 15 +++- neo/store/xun_test.go | 158 ++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 242 insertions(+), 1 deletion(-) diff --git a/neo/store/README.md b/neo/store/README.md index a4c00612..ae64c1dc 100644 --- a/neo/store/README.md +++ b/neo/store/README.md @@ -299,6 +299,9 @@ CREATE TABLE neo_attachment ( gzip BOOLEAN DEFAULT FALSE INDEX, -- Compression flag bytes BIGINT INDEX, -- File size collection_id VARCHAR(200) INDEX, -- Associated knowledge collection + status ENUM('uploading', 'uploaded', 'indexing', 'indexed', 'upload_failed', 'index_failed') DEFAULT 'uploading' INDEX, -- Processing status + progress VARCHAR(200), -- Progress information (nullable) + error VARCHAR(600), -- Error message (nullable) created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP INDEX, updated_at TIMESTAMP INDEX ); @@ -358,6 +361,26 @@ type AssistantFilter struct { } ``` +#### AttachmentFilter + +```go +type AttachmentFilter struct { + UID string `json:"uid,omitempty"` // Filter by user ID + Guest *bool `json:"guest,omitempty"` // Filter by guest status + Manager string `json:"manager,omitempty"` // Filter by upload manager + ContentType string `json:"content_type,omitempty"` // Filter by content type + Name string `json:"name,omitempty"` // Filter by filename + Public *bool `json:"public,omitempty"` // Filter by public status + Gzip *bool `json:"gzip,omitempty"` // Filter by gzip compression + CollectionID string `json:"collection_id,omitempty"` // Filter by knowledge collection ID + Status string `json:"status,omitempty"` // Filter by processing status + Keywords string `json:"keywords,omitempty"` // Search in filename + Page int `json:"page,omitempty"` // Page number + PageSize int `json:"pagesize,omitempty"` // Items per page + Select []string `json:"select,omitempty"` // Fields to return +} +``` + ## Usage Examples ### 1. Chat Management @@ -439,19 +462,65 @@ attachment := map[string]interface{}{ "bytes": 102400, "collection_id": "knowledge456", "scope": []string{"user", "admin"}, + "status": "uploaded", // Status: uploading, uploaded, indexing, indexed, upload_failed, index_failed + "progress": "Upload completed", // Progress information (optional) + "error": nil, // Error message (optional, for failed statuses) } fileID, err := store.SaveAttachment(attachment) +// Update attachment status during processing workflow +attachment["status"] = "indexing" +attachment["progress"] = "Processing file for indexing..." +_, err = store.SaveAttachment(attachment) + +// Handle failed upload +attachment["status"] = "upload_failed" +attachment["progress"] = nil +attachment["error"] = "Network connection timeout" +_, err = store.SaveAttachment(attachment) + +// Complete indexing +attachment["status"] = "indexed" +attachment["progress"] = "File indexed successfully" +attachment["error"] = nil +_, err = store.SaveAttachment(attachment) + // Get attachments with filtering filter := AttachmentFilter{ UID: "user123", ContentType: "image/jpeg", + Status: "indexed", // Filter by status Page: 1, PageSize: 20, } attachments, err := store.GetAttachments(filter) + +// Get all failed uploads +failedFilter := AttachmentFilter{ + UID: "user123", + Status: "upload_failed", + Page: 1, + PageSize: 10, +} +failedUploads, err := store.GetAttachments(failedFilter) ``` +#### Attachment Status Workflow + +The attachment system supports a complete file processing workflow with the following status values: + +- **`uploading`** (default): File upload is in progress +- **`uploaded`**: File upload completed successfully +- **`indexing`**: File is being processed for search indexing +- **`indexed`**: File has been indexed and is ready for use +- **`upload_failed`**: File upload failed (check `error` field for details) +- **`index_failed`**: File indexing failed (check `error` field for details) + +#### Additional Fields + +- **`progress`**: Human-readable progress information (string, nullable) +- **`error`**: Error message for failed operations (string, nullable, max 600 characters) + ### 4. Knowledge Collection Management ```go diff --git a/neo/store/types.go b/neo/store/types.go index e5f3d344..390d3f11 100644 --- a/neo/store/types.go +++ b/neo/store/types.go @@ -235,6 +235,7 @@ type AttachmentFilter struct { Public *bool `json:"public,omitempty"` // Filter by public status Gzip *bool `json:"gzip,omitempty"` // Filter by gzip compression CollectionID string `json:"collection_id,omitempty"` // Filter by knowledge collection ID + Status string `json:"status,omitempty"` // Filter by processing status (uploading, uploaded, indexing, indexed, upload_failed, index_failed) Keywords string `json:"keywords,omitempty"` // Search in filename Page int `json:"page,omitempty"` // Page number, starting from 1 PageSize int `json:"pagesize,omitempty"` // Items per page diff --git a/neo/store/xun.go b/neo/store/xun.go index d2390e6f..eabd8176 100644 --- a/neo/store/xun.go +++ b/neo/store/xun.go @@ -367,6 +367,9 @@ func (conv *Xun) initAttachmentTable() error { table.Boolean("gzip").SetDefault(false).Index() table.BigInteger("bytes").Index() table.String("collection_id", 200).Null().Index() + table.Enum("status", []string{"uploading", "uploaded", "indexing", "indexed", "upload_failed", "index_failed"}).SetDefault("uploading").Index() // Status field enum + table.String("progress", 200).Null() // Progress information + table.String("error", 600).Null() // Error information table.TimestampTz("created_at").SetDefaultRaw("CURRENT_TIMESTAMP").Index() table.TimestampTz("updated_at").Null().Index() }) @@ -383,7 +386,7 @@ func (conv *Xun) initAttachmentTable() error { return err } - fields := []string{"id", "file_id", "uid", "guest", "manager", "content_type", "name", "public", "scope", "gzip", "bytes", "collection_id", "created_at", "updated_at"} + fields := []string{"id", "file_id", "uid", "guest", "manager", "content_type", "name", "public", "scope", "gzip", "bytes", "collection_id", "status", "progress", "error", "created_at", "updated_at"} for _, field := range fields { if !tab.HasColumn(field) { return fmt.Errorf("%s is required", field) @@ -1789,6 +1792,11 @@ func (conv *Xun) GetAttachments(filter AttachmentFilter, locale ...string) (*Att qb.Where("collection_id", filter.CollectionID) } + // Apply status filter if provided + if filter.Status != "" { + qb.Where("status", filter.Status) + } + // Apply keyword filter if provided if filter.Keywords != "" { qb.Where("name", "like", fmt.Sprintf("%%%s%%", filter.Keywords)) @@ -1945,6 +1953,11 @@ func (conv *Xun) DeleteAttachments(filter AttachmentFilter) (int64, error) { qb.Where("collection_id", filter.CollectionID) } + // Apply status filter if provided + if filter.Status != "" { + qb.Where("status", filter.Status) + } + // Apply keyword filter if provided if filter.Keywords != "" { qb.Where("name", "like", fmt.Sprintf("%%%s%%", filter.Keywords)) diff --git a/neo/store/xun_test.go b/neo/store/xun_test.go index 7c6450a3..aafc2e96 100644 --- a/neo/store/xun_test.go +++ b/neo/store/xun_test.go @@ -1289,6 +1289,9 @@ func TestXunAttachmentCRUD(t *testing.T) { "gzip": false, "bytes": 102400, "scope": []string{"user", "admin"}, + "status": "uploaded", + "progress": "100%", + "error": nil, } v, err := store.SaveAttachment(attachment) @@ -1307,10 +1310,16 @@ func TestXunAttachmentCRUD(t *testing.T) { assert.Equal(t, "test-image.jpg", attachmentData["name"]) assert.Equal(t, int64(1), attachmentData["public"]) assert.Equal(t, []interface{}{"user", "admin"}, attachmentData["scope"]) + assert.Equal(t, "uploaded", attachmentData["status"]) + assert.Equal(t, "100%", attachmentData["progress"]) + assert.Nil(t, attachmentData["error"]) // Test SaveAttachment (Update) attachment["name"] = "updated-image.jpg" attachment["bytes"] = 204800 + attachment["status"] = "indexing" + attachment["progress"] = "Processing..." + attachment["error"] = "Connection timeout" v, err = store.SaveAttachment(attachment) assert.Nil(t, err) assert.Equal(t, "test-file-123", v.(string)) @@ -1320,6 +1329,9 @@ func TestXunAttachmentCRUD(t *testing.T) { assert.Nil(t, err) assert.Equal(t, "updated-image.jpg", attachmentData["name"]) assert.Equal(t, int64(204800), attachmentData["bytes"]) + assert.Equal(t, "indexing", attachmentData["status"]) + assert.Equal(t, "Processing...", attachmentData["progress"]) + assert.Equal(t, "Connection timeout", attachmentData["error"]) // Test GetAttachments with filters resp, err := store.GetAttachments(AttachmentFilter{ @@ -1710,3 +1722,149 @@ func TestXunAttachmentFiltering(t *testing.T) { _, err = store.DeleteAttachments(AttachmentFilter{}) assert.Nil(t, err) } + +func TestXunAttachmentStatusFields(t *testing.T) { + test.Prepare(t, config.Conf) + defer test.Clean() + defer capsule.Schema().DropTableIfExists("__unit_test_conversation_attachment") + + // Drop attachment table before test + err := capsule.Schema().DropTableIfExists("__unit_test_conversation_attachment") + if err != nil { + t.Fatal(err) + } + + // Add a small delay to ensure table is created + time.Sleep(100 * time.Millisecond) + + store, err := NewXun(Setting{ + Connector: "default", + Prefix: "__unit_test_conversation_", + }) + if err != nil { + t.Fatal(err) + } + + // Clean up any existing data + _, err = store.DeleteAttachments(AttachmentFilter{}) + assert.Nil(t, err) + + // Test all possible enum status values + statusValues := []string{"uploading", "uploaded", "indexing", "indexed", "upload_failed", "index_failed"} + + for i, status := range statusValues { + // Create attachment with specific status + attachment := map[string]interface{}{ + "file_id": fmt.Sprintf("test-file-%s-%d", status, i), + "uid": "user-123", + "manager": "local", + "content_type": "image/jpeg", + "name": fmt.Sprintf("test-%s.jpg", status), + "guest": false, + "public": true, + "gzip": false, + "bytes": 102400, + "status": status, + "progress": fmt.Sprintf("%s in progress", status), + "error": nil, + } + + // Set error message for failed statuses + if status == "upload_failed" || status == "index_failed" { + attachment["error"] = fmt.Sprintf("%s error occurred", status) + } + + v, err := store.SaveAttachment(attachment) + assert.Nil(t, err) + fileID := v.(string) + + // Verify the attachment was saved with correct status + attachmentData, err := store.GetAttachment(fileID) + assert.Nil(t, err) + assert.Equal(t, status, attachmentData["status"]) + assert.Equal(t, fmt.Sprintf("%s in progress", status), attachmentData["progress"]) + + if status == "upload_failed" || status == "index_failed" { + assert.Equal(t, fmt.Sprintf("%s error occurred", status), attachmentData["error"]) + } else { + assert.Nil(t, attachmentData["error"]) + } + } + + // Test default status value (should be "uploading") + attachmentWithoutStatus := map[string]interface{}{ + "file_id": "test-file-default", + "uid": "user-123", + "manager": "local", + "content_type": "image/jpeg", + "name": "test-default.jpg", + "guest": false, + "public": true, + "gzip": false, + "bytes": 102400, + // status not specified - should use default + } + + v, err := store.SaveAttachment(attachmentWithoutStatus) + assert.Nil(t, err) + fileID := v.(string) + + // Verify default status + attachmentData, err := store.GetAttachment(fileID) + assert.Nil(t, err) + assert.Equal(t, "uploading", attachmentData["status"]) // Should be default value + assert.Nil(t, attachmentData["progress"]) // Should be null + assert.Nil(t, attachmentData["error"]) // Should be null + + // Test updating status workflow: uploading -> uploaded -> indexing -> indexed + workflowAttachment := map[string]interface{}{ + "file_id": "test-file-workflow", + "uid": "user-123", + "manager": "local", + "content_type": "text/plain", + "name": "workflow-test.txt", + "status": "uploading", + "progress": "Starting upload...", + } + + v, err = store.SaveAttachment(workflowAttachment) + assert.Nil(t, err) + workflowFileID := v.(string) + + // Update to uploaded + workflowAttachment["status"] = "uploaded" + workflowAttachment["progress"] = "Upload completed, starting indexing..." + _, err = store.SaveAttachment(workflowAttachment) + assert.Nil(t, err) + + attachmentData, err = store.GetAttachment(workflowFileID) + assert.Nil(t, err) + assert.Equal(t, "uploaded", attachmentData["status"]) + assert.Equal(t, "Upload completed, starting indexing...", attachmentData["progress"]) + + // Update to indexing + workflowAttachment["status"] = "indexing" + workflowAttachment["progress"] = "Indexing in progress..." + _, err = store.SaveAttachment(workflowAttachment) + assert.Nil(t, err) + + attachmentData, err = store.GetAttachment(workflowFileID) + assert.Nil(t, err) + assert.Equal(t, "indexing", attachmentData["status"]) + assert.Equal(t, "Indexing in progress...", attachmentData["progress"]) + + // Update to indexed (final state) + workflowAttachment["status"] = "indexed" + workflowAttachment["progress"] = "Indexing completed" + _, err = store.SaveAttachment(workflowAttachment) + assert.Nil(t, err) + + attachmentData, err = store.GetAttachment(workflowFileID) + assert.Nil(t, err) + assert.Equal(t, "indexed", attachmentData["status"]) + assert.Equal(t, "Indexing completed", attachmentData["progress"]) + + // Clean up test data + _, err = store.DeleteAttachments(AttachmentFilter{}) + assert.Nil(t, err) +}