Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
121 changes: 121 additions & 0 deletions contributions/65559.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,121 @@
---
pr-url: https://github.com/nodejs/node/pull/65559
---

# stream: fix pipeline function tail deadlock

## 문제 내용

`stream/promises.pipeline()`에서 마지막 entry가 함수이고 해당 함수가 `AsyncIterable`을 반환하는 경우, 내부에서 생성되는 `PassThrough`의 readable side가 소비되지 않아 pipeline이 끝나지 않는 문제를 분석했다.

데이터가 계속 쌓이면서 내부 buffer가 `highWaterMark`에 도달하면 Backpressure가 발생하고, 결국 pipeline이 완료되지 않아 반환된 Promise도 resolve되지 않는 상황이었다.

관련 Issue: #40685

## 원인 분석

Promise 기반 `pipeline()`은 callback 기반 API와 달리 내부에서 생성된 proxy Stream을 사용자에게 직접 노출하지 않는다.

따라서 마지막 함수가 `AsyncIterable`을 반환하면서 내부 `PassThrough`가 만들어진 경우 해당 readable side를 소비할 주체가 없어질 수 있었다.

문제 흐름은 다음과 같았다.

```text
마지막 entry가 함수
→ 내부 PassThrough 생성
→ readable side 미소비
→ buffer 증가
→ highWaterMark 도달
→ Backpressure 발생
→ pipeline 완료 불가
→ Promise가 resolve되지 않음
```

## 해결 과정

Promise 기반 `pipeline()`에서 실제 마지막 entry를 확인하고, 마지막 entry가 함수이면서 반환된 Stream이 readable인 경우 내부 proxy를 `resume()`하도록 수정했다.

배열 형태의 `pipeline()` 호출도 지원하기 위해 다음 두 경우를 모두 처리했다.

```js
pipeline(a, b);
pipeline([a, b]);
```

핵심 수정은 다음과 같다.

```js
let lastStream = streams[streams.length - 1];

if (streams.length === 1 && ArrayIsArray(streams[0])) {
lastStream = streams[0][streams[0].length - 1];
}

const stream = pl(streams, (err, value) => {
if (err) {
reject(err);
} else {
resolve(value);
}
}, { signal, end });

if (typeof lastStream === 'function' && stream.readable) {
stream.resume();
}
```

## 테스트 및 검증

`AsyncIterable`을 반환하는 함수를 마지막 entry로 사용하고 충분한 데이터를 전달해 Backpressure가 발생할 수 있는 회귀 테스트를 추가했다.

두 가지 호출 방식을 모두 검증했다.

```js
pipelinep(...streams()).then(common.mustCall());
pipelinep(streams()).then(common.mustCall());
```

로컬에서는 다음 검증을 수행했다.

```bash
./node test/parallel/test-stream-pipeline.js
make lint
```

Codecov에서도 수정된 coverable line들이 테스트에 포함되는 것을 확인했다.

## 리뷰 과정

처음에는 `ronag`에게 Approved를 받았지만 이후 `mcollina`에게 `.resume()`이 destination 없이 사용될 경우 데이터를 소비하고 버릴 수 있다는 우려와 함께 Changes Requested를 받았다.

이에 Promise 기반 `pipeline()`에서는 callback 기반 API와 달리 내부 proxy Stream이 사용자에게 노출되지 않는다는 점을 다시 분석했고, 마지막 함수가 `AsyncIterable`을 반환하는 경우의 의도된 semantics에 대해 maintainer에게 질문했다.

이후 follow-up을 통해 방향을 확인했고 `mcollina`에게 `LGTM`과 Approved를 받았다.

추가로 `jasnell`의 Approved도 받았다.

리뷰 흐름은 다음과 같았다.

```text
첫 Approved
→ Changes Requested
→ Promise variant semantics 재검토
→ maintainer에게 질문 및 follow-up
→ mcollina Approved
→ jasnell Approved
```

## 결과

PR #65559는 최종적으로 Node.js Core에 Merge되었다.

```text
Status: Merged
Merge commit: 31e841f05633513db165eeaa8f51986d6932f8ca
```

## 배운 점

이번 기여를 통해 버그 현상을 없애는 것과 기존 API semantics를 보존하면서 문제를 해결하는 것은 다를 수 있다는 점을 배웠다.

또한 Stream에서는 Producer뿐 아니라 Consumer의 존재 여부와 Backpressure까지 함께 확인해야 하며, Core 수정에서는 테스트 통과뿐 아니라 기존 동작에 미치는 영향과 reviewer가 기대하는 semantics까지 검토해야 한다는 점을 경험했다.
Loading