addressing test failures

This commit is contained in:
jacob 2023-04-05 12:24:25 -05:00
parent e24c73e883
commit 9209643ebd
11 changed files with 60 additions and 104 deletions

View file

@ -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:

View file

@ -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
}

View file

@ -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),

View file

@ -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

View file

@ -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)
}

View file

@ -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

View file

@ -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

View file

@ -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

View file

@ -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)

View file

@ -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 {

View file

@ -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