/
githubmirror
/
ydb-go-sdk
Обзор
Документация
Войти
/
githubmirror
/
ydb-go-sdk
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
internal/query/execute_query_test.go
1 344 строки
36 KB
Aleksey Myasnikov
Deprecate WithConcurrentResultSets and always enable it in Query (#2216)
29 июн 2026, 15:49
Не верифицирован
29 июн 2026, 15:49
3810a99
Код
Авторство
О чём код?
package query import ( "context" "io" "testing" "time" "github.com/stretchr/testify/require" "github.com/ydb-platform/ydb-go-genproto/Ydb_Query_V1" "github.com/ydb-platform/ydb-go-genproto/protos/Ydb" "github.com/ydb-platform/ydb-go-genproto/protos/Ydb_Query" "go.uber.org/mock/gomock" "google.golang.org/grpc" grpcCodes "google.golang.org/grpc/codes" "google.golang.org/grpc/metadata" grpcStatus "google.golang.org/grpc/status" "github.com/ydb-platform/ydb-go-sdk/v3/internal/params" "github.com/ydb-platform/ydb-go-sdk/v3/internal/query/options" "github.com/ydb-platform/ydb-go-sdk/v3/internal/xcontext" "github.com/ydb-platform/ydb-go-sdk/v3/internal/xerrors" "github.com/ydb-platform/ydb-go-sdk/v3/query" ) func TestExecute(t *testing.T) { t.Run("HappyWay", func(t *testing.T) { ctx := t.Context() ctrl := gomock.NewController(t) stream := happyWayStream(ctrl) client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).Return(stream, nil) var txID string r, err := execute(ctx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered, onTxMeta(func(txMeta *Ydb_Query.TransactionMeta) { txID = txMeta.GetId() }), ) require.NoError(t, err) defer r.Close(ctx) require.EqualValues(t, "456", txID) require.EqualValues(t, -1, r.resultSetIndex) { t.Log("nextResultSet") rs, err := r.nextResultSet(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.index) { t.Log("next (row=1)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.rowIndex) } { t.Log("next (row=2)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.rowIndex) } { t.Log("next (row=3)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 2, rs.rowIndex) } { t.Log("next (row=4)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.rowIndex) } { t.Log("next (row=5)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.rowIndex) } { t.Log("next (row=6)") _, err := rs.nextRow(ctx) require.ErrorIs(t, err, io.EOF) } } { t.Log("nextResultSet") rs, err := r.nextResultSet(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.index) } { t.Log("nextResultSet") rs, err := r.nextResultSet(ctx) require.NoError(t, err) require.EqualValues(t, 2, rs.index) { t.Log("next (row=1)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.rowIndex) } { t.Log("next (row=2)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.rowIndex) } { t.Log("next (row=3)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.rowIndex) } { t.Log("next (row=4)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.rowIndex) } { t.Log("next (row=5)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 2, rs.rowIndex) } { t.Log("next (row=6)") _, err := rs.nextRow(ctx) require.ErrorIs(t, err, io.EOF) } } { t.Log("close result") r.Close(t.Context()) } { t.Log("nextResultSet") _, err := r.nextResultSet(t.Context()) require.ErrorIs(t, err, io.EOF) } }) t.Run("TransportError", func(t *testing.T) { t.Run("OnCall", func(t *testing.T) { ctx := t.Context() ctrl := gomock.NewController(t) client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).Return(nil, grpcStatus.Error(grpcCodes.Unavailable, "")) t.Log("execute") _, err := execute(ctx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered) require.Error(t, err) require.True(t, xerrors.IsTransportError(err, grpcCodes.Unavailable)) }) t.Run("OnStream", func(t *testing.T) { ctx := t.Context() ctrl := gomock.NewController(t) stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, TxMeta: &Ydb_Query.TransactionMeta{ Id: "456", }, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{ Columns: []*Ydb.Column{ { Name: "a", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UINT64, }, }, }, { Name: "b", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, }, }, Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 1, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "1", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 2, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "2", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 3, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "3", }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{ Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 4, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "4", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 5, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "5", }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(nil, grpcStatus.Error(grpcCodes.Unavailable, "")) client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).Return(stream, nil) t.Log("execute") var txID string r, err := execute(ctx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered, onTxMeta(func(txMeta *Ydb_Query.TransactionMeta) { txID = txMeta.GetId() }), ) require.NoError(t, err) defer r.Close(ctx) require.EqualValues(t, "456", txID) require.EqualValues(t, -1, r.resultSetIndex) { t.Log("nextResultSet") rs, err := r.nextResultSet(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.index) { t.Log("next (row=1)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.rowIndex) } { t.Log("next (row=2)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.rowIndex) } { t.Log("next (row=3)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 2, rs.rowIndex) } { t.Log("next (row=4)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.rowIndex) } { t.Log("next (row=5)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.rowIndex) } { t.Log("next (row=6)") _, err := rs.nextRow(ctx) require.Error(t, err) require.True(t, xerrors.IsTransportError(err, grpcCodes.Unavailable)) } } }) }) t.Run("OperationError", func(t *testing.T) { t.Run("OnCall", func(t *testing.T) { ctx := t.Context() ctrl := gomock.NewController(t) stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().Return(nil, xerrors.Operation(xerrors.WithStatusCode( Ydb.StatusIds_UNAVAILABLE, ))) client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).Return(stream, nil) t.Log("execute") _, err := execute(ctx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered) require.Error(t, err) require.True(t, xerrors.IsOperationError(err, Ydb.StatusIds_UNAVAILABLE)) }) t.Run("OnStream", func(t *testing.T) { ctx := t.Context() ctrl := gomock.NewController(t) stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, TxMeta: &Ydb_Query.TransactionMeta{ Id: "456", }, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{ Columns: []*Ydb.Column{ { Name: "a", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UINT64, }, }, }, { Name: "b", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, }, }, Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 1, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "1", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 2, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "2", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 3, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "3", }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(nil, xerrors.Operation(xerrors.WithStatusCode( Ydb.StatusIds_UNAVAILABLE, ))) client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).Return(stream, nil) t.Log("execute") var txID string r, err := execute(ctx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered, onTxMeta(func(txMeta *Ydb_Query.TransactionMeta) { txID = txMeta.GetId() }), ) require.NoError(t, err) defer r.Close(ctx) require.EqualValues(t, "456", txID) require.EqualValues(t, -1, r.resultSetIndex) { t.Log("nextResultSet") rs, err := r.nextResultSet(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.index) { t.Log("next (row=1)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 0, rs.rowIndex) } { t.Log("next (row=2)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 1, rs.rowIndex) } { t.Log("next (row=3)") _, err := rs.nextRow(ctx) require.NoError(t, err) require.EqualValues(t, 2, rs.rowIndex) } { t.Log("next (row=4)") _, err := rs.nextRow(ctx) require.Error(t, err) require.True(t, xerrors.IsOperationError(err, Ydb.StatusIds_UNAVAILABLE)) } } }) }) t.Run("ContextCancellation", func(t *testing.T) { t.Run("CancelWhileExecute", func(t *testing.T) { ctrl := gomock.NewController(t) ctx, cancel := context.WithCancel(t.Context()) var executeCtx context.Context stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().DoAndReturn(func() (*Ydb_Query.ExecuteQueryResponsePart, error) { cancel() // canceling happen in the beginning of the Recv() call <-executeCtx.Done() return nil, executeCtx.Err() }) client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).DoAndReturn( func(ctx context.Context, _ *Ydb_Query.ExecuteQueryRequest, _ ...grpc.CallOption) ( Ydb_Query_V1.QueryService_ExecuteQueryClient, error, ) { executeCtx = ctx return stream, nil }) // When execute() with context, cancelled in progress _, err := execute(ctx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered) // Then context cancellation error is returned require.ErrorIs(t, err, context.Canceled) }) t.Run("CancelAfterExecute", func(t *testing.T) { ctrl := gomock.NewController(t) stream := happyWayStream(ctrl) var streamCtx context.Context client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).DoAndReturn( func(ctx context.Context, _ *Ydb_Query.ExecuteQueryRequest, _ ...grpc.CallOption) ( Ydb_Query_V1.QueryService_ExecuteQueryClient, error, ) { streamCtx = ctx return stream, nil }) executeCtx, cancelExecuteCtx := context.WithCancel(t.Context()) r, err := execute(executeCtx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered) require.NoError(t, err) cancelExecuteCtx() _, err = readResultSet(t.Context(), r) require.NoError(t, err) _, err = readResultSet(t.Context(), r) require.NoError(t, err) _, err = readResultSet(t.Context(), r) require.NoError(t, err) // check here because the last `readResultSet()` closes stream with stream cancellation require.NoError(t, streamCtx.Err()) _, err = readResultSet(t.Context(), r) require.ErrorIs(t, err, io.EOF) }) // Regression test for https://github.com/ydb-platform/ydb-go-sdk/issues/2081. // // When the parent ctx is cancelled inside ExecuteQuery (simulating session // expiry while the gRPC stream is still open), execute() must surface that // cancellation as context.Canceled to the caller, regardless of whether the // stream's Recv() is reached. // // Two paths can lead to context.Canceled, depending on how the // context.AfterFunc that forwards the parent ctx to executeCtx races with // execute() proceeding into newResult: // - the AfterFunc has already cancelled executeCtx by the time // nextPart() runs, so the ctx.Err() check at the top of nextPart() // returns context.Canceled without invoking Recv(); // - the AfterFunc has not fired yet, nextPart() proceeds to Recv() and // the mock returns ctx.Err() (parent ctx is already cancelled). // Either way the final error must be context.Canceled. The mock therefore // allows Recv() any number of times - including zero - so the test stays // deterministic across both paths. t.Run("CancelParentContextAfterStreamOpen", func(t *testing.T) { t.Run("idempotent=true", func(t *testing.T) { ctrl := gomock.NewController(t) ctx, cancel := context.WithCancel(t.Context()) // Recv() may or may not be reached depending on the AfterFunc race // described above; on either path it returns the parent ctx's // cancellation error. stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().DoAndReturn(func() (*Ydb_Query.ExecuteQueryResponsePart, error) { return nil, ctx.Err() }).AnyTimes() client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).DoAndReturn( func(_ context.Context, _ *Ydb_Query.ExecuteQueryRequest, _ ...grpc.CallOption) ( Ydb_Query_V1.QueryService_ExecuteQueryClient, error, ) { // Simulate session expiry: canceling ctx closes ctx.Done() so that // the non-blocking check in execute() fires after ExecuteQuery returns. cancel() return stream, nil }) _, err := execute(xcontext.WithIdempotent(ctx, true), "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeOrdered, ) require.Error(t, err) require.ErrorIs(t, err, context.Canceled) }) t.Run("idempotent=false", func(t *testing.T) { ctrl := gomock.NewController(t) ctx, cancel := context.WithCancel(t.Context()) // Same race-tolerant expectation as the idempotent=true case: // Recv() may be invoked zero or more times depending on whether // the AfterFunc forwarding parent ctx → executeCtx fires before // nextPart()'s ctx.Err() check. stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().DoAndReturn(func() (*Ydb_Query.ExecuteQueryResponsePart, error) { return nil, ctx.Err() }).AnyTimes() client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).DoAndReturn( func(_ context.Context, _ *Ydb_Query.ExecuteQueryRequest, _ ...grpc.CallOption) ( Ydb_Query_V1.QueryService_ExecuteQueryClient, error, ) { // Simulate session expiry: canceling ctx closes ctx.Done() so that // the non-blocking check in execute() fires after ExecuteQuery returns. cancel() return stream, nil }) _, err := execute(ctx, "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeConcurrent, ) require.Error(t, err) require.ErrorIs(t, err, context.Canceled) }) }) // Per-call ctx cancellation is checked before stream.Recv() only. See // per_call_ctx_recv_test.go for blocked-Recv behavior and Close unblocking. t.Run("CancelCallCtxReturnsWithoutCancelingExecuteStream", func(t *testing.T) { ctrl := gomock.NewController(t) stream := newExecuteQueryStreamMock(ctrl) gomock.InOrder( stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, TxMeta: &Ydb_Query.TransactionMeta{Id: "456"}, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{}, }, nil), // Close(background) must still drain the stream after per-call ctx cancel. stream.EXPECT().Recv().Return(nil, io.EOF), ) client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).DoAndReturn( func(ctx context.Context, _ *Ydb_Query.ExecuteQueryRequest, _ ...grpc.CallOption) ( Ydb_Query_V1.QueryService_ExecuteQueryClient, error, ) { return stream, nil }) r, err := execute(t.Context(), "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeConcurrent) require.NoError(t, err) callCtx, callCancel := context.WithCancel(t.Context()) callCancel() _, err = r.nextPart(callCtx) require.ErrorIs(t, err, context.Canceled) require.NoError(t, r.lastErr) require.NoError(t, r.Close(t.Context())) }) t.Run("CancelCallCtxWhileRecvBlockedViaExecute", func(t *testing.T) { ctrl := gomock.NewController(t) var executeCtx context.Context recvEntered := make(chan struct{}) stream := NewMockQueryService_ExecuteQueryClient(ctrl) gomock.InOrder( stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, TxMeta: &Ydb_Query.TransactionMeta{Id: "456"}, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{}, }, nil), stream.EXPECT().Recv().DoAndReturn(func() (*Ydb_Query.ExecuteQueryResponsePart, error) { close(recvEntered) <-executeCtx.Done() return nil, executeCtx.Err() }), ) // Drain Recv during Close is optional: if callCtx cancel already // propagated to executeCtx via withStreamCancel, Close returns early. stream.EXPECT().Recv().Return(nil, io.EOF).AnyTimes() client := NewMockQueryServiceClient(ctrl) client.EXPECT().ExecuteQuery(gomock.Any(), gomock.Any()).DoAndReturn( func(ctx context.Context, _ *Ydb_Query.ExecuteQueryRequest, _ ...grpc.CallOption) ( Ydb_Query_V1.QueryService_ExecuteQueryClient, error, ) { executeCtx = ctx stubExecuteQueryStreamContext(ctx, stream) return stream, nil }) r, err := execute(t.Context(), "123", client, "", options.ExecuteSettings(), options.ResultSetsTypeConcurrent, withStreamResultCloseTimeout(50*time.Millisecond), ) require.NoError(t, err) callCtx, callCancel := context.WithCancel(t.Context()) iterDone := make(chan error, 1) go func() { _, err := r.nextPart(callCtx) iterDone <- err }() <-recvEntered callCancel() // execute() wires withStreamCancel(executeCancel), so per-call ctx // cancel forwards to the gRPC stream context asynchronously via // context.AfterFunc — not synchronously at callCancel() time. require.Eventually(t, func() bool { return executeCtx.Err() != nil }, time.Second, time.Millisecond, "callCtx cancel must propagate to execute stream") require.ErrorIs(t, executeCtx.Err(), context.Canceled) start := time.Now() closeErr := r.Close(t.Context()) require.Less(t, time.Since(start), time.Second) select { case err := <-iterDone: require.Error(t, err) case <-time.After(time.Second): t.Fatal("nextPart still blocked after Close") } // Close may return nil if the stream was already torn down by // streamCancel, or DeadlineExceeded if drain hit closeTimeout first. if closeErr != nil { require.ErrorIs(t, closeErr, context.DeadlineExceeded) } }) }) } // TestNewResult_DecoupledExecuteCtx verifies the semantic contract that // newResult succeeds when passed executeCtx (derived from a cancelled parent // via xcontext.ValueOnly), which is exactly what execute() does after the fix // for https://github.com/ydb-platform/ydb-go-sdk/issues/2081. func TestNewResult_DecoupledExecuteCtx(t *testing.T) { t.Run("CancelledParentCtxCausesImmediateError", func(t *testing.T) { // With the parent ctx passed directly, a cancelled ctx makes newResult // fail before it ever calls Recv(). This is the old (buggy) behavior // that the fix addresses at the execute() call-site. parentCtx, parentCancel := context.WithCancel(t.Context()) parentCancel() ctrl := gomock.NewController(t) stream := newExecuteQueryStreamMock(ctrl) // Recv must NOT be called — the cancelled ctx short-circuits newResult. _, err := newResult(parentCtx, stream) require.ErrorIs(t, err, context.Canceled) }) t.Run("DecoupledExecuteCtxSucceedsWhenParentCancelled", func(t *testing.T) { // Create executeCtx exactly as execute() does: strip cancellation from // parent via xcontext.ValueOnly, then add an independent cancel. // Even though parentCtx is already cancelled, executeCtx is not — so // newResult can proceed to Recv() and return the first response part. parentCtx, parentCancel := context.WithCancel(t.Context()) parentCancel() executeCtx, executeCancel := xcontext.WithCancel(xcontext.ValueOnly(parentCtx)) defer executeCancel() ctrl := gomock.NewController(t) stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, TxMeta: &Ydb_Query.TransactionMeta{ Id: "456", }, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{}, }, nil) stream.EXPECT().Recv().Return(nil, io.EOF) r, err := newResult(executeCtx, stream) require.NoError(t, err) if r != nil { r.Close(t.Context()) } }) } func TestExecuteQueryRequest(t *testing.T) { for _, tt := range []struct { name string opts []options.Execute request *Ydb_Query.ExecuteQueryRequest callOptions []grpc.CallOption }{ { name: "WithoutOptions", request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithoutOptions", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXECUTE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithoutOptions", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_NONE, ConcurrentResultSets: false, }, }, { name: "WithTxControl", opts: []options.Execute{ options.WithTxControl(query.SerializableReadWriteTxControl(query.CommitTx())), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithTxControl", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXECUTE, TxControl: &Ydb_Query.TransactionControl{ TxSelector: &Ydb_Query.TransactionControl_BeginTx{ BeginTx: &Ydb_Query.TransactionSettings{ TxMode: &Ydb_Query.TransactionSettings_SerializableReadWrite{ SerializableReadWrite: &Ydb_Query.SerializableModeSettings{}, }, }, }, CommitTx: true, }, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithTxControl", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_NONE, ConcurrentResultSets: false, }, }, { name: "WithParams", opts: []options.Execute{ options.WithParameters( params.Builder{}. Param("$a").Text("A"). Param("$b").Text("B"). Param("$c").Text("C"). Build(), ), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithParams", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXECUTE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithParams", }, }, Parameters: map[string]*Ydb.TypedValue{ "$a": { Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, Value: &Ydb.Value{ Value: &Ydb.Value_TextValue{ TextValue: "A", }, }, }, "$b": { Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, Value: &Ydb.Value{ Value: &Ydb.Value_TextValue{ TextValue: "B", }, }, }, "$c": { Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, Value: &Ydb.Value{ Value: &Ydb.Value_TextValue{ TextValue: "C", }, }, }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_NONE, ConcurrentResultSets: false, }, }, { name: "WithExplain", opts: []options.Execute{ options.WithExecMode(options.ExecModeExplain), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithExplain", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXPLAIN, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithExplain", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_NONE, ConcurrentResultSets: false, }, }, { name: "WithValidate", opts: []options.Execute{ options.WithExecMode(options.ExecModeValidate), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithValidate", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_VALIDATE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithValidate", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_NONE, ConcurrentResultSets: false, }, }, { name: "WithValidate", opts: []options.Execute{ options.WithExecMode(options.ExecModeParse), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithValidate", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_PARSE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithValidate", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_NONE, ConcurrentResultSets: false, }, }, { name: "WithStatsFull", opts: []options.Execute{ options.WithStatsMode(options.StatsModeFull, nil), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithStatsFull", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXECUTE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithStatsFull", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_FULL, ConcurrentResultSets: false, }, }, { name: "WithStatsBasic", opts: []options.Execute{ options.WithStatsMode(options.StatsModeBasic, nil), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithStatsBasic", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXECUTE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithStatsBasic", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_BASIC, ConcurrentResultSets: false, }, }, { name: "WithStatsProfile", opts: []options.Execute{ options.WithStatsMode(options.StatsModeProfile, nil), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithStatsProfile", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXECUTE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithStatsProfile", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_PROFILE, ConcurrentResultSets: false, }, }, { name: "WithGrpcCallOptions", opts: []options.Execute{ options.WithCallOptions(grpc.Header(&metadata.MD{ "ext-header": []string{"test"}, })), }, request: &Ydb_Query.ExecuteQueryRequest{ SessionId: "WithGrpcCallOptions", ExecMode: Ydb_Query.ExecMode_EXEC_MODE_EXECUTE, Query: &Ydb_Query.ExecuteQueryRequest_QueryContent{ QueryContent: &Ydb_Query.QueryContent{ Syntax: Ydb_Query.Syntax_SYNTAX_YQL_V1, Text: "WithGrpcCallOptions", }, }, StatsMode: Ydb_Query.StatsMode_STATS_MODE_NONE, ConcurrentResultSets: false, }, callOptions: []grpc.CallOption{ grpc.Header(&metadata.MD{ "ext-header": []string{"test"}, }), }, }, } { t.Run(tt.name, func(t *testing.T) { request, callOptions, err := executeQueryRequest( tt.name, tt.name, options.ExecuteSettings(tt.opts...), options.ResultSetsTypeOrdered, ) require.NoError(t, err) require.Equal(t, request.String(), tt.request.String()) require.Equal(t, tt.callOptions, callOptions) }) } } func happyWayStream(ctrl *gomock.Controller) Ydb_Query_V1.QueryService_ExecuteQueryClient { stream := newExecuteQueryStreamMock(ctrl) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, TxMeta: &Ydb_Query.TransactionMeta{ Id: "456", }, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{ Columns: []*Ydb.Column{ { Name: "a", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UINT64, }, }, }, { Name: "b", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, }, }, Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 1, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "1", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 2, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "2", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 3, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "3", }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, ResultSetIndex: 0, ResultSet: &Ydb.ResultSet{ Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 4, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "4", }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 5, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "5", }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, ResultSetIndex: 1, ResultSet: &Ydb.ResultSet{ Columns: []*Ydb.Column{ { Name: "c", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UINT64, }, }, }, { Name: "d", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, }, { Name: "e", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_BOOL, }, }, }, }, Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 1, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "1", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: true, }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 2, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "2", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: false, }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, ResultSetIndex: 1, ResultSet: &Ydb.ResultSet{ Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 3, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "3", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: true, }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 4, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "4", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: false, }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 5, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "5", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: false, }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, ResultSetIndex: 2, ResultSet: &Ydb.ResultSet{ Columns: []*Ydb.Column{ { Name: "c", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UINT64, }, }, }, { Name: "d", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_UTF8, }, }, }, { Name: "e", Type: &Ydb.Type{ Type: &Ydb.Type_TypeId{ TypeId: Ydb.Type_BOOL, }, }, }, }, Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 1, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "1", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: true, }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 2, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "2", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: false, }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(&Ydb_Query.ExecuteQueryResponsePart{ Status: Ydb.StatusIds_SUCCESS, ResultSetIndex: 2, ResultSet: &Ydb.ResultSet{ Rows: []*Ydb.Value{ { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 3, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "3", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: true, }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 4, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "4", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: false, }, }}, }, { Items: []*Ydb.Value{{ Value: &Ydb.Value_Uint64Value{ Uint64Value: 5, }, }, { Value: &Ydb.Value_TextValue{ TextValue: "5", }, }, { Value: &Ydb.Value_BoolValue{ BoolValue: false, }, }}, }, }, }, }, nil) stream.EXPECT().Recv().Return(nil, io.EOF) return stream }