Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,7 @@ ci-tests-race: test-deps
# run diff lint like in pipeline
.lint:
$(info Running lint...)
GOBIN=$(LOCAL_BIN) go run github.com/golangci/golangci-lint/v2/cmd/golangci-lint@v2.12.1 run \
GOBIN=$(LOCAL_BIN) go run github.com/golangci/golangci-lint/v2/cmd/golangci-lint@v2.14.0 run \
--config=.golangci.yaml ./...

.PHONY: lint
Expand Down
53 changes: 53 additions & 0 deletions pkg/seqproxyapi/v1/mappings.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"fmt"

"github.com/ozontech/seq-db/asyncsearcher"
"github.com/ozontech/seq-db/query"
"github.com/ozontech/seq-db/seq"
)

Expand Down Expand Up @@ -128,3 +129,55 @@ func AsyncSearchStatusFromString(s string) (AsyncSearchStatus, error) {

return 0, fmt.Errorf("unknown status")
}

var typeMappings = []DataType{
query.DataTypeBytes: DataType_BYTES,
query.DataTypeSeqID: DataType_SEQ_ID,
query.DataTypeDocument: DataType_RAW_DOCUMENT,
query.DataTypeString: DataType_STRING,
query.DataTypeUint32: DataType_UINT32,
query.DataTypeUint64: DataType_UINT64,
query.DataTypeInt32: DataType_INT32,
query.DataTypeInt64: DataType_INT64,
query.DataTypeFloat64: DataType_FLOAT64,
query.DataTypeFloat64Array: DataType_FLOAT64_ARRAY,
query.DataTypeStringArray: DataType_STRING_ARRAY,
}

var typeMappingsPb = func() []query.DataType {
mappings := make([]query.DataType, len(typeMappings))
for from, to := range typeMappings {
mappings[to] = query.DataType(from)
}
return mappings
}()

func (t DataType) ToQueryDataType() (query.DataType, error) {
if int(t) >= len(typeMappingsPb) || t < 0 {
return 0, fmt.Errorf("unknown data type: %d", t)
}
return typeMappingsPb[t], nil
}

func (t DataType) MustQueryDataType() query.DataType {
v, err := t.ToQueryDataType()
if err != nil {
panic(err)
}
return v
}

func ToProtoDataType(t query.DataType) (DataType, error) {
if int(t) >= len(typeMappings) {
return 0, fmt.Errorf("unknown data type: %d", t)
}
return typeMappings[t], nil
}

func MustProtoDataType(t query.DataType) DataType {
v, err := ToProtoDataType(t)
if err != nil {
panic(err)
}
return v
}
53 changes: 53 additions & 0 deletions pkg/storeapi/mappings.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"fmt"

"github.com/ozontech/seq-db/asyncsearcher"
"github.com/ozontech/seq-db/query"
"github.com/ozontech/seq-db/seq"
)

Expand Down Expand Up @@ -142,3 +143,55 @@ func MustProtoAsyncSearchStatus(s asyncsearcher.AsyncSearchStatus) AsyncSearchSt
}
return v
}

var typeMappings = []DataType{
query.DataTypeBytes: DataType_BYTES,
query.DataTypeSeqID: DataType_SEQ_ID,
query.DataTypeDocument: DataType_RAW_DOCUMENT,
query.DataTypeString: DataType_STRING,
query.DataTypeUint32: DataType_UINT32,
query.DataTypeUint64: DataType_UINT64,
query.DataTypeInt32: DataType_INT32,
query.DataTypeInt64: DataType_INT64,
query.DataTypeFloat64: DataType_FLOAT64,
query.DataTypeFloat64Array: DataType_FLOAT64_ARRAY,
query.DataTypeStringArray: DataType_STRING_ARRAY,
}

var typeMappingsPb = func() []query.DataType {
mappings := make([]query.DataType, len(typeMappings))
for from, to := range typeMappings {
mappings[to] = query.DataType(from)
}
return mappings
}()

func (t DataType) ToQueryDataType() (query.DataType, error) {
if int(t) >= len(typeMappingsPb) || t < 0 {
return 0, fmt.Errorf("unknown data type: %d", t)
}
return typeMappingsPb[t], nil
}

func (t DataType) MustQueryDataType() query.DataType {
v, err := t.ToQueryDataType()
if err != nil {
panic(err)
}
return v
}

func ToProtoDataType(t query.DataType) (DataType, error) {
if int(t) >= len(typeMappings) {
return 0, fmt.Errorf("unknown data type: %d", t)
}
return typeMappings[t], nil
}

func MustProtoDataType(t query.DataType) DataType {
v, err := ToProtoDataType(t)
if err != nil {
panic(err)
}
return v
}
74 changes: 69 additions & 5 deletions proxy/search/stream_search.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,11 @@ func (si *Ingestor) StreamSearch(
}
}

