-
Notifications
You must be signed in to change notification settings - Fork 1
/
client.go
105 lines (93 loc) · 2.81 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
package api
import (
"context"
"fmt"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/protobuf/types/known/timestamppb"
"github.com/seal-io/meta-api/schema"
)
// Client holds the actions for receiving from the exposing service.
type Client interface {
// Ingest ingests specified type dataset from the exposing service,
// and parses dataset with the given IngestParser.
Ingest(ctx context.Context, typ schema.DatasetIngestRequestType, since time.Time, parse IngestParser) (err error)
// IngestAll ingests all types dataset from the exposing service,
// and parses dataset with the given IngestParser.
IngestAll(ctx context.Context, since time.Time, parse IngestParser) (err error)
// Close closes the client.
Close() error
}
// GetClient returns the Client.
func GetClient(ctx context.Context, listenOn string) (Client, error) {
var opts = []grpc.DialOption{
grpc.WithBlock(),
grpc.WithDefaultCallOptions(grpc.MaxCallRecvMsgSize(128 * 1024 * 1024)),
grpc.WithTransportCredentials(insecure.NewCredentials()),
}
var cc, err = grpc.DialContext(ctx, listenOn, opts...)
if err != nil {
return nil, fmt.Errorf("error dialing %s: %w", listenOn, err)
}
var cli = &client{
cc: cc,
}
return cli, nil
}
type client struct {
cc *grpc.ClientConn
}
// IngestParser is the parser to parse the given api.DatasetIngestResponseBody.
type IngestParser func(currentWindow int32, body schema.DatasetIngestResponseBody) error
func (in *client) Ingest(ctx context.Context, typ schema.DatasetIngestRequestType, since time.Time, parse IngestParser) error {
var cli, err = schema.NewDatasetServiceClient(in.cc).Ingest(ctx)
if err != nil {
return fmt.Errorf("error creating ingest client: %w", err)
}
var window int32
for window >= 0 {
var req = &schema.DatasetIngestRequest{
Window: window,
Type: typ,
}
if !since.IsZero() {
req.Since = timestamppb.New(since)
}
err = cli.Send(req)
if err != nil {
return fmt.Errorf("error sending ingest request: %w", err)
}
var resp *schema.DatasetIngestResponse
resp, err = cli.Recv()
if err != nil {
return fmt.Errorf("error receiving ingest response: %w", err)
}
if parse != nil && resp.GetBody() != nil {
err = parse(window, resp.GetBody())
if err != nil {
return fmt.Errorf("error parsing ingest response: %w", err)
}
}
window = resp.GetNextWindow()
if resp.NextWindow == nil {
window = -1
}
}
return nil
}
func (in *client) IngestAll(ctx context.Context, since time.Time, parse IngestParser) error {
for typ := 0; typ < len(schema.DatasetIngestRequestType_name); typ++ {
var err = in.Ingest(ctx, schema.DatasetIngestRequestType(typ), since, parse)
if err != nil {
return err
}
}
return nil
}
func (in *client) Close() error {
if in.cc != nil {
return in.cc.Close()
}
return nil
}