forked from vitessio/vitess
/
client.go
122 lines (109 loc) · 3.2 KB
/
client.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
// Copyright 2015, Google Inc. All rights reserved.
// Use of this source code is governed by a BSD-style
// license that can be found in the LICENSE file.
package framework
import (
"errors"
"github.com/youtube/vitess/go/sqltypes"
"github.com/youtube/vitess/go/vt/callerid"
querypb "github.com/youtube/vitess/go/vt/proto/query"
"github.com/youtube/vitess/go/vt/tabletserver"
"github.com/youtube/vitess/go/vt/tabletserver/querytypes"
"golang.org/x/net/context"
vtrpcpb "github.com/youtube/vitess/go/vt/proto/vtrpc"
)
// QueryClient provides a convenient wrapper for TabletServer's query service.
// It's not thread safe, but you can create multiple clients that point to the
// same server.
type QueryClient struct {
ctx context.Context
target querypb.Target
server *tabletserver.TabletServer
transactionID int64
}
// NewClient creates a new client for Server.
func NewClient() *QueryClient {
return &QueryClient{
ctx: callerid.NewContext(
context.Background(),
&vtrpcpb.CallerID{},
&querypb.VTGateCallerID{Username: "dev"},
),
target: Target,
server: Server,
}
}
// Begin begins a transaction.
func (client *QueryClient) Begin() error {
if client.transactionID != 0 {
return errors.New("already in transaction")
}
transactionID, err := client.server.Begin(client.ctx, &client.target)
if err != nil {
return err
}
client.transactionID = transactionID
return nil
}
// Commit commits the current transaction.
func (client *QueryClient) Commit() error {
defer func() { client.transactionID = 0 }()
return client.server.Commit(client.ctx, &client.target, client.transactionID)
}
// Rollback rolls back the current transaction.
func (client *QueryClient) Rollback() error {
defer func() { client.transactionID = 0 }()
return client.server.Rollback(client.ctx, &client.target, client.transactionID)
}
// Execute executes a query.
func (client *QueryClient) Execute(query string, bindvars map[string]interface{}) (*sqltypes.Result, error) {
return client.server.Execute(
client.ctx,
&client.target,
query,
bindvars,
client.transactionID,
)
}
// StreamExecute executes a query & streams the results.
func (client *QueryClient) StreamExecute(query string, bindvars map[string]interface{}) (*sqltypes.Result, error) {
result := &sqltypes.Result{}
err := client.server.StreamExecute(
client.ctx,
&client.target,
query,
bindvars,
func(res *sqltypes.Result) error {
if result.Fields == nil {
result.Fields = res.Fields
}
result.Rows = append(result.Rows, res.Rows...)
result.RowsAffected += uint64(len(res.Rows))
return nil
},
)
if err != nil {
return nil, err
}
return result, nil
}
// Stream streams the resutls of a query.
func (client *QueryClient) Stream(query string, bindvars map[string]interface{}, sendFunc func(*sqltypes.Result) error) error {
return client.server.StreamExecute(
client.ctx,
&client.target,
query,
bindvars,
sendFunc,
)
}
// ExecuteBatch executes a batch of queries.
func (client *QueryClient) ExecuteBatch(queries []querytypes.BoundQuery, asTransaction bool) ([]sqltypes.Result, error) {
return client.server.ExecuteBatch(
client.ctx,
&client.target,
queries,
asTransaction,
client.transactionID,
)
}