-
Notifications
You must be signed in to change notification settings - Fork 9
/
ndjson_record_io_stream.go
104 lines (82 loc) · 2.08 KB
/
ndjson_record_io_stream.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
package ndjsonrecordiostream
import (
"encoding/json"
"fmt"
"io"
"reflect"
"github.com/inklabs/rangedb"
"github.com/inklabs/rangedb/provider/jsonrecordserializer"
)
type ndJSONRecordIoStream struct {
eventTypes map[string]reflect.Type
}
// New constructs an ndjson implementation of rangedb.RecordIoStream.
func New() *ndJSONRecordIoStream {
return &ndJSONRecordIoStream{
eventTypes: map[string]reflect.Type{},
}
}
func (s *ndJSONRecordIoStream) Write(writer io.Writer, recordIterator rangedb.RecordIterator) <-chan error {
errors := make(chan error)
go func() {
defer close(errors)
totalSaved := 0
for recordIterator.Next() {
if recordIterator.Err() != nil {
errors <- recordIterator.Err()
return
}
if totalSaved > 0 {
_, _ = fmt.Fprint(writer, "\n")
}
data, err := json.Marshal(recordIterator.Record())
if err != nil {
errors <- fmt.Errorf("failed marshalling event: %v", err)
return
}
_, _ = fmt.Fprintf(writer, "%s", data)
totalSaved++
}
}()
return errors
}
func (s *ndJSONRecordIoStream) Read(reader io.Reader) rangedb.RecordIterator {
resultRecords := make(chan rangedb.ResultRecord)
go func() {
defer close(resultRecords)
decoder := json.NewDecoder(reader)
decoder.UseNumber()
for decoder.More() {
record, err := jsonrecordserializer.UnmarshalRecord(decoder, s)
if err != nil {
resultRecords <- rangedb.ResultRecord{
Record: nil,
Err: err,
}
return
}
// TODO: Add cancel context to avoid deadlock
resultRecords <- rangedb.ResultRecord{
Record: record,
Err: nil,
}
}
}()
return rangedb.NewRecordIterator(resultRecords)
}
func (s *ndJSONRecordIoStream) Bind(events ...rangedb.Event) {
for _, e := range events {
s.eventTypes[e.EventType()] = getType(e)
}
}
func (s *ndJSONRecordIoStream) EventTypeLookup(eventTypeName string) (r reflect.Type, b bool) {
eventType, ok := s.eventTypes[eventTypeName]
return eventType, ok
}
func getType(object interface{}) reflect.Type {
t := reflect.TypeOf(object)
if t.Kind() == reflect.Ptr {
t = t.Elem()
}
return t
}