if len(streams) == 0 {
// nothing to read from, return an empty stream
return exec.NewNMergedProducers(nil, query.SeqIDColumn(0), sr.Order), newControlBroadcaster(streams), partialRespErr
}

broadcaster := newControlBroadcaster(streams)
producers := make([]query.RecordProducer, 0, len(streams))
for _, s := range streams {
Expand All @@ -88,12 +93,28 @@ func (si *Ingestor) StreamSearch(
offset = 0
}

// TODO: store can change docs' schema in case of filter pipe, need to handle it on proxy.
var expectedSchema *query.Schema
if sr.Agg != nil {
expectedSchema = query.AggsSchema
}
// Shard schemas come off the wire, so validate them before any consumer reads records.
if err := validateShardSchemas(streams, expectedSchema); err != nil {
closeStreams(streams)
return nil, nil, err
}

var mergedStream query.RecordProducer
if sr.Agg != nil {
mergedStream = exec.NewDistributedAggregator(producers, sr.Agg.Func, sr.Agg.Quantiles)
} else {
const seqIdColIdx = 0
mergedDocsStream := exec.NewNMergedProducers(producers, query.SeqIDColumn(seqIdColIdx), sr.Order)
// streams[0] is safe - streams len is already checked
idCol, err := streams[0].OutSchema().Column[seq.ID](query.DocsIDCol)
if err != nil {
closeStreams(streams)
return nil, nil, err
}
mergedDocsStream := exec.NewNMergedProducers(producers, idCol, sr.Order)
mergedStream = exec.NewLimiter(mergedDocsStream, uint32(sr.Size), uint32(offset))
}

Expand Down Expand Up @@ -267,6 +288,40 @@ func closeStreams(streams []*StreamSearchIterator) {
}
}

func schemaFromTyping(typing []*storeapi.Typing) (*query.Schema, []query.DataType, error) {
cols := make([]query.ColumnDesc, 0, len(typing))
types := make([]query.DataType, 0, len(typing))
for _, t := range typing {
dt, err := t.GetType().ToQueryDataType()
if err != nil {
return nil, nil, err
}
cols = append(cols, query.ColumnDesc{Name: t.GetTitle(), Type: dt})
types = append(types, dt)
}
schema, err := query.NewSchema(cols...)
if err != nil {
return nil, nil, err
}
return schema, types, nil
}

func validateShardSchemas(streams []*StreamSearchIterator, expected *query.Schema) error {
if len(streams) == 0 {
return nil
}
base := streams[0].OutSchema()
for _, s := range streams[1:] {
if !s.OutSchema().Equal(base) {
return fmt.Errorf("shard schemas mismatch: %v vs %v", base.Cols(), s.OutSchema().Cols())
}
}
if expected != nil && !base.Equal(expected) {
return fmt.Errorf("shard schema mismatch: got %v, want %v", base.Cols(), expected.Cols())
}
return nil
}

