From 901cf0f96e8b48c2788063bd025dfd272c5eb41c Mon Sep 17 00:00:00 2001 From: Christopher Lowenthal Date: Tue, 11 Oct 2022 10:59:31 -0400 Subject: [PATCH] extracts transaction function --- api_import_worker.go | 101 ++++++++++++++++++++++-------------------- etcd/leasedkv_test.go | 1 - executor_test.go | 3 +- 3 files changed, 53 insertions(+), 52 deletions(-) diff --git a/api_import_worker.go b/api_import_worker.go index 043f7eb61..aa4145e04 100644 --- a/api_import_worker.go +++ b/api_import_worker.go @@ -48,57 +48,60 @@ func importWorkerFunc(job importJob) error { } } - if err := func() (err1 error) { - tx, finisher, err := job.qcx.GetTx(Txo{Write: writable, Index: job.field.idx, Shard: job.shard}) - if err != nil { - return err - } - defer finisher(&err1) - - var doClear bool - switch doAction { - case RequestActionOverwrite: - err := job.field.importRoaringOverwrite(job.ctx, tx, viewData, job.shard, viewName, job.req.Block) - if err != nil { - return errors.Wrap(err, "importing roaring as overwrite") - } - case RequestActionClear: - doClear = true - fallthrough - case RequestActionSet: - fileMagic := uint32(binary.LittleEndian.Uint16(viewData[0:2])) - data := viewData - if fileMagic != roaring.MagicNumber { - // if the view data arrives is in the "standard" roaring format, we must - // make a copy of data in order allow for the conversion to the pilosa roaring run format - // in field.importRoaring - data = make([]byte, len(viewData)) - copy(data, viewData) - } - if job.req.UpdateExistence { - if ef := job.field.idx.existenceField(); ef != nil { - existence, err := combineForExistence(data) - if err != nil { - return errors.Wrap(err, "merging existence on roaring import") - } - - err = ef.importRoaring(job.ctx, tx, existence, job.shard, "standard", false) - if err != nil { - return errors.Wrap(err, "updating existence on roaring import") - } - } - } - - err := job.field.importRoaring(job.ctx, tx, data, job.shard, viewName, doClear) - - if err != nil { - return errors.Wrap(err, "importing standard roaring") - } - } - return nil - }(); err != nil { + if err := importWorkerTx(job, doAction, viewName, viewData); err != nil { return err } } return nil } + +func importWorkerTx(job importJob, doAction string, viewName string, viewData []byte) (txErr error) { + tx, finisher, err := job.qcx.GetTx(Txo{Write: writable, Index: job.field.idx, Shard: job.shard}) + if err != nil { + return err + } + defer finisher(&txErr) + + var doClear bool + switch doAction { + case RequestActionOverwrite: + err := job.field.importRoaringOverwrite(job.ctx, tx, viewData, job.shard, viewName, job.req.Block) + if err != nil { + return errors.Wrap(err, "importing roaring as overwrite") + } + case RequestActionClear: + doClear = true + fallthrough + case RequestActionSet: + fileMagic := uint32(binary.LittleEndian.Uint16(viewData[0:2])) + data := viewData + if fileMagic != roaring.MagicNumber { + // if the view data arrives is in the "standard" roaring format, we must + // make a copy of data in order allow for the conversion to the pilosa roaring run format + // in field.importRoaring + data = make([]byte, len(viewData)) + copy(data, viewData) + } + if job.req.UpdateExistence { + if ef := job.field.idx.existenceField(); ef != nil { + existence, err := combineForExistence(data) + if err != nil { + return errors.Wrap(err, "merging existence on roaring import") + } + + err = ef.importRoaring(job.ctx, tx, existence, job.shard, "standard", false) + if err != nil { + return errors.Wrap(err, "updating existence on roaring import") + } + } + } + + err := job.field.importRoaring(job.ctx, tx, data, job.shard, viewName, doClear) + + if err != nil { + return errors.Wrap(err, "importing standard roaring") + } + } + return nil + +} diff --git a/etcd/leasedkv_test.go b/etcd/leasedkv_test.go index b87789d06..4c9ec9fb8 100644 --- a/etcd/leasedkv_test.go +++ b/etcd/leasedkv_test.go @@ -12,7 +12,6 @@ import ( "github.com/featurebasedb/featurebase/v3/logger" "github.com/featurebasedb/featurebase/v3/testhook" "github.com/pkg/errors" - "github.com/stretchr/testify/assert" "go.etcd.io/etcd/server/v3/embed" "go.etcd.io/etcd/server/v3/etcdserver/api/v3client" diff --git a/executor_test.go b/executor_test.go index 9ee6f8473..ee3ce29ec 100644 --- a/executor_test.go +++ b/executor_test.go @@ -31,8 +31,7 @@ import ( "github.com/featurebasedb/featurebase/v3/server" "github.com/featurebasedb/featurebase/v3/test" "github.com/featurebasedb/featurebase/v3/testhook" - . "github.com/featurebasedb/featurebase/v3/vprint" // nolint:staticcheck - . "github.com/featurebasedb/featurebase/v3/vprint" // nolint:staticcheck + _ "github.com/featurebasedb/featurebase/v3/vprint" // nolint:staticcheck "github.com/google/go-cmp/cmp" "github.com/pkg/errors" )