mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-06 19:07:50 +00:00
* use t.Fatal(f) to abort tests, not panic * make perf_able run at all, make it debug a bit better switch perf-able to using same node type we use for other spot instances, because otherwise it never finds any available capacity. we switch the perf-able script to use the standard get_value function instead of direct jq calls. we try to grab server logs if the restore fails in the hopes of finding out why the restore very occasionally fails. * Fix some issues with running IDK tests in docker. (#2248) *Stop running TestKafkaSourceIntegration with t.Parallel() This test can't be run in parallel as it's currently written. Doing so allows for interleaving of messages to the same kafka topic between tests. I didn't attempt to modify the test so it could be run in parallel. That could be done, but left for someone more ambitious. * Remove idk/testenv/certs which got accidentally committed. also update .gitignore to include those. * changes to add bool support in idk (#2240) * initial changes to add bool support in idk * modifying some default parameters for testing, will revert them later * adding support for bool in making fragments function * boolean values implementation without supporting empty or null values at this point * Implement bool support in batch using a map (and a slice for nulls) (#2247) * Implement bool support in batch using a map (and a slice for nulls) * Keep the PackBools default for now But set it explicity in the ingest tests which rely on it. * Modify batch to construct bool update like mutex The code in API.ImportRoaringShard has a switch statement which causes bool fields to be handled like mutex fields. This means, that the viewUpdate.Clear value should only contain data in the first "row" of the fragment, which it will treat as records to clear for *all* rows. This makes more sense for mutex fields; for bool fields, there's only one other row to clear. But since the code is currently handling them the same, we need to construct viewUpdate.Clear such that it conforms to that pattern. This commit also adds a test which covers this logic. * Remove commented code; revert config for testing This commit also removes the DELETE_SENTINEL case for non-packed bools, since that isn't supported anyway. * Revert default setting * remove inconsistent type scope * correcting the logic of string converstion to bool * resolving an error in a test * adding tests to cover code related to bool support in batch.go file and interface.go files * modifying interfaces test * added one more test case Co-authored-by: Travis Turner <travis@pilosa.com> Co-authored-by: Travis Turner <travis@molecula.com> * resolving bool null field ingestion error (#2254) * resolving bool null field ingestion error * testing issues * adding null support for bools * updating the null bool field ingestion * trying to resolve issue when ingesting null value for bool type * adding a clearing support for bool type * resolving issues with bool null value ingestion * updating the jwt go package version and removing changes made in docker compose file * reverting jwt go version * removing v4 of jwt * adding a comment in test file to see if sonar cloud accepts this file * don't obtain stack traces on rbf.Tx creation We thought stack traces were mildly expensive. We were very wrong. Due to a complicated issue in the Go runtime, simultaneous requests for stack traces end up contending on a lock even when they're not actually contending on any resources. I've filed a ticket in the Go issue tracker for this: https://github.com/golang/go/issues/56400 In the mean time: Under some workloads, we were seeing 85% of all CPU time go into the stack backtraces, of which 81% went into the contention on those locks. But even if you take away the contention, that leaves us with 4/19 of all CPU time in our code going into building those stack backtraces. That's a lot of overhead for a feature we virtually never use. We might consider adding a backtrace functionality here, possibly using `runtime.Callers` which is much lower overhead, and allows us to generate a backtrace on demand (no argument values available, but then, we never read those because they're unformatted hex values), but I don't think it's actually very informative to know what the stack traces were of the Tx; they don't necessarily reflect the current state of any ongoing use of the Tx, so we can't necessarily correlate them to goroutine stack dumps, and so on. * fb-1729 Enriched Table Metadata (#2255) enriched metadata for tables added support for the concept of a table and field owners in metadata; mechanism to derive owner from http request metadata; metadata for table description * tightened up is/is not null filter expressions (FB-1741) (#2260) Covers tightening up handling filter expressions that contain is/is not null ops. These filters may have to be translated into PQL calls to be passed to the executor and even though sql3 language supports nullability for any data type, currently only BSI fields are nullable at the storage engine level (there is a ticket to add support for non-BSI field here FB-1689: IS SQL Argument returns incorrect error) so when these fields are used in filter conditions we need to handle BSI and non-BSI fields differently. * added a test to cover the keyword replace as being synonymous with insert (#2261) * update molecula references to featurebase (#2262) Co-authored-by: Seebs <seebs@molecula.com> Co-authored-by: Travis Turner <travis@pilosa.com> Co-authored-by: Pranitha-malae <56414132+Pranitha-malae@users.noreply.github.com> Co-authored-by: Travis Turner <travis@molecula.com> Co-authored-by: pokeeffe-molecula <85502298+pokeeffe-molecula@users.noreply.github.com> Co-authored-by: Stephanie Yang <stephanie@pilosa.com>
1156 lines
33 KiB
Go
1156 lines
33 KiB
Go
//go:build !kafka_sasl
|
|
// +build !kafka_sasl
|
|
|
|
package kafka
|
|
|
|
import (
|
|
"context"
|
|
"crypto/tls"
|
|
"encoding/binary"
|
|
"encoding/json"
|
|
"fmt"
|
|
"math/big"
|
|
"math/rand"
|
|
"os"
|
|
"reflect"
|
|
"strings"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/confluentinc/confluent-kafka-go/kafka"
|
|
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/idk/kafka/csrc"
|
|
"github.com/featurebasedb/featurebase/v3/logger"
|
|
"github.com/glycerine/vprint"
|
|
"github.com/go-avro/avro"
|
|
liavro "github.com/linkedin/goavro/v2"
|
|
)
|
|
|
|
func configureSourceTestFlags(source *Source) {
|
|
source.KafkaBootstrapServers = []string{kafkaHost}
|
|
source.SchemaRegistryURL = registryHost
|
|
}
|
|
|
|
func TestAvroToPDKSchema(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
tests := []struct {
|
|
schemaFile string
|
|
exp []idk.Field
|
|
expErr string
|
|
}{
|
|
{
|
|
schemaFile: "simple.json",
|
|
exp: expectedSchemas["simple.json"],
|
|
},
|
|
{
|
|
schemaFile: "stringtypes.json",
|
|
exp: expectedSchemas["stringtypes.json"],
|
|
},
|
|
{
|
|
schemaFile: "decimal.json",
|
|
exp: expectedSchemas["decimal.json"],
|
|
},
|
|
{
|
|
schemaFile: "othertypes.json",
|
|
exp: expectedSchemas["othertypes.json"],
|
|
},
|
|
{
|
|
schemaFile: "unions.json",
|
|
exp: expectedSchemas["unions.json"],
|
|
},
|
|
{
|
|
schemaFile: "floatscale.json",
|
|
exp: expectedSchemas["floatscale.json"],
|
|
},
|
|
{
|
|
schemaFile: "notarecord.json",
|
|
expErr: "unsupported Avro Schema type",
|
|
},
|
|
{
|
|
schemaFile: "fieldisrecord.json",
|
|
expErr: "nested fields are not currently supported",
|
|
},
|
|
{
|
|
schemaFile: "timestamp.json",
|
|
exp: expectedSchemas["timestamp.json"],
|
|
},
|
|
}
|
|
|
|
// check that we've covered all the test schemas
|
|
files, err := os.ReadDir("./testdata/schemas")
|
|
if err != nil {
|
|
t.Fatalf("reading directory: %v", err)
|
|
}
|
|
if len(files) != len(tests)+4 { // +4 because we aren't testing bigschema.json, the two delete ones, or the ID allocation one here.
|
|
t.Errorf("have different number of schemas and tests: %d and %d\n%+v", len(files), len(tests), files)
|
|
}
|
|
|
|
for _, test := range tests {
|
|
test := test
|
|
t.Run(test.schemaFile, func(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
codec := decodeTestSchema(t, test.schemaFile)
|
|
schema, err := avroToPDKSchema(codec)
|
|
if err != nil && test.expErr == "" {
|
|
t.Fatalf("unexpected error: %v", err)
|
|
}
|
|
if test.expErr != "" && err == nil {
|
|
t.Fatalf("expected error")
|
|
}
|
|
if test.expErr != "" && !strings.Contains(err.Error(), test.expErr) {
|
|
t.Fatalf("error expected/got\n%s\n%v", test.expErr, err.Error())
|
|
}
|
|
if !reflect.DeepEqual(test.exp, schema) {
|
|
t.Fatalf("schema exp/got\n%+v\n%+v", test.exp, schema)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
func TestAvroToPDKField(t *testing.T) {
|
|
tests := []struct {
|
|
name string
|
|
schemaField *avro.SchemaField
|
|
expField idk.Field
|
|
expErr string
|
|
}{
|
|
{
|
|
name: "string ttl",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "string-ttl",
|
|
Type: &avro.StringSchema{},
|
|
Properties: map[string]interface{}{
|
|
"ttl": "30s",
|
|
"quantum": "YMD",
|
|
},
|
|
},
|
|
expField: idk.StringField{NameVal: "string-ttl", Quantum: "YMD", TTL: "30s"},
|
|
expErr: "nil",
|
|
},
|
|
{
|
|
name: "string ttl without quantum",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "string-ttl",
|
|
Type: &avro.StringSchema{},
|
|
Properties: map[string]interface{}{
|
|
"ttl": "30s",
|
|
},
|
|
},
|
|
expField: idk.StringField{NameVal: "string-ttl"},
|
|
expErr: "nil",
|
|
},
|
|
{
|
|
name: "bytes ttl",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "bytes-ttl",
|
|
Type: &avro.BytesSchema{},
|
|
Properties: map[string]interface{}{
|
|
"ttl": "30s",
|
|
"quantum": "YMD",
|
|
},
|
|
},
|
|
expField: idk.StringField{NameVal: "bytes-ttl", Quantum: "YMD", TTL: "30s"},
|
|
expErr: "nil",
|
|
},
|
|
{
|
|
name: "bytes ttl without quantum",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "bytes-ttl",
|
|
Type: &avro.BytesSchema{},
|
|
Properties: map[string]interface{}{
|
|
"ttl": "30s",
|
|
},
|
|
},
|
|
expField: idk.StringField{NameVal: "bytes-ttl"},
|
|
expErr: "nil",
|
|
},
|
|
{
|
|
name: "string array ttl",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "array-ttl",
|
|
Type: &avro.ArraySchema{
|
|
Items: &avro.StringSchema{
|
|
Properties: map[string]interface{}{
|
|
"quantum": "YMD",
|
|
"ttl": "60s",
|
|
},
|
|
},
|
|
},
|
|
},
|
|
expField: idk.StringArrayField{NameVal: "array-ttl", Quantum: "YMD", TTL: "60s"},
|
|
expErr: "nil",
|
|
},
|
|
{
|
|
name: "string array ttl without quantum",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "array-ttl",
|
|
Type: &avro.ArraySchema{
|
|
Items: &avro.StringSchema{
|
|
Properties: map[string]interface{}{
|
|
"ttl": "60s",
|
|
},
|
|
},
|
|
},
|
|
},
|
|
expField: idk.StringArrayField{NameVal: "array-ttl"},
|
|
expErr: "nil",
|
|
},
|
|
{
|
|
name: "Int ttl",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "int-ttl",
|
|
Type: &avro.IntSchema{},
|
|
Properties: map[string]interface{}{
|
|
"fieldType": "id",
|
|
"ttl": "30s",
|
|
"quantum": "YMD",
|
|
},
|
|
},
|
|
expField: idk.IDField{NameVal: "int-ttl", Quantum: "YMD", TTL: "30s"},
|
|
expErr: "nil",
|
|
},
|
|
{
|
|
name: "Int ttl without quantum",
|
|
schemaField: &avro.SchemaField{
|
|
Name: "int-ttl",
|
|
Type: &avro.IntSchema{},
|
|
Properties: map[string]interface{}{
|
|
"fieldType": "id",
|
|
"ttl": "30s",
|
|
},
|
|
},
|
|
expField: idk.IDField{NameVal: "int-ttl"},
|
|
expErr: "nil",
|
|
},
|
|
}
|
|
|
|
for _, test := range tests {
|
|
t.Run(test.name, func(t *testing.T) {
|
|
field, err := avroToPDKField(test.schemaField)
|
|
if !reflect.DeepEqual(field, test.expField) {
|
|
t.Errorf("expected field: '%v', got: '%v'", test.expField, field)
|
|
}
|
|
if err != nil && !strings.Contains(err.Error(), test.expErr) {
|
|
t.Errorf("expected error: '%s', got: '%s'", test.expErr, err.Error())
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
func decodeTestSchema(t *testing.T, filename string) avro.Schema {
|
|
codec, err := avro.ParseSchema(readTestSchema(t, filename))
|
|
if err != nil {
|
|
t.Fatalf("parsing schema: %v", err)
|
|
}
|
|
return codec
|
|
}
|
|
|
|
func readTestSchema(t *testing.T, filename string) string {
|
|
bytes, err := os.ReadFile("./testdata/schemas/" + filename)
|
|
if err != nil {
|
|
t.Fatalf("reading schema file: %v", err)
|
|
}
|
|
return string(bytes)
|
|
}
|
|
|
|
func liDecodeTestSchema(t *testing.T, filename string) *liavro.Codec {
|
|
codec, err := liavro.NewCodec(readTestSchema(t, filename))
|
|
if err != nil {
|
|
t.Fatalf("li parsing schema: %v", err)
|
|
}
|
|
return codec
|
|
}
|
|
|
|
var tests = []struct {
|
|
data []map[string]interface{}
|
|
schemaFile string
|
|
exp [][]interface{}
|
|
}{
|
|
{
|
|
schemaFile: "simple.json",
|
|
data: []map[string]interface{}{{"first": "hello", "last": "goodbye"}, {"first": "one", "last": "two"}},
|
|
exp: [][]interface{}{{"hello", "goodbye"}, {"one", "two"}},
|
|
},
|
|
{
|
|
schemaFile: "stringtypes.json",
|
|
data: []map[string]interface{}{{"first": "blah", "last": "goodbye", "middle": "123456789"}},
|
|
exp: [][]interface{}{{"blah", []byte("goodbye"), []byte("123456789")}},
|
|
},
|
|
{
|
|
schemaFile: "decimal.json",
|
|
data: []map[string]interface{}{{"somenum": &big.Rat{}}, {"somenum": big.NewRat(10, 1)}, {"somenum": big.NewRat(1, 1)}, {"somenum": big.NewRat(5, 2)}, {"somenum": big.NewRat(1234567890, 1)}},
|
|
exp: [][]interface{}{{[]byte{0}}, {[]byte{0x3, 0xE8}}, {[]byte{100}}, {[]byte{0, 250}}, {[]byte{0x1C, 0xBE, 0x99, 0x1A, 0x08}}},
|
|
},
|
|
{
|
|
schemaFile: "othertypes.json",
|
|
data: []map[string]interface{}{{"first": "a", "second": []string{"b", "c"}, "third": -8, "fourth": 99, "fifth": 99.9, "sixth": 101.1, "seventh": true, "eighth": map[string]interface{}{"string": "a"}, "ninth": map[string]interface{}{"string": "b"}}},
|
|
exp: [][]interface{}{{"a", []interface{}{"b", "c"}, int32(-8), int64(99), float32(99.9), float64(101.1), true, "a", "b"}},
|
|
},
|
|
{
|
|
schemaFile: "unions.json",
|
|
data: []map[string]interface{}{
|
|
{"first": map[string]interface{}{"string": "a"}, "second": map[string]interface{}{"boolean": true}, "third": map[string]interface{}{"long": 101}, "fourth": map[string]interface{}{"bytes.decimal": big.NewRat(5, 2)}, "fifth": map[string]interface{}{"double": float64(9.4921)}},
|
|
{"first": nil, "second": nil, "third": map[string]interface{}{"null": nil}, "fourth": nil, "fifth": nil},
|
|
},
|
|
exp: [][]interface{}{
|
|
{"a", true, int64(101), []byte{9, 196}, float64(9.4921)},
|
|
{nil, nil, nil, nil, nil}},
|
|
},
|
|
{
|
|
schemaFile: "floatscale.json",
|
|
data: []map[string]interface{}{{"first": 23.12345}},
|
|
exp: [][]interface{}{{float32(23.12345)}},
|
|
},
|
|
{
|
|
schemaFile: "timestamp.json",
|
|
data: []map[string]interface{}{{"time_samir": []byte("5000")}},
|
|
exp: [][]interface{}{{[]byte("5000")}},
|
|
},
|
|
}
|
|
|
|
func TestKafkaSourceLocal(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
// this is not an integration test, so we'll take steps to avoid
|
|
// actually connecting to Kafka or Schema Registry.
|
|
|
|
for i, test := range tests {
|
|
i, test := i, test
|
|
t.Run(test.schemaFile, func(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
schema := liDecodeTestSchema(t, test.schemaFile)
|
|
|
|
src := NewSource()
|
|
defer src.Close()
|
|
configureSourceTestFlags(src)
|
|
cfg, err := common.SetupConfluent(&src.ConfluentCommand)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
src.ConfigMap = cfg
|
|
// note: we will not call Open on the source which would connect
|
|
// to Kafka. Instead, we'll set the reader manually so we
|
|
// can inject messages.
|
|
go func() {
|
|
topic := string("test")
|
|
for j, record := range test.data {
|
|
buf := make([]byte, 5, 1000)
|
|
buf[0] = 0
|
|
binary.BigEndian.PutUint32(buf[1:], uint32(i))
|
|
buf, err := schema.BinaryFromNative(buf, record)
|
|
if err != nil {
|
|
t.Errorf("encoding:\n%+v\nerr: %v", record, err)
|
|
}
|
|
e := &confluent.Message{
|
|
TopicPartition: confluent.TopicPartition{
|
|
Topic: &topic,
|
|
Partition: 0,
|
|
Offset: confluent.Offset(j),
|
|
},
|
|
Timestamp: time.Now(),
|
|
Value: buf,
|
|
}
|
|
|
|
src.recordChannel <- recordWithError{Record: e}
|
|
}
|
|
close(src.recordChannel)
|
|
}()
|
|
|
|
// prefill the schema cache so the registry isn't contacted.
|
|
src.cache[int32(i)] = decodeTestSchema(t, test.schemaFile)
|
|
for j, expect := range test.exp {
|
|
pdkRec, err := src.Record()
|
|
if j == 0 {
|
|
if err != idk.ErrSchemaChange {
|
|
t.Errorf("expected schema changed signal, got: %v", err)
|
|
}
|
|
gotSchema := src.Schema()
|
|
if !reflect.DeepEqual(gotSchema, expectedSchemas[test.schemaFile]) {
|
|
t.Errorf("unexpected schema exp/got:\n%+v\n%+v", expectedSchemas[test.schemaFile], gotSchema)
|
|
}
|
|
} else if err != nil {
|
|
t.Fatalf("unexpected error getting record: %v", err)
|
|
}
|
|
if pdkRec == nil {
|
|
t.Fatalf("should have a record")
|
|
}
|
|
data := pdkRec.Data()
|
|
if !reflect.DeepEqual(data, expect) {
|
|
t.Errorf("data mismatch exp/got:\n%+v\n%+v", expect, data)
|
|
if len(data) != len(expect) {
|
|
t.Fatalf("mismatched lengths exp/got %d/%d", len(expect), len(data))
|
|
}
|
|
for k, exp := range expect {
|
|
got := data[k]
|
|
if !reflect.DeepEqual(exp, got) {
|
|
t.Errorf("Mismatch at %d, exp/got\n%v of %[2]T\n%v of %[3]T", k, exp, got)
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
err = src.Close()
|
|
if err != nil {
|
|
t.Fatalf("failed to close source: %v", err)
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
// TestKafkaSourceSchemaChangeCommitRegresion is a regression test for a bug where a failure during a schema change could cause a record to be lost.
|
|
func TestKafkaSourceSchemaChangeCommitRegression(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
src := NewSource()
|
|
defer src.Close()
|
|
configureSourceTestFlags(src)
|
|
cfg, err := common.SetupConfluent(&src.ConfluentCommand)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
src.ConfigMap = cfg
|
|
|
|
// prefill the schema cache so the registry isn't contacted.
|
|
t1, t2 := tests[0], tests[1]
|
|
schema1 := liDecodeTestSchema(t, t1.schemaFile)
|
|
schema2 := liDecodeTestSchema(t, t2.schemaFile)
|
|
schemaID1, schemaID2 := int32(0), int32(1)
|
|
src.cache[schemaID1] = decodeTestSchema(t, t1.schemaFile)
|
|
src.cache[schemaID2] = decodeTestSchema(t, t2.schemaFile)
|
|
|
|
// synthesize fake Kakfa messages
|
|
data1, err := endcodeAvro(int(schemaID1), schema1, t1.data[0])
|
|
if err != nil {
|
|
t.Fatalf("encoding record: %v", err)
|
|
}
|
|
data2, err := endcodeAvro(int(schemaID2), schema2, t2.data[0])
|
|
if err != nil {
|
|
t.Fatalf("encoding record: %v", err)
|
|
}
|
|
|
|
go func() {
|
|
|
|
topic := "xyzzy"
|
|
src.recordChannel <- recordWithError{Record: &confluent.Message{
|
|
TopicPartition: confluent.TopicPartition{
|
|
Topic: &topic,
|
|
Partition: 0,
|
|
Offset: confluent.Offset(0),
|
|
},
|
|
Timestamp: time.Now(),
|
|
Value: data1,
|
|
},
|
|
}
|
|
src.recordChannel <- recordWithError{Record: &confluent.Message{
|
|
TopicPartition: confluent.TopicPartition{
|
|
Topic: &topic,
|
|
Partition: 0,
|
|
Offset: confluent.Offset(1),
|
|
},
|
|
Timestamp: time.Now(),
|
|
Value: data2,
|
|
},
|
|
}
|
|
close(src.recordChannel)
|
|
}()
|
|
|
|
rec1, err := src.Record()
|
|
if err != idk.ErrSchemaChange {
|
|
t.Fatalf("expected schema change but got: %v", err)
|
|
}
|
|
if rec1 == nil {
|
|
t.Fatal("first record is missing")
|
|
}
|
|
rec2, err := src.Record()
|
|
if err != idk.ErrSchemaChange {
|
|
t.Fatalf("expected schema change but got: %v", err)
|
|
}
|
|
if rec2 == nil {
|
|
t.Fatal("second record is missing")
|
|
}
|
|
}
|
|
|
|
func TestKafkaSourceTimeout(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
src := NewSource()
|
|
defer src.Close()
|
|
cfg, err := common.SetupConfluent(&src.ConfluentCommand)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
src.Timeout = time.Millisecond
|
|
src.KafkaSocketTimeoutMs = 10 // lowest value you can set librdkafka
|
|
|
|
src.ConfigMap = cfg
|
|
|
|
schema := liDecodeTestSchema(t, "simple.json")
|
|
|
|
// prefill the schema cache so the registry isn't contacted.
|
|
schemaID := int32(0)
|
|
src.cache[schemaID] = decodeTestSchema(t, "simple.json")
|
|
|
|
// synthesize fake Kakfa message
|
|
record := map[string]interface{}{"first": "hello", "last": "goodbye"}
|
|
buf := make([]byte, 5, 1000)
|
|
buf[0] = 0
|
|
binary.BigEndian.PutUint32(buf[1:], uint32(schemaID))
|
|
buf, err = schema.BinaryFromNative(buf, record)
|
|
if err != nil {
|
|
t.Errorf("encoding:\n%+v\nerr: %v", record, err)
|
|
}
|
|
|
|
go func() {
|
|
topic := "test"
|
|
src.recordChannel <- recordWithError{Record: &confluent.Message{
|
|
TopicPartition: confluent.TopicPartition{
|
|
Topic: &topic,
|
|
},
|
|
Value: buf,
|
|
},
|
|
}
|
|
|
|
}()
|
|
|
|
// ensure we can get a message if one is available
|
|
if rec, err := src.Record(); err != nil && err != idk.ErrSchemaChange {
|
|
t.Fatalf("expected record, but got error %v", err)
|
|
} else if rec.Data()[0].(string) != "hello" || rec.Data()[1].(string) != "goodbye" {
|
|
t.Fatalf("unexpected record: %v", rec.Data())
|
|
}
|
|
|
|
// ensure timeout works and we get ErrFlush when no message is available
|
|
start := time.Now()
|
|
if rec, err := src.Record(); err != idk.ErrFlush {
|
|
t.Fatalf("expected Flush, but got %+v err: %v", rec, err)
|
|
} else if dur := time.Since(start); dur < time.Millisecond {
|
|
t.Fatalf("expected to wait at least 1ms before getting Flush, but only waited %v", dur)
|
|
}
|
|
close(src.recordChannel)
|
|
}
|
|
|
|
// TestKafkaSource uses a real Kafka and Schema Registry. I downloaded
|
|
// the tar archive of the Confluent Platform (self managed software)
|
|
// from confluent.io/download (I got version 5.3.1). I ran `tar xzf`
|
|
// on the file, changed into the directory, ran `curl -L
|
|
// https://cnfl.io/cli | sh -s -- -b /Users/jaffee/bin` (that
|
|
// directory is on my PATH), then ran `confluent local start
|
|
// schema-registry`.
|
|
//
|
|
// I find that this test runs much faster after a `confluent local
|
|
// destroy` followed by `confluent local start schema-registry`. The
|
|
// difference is stark—10s of seconds—and I don't know why this should
|
|
// be, but I think it has something to do with kafka rebalancing
|
|
// itself when a new client joins.
|
|
func TestKafkaSourceIntegration(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip()
|
|
}
|
|
src := NewSource()
|
|
defer src.Close()
|
|
configureSourceTestFlags(src)
|
|
cfg, err := common.SetupConfluent(&src.ConfluentCommand)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
src.ConfigMap = cfg
|
|
|
|
src.Topics = []string{"testKafkaSourceIntegration"}
|
|
src.Group = "group0"
|
|
rnd := rand.New(rand.NewSource(time.Now().UnixNano()))
|
|
|
|
// Create Producer instance
|
|
p, err := confluent.NewProducer(&confluent.ConfigMap{
|
|
"bootstrap.servers": kafkaHost,
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("Failed to create producer: %s", err)
|
|
}
|
|
defer p.Close()
|
|
tCreateTopic(t, src.Topics[0], p)
|
|
/* "sasl.mechanisms": conf["sasl.mechanisms"],
|
|
"security.protocol": conf["security.protocol"],
|
|
"sasl.username": conf["sasl.username"],
|
|
"sasl.password": conf["sasl.password"]})*/
|
|
err = src.Open()
|
|
if err != nil {
|
|
t.Fatalf("Failed to open src:%v", err)
|
|
}
|
|
|
|
key := fmt.Sprintf("%d", rnd.Int())
|
|
for i, test := range tests {
|
|
schemaID := postSchema(t, test.schemaFile, fmt.Sprintf("schema%d", i), registryHost, nil)
|
|
schema := liDecodeTestSchema(t, test.schemaFile)
|
|
t.Run(test.schemaFile, func(t *testing.T) {
|
|
for j, record := range test.data {
|
|
tPutRecordsKafka(t, p, src.Topics[0], schemaID, schema, key, record)
|
|
pdkRec, err := src.Record()
|
|
if j == 0 {
|
|
if err != idk.ErrSchemaChange {
|
|
t.Errorf("expected schema changed signal, got: %v", err)
|
|
}
|
|
gotSchema := src.Schema()
|
|
if !reflect.DeepEqual(gotSchema, expectedSchemas[test.schemaFile]) {
|
|
t.Errorf("unexpected schema got/exp:\n%+v\n%+v", gotSchema, expectedSchemas[test.schemaFile])
|
|
}
|
|
} else if err != nil {
|
|
t.Fatalf("unexpected error getting record: %v", err)
|
|
}
|
|
if pdkRec == nil {
|
|
t.Fatalf("should have a record")
|
|
}
|
|
data := pdkRec.Data()
|
|
if !reflect.DeepEqual(data, test.exp[j]) {
|
|
t.Errorf("data mismatch exp/got:\n%+v\n%+v", test.exp[j], data)
|
|
if len(data) != len(test.exp[j]) {
|
|
t.Fatalf("mismatched lengths exp/got %d/%d", len(test.exp[j]), len(data))
|
|
}
|
|
for k := range test.exp[j] {
|
|
if !reflect.DeepEqual(test.exp[j][k], data[k]) {
|
|
t.Errorf("Mismatch at %d, exp/got\n%v of %[2]T\n%v of %[3]T", k, test.exp[j][k], data[k])
|
|
}
|
|
}
|
|
|
|
}
|
|
}
|
|
})
|
|
}
|
|
|
|
}
|
|
|
|
func mustNewProducer(t *testing.T, kafkaHost string) *confluent.Producer {
|
|
t.Helper()
|
|
p, err := confluent.NewProducer(&confluent.ConfigMap{
|
|
"bootstrap.servers": kafkaHost,
|
|
})
|
|
if err != nil {
|
|
t.Fatalf("Failed to create producer: %s", err)
|
|
}
|
|
return p
|
|
}
|
|
|
|
func TestKafkaSourceSkipOld(t *testing.T) {
|
|
if testing.Short() {
|
|
t.Skip()
|
|
}
|
|
src := NewSource()
|
|
defer src.Close()
|
|
configureSourceTestFlags(src)
|
|
src.Topics = []string{"testKafkaSourceSkipOld"}
|
|
src.Group = "group1"
|
|
src.SkipOld = true
|
|
src.Timeout = 10 * time.Second
|
|
src.Verbose = true
|
|
src.KafkaSocketTimeoutMs = int(src.Timeout / 2 / time.Millisecond)
|
|
src.Log = logger.NewVerboseLogger(os.Stdout)
|
|
|
|
// Create Producer instance
|
|
p := mustNewProducer(t, kafkaHost)
|
|
defer p.Close()
|
|
tCreateTopic(t, src.Topics[0], p)
|
|
|
|
schemaID := postSchema(t, "simple.json", "skipOldSchema", registryHost, nil)
|
|
schema := liDecodeTestSchema(t, "simple.json")
|
|
|
|
rnd := rand.New(rand.NewSource(time.Now().UnixNano()))
|
|
key := fmt.Sprintf("%d", rnd.Int())
|
|
|
|
// 1. put a record into kafka
|
|
tPutRecordsKafka(t, p, src.Topics[0], schemaID, schema, key, map[string]interface{}{"first": "hello", "last": "goodbye"})
|
|
|
|
// src.Open... this starts actually reading from Kafka asynchronously. With skipOld, it should not see the record.
|
|
if err := src.Open(); err != nil {
|
|
t.Fatalf("opening source: %v", err)
|
|
}
|
|
|
|
// grab the record which has been read from Kafka by the goroutine started in src.Open.... or timeout (we're expecting timeout)
|
|
if rec, err := src.Record(); err == nil {
|
|
t.Fatalf("expected error, got record: %v", rec)
|
|
} else if err != idk.ErrFlush {
|
|
t.Fatalf("expected flush error, got %v", err)
|
|
}
|
|
|
|
// put another record into Kafka
|
|
tPutRecordsKafka(t, p, src.Topics[0], schemaID, schema, key, map[string]interface{}{"first": "a", "last": "b"})
|
|
|
|
// grab that record or timeout (we're expecting to get it this
|
|
// time)... but we've seen weird issues where if src.Timeout is
|
|
// too low (like 5s) we don't see this record (on confluent 1.4.2)
|
|
rec, err := src.Record()
|
|
switch err {
|
|
case nil:
|
|
t.Fatal("expected schema change")
|
|
case idk.ErrSchemaChange:
|
|
default:
|
|
t.Fatalf("expected schema change, got %v", err)
|
|
}
|
|
expect := []interface{}{"a", "b"}
|
|
if got := rec.Data(); !reflect.DeepEqual(got, []interface{}{"a", "b"}) {
|
|
t.Fatalf("expected %v but got %v", expect, got)
|
|
}
|
|
}
|
|
|
|
func TestKafkaSourceNotAutoCommitting(t *testing.T) {
|
|
t.Parallel()
|
|
if testing.Short() {
|
|
t.Skip()
|
|
}
|
|
rand.Seed(time.Now().UnixNano())
|
|
src := NewSource()
|
|
defer src.Close()
|
|
configureSourceTestFlags(src)
|
|
src.Topics = []string{"testKafkaSourceNotAutoCommit"}
|
|
src.Group = fmt.Sprintf("group%d", rand.Int())
|
|
|
|
// Create Producer instance
|
|
p := mustNewProducer(t, kafkaHost)
|
|
defer p.Close()
|
|
tCreateTopic(t, src.Topics[0], p)
|
|
|
|
schemaID := postSchema(t, "simple.json", "skipOldSchema", registryHost, nil)
|
|
schema := liDecodeTestSchema(t, "simple.json")
|
|
|
|
rnd := rand.New(rand.NewSource(time.Now().UnixNano()))
|
|
key := fmt.Sprintf("%d", rnd.Int())
|
|
|
|
if err := src.Open(); err != nil {
|
|
t.Fatalf("opening source: %v", err)
|
|
}
|
|
partition := int32(1)
|
|
|
|
// 1. put records into Kafka
|
|
numRecords := int64(7)
|
|
for i := int64(0); i < numRecords; i++ {
|
|
tPutRecordsKafkaPartition(t, p, src.Topics[0], schemaID, schema, key, partition, map[string]interface{}{"first": fmt.Sprintf("%d", i), "last": fmt.Sprintf("%d", i+1)})
|
|
}
|
|
|
|
// 2. read records from Kafka (but don't commit!)
|
|
var rec idk.Record
|
|
var err error
|
|
for i := 0; i < 7; i++ {
|
|
if rec, err = src.Record(); err != nil && err != idk.ErrSchemaChange {
|
|
t.Fatalf("expected schema change or no error, got: %v", err)
|
|
}
|
|
}
|
|
|
|
// 3. wait for it...
|
|
time.Sleep(time.Second * 6) // default autocommit interval is 5s, so we have to wait if we're going to catch that
|
|
|
|
// 4. check to see if any offsets got committed. If
|
|
// enable.auto.commit somehow gets set to true (which is the
|
|
// default), then this should fail.
|
|
if offsets, err := src.client.Committed([]confluent.TopicPartition{{Topic: &src.Topics[0], Partition: partition}}, 100); err != nil {
|
|
t.Fatalf("getting committed offsets: %v", err)
|
|
} else if len(offsets) != 1 {
|
|
t.Fatalf("unexpected number of offsets in response: %v", offsets)
|
|
} else if off := offsets[0]; off.Offset > 0 {
|
|
t.Fatalf("initial committed offset is greater than 0: %d", off.Offset)
|
|
}
|
|
|
|
// 5. now commit
|
|
err = rec.Commit(context.Background())
|
|
if err != nil {
|
|
t.Fatalf("committing: %v", err)
|
|
}
|
|
|
|
// 6. now verify commit worked
|
|
if offsets, err := src.client.Committed([]confluent.TopicPartition{{Topic: &src.Topics[0], Partition: partition}}, 100); err != nil {
|
|
t.Fatalf("getting committed offsets: %v", err)
|
|
} else if len(offsets) != 1 {
|
|
t.Fatalf("unexpected number of offsets in response: %v", offsets)
|
|
} else if off := offsets[0]; int64(off.Offset) != numRecords {
|
|
t.Fatalf("after commit, offset is not %d: %d", numRecords-1, off.Offset)
|
|
}
|
|
|
|
}
|
|
|
|
func TestRegistryURLParsing(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
tests := []struct {
|
|
name string
|
|
registryURL string
|
|
expectedCodecURL string
|
|
}{
|
|
{
|
|
name: "host:port",
|
|
registryURL: "localhost:8081",
|
|
expectedCodecURL: "http://localhost:8081/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "host",
|
|
registryURL: "localhost",
|
|
expectedCodecURL: "http://localhost/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "scheme:host:port",
|
|
registryURL: "http://localhost:8081",
|
|
expectedCodecURL: "http://localhost:8081/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "scheme:host",
|
|
registryURL: "http://localhost",
|
|
expectedCodecURL: "http://localhost/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "scheme:host",
|
|
registryURL: "http://localhost",
|
|
expectedCodecURL: "http://localhost/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "scheme:host/path",
|
|
registryURL: "http://subsub.subdomain.domain/kafka-schema-registry",
|
|
expectedCodecURL: "http://subsub.subdomain.domain/kafka-schema-registry/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "scheme:host:port/path",
|
|
registryURL: "http://subsub.subdomain.domain:8081/kafka-schema-registry",
|
|
expectedCodecURL: "http://subsub.subdomain.domain:8081/kafka-schema-registry/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "https:host:port/path",
|
|
registryURL: "https://subsub.subdomain.domain:8081/kafka-schema-registry",
|
|
expectedCodecURL: "https://subsub.subdomain.domain:8081/kafka-schema-registry/schemas/ids/1",
|
|
},
|
|
{
|
|
name: "https:host:port",
|
|
registryURL: "https://localhost:8081",
|
|
expectedCodecURL: "https://localhost:8081/schemas/ids/1",
|
|
},
|
|
}
|
|
|
|
for i, test := range tests {
|
|
t.Run(fmt.Sprintf("%s-%d", test.name, i), func(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
src := &Source{}
|
|
src.SchemaRegistryURL = test.registryURL
|
|
|
|
err := src.cleanRegistryURL()
|
|
if err != nil {
|
|
t.Fatalf("cleaning registry URL: %v", err)
|
|
}
|
|
|
|
codecURL, err := src.codecURL(1, "schemas/ids/%d")
|
|
if err != nil {
|
|
t.Fatalf("codec URL: %v", err)
|
|
}
|
|
if codecURL != test.expectedCodecURL {
|
|
t.Errorf("codec URL exp:\n%s\ngot:\n%s", test.expectedCodecURL, codecURL)
|
|
}
|
|
|
|
})
|
|
}
|
|
|
|
}
|
|
|
|
func postSchema(t *testing.T, schemaFile, subj, regURL string, tlsConfig *tls.Config) (schemaID int) {
|
|
schemaClient := csrc.NewClient(regURL, tlsConfig, nil)
|
|
schemaStr := readTestSchema(t, schemaFile)
|
|
resp, err := schemaClient.PostSubjects(subj, schemaStr)
|
|
if err != nil {
|
|
t.Fatalf("posting schema: %v", err)
|
|
}
|
|
return resp.ID
|
|
}
|
|
|
|
func tCreateTopic(t *testing.T, topic string, p *confluent.Producer) {
|
|
t.Helper()
|
|
a, err := confluent.NewAdminClientFromProducer(p)
|
|
if err != nil {
|
|
t.Fatalf("Failed to create new admin client from producer: %s", err)
|
|
}
|
|
// Contexts are used to abort or limit the amount of time
|
|
// the Admin call blocks waiting for a result.
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
// Create topics on cluster.
|
|
// Set Admin options to wait up to 60s for the operation to finish on the remote cluster
|
|
maxDur, err := time.ParseDuration("60s")
|
|
if err != nil {
|
|
t.Fatalf("ParseDuration(60s): %s", err)
|
|
}
|
|
results, err := a.CreateTopics(
|
|
ctx,
|
|
// Multiple topics can be created simultaneously
|
|
// by providing more TopicSpecification structs here.
|
|
[]confluent.TopicSpecification{{
|
|
Topic: topic,
|
|
NumPartitions: 64,
|
|
ReplicationFactor: 1}},
|
|
// Admin options
|
|
confluent.SetAdminOperationTimeout(maxDur))
|
|
if err != nil {
|
|
t.Fatalf("Admin Client request error: %v\n", err)
|
|
}
|
|
for _, result := range results {
|
|
if result.Error.Code() != confluent.ErrNoError && result.Error.Code() != confluent.ErrTopicAlreadyExists {
|
|
t.Fatalf("Failed to create topic: %v\n", result.Error)
|
|
}
|
|
}
|
|
a.Close()
|
|
|
|
}
|
|
|
|
func tPutRecordsKafka(t *testing.T, p *confluent.Producer, topic string, schemaID int, schema *liavro.Codec, key string, records ...map[string]interface{}) {
|
|
tPutRecordsKafkaPartition(t, p, topic, schemaID, schema, key, confluent.PartitionAny, records...)
|
|
}
|
|
|
|
func tPutRecordsKafkaPartition(t *testing.T, p *confluent.Producer, topic string, schemaID int, schema *liavro.Codec, key string, partition int32, records ...map[string]interface{}) {
|
|
t.Helper()
|
|
delivery_chan := make(chan kafka.Event, 10000)
|
|
for _, record := range records {
|
|
data, err := endcodeAvro(schemaID, schema, record)
|
|
if err != nil {
|
|
t.Fatalf("encoding record: %v", err)
|
|
}
|
|
err = p.Produce(&kafka.Message{
|
|
TopicPartition: kafka.TopicPartition{Topic: &topic, Partition: partition},
|
|
Key: []byte(key),
|
|
Value: data,
|
|
}, delivery_chan)
|
|
if err != nil {
|
|
t.Fatalf("producing record: %v", err)
|
|
}
|
|
e := <-delivery_chan
|
|
m := e.(*kafka.Message)
|
|
|
|
if m.TopicPartition.Error != nil {
|
|
t.Fatalf("Delivery failed: %v\n", m.TopicPartition.Error)
|
|
} else {
|
|
vprint.VV("Delivered message to topic %s [%d] at offset %v\n",
|
|
*m.TopicPartition.Topic, m.TopicPartition.Partition, m.TopicPartition.Offset)
|
|
}
|
|
}
|
|
close(delivery_chan)
|
|
}
|
|
|
|
var expectedSchemas = map[string][]idk.Field{
|
|
"simple.json": {
|
|
idk.StringField{NameVal: "first"},
|
|
idk.StringField{NameVal: "last"},
|
|
},
|
|
"stringtypes.json": {
|
|
idk.StringField{NameVal: "first", Mutex: true},
|
|
idk.StringField{NameVal: "last"},
|
|
idk.StringField{NameVal: "middle"},
|
|
},
|
|
"decimal.json": {
|
|
idk.DecimalField{NameVal: "somenum", Scale: 2},
|
|
},
|
|
"unions.json": {
|
|
idk.StringField{NameVal: "first"},
|
|
idk.BoolField{NameVal: "second"},
|
|
idk.IntField{NameVal: "third"},
|
|
idk.DecimalField{NameVal: "fourth", Scale: 3},
|
|
idk.DecimalField{NameVal: "fifth", Scale: 2},
|
|
},
|
|
"othertypes.json": {
|
|
idk.StringField{NameVal: "first", Mutex: true},
|
|
idk.StringArrayField{NameVal: "second"},
|
|
idk.IntField{NameVal: "third"},
|
|
idk.IntField{NameVal: "fourth"},
|
|
idk.DecimalField{NameVal: "fifth"},
|
|
idk.DecimalField{NameVal: "sixth"},
|
|
idk.BoolField{NameVal: "seventh"},
|
|
idk.StringField{NameVal: "eighth", Mutex: true},
|
|
idk.StringField{NameVal: "ninth", Mutex: true},
|
|
},
|
|
"floatscale.json": {
|
|
idk.DecimalField{NameVal: "first", Scale: 4},
|
|
},
|
|
"timestamp.json": {
|
|
idk.TimestampField{NameVal: "time_samir", Granularity: "s", Unit: "ms"},
|
|
},
|
|
}
|
|
|
|
// Ensure the avro schemas produce the exptected JSON.
|
|
func TestAvroSchemaFieldToJSON(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
tests := []struct {
|
|
schemaField *avro.SchemaField
|
|
exp string
|
|
}{
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Type: &avro.BytesSchema{},
|
|
},
|
|
exp: `{"type":"bytes"}`,
|
|
},
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Name: "union-null-bytes",
|
|
Type: &avro.UnionSchema{
|
|
Types: []avro.Schema{
|
|
&avro.NullSchema{},
|
|
&avro.BytesSchema{},
|
|
},
|
|
},
|
|
Default: nil,
|
|
},
|
|
exp: `{"name":"union-null-bytes","default":null,"type":["null","bytes"]}`,
|
|
},
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Name: "decimal",
|
|
Type: &avro.BytesSchema{
|
|
Properties: map[string]interface{}{
|
|
"logicalType": "decimal",
|
|
"scale": 2,
|
|
"precision": 5,
|
|
},
|
|
},
|
|
},
|
|
exp: `{"name":"decimal","type":{"type":"bytes","logicalType":"decimal","scale":2,"precision":5}}`,
|
|
},
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Name: "union-null-decimal",
|
|
Type: &avro.UnionSchema{
|
|
Types: []avro.Schema{
|
|
&avro.NullSchema{},
|
|
&avro.BytesSchema{
|
|
Properties: map[string]interface{}{
|
|
"logicalType": "decimal",
|
|
"scale": 2,
|
|
"precision": 5,
|
|
},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
exp: `{"name":"union-null-decimal","type":["null",{"type":"bytes","logicalType":"decimal","scale":2,"precision":5}],"default":null}`,
|
|
},
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Name: "float",
|
|
Type: &avro.FloatSchema{},
|
|
},
|
|
exp: `{"name":"float","type":"float"}`,
|
|
},
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Name: "float-props",
|
|
Type: &avro.FloatSchema{
|
|
Properties: map[string]interface{}{
|
|
"scale": 2,
|
|
},
|
|
},
|
|
},
|
|
exp: `{"name":"float-props","type":{"type":"float","scale":2}}`,
|
|
},
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Name: "union-null-float",
|
|
Type: &avro.UnionSchema{
|
|
Types: []avro.Schema{
|
|
&avro.NullSchema{},
|
|
&avro.FloatSchema{
|
|
Properties: map[string]interface{}{
|
|
"scale": 2,
|
|
},
|
|
},
|
|
},
|
|
},
|
|
},
|
|
exp: `{"name":"union-null-float","type":["null",{"type":"float","scale":2}],"default":null}`,
|
|
},
|
|
{
|
|
schemaField: &avro.SchemaField{
|
|
Name: "timestamp",
|
|
Type: &avro.BytesSchema{
|
|
Properties: map[string]interface{}{
|
|
"logicalType": "timestamp",
|
|
"granularity": "s",
|
|
"unit": "ms",
|
|
},
|
|
},
|
|
},
|
|
exp: `{"name":"timestamp","type":{"granularity":"s","logicalType":"timestamp","type":"bytes","unit":"ms"}}`,
|
|
},
|
|
}
|
|
|
|
for i, test := range tests {
|
|
j, err := json.Marshal(test.schemaField)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if eq, err := equalJSON(string(j), test.exp); err != nil {
|
|
t.Fatal(err)
|
|
} else if !eq {
|
|
t.Fatalf("test %d: \nexp: %s, \ngot: %s", i, test.exp, j)
|
|
}
|
|
}
|
|
}
|
|
|
|
func equalJSON(s1, s2 string) (bool, error) {
|
|
var o1 interface{}
|
|
var o2 interface{}
|
|
|
|
var err error
|
|
err = json.Unmarshal([]byte(s1), &o1)
|
|
if err != nil {
|
|
return false, fmt.Errorf("Error mashalling string 1 :: %s", err.Error())
|
|
}
|
|
err = json.Unmarshal([]byte(s2), &o2)
|
|
if err != nil {
|
|
return false, fmt.Errorf("Error mashalling string 2 :: %s", err.Error())
|
|
}
|
|
|
|
return reflect.DeepEqual(o1, o2), nil
|
|
}
|
|
|
|
func TestGetKeyFunc(t *testing.T) {
|
|
tests := []struct {
|
|
rec []interface{}
|
|
expKey uint64
|
|
}{
|
|
{
|
|
rec: []interface{}{"aaa"},
|
|
expKey: uint64(162),
|
|
},
|
|
{
|
|
rec: []interface{}{"aaaaaaaaaaaaaaa"},
|
|
expKey: uint64(86),
|
|
},
|
|
}
|
|
|
|
p, err := NewPutSource()
|
|
if err != nil {
|
|
t.Fatalf("failed to return new putsource: %v", err)
|
|
}
|
|
|
|
p.FBPrimaryKeyFields = []string{"uuid"}
|
|
p.FBIDField = "uuid"
|
|
p.FBIndexName = "test"
|
|
|
|
schema := []idk.Field{
|
|
idk.StringField{NameVal: "uuid"},
|
|
}
|
|
|
|
getKeyFunc, err := p.getKeyFunc(schema)
|
|
if err != nil {
|
|
t.Fatalf("failed to getKeyFunc: %v", err)
|
|
}
|
|
|
|
for i, test := range tests {
|
|
tc := test
|
|
t.Run(fmt.Sprintf("test%d", i), func(t *testing.T) {
|
|
key, err := getKeyFunc(tc.rec)
|
|
if err != nil {
|
|
t.Fatalf("failed to get key for record: %v", err)
|
|
}
|
|
keyInt := uint64(binary.BigEndian.Uint64(key))
|
|
if tc.expKey != keyInt {
|
|
t.Fatalf("expected %d, got %d", tc.expKey, keyInt)
|
|
}
|
|
})
|
|
}
|
|
}
|