mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-09-07 09:05:55 +00:00
Testing unicode is great, but we appear to have had a couple
of cases where we were using strings that weren't valid UTF8.
Weirdly, other instances of these strings work -- I think because
they're in raw quotes (backticks) rather than strings. Anyway,
this is what SonarCloud fusses about.
(cherry picked from commit ab543adbe3)
733 lines
20 KiB
Go
733 lines
20 KiB
Go
//go:build !kafka_sasl
|
|
// +build !kafka_sasl
|
|
|
|
package kafka_static
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"math/rand"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/pkg/errors"
|
|
"github.com/segmentio/kafka-go"
|
|
|
|
pilosaclient "github.com/featurebasedb/featurebase/v3/client"
|
|
"github.com/featurebasedb/featurebase/v3/idk/idktest"
|
|
)
|
|
|
|
var pilosaHost string
|
|
var pilosaTLSHost string
|
|
var pilosaGrpcHost string
|
|
var kafkaHost string
|
|
var certPath string
|
|
|
|
func init() {
|
|
var ok bool
|
|
if pilosaHost, ok = os.LookupEnv("IDK_TEST_PILOSA_HOST"); !ok {
|
|
pilosaHost = "pilosa:10101"
|
|
}
|
|
if pilosaTLSHost, ok = os.LookupEnv("IDK_TEST_PILOSA_TLS_HOST"); !ok {
|
|
pilosaTLSHost = "https://pilosa-tls:10111"
|
|
}
|
|
if pilosaGrpcHost, ok = os.LookupEnv("IDK_TEST_PILOSA_GRPC_HOST"); !ok {
|
|
pilosaGrpcHost = "pilosa:20101"
|
|
}
|
|
if kafkaHost, ok = os.LookupEnv("IDK_TEST_KAFKA_HOST"); !ok {
|
|
kafkaHost = "kafka:9092"
|
|
}
|
|
if certPath, ok = os.LookupEnv("IDK_TEST_CERT_PATH"); !ok {
|
|
certPath = "/certs"
|
|
}
|
|
}
|
|
|
|
func configureTestFlags(main *Main) {
|
|
main.PilosaHosts = []string{pilosaHost}
|
|
main.PilosaGRPCHosts = []string{pilosaGrpcHost}
|
|
main.KafkaHosts = []string{kafkaHost}
|
|
// intentionally low timeout — if this gets triggered it shouldn't
|
|
// have any negative effects
|
|
main.Timeout = time.Millisecond * 20
|
|
main.Stats = ""
|
|
_, main.Verbose = os.LookupEnv("IDK_TEST_VERBOSE")
|
|
}
|
|
|
|
func TestFieldTypes(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
fieldNames := []string{"i", "d", "t", "@s", "@st", "unixtime", "e"}
|
|
records := [][]interface{}{
|
|
{100, 34.0404, "2021-02-01", "apple", "egg", 1617246530, "Lorem ipsum dolor sit amet"}, // Thu Apr 01 03:08:50 2021 UTC
|
|
}
|
|
sPilosaName := "s"
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
a := rand.Int()
|
|
topic := "xyz" + strconv.Itoa(a)
|
|
|
|
// create Main and run with MaxMsgs
|
|
m := NewMain()
|
|
configureTestFlags(m)
|
|
m.Index = fmt.Sprintf("cmd_test_index223ij%s", topic)
|
|
m.AutoGenerate = true
|
|
m.Header = "testdata/TestFieldTypes.json"
|
|
m.PackBools = "bools"
|
|
m.BatchSize = 1
|
|
m.Topics = []string{topic}
|
|
m.MaxMsgs = uint64(len(records))
|
|
m.PilosaHosts = []string{pilosaHost}
|
|
m.Timeout = time.Minute
|
|
m.LookupDBDSN = "postgresql://postgres:password@postgres:5432/postgres?sslmode=disable"
|
|
m.LookupBatchSize = 1
|
|
|
|
// put records in kafka
|
|
addr := kafka.TCP(kafkaHost)
|
|
tCreateTopic(t, topic, addr)
|
|
writer := &kafka.Writer{
|
|
Addr: addr,
|
|
Topic: topic,
|
|
Balancer: &kafka.Hash{},
|
|
BatchTimeout: time.Nanosecond,
|
|
}
|
|
defer writer.Close()
|
|
|
|
for _, vals := range records {
|
|
rec := makeRecordString(t, fieldNames, vals)
|
|
tPutStringKafka(t, writer, "akey", rec)
|
|
}
|
|
|
|
err := m.Run()
|
|
if err != nil {
|
|
t.Fatalf("running main: %v", err)
|
|
}
|
|
|
|
client := m.PilosaClient()
|
|
schema, err := client.Schema()
|
|
if err != nil {
|
|
t.Fatalf("getting client: %v", err)
|
|
}
|
|
index := schema.Index(m.Index)
|
|
defer func() {
|
|
err := client.DeleteIndex(index)
|
|
if err != nil {
|
|
t.Logf("deleting index: %v", err)
|
|
}
|
|
}()
|
|
|
|
if status, body, err := client.HTTPRequest("POST", "/recalculate-caches", nil, nil); err != nil {
|
|
t.Fatalf("recalculating cache: status: %d, response: %s, error: %s", status, body, err)
|
|
}
|
|
|
|
// check data in Pilosa
|
|
if !index.HasField(sPilosaName) {
|
|
t.Fatalf("don't have field '%s'", sPilosaName)
|
|
}
|
|
|
|
fields := index.Field(sPilosaName)
|
|
stPilosaName := "st"
|
|
fieldSt := index.Field(stPilosaName)
|
|
|
|
qr, err := client.Query(index.Count(fields.Row("apple")))
|
|
if err != nil {
|
|
t.Errorf("querying: %v", err)
|
|
}
|
|
if qr.Result().Count() != 1 {
|
|
t.Errorf("wrong count for field '%s', %d is not 1", sPilosaName, qr.Result().Count())
|
|
}
|
|
|
|
qr, err = client.Query(index.Count(fieldSt.Row("egg")))
|
|
if err != nil {
|
|
t.Fatalf("querying time range for egg: %v", err)
|
|
}
|
|
if qr.Result().Count() != 1 {
|
|
t.Errorf("wrong count for field '%s', %d is not 1", stPilosaName, qr.Result().Count())
|
|
}
|
|
qr, err = client.Query(index.Count(fieldSt.Range("egg", time.Unix(1617145530, 0), time.Unix(1617348530, 0))))
|
|
if err != nil {
|
|
t.Fatalf("querying time range for egg: %v", err)
|
|
}
|
|
if qr.Result().Count() != 1 {
|
|
t.Errorf("wrong count for field '%s' with time range, %d is not 1", stPilosaName, qr.Result().Count())
|
|
}
|
|
|
|
if !index.HasField("i") {
|
|
t.Fatal("don't have field 'i'")
|
|
}
|
|
fieldi := index.Field("i")
|
|
|
|
qr, err = client.Query(index.Count(fieldi.Row(100)))
|
|
if err != nil {
|
|
t.Errorf("querying: %v", err)
|
|
}
|
|
if qr.Result().Count() != 1 {
|
|
t.Errorf("wrong count for field 'i', %d is not 1", qr.Result().Count())
|
|
}
|
|
|
|
if !index.HasField("d") {
|
|
t.Fatal("don't have field 'd'")
|
|
}
|
|
fieldd := index.Field("d")
|
|
|
|
qr, err = client.Query(index.Count(fieldd.Row(34.0404)))
|
|
if err != nil {
|
|
t.Errorf("querying: %v", err)
|
|
}
|
|
if qr.Result().Count() != 1 {
|
|
t.Errorf("wrong count for field 'd', %d is not 1", qr.Result().Count())
|
|
}
|
|
|
|
}
|
|
|
|
func TestLookupFieldIdNameDisallowed(t *testing.T) {
|
|
t.Parallel()
|
|
tcs := []lookupTestCase{
|
|
{name: "junk", text: "foo", expText: "foo"},
|
|
}
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
a := rand.Int()
|
|
topic := "xyz" + strconv.Itoa(a)
|
|
|
|
// create Main and run with MaxMsgs
|
|
m := NewMain()
|
|
configureTestFlags(m)
|
|
m.Index = fmt.Sprintf("cmd_test_index223ij%s", topic)
|
|
m.AutoGenerate = true
|
|
m.PackBools = "bools"
|
|
m.BatchSize = 1
|
|
m.Topics = []string{topic}
|
|
m.PilosaHosts = []string{pilosaHost}
|
|
m.Timeout = time.Minute
|
|
m.LookupDBDSN = "postgresql://postgres:password@postgres:5432/postgres?sslmode=disable"
|
|
m.Header = "testdata/LookupId.json"
|
|
|
|
// lookupClient.Setup() can't get called until after initialFetch, so need to insert
|
|
// some junk into kafka to test this at this level.
|
|
m.MaxMsgs = uint64(len(tcs)) // make source wait for TimeOut after last message
|
|
m.LookupBatchSize = 1
|
|
|
|
ingesterErrs := make(chan error, 1)
|
|
go func() {
|
|
err := m.Run()
|
|
ingesterErrs <- err
|
|
}()
|
|
|
|
// put records from all testcases into kafka
|
|
addr := kafka.TCP(kafkaHost)
|
|
tCreateTopic(t, topic, addr)
|
|
writer := &kafka.Writer{
|
|
Addr: addr,
|
|
Topic: topic,
|
|
Balancer: &kafka.Hash{},
|
|
BatchTimeout: time.Nanosecond,
|
|
}
|
|
defer writer.Close()
|
|
for _, tc := range tcs {
|
|
msg := messageStringFromTestcase(t, tc)
|
|
tPutStringKafka(t, writer, "bkey", msg)
|
|
}
|
|
|
|
// Wait for ingestion to finish.
|
|
err := <-ingesterErrs
|
|
|
|
if !strings.Contains(err.Error(), "field name 'id' not allowed for LookupText fields") {
|
|
t.Fatalf("invalid field name 'id' not detected: %s", err)
|
|
}
|
|
}
|
|
|
|
// lookupTestCase consolidates definitions for:
|
|
// - messages sent to Kafka
|
|
// - record IDs generated by Pilosa and retrieved within a test
|
|
// - values to check against Postgres
|
|
// - values to check against Pilosa
|
|
// name: testcase name
|
|
// uniquePilosaVal: unqiue integer value sent to Pilosa. Corresponds to `int` field in Lookup.json
|
|
// externalId: Pilosa record ID, allocated by Pilosa, looked up by test, used as Postgres ID as well
|
|
// text: raw text sent to Postgres. Corresponds to `text` field in Lookup.josn
|
|
// expText: text after retrieving from Postges (distinct from `text` due to escape characters)
|
|
// missing: true if `text` should NOT be present in the Kafka message
|
|
type lookupTestCase struct {
|
|
name string
|
|
uniquePilosaVal uint64
|
|
externalId uint64
|
|
text string
|
|
expText string
|
|
missing bool
|
|
}
|
|
|
|
// messageStringFromTestCase defines a Kafka json message string, to be
|
|
// sent to Kafka. Matches Lookup.json.
|
|
func messageStringFromTestcase(t *testing.T, tc lookupTestCase) string {
|
|
if tc.missing {
|
|
return fmt.Sprintf(`{"int":%d}`, tc.uniquePilosaVal)
|
|
} else {
|
|
return fmt.Sprintf(`{"int":%d,"text":"%s"}`, tc.uniquePilosaVal, tc.text)
|
|
}
|
|
|
|
}
|
|
|
|
// TestLookupFieldWithExternalId checks that the external
|
|
// lookup (postgres) feature works as expected,
|
|
// - using ExternalGenerate to use pilosa to generate IDs
|
|
// - with bad string input
|
|
// - with missing data
|
|
// This test is intended to match THR's use case.
|
|
func TestLookupFieldWithExternalId(t *testing.T) {
|
|
// t.Parallel()
|
|
tcs := []lookupTestCase{
|
|
// When using ExternalGenerate (pilosa Nexter), then the only kafka setting that
|
|
// makes sense is at-most-once delivery.
|
|
// That means this postgres client/batcher should only be used in its current state
|
|
// for this very specific use case.
|
|
|
|
// Testing for duplicate records and overwriting behavior would be done here by defining
|
|
// multiple testcases with overlapping ID values.
|
|
// This is not done because that doesn't make sense when using ExternalGenerate:
|
|
// - key is not present in Kafka message, so it is generated by Pilosa
|
|
// - identifying duplicates is not possible without a primary key in the message
|
|
{name: "missing-data", missing: true},
|
|
{name: "normal-write", text: "D", expText: "D"},
|
|
{name: "weirdstring-1", text: "Æ漢д ☮♬ ♞🜻💣", expText: "Æ漢д ☮♬ ♞🜻💣"},
|
|
{name: "weirdstring-2", text: "", expText: ""},
|
|
{name: "weirdstring-3", text: "'", expText: "'"},
|
|
{name: "doublequotes-1", text: `\"`, expText: `"`}, // is this sensible?
|
|
{name: "doublequotes-2", text: `{\"log\": \"message\", 'with': 'whatever', ` + "`weird`: `syntax`}", expText: `{"log": "message", 'with': 'whatever', ` + "`weird`: `syntax`}"},
|
|
}
|
|
|
|
// lookupTestCase.uniquePilosaVal is required to be unique,
|
|
// so it can be used to correlate testcases with IDs allocated by Pilosa,
|
|
// in retrieveTestCaseIds. Set it automatically here.
|
|
lookupRecordCount := 0
|
|
for n := range tcs {
|
|
tcs[n].uniquePilosaVal = uint64(100 * (n + 1))
|
|
if !tcs[n].missing {
|
|
lookupRecordCount++
|
|
}
|
|
}
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
a := rand.Int()
|
|
topic := "xyz" + strconv.Itoa(a)
|
|
|
|
// create Main and run with MaxMsgs
|
|
m := NewMain()
|
|
configureTestFlags(m)
|
|
m.Index = fmt.Sprintf("cmd_test_index223ij%s", topic)
|
|
m.AutoGenerate = true
|
|
m.PackBools = "bools"
|
|
m.BatchSize = 1
|
|
m.Topics = []string{topic}
|
|
m.PilosaHosts = []string{pilosaHost}
|
|
m.Timeout = time.Minute
|
|
m.Header = "testdata/Lookup.json"
|
|
pilosaFieldName := "int" // matches Lookup.json
|
|
postgresColumnName := "text" // matches Lookup.json
|
|
|
|
// lookup+external specific settings
|
|
m.ExternalGenerate = true
|
|
m.LookupDBDSN = "postgresql://postgres:password@postgres:5432/postgres?sslmode=disable"
|
|
m.AllowMissingFields = true // needed for missing data testcases
|
|
m.MaxMsgs = uint64(len(tcs)) // make source wait for TimeOut after last message
|
|
m.LookupBatchSize = lookupRecordCount
|
|
|
|
ingesterErrs := make(chan error, 1)
|
|
go func() {
|
|
err := m.Run()
|
|
ingesterErrs <- err
|
|
}()
|
|
|
|
// put records from all testcases into kafka
|
|
addr := kafka.TCP(kafkaHost)
|
|
tCreateTopic(t, topic, addr)
|
|
writer := &kafka.Writer{
|
|
Addr: addr,
|
|
Topic: topic,
|
|
Balancer: &kafka.Hash{},
|
|
BatchTimeout: time.Nanosecond,
|
|
}
|
|
defer writer.Close()
|
|
for _, tc := range tcs {
|
|
msg := messageStringFromTestcase(t, tc)
|
|
tPutStringKafka(t, writer, "akey", msg)
|
|
}
|
|
|
|
// Wait for ingestion to finish.
|
|
if err := <-ingesterErrs; err != nil {
|
|
t.Fatalf("running main: %v", err)
|
|
}
|
|
|
|
err := retrieveTestCaseIds(tcs, m.Index, pilosaFieldName)
|
|
if err != nil {
|
|
t.Fatalf("looking up ExternalIds: %v", err)
|
|
}
|
|
|
|
lookupClient, err := m.NewLookupClient()
|
|
if err != nil {
|
|
t.Fatal("creating lookup client")
|
|
}
|
|
defer lookupClient.Close()
|
|
|
|
// check final postgres values
|
|
for _, tc := range tcs {
|
|
t.Run(tc.name+"-postgres", func(t *testing.T) {
|
|
if tc.missing {
|
|
if present, err := lookupClient.RowExists(tc.externalId); err != nil {
|
|
t.Fatalf("querying postgres: %s", err)
|
|
} else if present {
|
|
t.Fatalf("present and shouldn't be")
|
|
}
|
|
} else {
|
|
if got, err := lookupClient.ReadString(tc.externalId, postgresColumnName); err != nil {
|
|
t.Fatalf("querying postgres: %s", err)
|
|
} else if tc.expText != got {
|
|
t.Errorf("wrong value from postgres, expected\n%s\n got\n%s\n", tc.expText, got)
|
|
}
|
|
}
|
|
})
|
|
}
|
|
|
|
// check postgres values via pilosa
|
|
for _, tc := range tcs {
|
|
t.Run(tc.name+"-pilosa", func(t *testing.T) {
|
|
pilosaCount, pilosaVal, err := lookupViaPilosa(tc.externalId, postgresColumnName, m.Index)
|
|
if err != nil {
|
|
t.Fatalf("querying pilosa: %v", err)
|
|
}
|
|
if tc.missing {
|
|
if pilosaCount > 0 {
|
|
t.Errorf("present and shouldn't be")
|
|
}
|
|
} else {
|
|
if pilosaVal != tc.expText {
|
|
t.Errorf("wrong value from pilosa, expected\n%s\n got\n%s\n", tc.expText, pilosaVal)
|
|
}
|
|
}
|
|
})
|
|
}
|
|
|
|
// delete data from pilosa
|
|
schema, err := m.PilosaClient().Schema()
|
|
if err != nil {
|
|
t.Errorf("getting client: %v", err)
|
|
}
|
|
index := schema.Index(m.Index)
|
|
|
|
err = m.PilosaClient().DeleteIndex(index)
|
|
if err != nil {
|
|
t.Errorf("deleting index: %v", err)
|
|
}
|
|
|
|
// delete data from postgres
|
|
err = lookupClient.DropTable()
|
|
if err != nil {
|
|
t.Errorf("dropping table: %v", err)
|
|
}
|
|
|
|
}
|
|
|
|
// lookupViaPilosa retrieves lookupText values from Postgres via the ExternalLookup
|
|
// Pilosa query.
|
|
// This was created as a helper function for TestLookupFieldWithExternalId.
|
|
func lookupViaPilosa(id uint64, column, index string) (int, string, error) {
|
|
pql := fmt.Sprintf(`ExternalLookup(ConstRow(columns=[%d]), query="select id, %s from %s where id = ANY($1)")`, id, "text", index)
|
|
|
|
eResp, err := idktest.DoExtractQuery(pql, index)
|
|
if err != nil {
|
|
return 0, "", err
|
|
}
|
|
|
|
if eResp.Results[0].Columns == nil {
|
|
// This itself is not an error condition; it is expected for the
|
|
// missing-data testcase.
|
|
return 0, "", nil
|
|
}
|
|
|
|
rows := eResp.Results[0].Columns[0].Rows
|
|
count := len(rows)
|
|
text := rows[0].(string)
|
|
return count, text, nil
|
|
}
|
|
|
|
// retrieveTestCaseIds populates the externalId field of each testcase
|
|
// by comparing the uniquePilosaValue in the testcase with the results of
|
|
// an Extract query which correlates uniquePilosaValue with its
|
|
// pilosa-allocated ID.
|
|
// This was created as a helper function for TestLookupFieldWithExternalId.
|
|
func retrieveTestCaseIds(tcs []lookupTestCase, index, field string) error {
|
|
pql := fmt.Sprintf("Extract(All(), Rows(%s))", field)
|
|
|
|
eResp, err := idktest.DoExtractQuery(pql, index)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
if eResp.Results[0].Columns == nil {
|
|
return errors.Errorf("no data in Extract response")
|
|
}
|
|
|
|
// Correlate PilosaVal to assign corresponding IDs.
|
|
for n, tc := range tcs {
|
|
for _, col := range eResp.Results[0].Columns {
|
|
// ?? panic: interface conversion: interface {} is float64, not uint64
|
|
if uint64(col.Rows[0].(float64)) == tc.uniquePilosaVal {
|
|
tcs[n].externalId = uint64(col.ColumnID)
|
|
break
|
|
}
|
|
}
|
|
if tcs[n].externalId == 0 {
|
|
// NOTE This assumes an ID of 0 will not be used by the Nexter.
|
|
return errors.Errorf("no externalID found for test case with uniquePilosaVal=%d", tc.uniquePilosaVal)
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func TestDuplicateFieldNameDisallowed(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
a := rand.Int()
|
|
topic := "xyz" + strconv.Itoa(a)
|
|
|
|
// create Main and run with MaxMsgs
|
|
m := NewMain()
|
|
configureTestFlags(m)
|
|
m.Index = fmt.Sprintf("cmd_test_index223ij%s", topic)
|
|
m.AutoGenerate = true
|
|
m.PackBools = "bools"
|
|
m.BatchSize = 1
|
|
m.Topics = []string{topic}
|
|
m.PilosaHosts = []string{pilosaHost}
|
|
m.Timeout = time.Minute
|
|
m.LookupDBDSN = "postgresql://postgres:password@postgres:5432/postgres?sslmode=disable"
|
|
m.Header = "testdata/LookupDuplicate.json"
|
|
m.MaxMsgs = uint64(0)
|
|
|
|
err := m.Run()
|
|
if !strings.Contains(err.Error(), "schema field 2 duplicates name of field 1 (text)") {
|
|
t.Fatalf("duplicate field name not detected: %s", err)
|
|
}
|
|
}
|
|
|
|
func TestPrimaryKeyFieldsMissing(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
fieldNames := []string{"i", "d", "t", "@s"}
|
|
records := [][]interface{}{
|
|
{2, 5.5, "2021-02-01", nil},
|
|
}
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
a := rand.Int()
|
|
topic := "xyz" + strconv.Itoa(a)
|
|
|
|
// create Main and run with MaxMsgs
|
|
m := NewMain()
|
|
configureTestFlags(m)
|
|
m.Index = fmt.Sprintf("cmd_test_index223ij%s", topic)
|
|
m.PrimaryKeyFields = []string{"s"}
|
|
m.AllowMissingFields = true
|
|
m.Header = "testdata/TestFieldTypes.json"
|
|
m.PackBools = "bools"
|
|
m.BatchSize = 1
|
|
m.Topics = []string{topic}
|
|
m.MaxMsgs = uint64(len(records))
|
|
m.PilosaHosts = []string{pilosaHost}
|
|
m.Timeout = time.Minute
|
|
m.Verbose = true
|
|
|
|
// put records in kafka
|
|
addr := kafka.TCP(kafkaHost)
|
|
tCreateTopic(t, topic, addr)
|
|
writer := &kafka.Writer{
|
|
Addr: addr,
|
|
Topic: topic,
|
|
Balancer: &kafka.Hash{},
|
|
BatchTimeout: time.Nanosecond,
|
|
}
|
|
defer writer.Close()
|
|
|
|
for _, vals := range records {
|
|
rec := makeRecordString(t, fieldNames, vals)
|
|
tPutStringKafka(t, writer, "akey", rec)
|
|
}
|
|
|
|
err := m.Run()
|
|
if err == nil {
|
|
t.Fatal("running main should have failed")
|
|
}
|
|
|
|
client := m.PilosaClient()
|
|
schema, err := client.Schema()
|
|
if err != nil {
|
|
t.Fatalf("getting client: %v", err)
|
|
}
|
|
index := schema.Index(m.Index)
|
|
|
|
err = client.DeleteIndex(index)
|
|
if err != nil {
|
|
t.Logf("deleting index: %v", err)
|
|
}
|
|
|
|
}
|
|
|
|
func TestIDFieldMissing(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
fieldNames := []string{"i", "d", "t", "@s"}
|
|
records := [][]interface{}{
|
|
{nil, 5.5, "2021-02-01", "apple"},
|
|
}
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
a := rand.Int()
|
|
topic := "xyz" + strconv.Itoa(a)
|
|
|
|
// create Main and run with MaxMsgs
|
|
m := NewMain()
|
|
configureTestFlags(m)
|
|
m.Index = fmt.Sprintf("cmd_test_index223ij%s", topic)
|
|
m.IDField = "i"
|
|
m.AllowMissingFields = true
|
|
m.Header = "testdata/TestFieldTypes.json"
|
|
m.PackBools = "bools"
|
|
m.BatchSize = 1
|
|
m.Topics = []string{topic}
|
|
m.MaxMsgs = uint64(len(records))
|
|
m.PilosaHosts = []string{pilosaHost}
|
|
m.Timeout = time.Minute
|
|
m.Verbose = true
|
|
|
|
// put records in kafka
|
|
addr := kafka.TCP(kafkaHost)
|
|
tCreateTopic(t, topic, addr)
|
|
writer := &kafka.Writer{
|
|
Addr: addr,
|
|
Topic: topic,
|
|
Balancer: &kafka.Hash{},
|
|
BatchTimeout: time.Nanosecond,
|
|
}
|
|
defer writer.Close()
|
|
|
|
for _, vals := range records {
|
|
rec := makeRecordString(t, fieldNames, vals)
|
|
tPutStringKafka(t, writer, "akey", rec)
|
|
}
|
|
|
|
err := m.Run()
|
|
if err == nil {
|
|
t.Fatal("running main should have failed")
|
|
}
|
|
|
|
client := m.PilosaClient()
|
|
schema, err := client.Schema()
|
|
if err != nil {
|
|
t.Fatalf("getting client: %v", err)
|
|
}
|
|
index := schema.Index(m.Index)
|
|
|
|
err = client.DeleteIndex(index)
|
|
if err != nil {
|
|
t.Logf("deleting index: %v", err)
|
|
}
|
|
|
|
}
|
|
|
|
func TestCmdAutoID(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
fieldNames := []string{"first"}
|
|
records := [][]interface{}{
|
|
{"a"},
|
|
{"b"},
|
|
{"c"},
|
|
}
|
|
|
|
rand.Seed(time.Now().UnixNano())
|
|
a := rand.Int()
|
|
topic := "xyz" + strconv.Itoa(a)
|
|
|
|
// create Main and run with MaxMsgs
|
|
m := NewMain()
|
|
configureTestFlags(m)
|
|
m.Index = fmt.Sprintf("cmd_test_auto_id223ij%s", topic)
|
|
m.AutoGenerate = true
|
|
m.ExternalGenerate = true
|
|
m.Header = "testdata/Flat.json"
|
|
m.BatchSize = 1
|
|
m.Topics = []string{topic}
|
|
m.MaxMsgs = uint64(len(records))
|
|
m.PilosaHosts = []string{pilosaHost}
|
|
m.Timeout = time.Minute
|
|
m.Verbose = true
|
|
|
|
// put records in kafka
|
|
addr := kafka.TCP(kafkaHost)
|
|
tCreateTopic(t, topic, addr)
|
|
writer := &kafka.Writer{
|
|
Addr: addr,
|
|
Topic: topic,
|
|
Balancer: &kafka.Hash{},
|
|
BatchTimeout: time.Nanosecond,
|
|
}
|
|
defer writer.Close()
|
|
|
|
for _, vals := range records {
|
|
rec := makeRecordString(t, fieldNames, vals)
|
|
tPutStringKafka(t, writer, "a", rec)
|
|
}
|
|
|
|
err := m.Run()
|
|
if err != nil {
|
|
t.Fatalf("running main: %v", err)
|
|
}
|
|
|
|
client := m.PilosaClient()
|
|
schema, err := client.Schema()
|
|
if err != nil {
|
|
t.Fatalf("getting client: %v", err)
|
|
}
|
|
index := schema.Index(m.Index)
|
|
defer func() {
|
|
err := client.DeleteIndex(index)
|
|
if err != nil {
|
|
t.Logf("deleting index: %v", err)
|
|
}
|
|
}()
|
|
|
|
qr, err := client.Query(index.Count(index.All()))
|
|
if err != nil {
|
|
t.Errorf("querying: %v", err)
|
|
}
|
|
if qr.Result().Count() != 3 {
|
|
t.Errorf("wrong count for columns, %d is not 3", qr.Result().Count())
|
|
}
|
|
|
|
qr, err = client.Query(pilosaclient.NewPQLBaseQuery(`Count(Distinct(All(), field="first"))`, index, nil))
|
|
if err != nil {
|
|
t.Errorf("querying: %v", err)
|
|
}
|
|
if qr.Result().Count() != 3 {
|
|
t.Errorf("wrong count for val, %d is not 3", qr.Result().Count())
|
|
}
|
|
}
|
|
|
|
func makeRecordString(t *testing.T, fields []string, vals []interface{}) string {
|
|
if len(fields) != len(vals) {
|
|
t.Fatalf("have %d fields and %d vals", len(fields), len(vals))
|
|
}
|
|
rec := make(map[string]interface{})
|
|
for i, field := range fields {
|
|
rec[field] = vals[i]
|
|
}
|
|
ret, err := json.Marshal(rec)
|
|
if err != nil {
|
|
t.Fatalf("error marshaling record to json")
|
|
}
|
|
return string(ret)
|
|
}
|