-
Notifications
You must be signed in to change notification settings - Fork 10
Expand file tree
/
Copy pathprotocol.go
More file actions
293 lines (248 loc) · 9.5 KB
/
Copy pathprotocol.go
File metadata and controls
293 lines (248 loc) · 9.5 KB
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
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
package airbyte
import (
"encoding/json"
"errors"
"io"
"time"
)
// Should conform to https://github.com/airbytehq/airbyte/blob/master/airbyte-protocol/models/src/main/resources/airbyte_protocol/airbyte_protocol.yaml
type cmd string
const (
cmdSpec cmd = "spec"
cmdCheck cmd = "check"
cmdDiscover cmd = "discover"
cmdRead cmd = "read"
)
type msgType string
const (
msgTypeRecord msgType = "RECORD"
msgTypeState msgType = "STATE"
msgTypeLog msgType = "LOG"
msgTypeConnectionStat msgType = "CONNECTION_STATUS"
msgTypeCatalog msgType = "CATALOG"
msgTypeSpec msgType = "SPEC"
)
var errInvalidTypePayload = errors.New("message type and payload are invalid")
type message struct {
Type msgType `json:"type"`
*record `json:"record,omitempty"`
*state `json:"state,omitempty"`
*logMessage `json:"log,omitempty"`
*ConnectorSpecification `json:"spec,omitempty"`
*connectionStatus `json:"connectionStatus,omitempty"`
*Catalog `json:"catalog,omitempty"`
}
// message MarshalJSON is a custom marshaller which validates the messageType with the sub-struct
func (m *message) MarshalJSON() ([]byte, error) {
switch m.Type {
case msgTypeRecord:
if m.record == nil ||
m.state != nil ||
m.logMessage != nil ||
m.connectionStatus != nil ||
m.Catalog != nil {
return nil, errInvalidTypePayload
}
case msgTypeState:
if m.state == nil ||
m.record != nil ||
m.logMessage != nil ||
m.connectionStatus != nil ||
m.Catalog != nil {
return nil, errInvalidTypePayload
}
case msgTypeLog:
if m.logMessage == nil ||
m.record != nil ||
m.state != nil ||
m.connectionStatus != nil ||
m.Catalog != nil {
return nil, errInvalidTypePayload
}
}
type m2 message
return json.Marshal(m2(*m))
}
// write emits data outbound from your src/destination to airbyte workers
func write(w io.Writer, m *message) error {
return json.NewEncoder(w).Encode(m)
}
// record defines a record as per airbyte - a "data point"
type record struct {
EmittedAt int64 `json:"emitted_at"`
Namespace string `json:"namespace"`
Data interface{} `json:"data"`
Stream string `json:"stream"`
}
// state is used to store data between syncs - useful for incremental syncs and state storage
type state struct {
Data interface{} `json:"data"`
}
// LogLevel defines the log levels that can be emitted with airbyte logs
type LogLevel string
const (
LogLevelFatal LogLevel = "FATAL"
LogLevelError LogLevel = "ERROR"
LogLevelWarn LogLevel = "WARN"
LogLevelInfo LogLevel = "INFO"
LogLevelDebug LogLevel = "DEBUG"
LogLevelTrace LogLevel = "TRACE"
)
type logMessage struct {
Level LogLevel `json:"level"`
Message string `json:"message"`
}
type checkStatus string
const (
checkStatusSuccess checkStatus = "SUCCEEDED"
checkStatusFailed checkStatus = "FAILED"
)
type connectionStatus struct {
Status checkStatus `json:"status"`
}
// Catalog defines the complete available schema you can sync with a source
// This should not be mistaken with ConfiguredCatalog which is the "selected" schema you want to sync
type Catalog struct {
Streams []Stream `json:"streams"`
}
// Stream defines a single "schema" you'd like to sync - think of this as a table, collection, topic, etc. In airbyte terminology these are "streams"
type Stream struct {
Name string `json:"name"`
JSONSchema Properties `json:"json_schema"`
SupportedSyncModes []SyncMode `json:"supported_sync_modes,omitempty"`
SourceDefinedCursor bool `json:"source_defined_cursor,omitempty"`
DefaultCursorField []string `json:"default_cursor_field,omitempty"`
SourceDefinedPrimaryKey [][]string `json:"source_defined_primary_key,omitempty"`
Namespace string `json:"namespace"`
}
// ConfiguredCatalog is the "selected" schema you want to sync
// This should not be mistaken with Catalog which represents the complete available schema to sync
type ConfiguredCatalog struct {
Streams []ConfiguredStream `json:"streams"`
}
// ConfiguredStream defines a single selected stream to sync
type ConfiguredStream struct {
Stream Stream `json:"stream"`
SyncMode SyncMode `json:"sync_mode"`
CursorField []string `json:"cursor_field"`
DestinationSyncMode DestinationSyncMode `json:"destination_sync_mode"`
PrimaryKey [][]string `json:"primary_key"`
}
// SyncMode defines the modes that your source is able to sync in
type SyncMode string
const (
// SyncModeFullRefresh means the data will be wiped and fully synced on each run
SyncModeFullRefresh SyncMode = "full_refresh"
// SyncModeIncremental is used for incremental syncs
SyncModeIncremental SyncMode = "incremental"
)
// DestinationSyncMode represents how the destination should interpret your data
type DestinationSyncMode string
var (
// DestinationSyncModeAppend is used for the destination to know it needs to append data
DestinationSyncModeAppend DestinationSyncMode = "append"
// DestinationSyncModeOverwrite is used to indicate the destination should overwrite data
DestinationSyncModeOverwrite DestinationSyncMode = "overwrite"
)
// ConnectorSpecification is used to define the connector wide settings. Every connection using your connector will comply to these settings
type ConnectorSpecification struct {
DocumentationURL string `json:"documentationUrl,omitempty"`
ChangeLogURL string `json:"changeLogUrl"`
SupportsIncremental bool `json:"supportsIncremental"`
SupportsNormalization bool `json:"supportsNormalization"`
SupportsDBT bool `json:"supportsDBT"`
SupportedDestinationSyncModes []DestinationSyncMode `json:"supported_destination_sync_modes"`
ConnectionSpecification ConnectionSpecification `json:"connectionSpecification"`
}
// https://json-schema.org/learn/getting-started-step-by-step.html
// Properties defines the property map which is used to define any single "field name" along with its specification
type Properties struct {
Properties map[PropertyName]PropertySpec `json:"properties"`
}
// PropertyName is a alias for a string to make it clear to the user that the "key" in the map is the name of the property
type PropertyName string
// ConnectionSpecification is used to define the settings that are configurable "per" instance of your connector
type ConnectionSpecification struct {
Title string `json:"title"`
Description string `json:"description"`
Properties
Type string `json:"type"` // should always be "object"
Required []PropertyName `json:"required"`
}
// PropType defines the property types any field can take. See more here: https://docs.airbyte.com/understanding-airbyte/supported-data-types
type PropType string
const (
String PropType = "string"
Number PropType = "number"
Integer PropType = "integer"
Object PropType = "object"
Array PropType = "array"
Null PropType = "null"
)
// AirbytePropType is used to define airbyte specific property types. See more here: https://docs.airbyte.com/understanding-airbyte/supported-data-types
type AirbytePropType string
const (
TimestampWithTZ AirbytePropType = "timestamp_with_timezone"
TimestampWOTZ AirbytePropType = "timestamp_without_timezone"
BigInteger AirbytePropType = "big_integer"
BigNumber AirbytePropType = "big_number"
)
// FormatType is used to define data type formats supported by airbyte where needed (usually for strings formatted as dates). See more here: https://docs.airbyte.com/understanding-airbyte/supported-data-types
type FormatType string
const (
Date FormatType = "date"
DateTime FormatType = "datetime"
)
type PropertyType struct {
Type []PropType `json:"type,omitempty"`
AirbyteType AirbytePropType `json:"airbyte_type,omitempty"`
}
type PropertySpec struct {
Description string `json:"description"`
PropertyType `json:",omitempty"`
Examples []string `json:"examples,omitempty"`
Items map[string]interface{} `json:"items,omitempty"`
Properties map[PropertyName]PropertySpec `json:"properties,omitempty"`
IsSecret bool `json:"airbyte_secret,omitempty"`
}
// LogWriter is exported for documentation purposes - only use this through LogTracker or MessageTracker
// to ensure thread-safe behavior with the writer
type LogWriter func(level LogLevel, s string) error
// StateWriter is exported for documentation purposes - only use this through MessageTracker
type StateWriter func(v interface{}) error
// RecordWriter is exported for documentation purposes - only use this through MessageTracker
type RecordWriter func(v interface{}, streamName string, namespace string) error
func newLogWriter(w io.Writer) LogWriter {
return func(lvl LogLevel, s string) error {
return write(w, &message{
Type: msgTypeLog,
logMessage: &logMessage{
Level: lvl,
Message: s,
},
})
}
}
func newStateWriter(w io.Writer) StateWriter {
return func(s interface{}) error {
return write(w, &message{
Type: msgTypeState,
state: &state{
Data: s,
},
})
}
}
func newRecordWriter(w io.Writer) RecordWriter {
return func(s interface{}, stream string, namespace string) error {
return write(w, &message{
Type: msgTypeRecord,
record: &record{
EmittedAt: time.Now().UnixMilli(),
Data: s,
Namespace: namespace,
Stream: stream,
},
})
}
}