From 9209643ebd6a7e82cfeac1f316572b64bdd36df2 Mon Sep 17 00:00:00 2001 From: jacob Date: Wed, 5 Apr 2023 12:24:25 -0500 Subject: [PATCH] addressing test failures --- .gitlab/.gitlab-ci.yml | 9 +++--- cli/cli_kafka_integration_test.go | 7 +++-- cli/kafka.go | 7 +---- cli/kafka/config.go | 49 +++++++++++++++++-------------- cli/kafka/runner.go | 5 ++-- idk/Makefile | 13 ++++++++ idk/ingest.go | 42 ++------------------------ idk/kafka/source.go | 2 +- idk/kafka_sasl/cmd_test.go | 1 + idk/kafka_sasl/source.go | 10 ++----- wire_response.go | 19 ------------ 11 files changed, 60 insertions(+), 104 deletions(-) diff --git a/.gitlab/.gitlab-ci.yml b/.gitlab/.gitlab-ci.yml index 3cf8c1676..3482f3123 100644 --- a/.gitlab/.gitlab-ci.yml +++ b/.gitlab/.gitlab-ci.yml @@ -375,10 +375,7 @@ run go tests cli kafka integration: script: - echo "running fbsql integration tests" - cd ./idk/ - - BRANCH_NAME=${CI_COMMIT_REF_SLUG} make start-pilosa - - BRANCH_NAME=${CI_COMMIT_REF_SLUG} make start-zookeeper - - BRANCH_NAME=${CI_COMMIT_REF_SLUG} make start-kafka - - BRANCH_NAME=${CI_COMMIT_REF_SLUG} make start-schema-registry + - BRANCH_NAME=${CI_COMMIT_REF_SLUG} make cli-startup - cd .. - go test -coverprofile=coverage-cli-kafka-integration.out -run -timeout=10m TestKafkaRunner ./cli after_script: @@ -532,6 +529,7 @@ package for linux amd64: script: - echo 'deb [trusted=yes] https://repo.goreleaser.com/apt/ /' | tee /etc/apt/sources.list.d/goreleaser.list - apt update && apt install nfpm=2.11.3 + - apt-get update && apt-get install -y -qq musl-tools build-essential - make package artifacts: paths: @@ -562,6 +560,8 @@ package for linux arm64: stage: build image: golang:$GOVERSION extends: .go-cache + tags: + - shell-arm64 rules: - if: '$CI_PIPELINE_SOURCE == "push" || $CI_PIPELINE_SOURCE == "schedule" || $CI_PIPELINE_SOURCE == "web"' variables: @@ -570,6 +570,7 @@ package for linux arm64: script: - echo 'deb [trusted=yes] https://repo.goreleaser.com/apt/ /' | tee /etc/apt/sources.list.d/goreleaser.list - apt update && apt install nfpm=2.11.3 + - apt-get update && apt-get install -y -qq build-essential - make package artifacts: paths: diff --git a/cli/cli_kafka_integration_test.go b/cli/cli_kafka_integration_test.go index c8ca6ec8c..bac2e7615 100644 --- a/cli/cli_kafka_integration_test.go +++ b/cli/cli_kafka_integration_test.go @@ -8,6 +8,7 @@ import ( "fmt" "io/ioutil" "os" + "sort" "strconv" "strings" "testing" @@ -340,9 +341,11 @@ func verifyQueryReponse(t *testing.T, wqr *featurebase.WireQueryResponse, expect for j, element := range line { switch newElement := element.(type) { case featurebase.StringSet: - newline[j] = newElement.SortedStringSlice() + sort.Strings(newElement) + newline[j] = newElement case featurebase.IDSet: - newline[j] = newElement.SortedInt64Slice() + sort.Slice(newElement, func(i, j int) bool { return newElement[i] < newElement[j] }) + newline[j] = newElement default: newline[j] = element } diff --git a/cli/kafka.go b/cli/kafka.go index 1392dbc9f..23dc841d2 100644 --- a/cli/kafka.go +++ b/cli/kafka.go @@ -10,7 +10,7 @@ func (cmd *Command) newKafkaRunner(cfgFile string) (*kafka.Runner, error) { cfg, err := kafka.ConfigFromFile(cfgFile) if err != nil { - return nil, err + return nil, errors.Wrap(err, "getting config from file") } if err := kafka.ValidateConfig(cfg); err != nil { @@ -49,11 +49,6 @@ func (cmd *Command) newKafkaRunner(cfgFile string) (*kafka.Runner, error) { return nil, errors.Wrap(err, "getting fields from config") } - // for avro, let the SchemaManager and IDK handle fields - if cfg.Encode == "avro" { - flds = nil - } - return kafka.NewRunner( idkCfg, batch.NewSQLBatcher(cmd, flds), diff --git a/cli/kafka/config.go b/cli/kafka/config.go index 4e09b9dc2..fa044ad06 100644 --- a/cli/kafka/config.go +++ b/cli/kafka/config.go @@ -35,7 +35,7 @@ type Config struct { Encode string `mapstructure:"encode" help:"Encoding format (currently supported formats: avro, json)"` AllowMissingFields bool `mapstructure:"allow-missing-fields" help:"allow missing fields in messages from kafka"` MaxMessages int `mapstructure:"max-messages" help:"max messages read from kakfka"` - ConfluentConfig string `mapstructure:"confluent-config" help:"max messages read from kakfka"` + ConfluentConfig string `mapstructure:"confluent-config" help:"path to JSON file mapping librdkafka consumer configurations to configuration values"` } // Field is a user-facing configuration field. @@ -117,40 +117,40 @@ func ValidateConfig(c Config) error { return validateConfigJSON(c) case encodingTypeAvro: return validateConfigAvro(c) + default: + return errors.Errorf("encode configuration value must be %s or %s: got %s", encodingTypeJSON, encodingTypeAvro, c.Encode) } - - return nil - } func validateConfigJSON(c Config) error { - if len(c.Fields) > 0 { + switch len(c.Fields) { + case 0: // We only need to do these checks if any fields are specified at all. // If no fields are specified, that's ok because then we default to // using fields based off the existing table. - if len(c.Fields) < 2 { - return errors.Errorf("at least two fields are required (one should be a primary key)") - } else { - var found int - for i := range c.Fields { - if c.Fields[i].PrimaryKey { - found++ - } - if c.Fields[i].Name == "" { - return errors.Errorf("a name attribute (which isn't equal to \"\") should exist for all fields") - } - if c.Fields[i].SourceType == "" { - return errors.Errorf("a source-type attribute (which isn't equal to \"\") should exist for all fields") - } + return nil + case 1: + return errors.Errorf("at least two fields are required (one should be a primary key)") + default: + var found int + for i := range c.Fields { + if c.Fields[i].PrimaryKey { + found++ } - if found < 1 { - return errors.Errorf("at least one primary key field is required") + if c.Fields[i].Name == "" { + return errors.Errorf("a name attribute (which isn't equal to \"\") should exist for all fields") + } + if c.Fields[i].SourceType == "" { + return errors.Errorf("a source-type attribute (which isn't equal to \"\") should exist for all fields") } } + if found < 1 { + return errors.Errorf("at least one primary key field is required") + } + return nil } - return nil } // Only primary key fields required @@ -282,6 +282,11 @@ func ConfigToFields(c Config, primaryKeys []string) ([]*dax.Field, error) { // capacity to `len(c.Fields)-1`. out := make([]*dax.Field, 0, len(c.Fields)) + // for avro, let the SchemaManager and IDK handle fields + if c.Encode == encodingTypeAvro { + return nil, nil + } + for _, fld := range c.Fields { // When we have a single primary key, don't also store that value as a // field in FeatureBase. However, when we have more than one primary diff --git a/cli/kafka/runner.go b/cli/kafka/runner.go index 5c92a4806..0c82a6056 100644 --- a/cli/kafka/runner.go +++ b/cli/kafka/runner.go @@ -51,9 +51,10 @@ func NewRunner(cfg ConfigForIDK, batcher fbbatch.Batcher, logWriter io.Writer) * } // NewSource should be set based on the encoding of the source (e.g. JSON, Avro) - if cfg.Encode == encodingTypeAvro { + switch cfg.Encode { + case encodingTypeAvro: kr.GetAvroNewSource(cfg) - } else if cfg.Encode == encodingTypeJSON { + default: kr.GetJSONNewSource(cfg) } diff --git a/idk/Makefile b/idk/Makefile index 237b2aae8..acafadc46 100644 --- a/idk/Makefile +++ b/idk/Makefile @@ -338,3 +338,16 @@ update-mocks-%: install-mock-generator $(eval AWS_SDK_VERSION := $(shell grep 'github.com/aws/aws-sdk-go' ../go.mod | cut -d ' ' -f 2)) echo Generating mock for AWS service $(AWS_SERVICE) and SDK version $(AWS_SDK_VERSION) && \ $(GOPATH)/bin/mockery --name $*API --output idktest/mocks --filename $(AWS_SERVICE).go --dir $(GOPATH)/pkg/mod/github.com/aws/aws-sdk-go@$(AWS_SDK_VERSION)/service/$(AWS_SERVICE)/$(AWS_SERVICE)iface + + +### CLI Test Services + +cli-startup: + apt-get update + apt-get install -y pip git + pip install docker-compose + git clone https://github.com/square/certstrap + cd certstrap + go build + $(MAKE) testenv vendor + $(MAKE) startup \ No newline at end of file diff --git a/idk/ingest.go b/idk/ingest.go index 1bb834d4f..5ffac6381 100644 --- a/idk/ingest.go +++ b/idk/ingest.go @@ -309,51 +309,13 @@ func (m *Main) clone() (*Main, error) { index = schema.Index(m.Index) - // use a copy (schema race condition issu) + // use a copy (schema race condition issues) mClone := *m mClone.index = index return &mClone, nil } -// func (m *Main) clone() (*Main, error) { -// var index *pilosaclient.Index -// noOpSchemaManager := false - -// // If you have a schema manager, it does it's thing. Otherwise, it's a no -// // opt manager. Then you get a default schema. If you get a default schema, -// // you get a default index. This seems fine for IDK. However, for the CLI / -// // SQL kafka runner, we use the m.index to create the dax.Table. If m.index -// // is set to default values, then keys is false even when we don't want it -// // to be. This was causing issue in buildBulkInsert in the batch package of -// // the CLI. -// switch m.SchemaManager.(type) { -// case *nopSchemaManager: -// noOpSchemaManager = true -// } - -// schema, err := m.SchemaManager.Schema() -// if err != nil { -// return nil, err -// } - -// if noOpSchemaManager && len(m.PrimaryKeyFields) > 0 { -// // if IDK has PrimaryKeyFields, it expects keyed index -// keys := pilosaclient.OptIndexKeys(true) -// // most queries don't work with this set to false so set to true -// exists := pilosaclient.OptIndexTrackExistence(true) -// index = schema.Index(m.Index, keys, exists) -// } else { -// index = schema.Index(m.Index) -// } - -// // use a copy (schema race condition issues) -// mClone := *m -// mClone.index = index - -// return &mClone, nil -// } - func (m *Main) runIngester(c int, l *msgCounter) error { m.log.Printf("start ingester %d", c) // TODO: actually implement cancellation and graceful shutdown @@ -2203,7 +2165,7 @@ func (m *Main) newBatch(clientFields []*pilosaclient.Field) (pilosabatch.RecordB ii := pilosaclient.FromClientIndex(m.index) tbl := pilosacore.IndexInfoToTable(ii) - // Fields. // this is giving bad fieldInfos + // Fields. fields := pilosaclient.FromClientFields(clientFields) // If a custom Batcher has been defined, use that. Otherwise default to diff --git a/idk/kafka/source.go b/idk/kafka/source.go index 00adf0bff..78abec914 100644 --- a/idk/kafka/source.go +++ b/idk/kafka/source.go @@ -268,7 +268,7 @@ func (r *Record) Data() []interface{} { func (s *Source) Open() error { cfg, err := common.SetupConfluent(&s.ConfluentCommand) if err != nil { - return err + return errors.Wrap(err, "setting up confluent command") } s.ConfigMap = cfg diff --git a/idk/kafka_sasl/cmd_test.go b/idk/kafka_sasl/cmd_test.go index 7314f1740..b248ed56a 100644 --- a/idk/kafka_sasl/cmd_test.go +++ b/idk/kafka_sasl/cmd_test.go @@ -1022,6 +1022,7 @@ func TestMaxMsgs(t *testing.T) { producer.Close() collector.reset() + fmt.Println("able to produce messages") err = m.Run() if err != nil { t.Fatalf("running main: %v", err) diff --git a/idk/kafka_sasl/source.go b/idk/kafka_sasl/source.go index feb063c4e..c18239545 100644 --- a/idk/kafka_sasl/source.go +++ b/idk/kafka_sasl/source.go @@ -14,7 +14,6 @@ import ( confluent "github.com/confluentinc/confluent-kafka-go/kafka" "github.com/featurebasedb/featurebase/v3/idk" - "github.com/featurebasedb/featurebase/v3/idk/common" "github.com/featurebasedb/featurebase/v3/logger" "github.com/pkg/errors" ) @@ -219,17 +218,12 @@ func (r *Record) Data() []interface{} { // Open initializes the kafka source. func (s *Source) Open() error { - cfg, err := common.SetupConfluent(&s.ConfluentCommand) - if err != nil { - return err - } - s.ConfigMap = cfg - if len(s.Header) == 0 && len(s.HeaderFields) == 0 { - return errors.New("needs header specification file (file or fields)") + return errors.New("needs header specification (from file or from existing fields)") } var headerData []byte + var err error if s.Header != "" { headerData, err = os.ReadFile(s.Header) if err != nil { diff --git a/wire_response.go b/wire_response.go index 90eac3a1e..dfb08bc00 100644 --- a/wire_response.go +++ b/wire_response.go @@ -5,7 +5,6 @@ import ( "encoding/json" "fmt" "log" - "sort" "strings" "time" @@ -211,15 +210,6 @@ func (ii IDSet) String() string { return sb.String() } -// SortedInt64Slice returns the values in a IDSet field in a int64 slice that is -// sorted -func (ii IDSet) SortedInt64Slice() []int64 { - var idSetSlice = make([]int64, len(ii)) - copy(idSetSlice, ii) - sort.Slice(idSetSlice, func(i, j int) bool { return idSetSlice[i] < idSetSlice[j] }) - return idSetSlice -} - // StringSet is a return type specific to SQLResponse types. type StringSet []string @@ -237,15 +227,6 @@ func (ss StringSet) String() string { return sb.String() } -// SortedStringSlice returns the values in a StringSet field in a string slice -// that is sorted -func (ss StringSet) SortedStringSlice() []string { - var stringSetSlice = make([]string, len(ss)) - copy(stringSetSlice, ss) - sort.Strings(stringSetSlice) - return stringSetSlice -} - // ShowColumnsResponse returns a structure which is specific to a `SHOW COLUMNS` // statement, derived from the results in the WireQueryResponse. This is kind of // a crude way to unmarshal a WireQueryResponse into a type which is specific to