Skip to content

Commit 6d44174

Browse files
authored
Backport upgrade diagnostics lifecycle fixes to v6.x (#5833)
* fix: complete upgrade diagnostics lifecycle Successful CONNECT and protocol upgrades stop before Undici marks the request complete. Diagnostics subscribers retain request state, and streamed request-body listeners remain attached to the upgraded socket. Publish the existing response lifecycle around accepted upgrades. Preserve the early HTTP/2 CONNECT handoff in v6, and terminate diagnostics if the stream fails before its response. Refs: #5783 Signed-off-by: Ruben Bridgewater <ruben.bridgewater@datadoghq.com> * fix: preserve abort after upgrade handoff Completing upgrade requests suppresses late request errors, but v6 retains the abort callback as the transport cleanup owner after handoff. Returning early leaves upgraded sockets and HTTP/2 CONNECT streams open. Signed-off-by: Ruben Bridgewater <ruben.bridgewater@datadoghq.com> --------- Signed-off-by: Ruben Bridgewater <ruben.bridgewater@datadoghq.com>
1 parent 2a91fc8 commit 6d44174

5 files changed

Lines changed: 712 additions & 28 deletions

File tree

‎docs/docs/api/DiagnosticsChannel.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,8 @@ diagnosticsChannel.channel('undici:request:bodySent').subscribe(({ request }) =>
4040

4141
## `undici:request:headers`
4242

43-
This message is published after the response headers have been received, i.e. the response has been completed.
43+
This message is published after the response headers have been received. This includes a successful CONNECT or
44+
protocol upgrade response.
4445

4546
```js
4647
import diagnosticsChannel from 'diagnostics_channel'
@@ -57,6 +58,7 @@ diagnosticsChannel.channel('undici:request:headers').subscribe(({ request, respo
5758
## `undici:request:trailers`
5859

5960
This message is published after the response body and trailers have been received, i.e. the response has been completed.
61+
After an upgraded socket is passed to the request handler, this message is published with an empty `trailers` array.
6062

6163
```js
6264
import diagnosticsChannel from 'diagnostics_channel'

‎lib/core/request.js‎

Lines changed: 68 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -263,11 +263,77 @@ class Request {
263263
}
264264
}
265265

266-
onUpgrade (statusCode, headers, socket) {
266+
/**
267+
* @param {number|null} statusCode
268+
* @param {Buffer[]|null} headers
269+
* @param {import('node:stream').Duplex} socket
270+
* @param {string} [statusText]
271+
*/
272+
onUpgrade (statusCode, headers, socket, statusText = '') {
273+
this.onFinally()
274+
267275
assert(!this.aborted)
268276
assert(!this.completed)
269277

270-
return this[kHandler].onUpgrade(statusCode, headers, socket)
278+
if (statusCode !== null) {
279+
this.#publishUpgradeHeaders(statusCode, headers, statusText)
280+
}
281+
282+
const result = this[kHandler].onUpgrade(statusCode, headers, socket)
283+
284+
if (!this.aborted) {
285+
this.completed = true
286+
if (statusCode !== null) {
287+
this.#publishUpgradeTrailers()
288+
}
289+
}
290+
291+
return result
292+
}
293+
294+
/**
295+
* @param {number} statusCode
296+
* @param {import('node:http2').IncomingHttpHeaders} headers
297+
* @param {(headers: import('node:http2').IncomingHttpHeaders) => Buffer[]} parseHeaders
298+
* @param {string} [statusText]
299+
*/
300+
onUpgradeResponse (statusCode, headers, parseHeaders, statusText = '') {
301+
assert(!this.aborted)
302+
assert(this.completed)
303+
304+
if (channels.headers.hasSubscribers) {
305+
this.#publishUpgradeHeaders(statusCode, parseHeaders(headers), statusText)
306+
}
307+
this.#publishUpgradeTrailers()
308+
}
309+
310+
/**
311+
* @param {Error} error
312+
*/
313+
onUpgradeError (error) {
314+
assert(!this.aborted)
315+
assert(this.completed)
316+
317+
if (channels.error.hasSubscribers) {
318+
channels.error.publish({ request: this, error })
319+
}
320+
}
321+
322+
/**
323+
* @param {number} statusCode
324+
* @param {Buffer[]} headers
325+
* @param {string} statusText
326+
*/
327+
#publishUpgradeHeaders (statusCode, headers, statusText) {
328+
if (channels.headers.hasSubscribers) {
329+
channels.headers.publish({ request: this, response: { statusCode, headers, statusText } })
330+
}
331+
}
332+
333+
#publishUpgradeTrailers () {
334+
if (channels.trailers.hasSubscribers) {
335+
channels.trailers.publish({ request: this, trailers: [] })
336+
}
271337
}
272338

