Compare commits

...
2 changed files with 79 additions and 3 deletions

No files matched your search

+26 -3
View File
@@ -173,7 +173,12 @@ 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)
failures, err := json.Marshal(resp.Failures)
if err != nil {
return fmt.Errorf("failed to marshal delete by query failures: %w", err)
}
return fmt.Errorf("failed to delete by query, failures: %s", failures)
}
return nil
@@ -210,11 +215,29 @@ 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()),
Params: opensearchgoAPI.BulkParams{Refresh: "wait_for"},
}); err != nil {
})
switch {
case err != nil:
return fmt.Errorf("failed to execute bulk operations: %w", err)
case resp.Errors:
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 failures: %w", err)
}
return fmt.Errorf("failed to execute bulk operations, failures: %s", failures)
}
bulkOperations = nil
@@ -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)
})
}