// Copyright 2021 Molecula Corp. All rights reserved. package ingest import ( "fmt" "sort" "strings" "testing" "time" "github.com/molecula/featurebase/v2/shardwidth" ) func TestStableTranslator(t *testing.T) { tr := newStableTranslator() m1, err := tr.TranslateKeys("a", "b") if err != nil { t.Fatalf("translation error on initial keys: %v", err) } m2, err := tr.TranslateIDs(m1["a"], m1["b"], 6) if err != nil { t.Fatalf("translation error on reverse lookup: %v", err) } m3, err := tr.TranslateKeys("a", "k-6") if err != nil { t.Fatalf("translation error on new keys: %v", err) } for k, v := range m3 { if m2[v] != k { t.Fatalf("expected round trip to equate %q and %d", k, v) } } } func TestMakeCodec(t *testing.T) { codec, _ := NewJSONCodec(nil) err := codec.AddSetField("set", nil) if err != nil { t.Fatalf("unexpected error creating field: %v", err) } err = codec.AddSetField("set", nil) if err == nil { t.Fatalf("expected error creating duplicate field, didn't get it") } } func TestEncode(t *testing.T) { codec, _ := NewJSONCodec(nil) _ = codec.AddSetField("set", nil) _ = codec.AddSetField("setkeys", newStableTranslator()) _ = codec.AddMutexField("mutex", nil) _ = codec.AddMutexField("mutexkeys", newStableTranslator()) _ = codec.AddTimeQuantumField("tq", nil) _ = codec.AddTimeQuantumField("tqkeys", newStableTranslator()) _ = codec.AddIntField("int", nil) _ = codec.AddIntField("intkeys", newStableTranslator()) epoch, err := time.Parse("2006-01-02", "2020-01-01") if err != nil { t.Fatalf("can't parse sample epoch time: %v", err) } _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") codecs := []*JSONCodec{codec} // redo all of that, only on a keyed translator codec, _ = NewJSONCodec(newStableTranslator()) _ = codec.AddSetField("set", nil) _ = codec.AddSetField("setkeys", newStableTranslator()) _ = codec.AddMutexField("mutex", nil) _ = codec.AddMutexField("mutexkeys", newStableTranslator()) _ = codec.AddTimeQuantumField("tq", nil) _ = codec.AddTimeQuantumField("tqkeys", newStableTranslator()) _ = codec.AddIntField("int", nil) _ = codec.AddIntField("intkeys", newStableTranslator()) _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") codecs = append(codecs, codec) encodeTests := []*Request{ { Ops: []*Operation{ { OpType: OpWrite, ClearRecordIDs: []uint64{0, 1, 2, 3, 5}, ClearFields: []string{"bool", "dec", "int", "intkeys", "mutex", "mutexkeys", "set", "setkeys", "tq", "tqkeys", "ts"}, FieldOps: map[string]*FieldOperation{ "int": { RecordIDs: []uint64{0, 1}, Signed: []int64{1, -3}, }, "intkeys": { RecordIDs: []uint64{0, 1}, Signed: []int64{1, 1}, }, "set": { RecordIDs: []uint64{0, 1, 2, 2}, Values: []uint64{1, 1, 0, 1}, }, "setkeys": { RecordIDs: []uint64{0, 1, 2, 2}, Values: []uint64{1, 1, 0, 1}, }, "mutex": { RecordIDs: []uint64{0, 1}, Values: []uint64{1, 2}, }, "mutexkeys": { RecordIDs: []uint64{0, 1}, Values: []uint64{1, 2}, }, "tq": { RecordIDs: []uint64{5, 5}, Values: []uint64{8, 9}, Signed: []int64{1234567890e9, 1234567890e9}, }, "tqkeys": { RecordIDs: []uint64{3, 3}, Values: []uint64{2, 4}, Signed: []int64{1234567890e9, 1234567890e9}, }, "ts": { RecordIDs: []uint64{0}, Signed: []int64{1}, }, "bool": { RecordIDs: []uint64{0, 1}, Values: []uint64{0, 1}, }, "dec": { RecordIDs: []uint64{0, 1, 2}, Signed: []int64{123, -123, 0}, }, }, }, { OpType: OpClear, Seq: 1, ClearRecordIDs: []uint64{6}, ClearFields: []string{"tq"}, }, }, }, { // this one needs to get filled in programmatically; see below Ops: []*Operation{ { OpType: OpSet, FieldOps: map[string]*FieldOperation{}, }, }, }, { Ops: []*Operation{ { OpType: OpRemove, FieldOps: map[string]*FieldOperation{ "set": {}, }, }, }, }, } // and now we populate encodeTests[1] with a larger pool of data const dataSize = 5000 shardCount := uint64(600) // shards we want to target passes := uint64(0) recordIDs := make([]uint64, dataSize) values := make([]uint64, dataSize) timeStamps := make([]int64, dataSize) signedValues := make([]int64, dataSize) for i := uint64(0); i < dataSize; i++ { if (i % shardCount) == 0 { passes++ } recordIDs[i] = ((i % shardCount) << shardwidth.Exponent) + passes values[i] = (i % 4) timeStamps[i] = int64(1234567890e9 + (i * 100e9)) signedValues[i] = (int64(i) % 16) // no negative values because they won't work with keys } // ensure record IDs are sorted, because other stuff might rely on this sort.Slice(recordIDs, func(i, j int) bool { return recordIDs[i] < recordIDs[j] }) op := encodeTests[1].Ops[0] op.FieldOps["tq"] = &FieldOperation{ RecordIDs: append([]uint64{}, recordIDs...), Values: append([]uint64{}, values...), Signed: append([]int64{}, timeStamps...), } op.FieldOps["tqkeys"] = &FieldOperation{ RecordIDs: append([]uint64{}, recordIDs...), Values: values, Signed: timeStamps, } op.FieldOps["int"] = &FieldOperation{ RecordIDs: append([]uint64{}, recordIDs...), Signed: append([]int64{}, signedValues...), } op.FieldOps["intkeys"] = &FieldOperation{ RecordIDs: recordIDs, Signed: signedValues, } // for sets, we want to shuffle things into fewer shards, and ensure // non-duplication of values within each record, but also have lots // of duplication of record IDs in the low shards recordIDs = make([]uint64, dataSize) values = make([]uint64, dataSize) valuesPerRecord := uint64(5) recordsPerShard := dataSize / valuesPerRecord / 30 if recordsPerShard < 1 { recordsPerShard = 1 } shard := uint64(0) nextID := uint64(0) nextValue := uint64(0) for i := uint64(0); i < dataSize; i++ { recordIDs[i] = nextID values[i] = nextValue + (i % valuesPerRecord) nextValue++ if nextValue == valuesPerRecord { nextValue = 0 nextID++ if nextID%(1< 1 { valuesPerRecord-- } } } } sort.Slice(recordIDs, func(i, j int) bool { return recordIDs[i] < recordIDs[j] }) op.FieldOps["set"] = &FieldOperation{ RecordIDs: append([]uint64{}, recordIDs...), Values: append([]uint64{}, values...), } op.FieldOps["setkeys"] = &FieldOperation{ RecordIDs: recordIDs, Values: values, } var buf []byte for i, tc := range encodeTests { for _, c := range codecs { data, err := c.AppendBytes(tc, buf[:0]) if err != nil { t.Fatalf("encode test %d: error encoding: %v", i, err) } // t.Logf("data:\n%s", data) req, err := c.ParseBytes(data) if err != nil { t.Logf("encode test %d: data:\n%s", i, data) t.Fatalf("encode test %d: error parsing: %v", i, err) } err = req.Compare(tc) if err != nil { t.Logf("encode test %d: data:\n%s", i, data) t.Fatalf("encode test %d: round-trip mismatch: %v", i, err) } data, err = c.AppendBytes(req, buf[:0]) if err != nil { t.Fatalf("encode test %d: error encoding: %v", i, err) } // t.Logf("data:\n%s", data) req2, err := c.ParseBytes(data) if err != nil { t.Logf("encode test %d: data:\n%s", i, data) t.Fatalf("encode test %d: error parsing: %v", i, err) } err = req2.Compare(tc) if err != nil { t.Logf("encode test %d: data:\n%s", i, data) t.Fatalf("encode test %d: round-trip mismatch: %v", i, err) } } } } func TestCodecErrors(t *testing.T) { codec, _ := NewJSONCodec(nil) _ = codec.AddSetField("set", nil) _ = codec.AddSetField("setkeys", newStableTranslator()) _ = codec.AddMutexField("mutex", nil) _ = codec.AddMutexField("mutexkeys", newStableTranslator()) _ = codec.AddTimeQuantumField("tq", nil) _ = codec.AddTimeQuantumField("tqkeys", newStableTranslator()) _ = codec.AddIntField("int", nil) _ = codec.AddIntField("intkeys", newStableTranslator()) epoch, err := time.Parse("2006-01-02", "2020-01-01") if err != nil { t.Fatalf("can't parse sample epoch time: %v", err) } _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") testCases := []struct { name string json []byte error string }{ { name: "no action", json: []byte(`[{"records":{"0":{"set":[0]}}}]`), error: "action not specified", }, { name: "unknown action", json: []byte(`[{"action":"yeet","records":{"0":{"set":[0]}}}]`), error: "unknown action", }, { name: "unknown field", json: []byte(`[{"action":"set","records":{"0":{"settee":[0]}}}]`), error: "field not found", }, { name: "unknown operation field", json: []byte(`[{"action":"set","yeet":false,"records":{"0":{"set":[0]}}}]`), error: "unknown operation field", }, { name: "expected operation", json: []byte(`[true]`), error: "expected operation", }, { name: "expecting key", json: []byte(`[{"action":"set","records":{"0":{"setkeys":0}}}]`), error: "expecting key", }, { name: "invalid int for bool", json: []byte(`[{"action":"set","records":{"0":{"bool":2}}}]`), error: "boolean should be", }, { name: "invalid number for bool", json: []byte(`[{"action":"set","records":{"0":{"bool":1.3}}}]`), error: "looks like Number", }, { name: "invalid string for bool", json: []byte(`[{"action":"set","records":{"0":{"bool":"truly"}}}]`), error: "expecting boolean", }, { name: "nonsense bool", json: []byte(`[{"action":"set","records":{"0":{"bool":[]}}}]`), error: "boolean should be", }, { name: "expecting numeric value", json: []byte(`[{"action":"set","records":{"0":{"set":0.1}}}]`), error: "invalid syntax", }, { name: "expecting value", json: []byte(`[{"action":"set","records":{"0":{"setkeys":true}}}]`), error: "expecting value", }, { name: "expecting array-key", json: []byte(`[{"action":"set","records":{"0":{"setkeys":[0]}}}]`), error: "expecting key", }, { name: "expecting array-value", json: []byte(`[{"action":"set","records":{"0":{"setkeys":[true]}}}]`), error: "expecting value", }, { name: "expecting numeric array-value", json: []byte(`[{"action":"set","records":{"0":{"set":[0.1]}}}]`), error: "invalid syntax", }, { name: "expecting numeric value", json: []byte(`[{"action":"set","records":{"0":{"int":0.1}}}]`), error: "invalid syntax", }, { name: "expecting int key", json: []byte(`[{"action":"set","records":{"0":{"intkeys":0}}}]`), error: "expecting string key", }, { name: "expecting int value", json: []byte(`[{"action":"set","records":{"0":{"int":[0]}}}]`), error: "expecting integer value", }, { name: "expecting string array-value", json: []byte(`[{"action":"set","records":{"0":{"intkeys":["a"]}}}]`), error: "expecting string key", }, { name: "expecting numeric mutex value", json: []byte(`[{"action":"set","records":{"0":{"mutex":0.1}}}]`), error: "invalid syntax", }, { name: "expecting mutex key", json: []byte(`[{"action":"set","records":{"0":{"mutexkeys":0}}}]`), error: "expecting string key", }, { name: "expecting mutex value", json: []byte(`[{"action":"set","records":{"0":{"mutex":[0]}}}]`), error: "expecting integer value", }, { name: "expecting mutex string value", json: []byte(`[{"action":"set","records":{"0":{"mutexkeys":["a"]}}}]`), error: "expecting string key", }, { name: "time quantum invalid time", json: []byte(`[{"action":"set","records":{"0":{"tq":{"time":[],"values":[3]}}}}]`), error: "expecting time", }, { name: "time stamp invalid integer", json: []byte(`[{"action":"set","records":{"0":{"ts":1.3}}}]`), error: "parsing numeric time", }, { name: "time stamp invalid string", json: []byte(`[{"action":"set","records":{"0":{"ts":"RFC3339"}}}]`), error: "parsing time", }, { name: "time stamp invalid type", json: []byte(`[{"action":"set","records":{"0":{"ts":[]}}}]`), error: "expecting time", }, { name: "invalid decimal", json: []byte(`[{"action":"set","records":{"0":{"dec":[]}}}]`), error: "expecting floating", }, { name: "duplicate record", json: []byte(`[{"action":"set","records":{"0":{"int":1},"0":{"set":0}}}]`), error: "duplicated in input", }, } for _, tc := range testCases { t.Run(tc.name, func(t *testing.T) { req, err := codec.ParseBytes(tc.json) if err == nil { req.Dump(t.Logf) t.Fatalf("expected error like %q, got request instead", tc.error) } else { msg := err.Error() if !strings.Contains(msg, tc.error) { t.Fatalf("expected error like %q, got %q", tc.error, msg) } } }) } } func TestSimpleCodec(t *testing.T) { codec, _ := NewJSONCodec(nil) _ = codec.AddSetField("set", nil) _ = codec.AddSetField("setkeys", newStableTranslator()) _ = codec.AddMutexField("mutex", nil) _ = codec.AddMutexField("mutexkeys", newStableTranslator()) _ = codec.AddTimeQuantumField("tq", nil) _ = codec.AddIntField("int", nil) _ = codec.AddIntField("intkeys", newStableTranslator()) epoch, err := time.Parse("2006-01-02", "2020-01-01") if err != nil { t.Fatalf("can't parse sample epoch time: %v", err) } _ = codec.AddTimestampField("ts", time.Millisecond, epoch.Unix()*1000) _ = codec.AddDecimalField("dec", 2) _ = codec.AddBoolField("bool") var nextShard = uint64(1<