-
Notifications
You must be signed in to change notification settings - Fork 179
Expand file tree
/
Copy pathintakeProxyMiddleware.ts
More file actions
347 lines (312 loc) · 9.98 KB
/
intakeProxyMiddleware.ts
File metadata and controls
347 lines (312 loc) · 9.98 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
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
import { createInflate, inflateSync } from 'node:zlib'
import https from 'node:https'
import type express from 'express'
import createBusboy from 'busboy'
import type { BrowserProfileEvent, BrowserProfilerTrace } from '@datadog/browser-rum/src/types/profiling'
import type { BrowserSegment, BrowserSegmentMetadata } from '@datadog/browser-rum/src/types/sessionReplay'
import type { LogsEvent } from '@datadog/browser-logs/src/logsEvent.types'
import type { RumEvent } from '@datadog/browser-rum-core/src/rumEvent.types'
import type { TelemetryEvent } from '@datadog/browser-core/src/domain/telemetry/telemetryEvent.types'
interface BaseIntakeRequest {
isBridge: boolean
encoding: string | null
transport: string | null
}
export type LogsIntakeRequest = {
intakeType: 'logs'
events: LogsEvent[]
} & BaseIntakeRequest
export type RumIntakeRequest = {
intakeType: 'rum'
events: Array<RumEvent | TelemetryEvent>
} & BaseIntakeRequest
export type ReplayIntakeRequest = {
intakeType: 'replay'
segment: BrowserSegment
metadata: BrowserSegmentMetadataAndSegmentSizes
segmentFile: {
filename: string
encoding: string
mimetype: string
}
} & BaseIntakeRequest
export type BrowserSegmentMetadataAndSegmentSizes = BrowserSegmentMetadata & {
raw_segment_size: number
compressed_segment_size: number
}
export type ProfileIntakeRequest = {
intakeType: 'profile'
event: BrowserProfileEvent
trace: BrowserProfilerTrace
traceFile: {
filename: string
encoding: string | null
mimetype: string
}
} & BaseIntakeRequest
export type DebuggerIntakeRequest = {
intakeType: 'debugger'
events: Array<Record<string, unknown>>
} & BaseIntakeRequest
export type IntakeRequest =
| LogsIntakeRequest
| RumIntakeRequest
| ReplayIntakeRequest
| ProfileIntakeRequest
| DebuggerIntakeRequest
interface IntakeRequestInfos {
isBridge: boolean
intakeType: IntakeRequest['intakeType']
encoding: string | null
transport: string | null
}
interface IntakeProxyOptions {
onRequest?: (request: IntakeRequest) => void
}
export function createIntakeProxyMiddleware(options: IntakeProxyOptions): express.RequestHandler {
return async (req, res) => {
const infos = computeIntakeRequestInfos(req)
try {
const [intakeRequest] = await Promise.all([
readIntakeRequest(req, infos),
!infos.isBridge && forwardIntakeRequestToDatadog(req),
])
options.onRequest?.(intakeRequest)
} catch (error) {
console.error('Error while processing request:', error)
}
res.end()
}
}
function computeIntakeRequestInfos(req: express.Request): IntakeRequestInfos {
const ddforward = req.query.ddforward as string | undefined
if (!ddforward) {
throw new Error('ddforward is missing')
}
const { pathname, searchParams } = new URL(ddforward, 'https://example.org')
const encoding = req.headers['content-encoding'] || searchParams.get('dd-evp-encoding')
const transport = searchParams.get('_dd.api')
if (req.query.bridge === 'true') {
const eventType = req.query.event_type
return {
isBridge: true,
encoding,
transport,
intakeType: eventType === 'log' ? 'logs' : eventType === 'record' ? 'replay' : 'rum',
}
}
let intakeType: IntakeRequest['intakeType']
// pathname = /api/v2/rum
const endpoint = pathname.split(/[/?]/)[3]
if (
endpoint === 'logs' ||
endpoint === 'rum' ||
endpoint === 'replay' ||
endpoint === 'profile' ||
endpoint === 'debugger'
) {
intakeType = endpoint
} else {
throw new Error("Can't find intake type")
}
return {
isBridge: false,
encoding,
transport,
intakeType,
}
}
function readIntakeRequest(req: express.Request, infos: IntakeRequestInfos): Promise<IntakeRequest> {
if (infos.intakeType === 'replay') {
return readReplayIntakeRequest(req, infos as IntakeRequestInfos & { intakeType: 'replay' })
}
if (infos.intakeType === 'profile') {
return readProfileIntakeRequest(req, infos as IntakeRequestInfos & { intakeType: 'profile' })
}
if (infos.intakeType === 'debugger') {
return readDebuggerIntakeRequest(req, infos as IntakeRequestInfos & { intakeType: 'debugger' })
}
return readRumOrLogsIntakeRequest(req, infos as IntakeRequestInfos & { intakeType: 'rum' | 'logs' })
}
async function readRumOrLogsIntakeRequest(
req: express.Request,
infos: IntakeRequestInfos & { intakeType: 'rum' | 'logs' }
): Promise<RumIntakeRequest | LogsIntakeRequest> {
const rawBody = await readStream(req)
const encodedBody = infos.encoding === 'deflate' ? inflateSync(rawBody) : rawBody
return {
...infos,
events: encodedBody
.toString('utf-8')
.split('\n')
.map((line): any => JSON.parse(line)),
}
}
async function readDebuggerIntakeRequest(
req: express.Request,
infos: IntakeRequestInfos & { intakeType: 'debugger' }
): Promise<DebuggerIntakeRequest> {
const rawBody = await readStream(req)
const encodedBody = infos.encoding === 'deflate' ? inflateSync(rawBody) : rawBody
return {
...infos,
events: encodedBody
.toString('utf-8')
.split('\n')
// eslint-disable-next-line @typescript-eslint/no-unsafe-return
.map((line): Record<string, unknown> => JSON.parse(line)),
}
}
function readReplayIntakeRequest(
req: express.Request,
infos: IntakeRequestInfos & { intakeType: 'replay' }
): Promise<ReplayIntakeRequest> {
return new Promise((resolve, reject) => {
if (infos.isBridge) {
readStream(req)
.then((rawBody) => {
resolve({
...infos,
segment: {
records: rawBody
.toString('utf-8')
.split('\n')
.map((line): unknown => JSON.parse(line)),
},
} as ReplayIntakeRequest)
})
.catch(reject)
return
}
let segmentPromise: Promise<{
encoding: string
filename: string
mimetype: string
segment: BrowserSegment
}>
let metadataPromise: Promise<BrowserSegmentMetadataAndSegmentSizes>
const busboy = createBusboy({ headers: req.headers })
busboy.on('file', (name, stream, info) => {
const { filename, encoding, mimeType } = info
if (name === 'segment') {
segmentPromise = readStream(stream.pipe(createInflate())).then((data) => ({
encoding,
filename,
mimetype: mimeType,
segment: JSON.parse(data.toString()),
}))
} else if (name === 'event') {
metadataPromise = readStream(stream).then(
(data) => JSON.parse(data.toString()) as BrowserSegmentMetadataAndSegmentSizes
)
}
})
busboy.on('finish', () => {
Promise.all([segmentPromise, metadataPromise])
.then(([{ segment, ...segmentFile }, metadata]) => ({
...infos,
segmentFile,
metadata,
segment,
}))
.then(resolve, reject)
})
req.pipe(busboy)
})
}
function readProfileIntakeRequest(
req: express.Request,
infos: IntakeRequestInfos & { intakeType: 'profile' }
): Promise<ProfileIntakeRequest> {
return new Promise((resolve, reject) => {
let eventPromise: Promise<BrowserProfileEvent>
let tracePromise: Promise<{
trace: BrowserProfilerTrace
encoding: string | null
filename: string
mimetype: string
}>
const busboy = createBusboy({ headers: req.headers })
busboy.on('file', (name, stream, info) => {
const { filename, mimeType } = info
if (name === 'event') {
eventPromise = readStream(stream).then((data) => JSON.parse(data.toString()) as BrowserProfileEvent)
} else if (name === 'wall-time.json') {
tracePromise = readStream(stream).then((data) => {
let encoding: string | null
if (isDeflateEncoded(data)) {
encoding = 'deflate'
data = inflateSync(data)
} else {
encoding = null
}
return {
trace: JSON.parse(data.toString()) as BrowserProfilerTrace,
encoding,
filename,
mimetype: mimeType,
}
})
} else {
// Skip other attachments
stream.resume()
}
})
busboy.on('finish', () => {
Promise.all([eventPromise, tracePromise])
.then(([event, { trace, ...traceFile }]) => ({
...infos,
event,
trace,
traceFile,
}))
.then(resolve, reject)
})
req.pipe(busboy)
})
}
function forwardIntakeRequestToDatadog(req: express.Request): Promise<any> {
return new Promise<void>((resolve) => {
const ddforward = req.query.ddforward! as string
if (!/^\/api\/v2\//.test(ddforward)) {
throw new Error(`Unsupported ddforward: ${ddforward}`)
}
const options = {
method: 'POST',
headers: {
'X-Forwarded-For': req.socket.remoteAddress,
'Content-Type': req.headers['content-type'],
'User-Agent': req.headers['user-agent'],
},
}
const datadogIntakeRequest = https.request(new URL(ddforward, 'https://browser-intake-datadoghq.com'), options)
req.pipe(datadogIntakeRequest)
datadogIntakeRequest.on('response', resolve)
datadogIntakeRequest.on('error', (error) => {
console.log('Error while forwarding request to Datadog:', error)
resolve()
})
})
}
function readStream(stream: NodeJS.ReadableStream): Promise<Buffer> {
return new Promise((resolve, reject) => {
const buffers: Buffer[] = []
stream.on('data', (data: Buffer) => {
buffers.push(data)
})
stream.on('error', (error) => {
// eslint-disable-next-line @typescript-eslint/prefer-promise-reject-errors
reject(error)
})
stream.on('end', () => {
resolve(Buffer.concat(buffers))
})
})
}
function isDeflateEncoded(buffer: Buffer): boolean {
// Check for deflate/zlib magic numbers
// 0x78 0x01 - No Compression/low
// 0x78 0x9C - Default Compression
// 0x78 0xDA - Best Compression
return buffer.length >= 2 && buffer[0] === 0x78 && (buffer[1] === 0x01 || buffer[1] === 0x9c || buffer[1] === 0xda)
}