From ac3c11a069e319fd29c760fc1d3e28250babef4c Mon Sep 17 00:00:00 2001 From: Dominik Schmidt Date: Thu, 16 Jul 2026 01:20:28 +0200 Subject: [PATCH 1/3] fix(search): report rejected documents when indexing to OpenSearch --- services/search/pkg/opensearch/batch.go | 13 ++++- services/search/pkg/opensearch/batch_test.go | 53 ++++++++++++++++++++ 2 files changed, 64 insertions(+), 2 deletions(-) create mode 100644 services/search/pkg/opensearch/batch_test.go diff --git a/services/search/pkg/opensearch/batch.go b/services/search/pkg/opensearch/batch.go index 6297b40156..16a33fd46d 100644 --- a/services/search/pkg/opensearch/batch.go +++ b/services/search/pkg/opensearch/batch.go @@ -204,10 +204,19 @@ func (b *Batch) Push() error { body.WriteString("\n") } - if _, err := b.client.Bulk(context.Background(), opensearchgoAPI.BulkReq{ + resp, err := b.client.Bulk(context.Background(), opensearchgoAPI.BulkReq{ Body: strings.NewReader(body.String()), - }); err != nil { + }) + switch { + case err != nil: return fmt.Errorf("failed to execute bulk operations: %w", err) + case resp.Errors: + items, err := json.Marshal(resp.Items) + if err != nil { + return fmt.Errorf("failed to marshal bulk response: %w", err) + } + + return fmt.Errorf("failed to execute bulk operations, response: %s", items) } bulkOperations = nil diff --git a/services/search/pkg/opensearch/batch_test.go b/services/search/pkg/opensearch/batch_test.go new file mode 100644 index 0000000000..efd29e5597 --- /dev/null +++ b/services/search/pkg/opensearch/batch_test.go @@ -0,0 +1,53 @@ +package opensearch_test + +import ( + "strings" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/opencloud-eu/opencloud/services/search/pkg/opensearch" + "github.com/opencloud-eu/opencloud/services/search/pkg/opensearch/internal/test" +) + +func TestBatch_Push(t *testing.T) { + tc := opensearchtest.NewDefaultTestClient(t, defaultConfig.Engine.OpenSearch.Client) + + t.Run("reports the documents the bulk API rejected", func(t *testing.T) { + indexName := "opencloud-test-batch-push-rejected" + tc.Require.IndicesReset([]string{indexName}) + defer tc.Require.IndicesDelete([]string{indexName}) + + // Name is a string, mapping it as a long makes every document fail to parse. + tc.Require.IndicesCreate(indexName, strings.NewReader(`{"mappings":{"properties":{"Name":{"type":"long"}}}}`)) + + batch, err := opensearch.NewBatch(tc.Client(), indexName, 10) + require.NoError(t, err) + + document := opensearchtest.Testdata.Resources.File + require.NoError(t, batch.Upsert(document.ID, document)) + + err = batch.Push() + require.Error(t, err) + require.ErrorContains(t, err, document.ID) + require.ErrorContains(t, err, "mapper_parsing_exception") + tc.Require.IndicesCount([]string{indexName}, nil, 0) + }) + + t.Run("pushes the documents the bulk API accepted", func(t *testing.T) { + indexName := "opencloud-test-batch-push-accepted" + tc.Require.IndicesReset([]string{indexName}) + defer tc.Require.IndicesDelete([]string{indexName}) + + tc.Require.IndicesCreate(indexName, strings.NewReader(opensearch.IndexManagerLatest.String())) + + batch, err := opensearch.NewBatch(tc.Client(), indexName, 10) + require.NoError(t, err) + + document := opensearchtest.Testdata.Resources.File + require.NoError(t, batch.Upsert(document.ID, document)) + require.NoError(t, batch.Push()) + + tc.Require.IndicesCount([]string{indexName}, nil, 1) + }) +} From 9beb00afb028d5211c5b74f93b92fc7809f16cb4 Mon Sep 17 00:00:00 2001 From: Dominik Schmidt Date: Thu, 16 Jul 2026 01:20:28 +0200 Subject: [PATCH 2/3] fix(search): print delete by query failures as json --- services/search/pkg/opensearch/batch.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/services/search/pkg/opensearch/batch.go b/services/search/pkg/opensearch/batch.go index 16a33fd46d..af878a235f 100644 --- a/services/search/pkg/opensearch/batch.go +++ b/services/search/pkg/opensearch/batch.go @@ -167,7 +167,7 @@ func (b *Batch) Purge(id string, onlyDeleted bool) error { case err != nil: return fmt.Errorf("failed to delete by query: %w", err) case len(resp.Failures) != 0: - return fmt.Errorf("failed to delete by query, failures: %v", resp.Failures) + return fmt.Errorf("failed to delete by query, failures: %s", resp.Failures) } return nil From 87dd68a04a00c22d5807d7e93d85da5f519dfa89 Mon Sep 17 00:00:00 2001 From: Dominik Schmidt Date: Thu, 16 Jul 2026 01:28:58 +0200 Subject: [PATCH 3/3] Apply suggestion from @dschmidt --- services/search/pkg/opensearch/batch.go | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/services/search/pkg/opensearch/batch.go b/services/search/pkg/opensearch/batch.go index af878a235f..8366a3a24f 100644 --- a/services/search/pkg/opensearch/batch.go +++ b/services/search/pkg/opensearch/batch.go @@ -211,12 +211,21 @@ func (b *Batch) Push() error { case err != nil: return fmt.Errorf("failed to execute bulk operations: %w", err) case resp.Errors: - items, err := json.Marshal(resp.Items) + var failed []opensearchgoAPI.BulkRespItem + for _, item := range resp.Items { + for _, result := range item { + if result.Error != nil { + failed = append(failed, result) + } + } + } + + failures, err := json.Marshal(failed) if err != nil { - return fmt.Errorf("failed to marshal bulk response: %w", err) + return fmt.Errorf("failed to marshal bulk failures: %w", err) } - return fmt.Errorf("failed to execute bulk operations, response: %s", items) + return fmt.Errorf("failed to execute bulk operations, failures: %s", failures) } bulkOperations = nil