forked from asonawalla/gazette
/
read_api.go
149 lines (129 loc) · 4.15 KB
/
read_api.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
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
package broker
import (
"context"
"io"
"io/ioutil"
"time"
"github.com/LiveRamp/gazette/v2/pkg/client"
"github.com/LiveRamp/gazette/v2/pkg/fragment"
pb "github.com/LiveRamp/gazette/v2/pkg/protocol"
log "github.com/sirupsen/logrus"
"google.golang.org/grpc"
)
// Read dispatches the JournalServer.Read API.
func (svc *Service) Read(req *pb.ReadRequest, stream pb.Journal_ReadServer) (err error) {
defer instrumentJournalServerOp("read", &err, time.Now())
if err = req.Validate(); err != nil {
return err
}
var res resolution
res, err = svc.resolver.resolve(resolveArgs{
ctx: stream.Context(),
journal: req.Journal,
mayProxy: !req.DoNotProxy,
requirePrimary: false,
requireFullAssignment: false,
proxyHeader: req.Header,
})
if err != nil {
return err
} else if res.status != pb.Status_OK {
err = stream.Send(&pb.ReadResponse{Status: res.status, Header: &res.Header})
return err
} else if !res.journalSpec.Flags.MayRead() {
err = stream.Send(&pb.ReadResponse{Status: pb.Status_NOT_ALLOWED, Header: &res.Header})
return err
} else if res.replica == nil {
req.Header = &res.Header // Attach resolved Header to |req|, which we'll forward.
err = proxyRead(stream, req, svc.jc)
return err
}
if err = serveRead(stream, req, &res.Header, res.replica.index); err == context.Canceled {
err = nil // Gracefully terminate RPC.
} else if err != nil {
log.WithFields(log.Fields{"err": err, "req": req}).Warn("failed to serve Read")
}
return err
}
// proxyRead forwards a ReadRequest to a resolved peer broker.
func proxyRead(stream grpc.ServerStream, req *pb.ReadRequest, jc pb.JournalClient) error {
var ctx = pb.WithDispatchRoute(stream.Context(), req.Header.Route, req.Header.ProcessId)
var client, err = jc.Read(ctx, req)
if err != nil {
return err
}
// Ignore CloseSend's error. Currently, gRPC will never return one. If the
// stream is broken, it *could* return io.EOF but we'd rather read the actual
// casual error with RecvMsg.
_ = client.CloseSend()
var resp = new(pb.ReadResponse)
for {
if err = client.RecvMsg(resp); err == io.EOF {
return nil
} else if err != nil {
return err
} else if err = stream.SendMsg(resp); err != nil {
return err
}
}
}
// serveRead evaluates a client's Read RPC against the local replica index.
func serveRead(stream grpc.ServerStream, req *pb.ReadRequest, hdr *pb.Header, index *fragment.Index) error {
var buffer = make([]byte, chunkSize)
var reader io.ReadCloser
for i := 0; true; i++ {
var resp, file, err = index.Query(stream.Context(), req)
if err != nil {
return err
}
// Send the Header with the first response message (only).
if i == 0 {
resp.Header = hdr
}
if err = stream.SendMsg(resp); err != nil {
return err
}
// Return after sending Metadata if the Fragment query failed,
// or we were only asked to send metadata, or the Fragment is
// remote and we're instructed to not proxy.
if resp.Status != pb.Status_OK || req.MetadataOnly || file == nil && req.DoNotProxy {
return nil
}
// Note Query may have resolved or updated req.Offset. For the remainder of
// this iteration, we update |req.Offset| to reference the next byte to read.
req.Offset = resp.Offset
if file != nil {
reader = ioutil.NopCloser(io.NewSectionReader(
file, req.Offset-resp.Fragment.Begin, resp.Fragment.End-req.Offset))
} else {
if reader, err = fragment.Open(stream.Context(), *resp.Fragment); err != nil {
return err
} else if reader, err = client.NewFragmentReader(reader, *resp.Fragment, req.Offset); err != nil {
return err
}
}
// Loop over chunks read from |reader|, sending each to the client.
var n int
var readErr error
for readErr == nil {
if n, readErr = reader.Read(buffer); n == 0 {
continue
}
if err = stream.SendMsg(&pb.ReadResponse{
Offset: req.Offset,
Content: buffer[:n],
}); err != nil {
return err
}
req.Offset += int64(n)
}
if readErr != io.EOF {
return readErr
} else if err = reader.Close(); err != nil {
return err
}
// Loop to query and read the next Fragment.
}
return nil
}
var chunkSize = 1 << 17 // 128K.