273339
onComplete (trailers) {

‎lib/dispatcher/client-h1.js‎

Lines changed: 18 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -432,7 +432,7 @@ class Parser {
432432
}
433433

434434
onUpgrade (head) {
435-
const { upgrade, client, socket, headers, statusCode } = this
435+
const { upgrade, client, socket, headers, statusCode, statusText } = this
436436

437437
assert(upgrade)
438438
assert(client[kSocket] === socket)
@@ -467,9 +467,10 @@ class Parser {
467467
client.emit('disconnect', client[kUrl], [client], new InformationalError('upgrade'))
468468

469469
try {
470-
request.onUpgrade(statusCode, headers, socket)
471-
} catch (err) {
472-
util.destroy(socket, err)
470+
request.onUpgrade(statusCode, headers, socket, statusText)
471+
} catch (error) {
472+
util.errorRequest(client, request, error)
473+
util.destroy(socket, error)
473474
}
474475

475476
client[kResume]()
@@ -1050,12 +1051,22 @@ function writeH1 (client, request) {
10501051
const socket = client[kSocket]
10511052
clearIdleSocketValidation(socket)
10521053

1053-
const abort = (err) => {
1054-
if (request.aborted || request.completed) {
1054+
/**
1055+
* @param {Error} [error]
1056+
*/
1057+
const abort = (error) => {
1058+
if (request.aborted) {
1059+
return
1060+
}
1061+
1062+
if (request.completed) {
1063+
if (request.upgrade || request.method === 'CONNECT') {
1064+
util.destroy(socket, new InformationalError('aborted'))
1065+
}
10551066
return
10561067
}
10571068

1058-
util.errorRequest(client, request, err || new RequestAbortedError())
1069+
util.errorRequest(client, request, error || new RequestAbortedError())
10591070

10601071
util.destroy(body)
10611072
util.destroy(socket, new InformationalError('aborted'))

‎lib/dispatcher/client-h2.js‎

Lines changed: 70 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
'use strict'
22

33
const assert = require('node:assert')
4+
const { errorMonitor } = require('node:events')
45
const { pipeline } = require('node:stream')
56
const util = require('../core/util.js')
67
const {
@@ -77,6 +78,15 @@ function parseH2Headers (headers) {
7778
return result
7879
}
7980

81+
/**
82+
* @param {import('node:http2').IncomingHttpHeaders} headers
83+
* @returns {Buffer[]}
84+
*/
85+
function parseH2ResponseHeaders (headers) {
86+
const { [HTTP2_HEADER_STATUS]: _statusCode, ...realHeaders } = headers
87+
return parseH2Headers(realHeaders)
88+
}
89+
8090
async function connectH2 (client, socket) {
8191
client[kSocket] = socket
8292

@@ -297,22 +307,32 @@ function writeH2 (client, request) {
297307
headers[HTTP2_HEADER_AUTHORITY] = host || `${hostname}${port ? `:${port}` : ''}`
298308
headers[HTTP2_HEADER_METHOD] = method
299309

300-
const abort = (err) => {
301-
if (request.aborted || request.completed) {
310+
/**
311+
* @param {Error} [error]
312+
*/
313+
const abort = (error) => {
314+
if (request.aborted) {
315+
return
316+
}
317+
318+
if (request.completed) {
319+
if (method === 'CONNECT' && stream != null) {
320+
util.destroy(stream, error || new RequestAbortedError())
321+
}
302322
return
303323
}
304324

305-
err = err || new RequestAbortedError()
325+
error = error || new RequestAbortedError()
306326

307-
util.errorRequest(client, request, err)
327+
util.errorRequest(client, request, error)
308328

309329
if (stream != null) {
310-
util.destroy(stream, err)
330+
util.destroy(stream, error)
311331
}
312332

313333
// We do not destroy the socket as we can continue using the session
314334
// the stream get's destroyed and the session remains to create new streams
315-
util.destroy(body, err)
335+
util.destroy(body, error)
316336
client[kQueue][client[kRunningIdx]++] = null
317337
client[kResume]()
318338
}
@@ -331,25 +351,57 @@ function writeH2 (client, request) {
331351

332352
if (method === 'CONNECT') {
333353
session.ref()
334-
// We are already connected, streams are pending, first request
335-
// will create a new stream. We trigger a request to create the stream and wait until
336-
// `ready` event is triggered
337354
// We disabled endStream to allow the user to write to the stream
338355
stream = session.request(headers, { endStream: false, signal })
356+
let upgradeResponseFinished = false
357+
358+
/**
359+
* @param {import('node:http2').IncomingHttpHeaders} headers
360+
*/
361+
const onResponse = (headers) => {
362+
upgradeResponseFinished = true
363+
stream.off(errorMonitor, onUpgradeError)
364+
request.onUpgradeResponse(Number(headers[HTTP2_HEADER_STATUS]), headers, parseH2ResponseHeaders)
365+
}
339366

340-
if (stream.id && !stream.pending) {
341-
request.onUpgrade(null, null, stream)
342-
++session[kOpenStreams]
343-
client[kQueue][client[kRunningIdx]++] = null
344-
} else {
345-
stream.once('ready', () => {
367+
/**
368+
* @param {Error} error
369+
*/
370+
const onUpgradeError = (error) => {
371+
upgradeResponseFinished = true
372+
stream.off('response', onResponse)
373+
request.onUpgradeError(error)
374+
}
375+
376+
const onReady = () => {
377+
try {
346378
request.onUpgrade(null, null, stream)
347-
++session[kOpenStreams]
348-
client[kQueue][client[kRunningIdx]++] = null
349-
})
379+
} catch (error) {
380+
stream.off('response', onResponse)
381+
abort(error)
382+
return
383+
}
384+
385+
if (request.aborted) {
386+
return
387+
}
388+
389+
stream.off('error', abort)
390+
stream.once(errorMonitor, onUpgradeError)
391+
client[kQueue][client[kRunningIdx]++] = null
350392
}
351393

394+
stream.once('response', onResponse)
395+
stream.once('error', abort)
396+
++session[kOpenStreams]
397+
onReady()
398+
352399
stream.once('close', () => {
400+
if (!upgradeResponseFinished && request.completed) {
401+
stream.off('response', onResponse)
402+
stream.off(errorMonitor, onUpgradeError)
403+
request.onUpgradeError(new InformationalError(`HTTP/2: "stream error" received - code ${stream.rstCode}`))
404+
}
353405
session[kOpenStreams] -= 1
354406
if (session[kOpenStreams] === 0) session.unref()
355407
})

0 commit comments

Comments
 (0)