-
Notifications
You must be signed in to change notification settings - Fork 19
/
file_client.go
123 lines (97 loc) · 3.01 KB
/
file_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
123
// Copyright (C) 2023 Gobalsky Labs Limited
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as
// published by the Free Software Foundation, either version 3 of the
// License, or (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package broker
import (
"encoding/binary"
"fmt"
"os"
"path/filepath"
"sync"
"google.golang.org/protobuf/proto"
"code.vegaprotocol.io/vega/core/events"
"code.vegaprotocol.io/vega/logging"
)
type FileClient struct {
log *logging.Logger
config *FileConfig
file *os.File
mut sync.RWMutex
seqNum uint64
}
const (
NumberOfSeqNumBytes = 8
NumberOfSizeBytes = 4
namedFileClientLogger = "file-client"
)
func NewFileClient(log *logging.Logger, config *FileConfig) (*FileClient, error) {
log = log.Named(namedFileClientLogger)
fc := &FileClient{
log: log,
config: config,
}
filePath, err := filepath.Abs(config.File)
if err != nil {
return nil, fmt.Errorf("unable to determine absolute path of file %s: %w", config.File, err)
}
fc.file, err = os.Create(filePath)
if err != nil {
return nil, fmt.Errorf("unable to create file %s: %w", filePath, err)
}
log.Infof("persisting events to: %s\n", filePath)
return fc, nil
}
func (fc *FileClient) SendBatch(evts []events.Event) error {
for _, evt := range evts {
if err := fc.Send(evt); err != nil {
return err
}
}
return nil
}
func (fc *FileClient) Send(event events.Event) error {
fc.mut.RLock()
defer fc.mut.RUnlock()
err := WriteToBufferFile(fc.file, fc.seqNum, event)
fc.seqNum++
if err != nil {
return fmt.Errorf("failed to write event to buffer file: %w", err)
}
return nil
}
func WriteToBufferFile(bufferFile *os.File, bufferSeqNum uint64, event events.Event) error {
rawEvent, err := proto.Marshal(event.StreamMessage())
if err != nil {
return fmt.Errorf("failed to marshal bus event:%w", err)
}
return WriteRawToBufferFile(bufferFile, bufferSeqNum, rawEvent)
}
func WriteRawToBufferFile(bufferFile *os.File, bufferSeqNum uint64, rawEvent []byte) error {
seqNumBytes := make([]byte, NumberOfSeqNumBytes)
sizeBytes := make([]byte, NumberOfSizeBytes)
size := NumberOfSeqNumBytes + uint32(len(rawEvent))
binary.BigEndian.PutUint64(seqNumBytes, bufferSeqNum)
binary.BigEndian.PutUint32(sizeBytes, size)
allBytes := append([]byte{}, sizeBytes...)
allBytes = append(allBytes, seqNumBytes...)
allBytes = append(allBytes, rawEvent...)
_, err := bufferFile.Write(allBytes)
if err != nil {
return fmt.Errorf("failed to write to buffer file:%w", err)
}
return nil
}
func (fc *FileClient) Close() error {
return fc.file.Close()
}