// NewStreamSearchIterator reads one message ahead after the header so that a
// summary-with-error sent immediately after the header (before any data) is
// detected on the open-stream phase and can trigger fail-fast in the
Expand All @@ -276,7 +331,11 @@ func NewStreamSearchIterator(
header *storeapi.ResponseHeader,
stream storeapi.StoreApi_StreamSearchClient,
) (*StreamSearchIterator, error) {
it := &StreamSearchIterator{tr: tr, typing: header.Typing, stream: stream}
schema, types, err := schemaFromTyping(header.Typing)
if err != nil {
return nil, fmt.Errorf("bad response header: %w", err)
}
it := &StreamSearchIterator{tr: tr, schema: schema, types: types, stream: stream}

msg, err := stream.Recv()
if errors.Is(err, io.EOF) {
Expand All @@ -295,7 +354,8 @@ func NewStreamSearchIterator(
type StreamSearchIterator struct {
tr *querytracer.Tracer

typing []*storeapi.Typing
schema *query.Schema
types []query.DataType
stream storeapi.StoreApi_StreamSearchClient

curBatch []*storeapi.Record
Expand Down Expand Up @@ -336,11 +396,15 @@ func (it *StreamSearchIterator) Next() *query.Record {

recordVals := make([]*query.RecordVals, 0, len(record.RawData))
for i, rawData := range record.RawData {
recordVals = append(recordVals, query.NewRecordVals(query.DataType(it.typing[i].Type), rawData))
recordVals = append(recordVals, query.NewRecordVals(it.types[i], rawData))
}
return query.NewRecord(recordVals)
}

func (it *StreamSearchIterator) OutSchema() *query.Schema {
return it.schema
}

// push handles a single message received from the store stream.
func (it *StreamSearchIterator) push(msg *storeapi.StreamSearchResponse) error {
switch v := msg.ResponseType.(type) {
Expand Down
87 changes: 87 additions & 0 deletions proxy/search/stream_search_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,87 @@
package search

import (
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

insaneJSON "github.com/ozontech/insane-json"

"github.com/ozontech/seq-db/pkg/storeapi"
"github.com/ozontech/seq-db/query"
"github.com/ozontech/seq-db/seq"
)

func TestSchemaFromTyping(t *testing.T) {
typing := []*storeapi.Typing{
{Title: "id", Type: storeapi.DataType_SEQ_ID},
{Title: "data", Type: storeapi.DataType_RAW_DOCUMENT},
}

schema, types, err := schemaFromTyping(typing)
require.NoError(t, err)
assert.Equal(t, 2, schema.Len())
assert.Equal(t, []query.DataType{query.DataTypeSeqID, query.DataTypeDocument}, types)
assert.Equal(t, 0, schema.MustColumn[seq.ID]("id").Idx())
assert.Equal(t, query.DataTypeDocument, schema.MustColumn[*insaneJSON.Root]("data").DataType())
}

func TestSchemaFromTypingDuplicateName(t *testing.T) {
typing := []*storeapi.Typing{
{Title: "id", Type: storeapi.DataType_SEQ_ID},
{Title: "id", Type: storeapi.DataType_SEQ_ID},
}
_, _, err := schemaFromTyping(typing)
assert.Error(t, err)
}

func TestValidateShardSchemas(t *testing.T) {
newIterator := func(typing []*storeapi.Typing) *StreamSearchIterator {
schema, types, err := schemaFromTyping(typing)
require.NoError(t, err)
return &StreamSearchIterator{schema: schema, types: types}
}

docsTyping := []*storeapi.Typing{
{Title: "id", Type: storeapi.DataType_SEQ_ID},
{Title: "data", Type: storeapi.DataType_RAW_DOCUMENT},
}
otherTyping := []*storeapi.Typing{
{Title: "id", Type: storeapi.DataType_SEQ_ID},
{Title: "payload", Type: storeapi.DataType_RAW_DOCUMENT},
}
aggsTyping := []*storeapi.Typing{
{Title: "token", Type: storeapi.DataType_STRING},
}

t.Run("no shards", func(t *testing.T) {
assert.NoError(t, validateShardSchemas(nil, nil))
})

t.Run("matching", func(t *testing.T) {
streams := []*StreamSearchIterator{newIterator(docsTyping), newIterator(docsTyping)}
assert.NoError(t, validateShardSchemas(streams, nil))
})

t.Run("mismatch", func(t *testing.T) {
streams := []*StreamSearchIterator{newIterator(docsTyping), newIterator(otherTyping)}
err := validateShardSchemas(streams, nil)
require.Error(t, err)
assert.Contains(t, err.Error(), "shard schemas mismatch")
})

t.Run("expected schema mismatch", func(t *testing.T) {
streams := []*StreamSearchIterator{newIterator(docsTyping), newIterator(docsTyping)}
err := validateShardSchemas(streams, query.AggsSchema)
require.Error(t, err)
assert.Contains(t, err.Error(), "shard schema mismatch")
})

t.Run("expected schema match", func(t *testing.T) {
streams := []*StreamSearchIterator{newIterator(aggsTyping), newIterator(aggsTyping)}
expected, _, err := schemaFromTyping(aggsTyping)
require.NoError(t, err)
assert.NoError(t, validateShardSchemas(streams, expected))
})
}
12 changes: 6 additions & 6 deletions proxyapi/grpc_complex_search.go
Original file line number Diff line number Diff line change
Expand Up @@ -171,11 +171,11 @@ func (g *grpcV1) useStreamSearch(
func readDocuments(storesStream query.RecordProducer) []*seqproxyapi.Document {
var docs []*seqproxyapi.Document
for r := storesStream.Next(); r != nil; r = storesStream.Next() {
id := r.Vals[0].AsSeqID()
id := docIDCol.Val(r)
docs = append(docs, &seqproxyapi.Document{
Id: id.String(),
Time: timestamppb.New(id.MID.Time()),
Data: r.Vals[1].RawData(),
Data: docDataCol.RawData(r),
})
}
return docs
Expand All @@ -185,13 +185,13 @@ func readAggregations(storesStream query.RecordProducer) []*seqproxyapi.Aggregat
buckets := make([]*seqproxyapi.Aggregation_Bucket, 0)
for r := storesStream.Next(); r != nil; r = storesStream.Next() {
bucket := &seqproxyapi.Aggregation_Bucket{
Key: r.Vals[0].AsString(),
Value: r.Vals[1].AsFloat64(),
Key: aggKeyCol.Val(r),
Value: aggValueCol.Val(r),
}
if ts := r.Vals[2].AsUint64(); ts != consts.DummyMID {
if ts := aggTsCol.Val(r); ts != consts.DummyMID {
bucket.Ts = timestamppb.New(seq.MID(ts).Time())
}
if quantiles := r.Vals[3].AsFloat64Array(); len(quantiles) > 0 {
if quantiles := aggQuantilesCol.Val(r); len(quantiles) > 0 {
bucket.Quantiles = quantiles
}
buckets = append(buckets, bucket)
Expand Down
Loading
Loading