mirror of
https://github.com/opencloud-eu/opencloud.git
synced 2026-09-08 11:53:07 -04:00
Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0cd6e84c47 | ||
|
|
3d079b90f0 | ||
|
|
d7b8feaa18 |
No files matched your search
@@ -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)
|
||||
})
|
||||
}
|
||||
Reference in new issue
Block a user