forked from facebookarchive/zk
-
Notifications
You must be signed in to change notification settings - Fork 0
/
util.go
95 lines (77 loc) · 2.66 KB
/
util.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
/*
* Copyright (c) Facebook, Inc. and its affiliates.
*
* This source code is licensed under the MIT license found in the
* LICENSE file in the root directory of this source tree.
*
*/
package zk
import (
"bytes"
"fmt"
"io"
"github.com/facebookincubator/zk/internal/proto"
"github.com/go-zookeeper/jute/lib/go/jute"
)
// WriteRecords takes in one or more RecordWriter instances, serializes them to a byte array
// and writes them to the provided io.Writer.
func WriteRecords(w io.Writer, generated ...jute.RecordWriter) error {
sendBuf := &bytes.Buffer{}
enc := jute.NewBinaryEncoder(sendBuf)
for _, generatedStruct := range generated {
if err := generatedStruct.Write(enc); err != nil {
return fmt.Errorf("could not encode struct: %w", err)
}
}
// copy encoded request bytes
requestBytes := append([]byte(nil), sendBuf.Bytes()...)
// use encoder to prepend request length to the request bytes
sendBuf.Reset()
if err := enc.WriteBuffer(requestBytes); err != nil {
return fmt.Errorf("could not write buffer: %w", err)
}
if err := enc.WriteEnd(); err != nil {
return fmt.Errorf("could not write buffer: %w", err)
}
if _, err := w.Write(sendBuf.Bytes()); err != nil {
return fmt.Errorf("error writing to io.Writer: %w", err)
}
return nil
}
// ReadRecord reads the request header and body depending on the opcode.
// It returns the serialized request header and body, or an error if it occurs.
func ReadRecord(r io.Reader) (*proto.RequestHeader, jute.RecordReader, error) {
dec, err := createDecoder(r)
if err != nil {
return nil, nil, fmt.Errorf("error reading request length: %w", err)
}
header := &proto.RequestHeader{}
if err = dec.ReadRecord(header); err != nil {
return nil, nil, fmt.Errorf("error reading RequestHeader: %w", err)
}
var req jute.RecordReader
switch header.Type {
case opGetData:
req = &proto.GetDataRequest{}
case opGetChildren:
req = &proto.GetChildrenRequest{}
default:
return nil, nil, fmt.Errorf("unrecognized header type: %d", header.Type)
}
if err := dec.ReadRecord(req); err != nil {
return nil, nil, fmt.Errorf("error reading request: %w", err)
}
return header, req, nil
}
// createDecoder reads a packet from io.Reader by reading N bytes from the packet header first,
// and then reading the remaining N bytes as per the Zookeeper protocol.
// It returns a jute.Decoder which can then be used to serialize the bytes into a valid struct.
func createDecoder(r io.Reader) (jute.Decoder, error) {
dec := jute.NewBinaryDecoder(r)
readBytes, err := dec.ReadBuffer()
if err != nil {
return nil, fmt.Errorf("error reading packet: %w", err)
}
dec = jute.NewBinaryDecoder(bytes.NewBuffer(readBytes))
return dec, nil
}