From fca7d2727316a8f17b0e087eb4c4acd7367683c8 Mon Sep 17 00:00:00 2001 From: toki Date: Tue, 2 Jun 2026 17:03:49 +0900 Subject: [PATCH] feat(typescript): implement inbound gateway and queue ordering - Add inbound_gateway.ts for handling inbound message routing - Add node_inbound_gateway.ts for Node.js-specific inbound processing - Add node_inbound_gateway_worker.ts for worker-based processing - Update communicator.ts with queue ordering support - Update node.ts with inbound gateway integration - Update communicator.test.ts with new test cases - Archive completed task files for 08_typescript_gateway --- .../code_review_cloud_G08_0.log | 166 ++++++++++++ .../code_review_cloud_G08_1.log | 207 ++++++++++++++ .../08_typescript_gateway/complete.log | 43 +++ .../plan_cloud_G08_0.log} | 0 .../plan_cloud_G08_1.log | 118 ++++++++ .../CODE_REVIEW-cloud-G08.md | 73 ----- typescript/src/communicator.ts | 116 ++++++++ typescript/src/inbound_gateway.ts | 220 +++++++++++++++ typescript/src/node.ts | 1 + typescript/src/node_inbound_gateway.ts | 168 ++++++++++++ typescript/src/node_inbound_gateway_worker.ts | 30 +++ typescript/test/communicator.test.ts | 252 ++++++++++++++++++ 12 files changed, 1321 insertions(+), 73 deletions(-) create mode 100644 agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_0.log create mode 100644 agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_1.log create mode 100644 agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/complete.log rename agent-task/{m-inbound-queue-ordering/08_typescript_gateway/PLAN-cloud-G08.md => archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_0.log} (100%) create mode 100644 agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_1.log delete mode 100644 agent-task/m-inbound-queue-ordering/08_typescript_gateway/CODE_REVIEW-cloud-G08.md create mode 100644 typescript/src/inbound_gateway.ts create mode 100644 typescript/src/node_inbound_gateway.ts create mode 100644 typescript/src/node_inbound_gateway_worker.ts diff --git a/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_0.log b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_0.log new file mode 100644 index 0000000..9e62fbb --- /dev/null +++ b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_0.log @@ -0,0 +1,166 @@ + + +# Code Review Reference - API + +> **[IMPLEMENTING AGENT — READ FIRST]** 구현 후 이 파일을 채우고 active 파일을 그대로 둔 채 리뷰를 요청한다. 최종화는 code-review 전용이다. + +## 개요 + +date=2026-06-02 +task=m-inbound-queue-ordering/08_typescript_gateway, plan=0, tag=API + +## Roadmap Targets + +- Milestone: `agent-roadmap/milestones/inbound-queue-ordering.md` +- Task ids: + - `typescript-gateway`: Node/browser worker gateway 후보를 적용한다. +- Completion mode: check-on-pass + +## 구현 항목별 완료 여부 + +| 항목 | 완료 여부 | +|------|---------| +| [API-1] TypeScript worker gateway | [x] | + +## 구현 체크리스트 + +- [x] TypeScript gateway input/output에 내부 `seq`를 포함하고 Node worker_threads 후보를 추가한다. +- [x] browser target은 Web Worker 또는 safe fallback을 사용하며 unsupported 환경에서 기존 coordinator path를 유지한다. +- [x] transferable `ArrayBuffer` 사용 여부를 검토하고 복사 비용을 줄인다. +- [x] worker result를 `seq` 기준으로 reorder한 뒤 기존 `enqueueInbound`/coordinator path로 전달한다. +- [x] TypeScript tests에 Node gateway on, fallback, out-of-order completion 검증을 추가한다. +- [x] 중간 검증으로 `cd typescript && npm run check && npm test -- communicator`를 실행한다. +- [x] 최종 검증으로 `bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all`을 실행한다. +- [x] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. + +## 코드리뷰 전용 체크리스트 + +- [x] verdict append, log archive, PASS complete.log/archive move를 수행한다. FAIL 판정이므로 PASS complete.log/archive move 대신 후속 active plan/review를 작성했다. + +## 계획 대비 변경 사항 + +- gateway 계약(`InboundFrame`/`DecodedFrame`/`DecodedEnvelope`/`InboundGateway`/`FrameReorderBuffer`)과 worker-pool gateway를 신규 파일 `typescript/src/inbound_gateway.ts`(browser-safe)로 분리했다. Go의 `inbound_gateway.go`, Kotlin의 동일 계약과 심볼/주석을 맞췄다. +- Node `worker_threads` 후보는 browser bundle을 깨지 않도록 별도 파일 `typescript/src/node_inbound_gateway.ts` + worker entry `typescript/src/node_inbound_gateway_worker.ts`에 두고, `node:` import가 browser-safe `index.ts`로 새지 않게 `src/node.ts`에서만 re-export 했다. plan의 "gateway 전용 파일(필요 시)" 항목을 사용한 것이며, browser-safe 표면(`index.ts`)에는 worker_threads import가 들어가지 않는다. +- plan 수정 파일 요약대로 transport(`tcp_client.ts`/`node_ws_client.ts`)는 건드리지 않았다. gateway는 opt-in이며 `Communicator.onReceivedFrame`/`enableInboundGateway`/`attachInboundGateway`로만 활성화된다(Go/Kotlin의 opt-in 패턴과 동일). 기존 transport 경로는 inline `enqueueInbound`를 유지한다. +- `npm run check` 통과를 위해 worker의 transfer 대상 buffer를 `ArrayBuffer` instance로 narrowing 했다(`ArrayBufferLike`에 `SharedArrayBuffer` 포함). + +## 주요 설계 결정 + +- **단일 reorder/serialized-sink 코어**: `WorkerGateway`는 최대 `workers`개 async slot이 동시 decode하고, `FrameReorderBuffer`가 `seq`로 재정렬, 단일 `sinkChain` Promise로 sink 호출을 직렬화한다. 그 결과 worker가 out-of-order로 완료해도 coordinator(`enqueueInbound`)에는 입력 순서로만 들어가며 stateful dispatch 소유권이 coordinator에 남는다. +- **safe fallback = inline decode**: `WorkerGateway`의 기본 `decode`는 호출 컨텍스트의 inline `decodePacketBaseEnvelope`이다. gateway 미설정 시 `onReceivedFrame`도 inline decode 후 coordinator로 전달하므로, worker 미지원(browser 포함) 환경에서 기존 coordinator path가 그대로 유지된다. browser는 Web Worker 기반 `decode`를, Node는 `worker_threads` 기반 `decode`를 주입하는 구조다. +- **Node worker entry 로딩**: 빌드 산출물(.js)에서는 worker `.js`를 직접 로드한다. worker `execArgv`가 `--import`를 허용하지 않아(Node 제약) 빌드 전(.ts, vitest/tsx) 환경에서는 worker 안에서 `tsx/esm/api`의 `register()`로 loader를 등록한 뒤 `.ts` worker를 동적 import하는 부트스트랩을 `eval`로 실행한다. tsx는 devDependency이며 이 경로는 빌드 전 환경 전용이다. +- **transferable ArrayBuffer**: worker는 decode 결과 payload(`env.data`, protobuf decode가 입력 frame buffer를 view로 재사용)를 owned일 때 transferList로 main thread에 넘겨 결과 path의 clone 복사를 없앤다. 입력 raw frame은 source view가 buffer 전체를 소유할 때만 transfer하고, pool/공유 view(Node socket Buffer 등)이면 owned buffer로 복사해 detach 부작용을 피한다. inline fallback 경로는 `Uint8Array`를 참조 전달하므로 zero-copy다. +- **frame error 경계**: gateway decode 오류는 sink가 아니라 `seq` 순서로 `setFrameErrorHandler`에 보고되어 한 프레임의 실패가 reorder window를 막지 않는다. transport가 inline path와 동일한 parse-error disconnect 정책을 적용할 수 있게 하되, gateway 자체는 transport를 종료하지 않는다(Go/Kotlin `onError` 계약과 동일). + +## 사용자 리뷰 요청 + +- 상태: 없음 +- 사유 유형: 없음 +- 결정 필요: 없음 +- 차단 근거: 없음 +- 실행한 검증/명령: 없음 +- 자동 후속 불가 이유: 없음 +- 재개 조건: 없음 + +## 리뷰어를 위한 체크포인트 + +- Node/browser worker 경계가 package/runtime import를 깨지 않는지 확인한다. +- fallback 모드가 기존 coordinator path를 유지하는지 확인한다. + +## 검증 결과 + +### API-1 중간 검증 +``` +$ cd typescript && npm run check && npm test -- communicator +> proto-socket@1.0.5 check +> tsc --noEmit +(에러 없음) + +> proto-socket@1.0.5 test +> vitest run communicator + + RUN v3.2.4 /config/workspace/proto-socket/typescript + ✓ test/communicator.test.ts (21 tests) 363ms + ✓ Communicator > worker gateway: out-of-order seq를 reorder한 뒤 입력 순서로 dispatch한다 + ✓ FrameReorderBuffer > out-of-order 결과를 contiguous prefix가 채워질 때만 seq 순서로 release한다 + ✓ Communicator inbound gateway > gateway on: worker가 out-of-order로 완료해도 입력 seq 순서로 dispatch한다 + ✓ Communicator inbound gateway > gateway 미설정 fallback: onReceivedFrame이 inline decode로 coordinator path를 유지한다 + ✓ Communicator inbound gateway > gateway decode 오류는 frame error handler로 보고되고 이후 프레임을 막지 않는다 + ✓ Communicator inbound gateway > Node worker_threads gateway: off-thread decode가 입력 순서를 보존한다 + + Test Files 1 passed (1) + Tests 21 passed (21) +``` + +### 최종 검증 +``` +$ bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all +(종료 코드 0) + +**Proto 동기화**: schema sync PASS +**동일언어**: Dart / Go / Kotlin / Python / TypeScript 모두 PASS +**언어 PASS 매트릭스**: 5 서버 × (Dart.io/Dart.web/Dart.web(WSS)/Go/Kotlin/Python/TypeScript) 전부 PASS +**크로스테스트 상세**: 일반 언어간 20방향 16/16, Dart.web/Dart.web(WSS) 10방향 2/2, FAIL lines 0 + +결과 기록 파일: agent-test/runs/20260602-065016-proto-socket-full-matrix.md +git ref: 84635c3 (working tree, 미커밋 변경 포함) +``` + +## 코드리뷰 결과 + +- 종합 판정: FAIL + +### 차원별 평가 + +| 차원 | 평가 | 근거 | +|------|------|------| +| correctness | Fail | `InboundGateway.close()` 계약은 추가 전달 중단인데 `WorkerGateway`가 close 이후 pending decode 결과를 sink/onError로 계속 방출한다. | +| completeness | Fail | gateway close/cancel 경계가 구현 체크리스트의 worker gateway 안정성 범위에서 누락되어 있다. | +| test coverage | Fail | close 이후 pending result가 dispatch되지 않는지 검증하는 테스트가 없다. | +| API contract | Fail | `InboundGateway.close()` 문서 계약과 실제 동작이 불일치한다. | +| code quality | Warn | lifecycle token/closed guard 없이 async chain이 계속 진행되어 reattach/reinitialize 시 stale delivery 위험이 있다. | +| plan deviation | Pass | 계획 범위 자체는 TypeScript gateway 구현과 테스트에 집중되어 있다. | +| verification trust | Pass | `cd typescript && npm run check && npm test -- communicator`는 재실행해 21/21 PASS를 확인했다. 단, 통과 테스트가 아래 lifecycle 결함을 커버하지 못한다. | + +### 발견된 문제 + +- Required: `typescript/src/inbound_gateway.ts:147`의 `close()`는 `closed = true`만 설정하고 queued/in-flight decode, reorder buffer, `sinkChain` 전달을 중단하지 않는다. 실제로 pending decode를 close 이후 resolve하면 `typescript/src/inbound_gateway.ts:172`-`190`의 collect/sink 경로가 계속 실행된다. 이 때문에 `typescript/src/communicator.ts:198`-`208`의 `attachInboundGateway()`/`closeGateway()`가 기존 gateway를 닫고 새 gateway를 붙이거나 `initialize()`가 기존 gateway를 닫은 뒤 상태를 재설정해도, 이전 gateway의 늦은 결과가 현재 communicator 상태로 들어올 수 있다. `close()` 이후 결과와 error 전달을 모두 무시하도록 세대 토큰 또는 closed guard를 collect/sink chain에 적용하고, queued jobs/reorder pending도 정리한 뒤 close-after-pending 테스트를 추가해야 한다. + +### 리뷰 검증 + +```bash +$ cd typescript && npm run check && npm test -- communicator +> proto-socket@1.0.5 check +> tsc --noEmit + +> proto-socket@1.0.5 test +> vitest run communicator + + RUN v3.2.4 /config/workspace/proto-socket/typescript + ✓ test/communicator.test.ts (21 tests) 395ms + Test Files 1 passed (1) + Tests 21 passed (21) +``` + +```bash +$ cd typescript && node --import tsx --input-type=module <<'EOF' +import { WorkerGateway } from './src/inbound_gateway.ts'; +let resolveDecode; +const delivered = []; +const gw = new WorkerGateway({ + workers: 1, + decode: () => new Promise((resolve) => { resolveDecode = resolve; }), + sink: (frame) => { delivered.push(frame.seq); }, +}); +gw.submit({ seq: 1, bytes: new Uint8Array([1]) }); +gw.close(); +resolveDecode({ typeName: 't', data: new Uint8Array(), incomingNonce: 0, responseNonce: 0 }); +await new Promise((r) => setTimeout(r, 0)); +console.log(JSON.stringify(delivered)); +EOF +[1] +``` + +### 다음 단계 + +- WARN/FAIL 후속: active 파일을 archive한 뒤 `close()` 이후 stale sink/onError 전달을 차단하고 테스트를 추가하는 후속 `PLAN-cloud-G08.md` / `CODE_REVIEW-cloud-G08.md`를 작성한다. diff --git a/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_1.log b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_1.log new file mode 100644 index 0000000..1c14616 --- /dev/null +++ b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_1.log @@ -0,0 +1,207 @@ + + +# Code Review Reference - REVIEW_API + +> **[IMPLEMENTING AGENT — READ FIRST] Filling in this file is the mandatory final step of implementation.** +> The task is NOT complete until every implementation-owned section below is filled in. +> Complete the `구현 체크리스트`; the final checklist item is mandatory before saving. +> Fill implementation-owned sections, then stop with active files in place and report ready for review. +> If implementation is blocked by a user-only decision, user-owned external environment prerequisite, or scope conflict, fill `사용자 리뷰 요청` with evidence and stop with active files in place; code-review decides whether to write `USER_REVIEW.md`. Evidence gaps that a follow-up agent can close by rerunning commands or collecting artifacts are normal follow-up issues, not user-review blockers by themselves. +> Finalization (`코드리뷰 결과`, log rename, `complete.log`, archive moves, `코드리뷰 전용 체크리스트`) is review-agent-only, even after compaction/resume. + +## 개요 + +date=2026-06-02 +task=m-inbound-queue-ordering/08_typescript_gateway, plan=1, tag=REVIEW_API + +## Roadmap Targets + +- Milestone: `agent-roadmap/milestones/inbound-queue-ordering.md` +- Task ids: + - `typescript-gateway`: Node/browser worker gateway 후보를 적용한다. +- Completion mode: check-on-pass + +## 이 파일을 읽는 리뷰 에이전트에게 + +> **[REVIEW AGENT ONLY]** 아래 종결 절차는 코드리뷰 에이전트 전용이다. 구현 에이전트는 이 섹션을 실행하지 않는다. + +각 항목의 구현을 실제 소스 파일과 대조하고, `검증 결과` 섹션의 출력이 코드와 일치하는지 확인하세요. + +--- + +## 구현 항목별 완료 여부 + +| 항목 | 완료 여부 | +|------|---------| +| [REVIEW_API-1] WorkerGateway close lifecycle | [x] | + +## 구현 체크리스트 + +- [x] `WorkerGateway.close()`가 close 이후 queued/in-flight decode 결과와 error를 sink/onError로 전달하지 않도록 lifecycle guard를 추가한다. +- [x] close 시 queued jobs와 reorder pending 상태를 정리해 stale result가 contiguous release를 유발하지 않게 한다. +- [x] `Communicator.attachInboundGateway()` 또는 재초기화 경계에서 이전 gateway의 늦은 결과가 현재 communicator로 들어오지 않는 테스트를 추가한다. +- [x] TypeScript tests에 `WorkerGateway.close()` 이후 pending decode result/error가 전달되지 않는 검증을 추가한다. +- [x] 중간 검증으로 `cd typescript && npm run check && npm test -- communicator`를 실행한다. +- [x] 최종 검증으로 `bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all`을 실행한다. +- [x] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. + +## 코드리뷰 전용 체크리스트 + +> **[REVIEW AGENT ONLY]** 이 체크리스트는 코드리뷰 에이전트만 사용한다. +> 구현 에이전트는 이 섹션을 수정하거나 체크하지 않는다. + +- [x] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다. +- [x] 판정과 `차원별 평가`, Required/Suggested/Nit 분류가 서로 일치한다. +- [x] active `CODE_REVIEW-*-G??.md`를 `code_review_{review_lane}_GNN_N.log`로 아카이브한다. +- [x] active `PLAN-*-G??.md`를 `plan_{build_lane}_GNN_M.log`로 아카이브한다. +- [x] `.gitignore`의 Agent-Ops 관리 block이 `agent-task/**/*.md`와 `agent-task/**/*.log`를 unignore하여 plan/review/archive 산출물이 추적 가능한지 확인한다. +- [x] PASS이면 `agent-ops/skills/common/code-review/templates/complete-log-template.md` 기준으로 `complete.log`를 작성하고 active `.md` 파일을 남기지 않는다. +- [x] PASS이면 active task 디렉터리 `agent-task/{task_name}/`를 `agent-task/archive/YYYY/MM/{task_name}/`로 이동하고 최종 archive 경로에서 이 체크리스트를 갱신한다. +- [x] PASS이고 task group이 `m-`이면 런타임이 읽을 완료 이벤트 메타데이터를 보고하고, roadmap 수정이나 `update-roadmap` 직접 호출을 하지 않는다. +- [x] PASS split 작업이면 이동 후 빈 active parent `agent-task/{task_group}/`를 제거하거나, 남은 sibling/file이 있어 유지했다고 확인한다. +- [ ] WARN/FAIL이고 user-review gate가 트리거되지 않았으면 다음 active `PLAN-{build_lane}-GNN.md`와 `CODE_REVIEW-{review_lane}-GNN.md`를 작성하고 `complete.log`를 작성하지 않는다. +- [ ] USER_REVIEW이면 `agent-ops/skills/common/code-review/templates/user-review-template.md` 기준으로 `USER_REVIEW.md`를 작성하고 active `PLAN-*.md`, `CODE_REVIEW-*.md`, `complete.log`를 남기지 않는다. +- [ ] USER_REVIEW가 사용자 결정으로 완료/PASS 해소되면 `USER_REVIEW.md`를 해소 상태로 갱신하고 `complete.log`를 작성한 뒤 task directory를 archive로 이동한다. + +## 계획 대비 변경 사항 + +- `code_review_cloud_G08_0.log`의 Required issue(`InboundGateway.close()` 계약 위반)만 해소했다. Node worker pool 구조, public export, transport 경로, wire schema는 변경하지 않았다. +- 수정 파일: `typescript/src/inbound_gateway.ts`(close lifecycle), `typescript/test/communicator.test.ts`(lifecycle 테스트 3건 추가). 계획의 `수정 파일 요약`과 일치한다. +- `FrameReorderBuffer`에 `clear()`를 추가해 close 시 reorder pending을 비울 수 있게 했다(계획의 "reorder pending 상태 정리" 항목 충족). + +## 주요 설계 결정 + +- **per-instance closed guard로 stale delivery 차단**: `WorkerGateway.close()`는 `closed = true` 설정에 더해 `jobs.length = 0`과 `reorder.clear()`로 queued/buffered 상태를 비운다. in-flight decode는 `workerLoop()`이 `await this.decode()` 직후 `this.closed`를 확인해 close된 경우 collect 없이 종료하고, `collect()`와 `sinkChain` 콜백 진입부에서도 `closed`를 재확인한다. 따라서 close 이후 완료/실패하는 decode는 sink·onError로 전달되지 않는다. +- **세대 토큰이 아닌 인스턴스 guard 선택 근거**: communicator가 `attachInboundGateway()`/`initialize()`로 gateway를 교체할 때 이전 gateway 인스턴스는 `closeGateway()`로 close된다. sink는 해당 인스턴스 생성 시점에 바인딩된 클로저이므로, 인스턴스별 `closed` guard만으로 이전 gateway의 늦은 결과가 현재 communicator(`enqueueInbound`)로 유입되는 것을 막는다. 별도 cross-instance 세대 토큰은 불필요하다. +- **close 이후 in-flight sink의 잔여 frame 중단**: `sinkChain` 콜백은 각 ready frame 전달 전에 `closed`를 확인해, 스케줄 이후 close되면 남은 batch 전달을 멈춘다. 이미 await 중인 단일 sink 호출은 완주하되 그 이후 frame은 전달되지 않는다. +- 기존 reorder/fallback/Node worker/decode-error 동작과 그 테스트는 모두 그대로 유지된다(회귀 없음). + +## 사용자 리뷰 요청 + +_기본값은 `없음`이다. 구현 중 사용자 결정, 사용자 소유 외부 환경/secret/서비스 준비, 또는 계획 범위 변경 없이는 안전하게 진행할 수 없으면 아래 항목을 실제 내용으로 교체하고, 구현을 중단한 뒤 active 파일을 그대로 둔 채 리뷰를 요청한다. 후속 에이전트가 명령 재실행이나 산출물 수집으로 해소할 수 있는 검증 증거 공백만으로는 사용자 리뷰 요청을 작성하지 않는다._ + +- 상태: 없음 +- 사유 유형: 없음 +- 결정 필요: 없음 +- 차단 근거: 없음 +- 실행한 검증/명령: 없음 +- 자동 후속 불가 이유: 없음 +- 재개 조건: 없음 + +## 리뷰어를 위한 체크포인트 + +- `WorkerGateway.close()` 이후 pending result/error가 sink/onError로 전달되지 않는지 확인한다. +- reattach/reinitialize 경계에서 이전 gateway의 늦은 결과가 현재 communicator 상태로 들어오지 않는지 확인한다. +- 기존 Node worker gateway, fallback, out-of-order reorder 테스트가 유지되는지 확인한다. + +## 검증 결과 + +### Required issue 재현 → 해소 확인 +``` +$ cd typescript && node --import tsx --input-type=module <<'EOF' +import { WorkerGateway } from './src/inbound_gateway.ts'; +let resolveDecode; +const delivered = []; +const gw = new WorkerGateway({ + workers: 1, + decode: () => new Promise((resolve) => { resolveDecode = resolve; }), + sink: (frame) => { delivered.push(frame.seq); }, +}); +gw.submit({ seq: 1, bytes: new Uint8Array([1]) }); +gw.close(); +resolveDecode({ typeName: 't', data: new Uint8Array(), incomingNonce: 0, responseNonce: 0 }); +await new Promise((r) => setTimeout(r, 0)); +console.log("delivered:", JSON.stringify(delivered)); +EOF +delivered: [] # 수정 전 [1] → 수정 후 [] (기대값) +``` + +### REVIEW_API-1 중간 검증 +``` +$ cd typescript && npm run check && npm test -- communicator +> proto-socket@1.0.5 check +> tsc --noEmit +(에러 없음) + +> proto-socket@1.0.5 test +> vitest run communicator + + RUN v3.2.4 /config/workspace/proto-socket/typescript + ✓ test/communicator.test.ts (24 tests) 409ms + ✓ WorkerGateway lifecycle > close() 이후 pending decode result는 sink로 전달되지 않는다 + ✓ WorkerGateway lifecycle > close() 이후 pending decode reject는 onError로 전달되지 않는다 + ✓ Communicator inbound gateway > attachInboundGateway 경계: 닫힌 이전 gateway의 늦은 결과가 현재 communicator로 들어오지 않는다 + + Test Files 1 passed (1) + Tests 24 passed (24) +``` +(기존 21건 + lifecycle 3건 = 24건 PASS) + +### 최종 검증 +``` +$ bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all +(종료 코드 0) + +**Proto 동기화**: schema sync PASS +**동일언어**: Dart / Go / Kotlin / Python / TypeScript 모두 PASS +**언어 PASS 매트릭스**: 5 서버 × (Dart.io/Dart.web/Dart.web(WSS)/Go/Kotlin/Python/TypeScript) 전부 PASS +**크로스테스트 상세**: 일반 언어간 20방향 16/16, Dart.web/Dart.web(WSS) 10방향 2/2, FAIL lines 0 + +결과 기록 파일: agent-test/runs/20260602-073737-proto-socket-full-matrix.md +``` + +## 코드리뷰 결과 + +- 종합 판정: PASS + +### 차원별 평가 + +| 차원 | 평가 | 근거 | +|------|------|------| +| correctness | Pass | close 이후 pending decode result/error가 sink/onError로 전달되지 않도록 `WorkerGateway` guard가 적용됐다. | +| completeness | Pass | 계획의 close lifecycle 항목, queued/reorder 정리, attach 경계 테스트, lifecycle 테스트가 모두 충족됐다. | +| test coverage | Pass | WorkerGateway close result/reject, Communicator attach stale result, 기존 gateway reorder/fallback/Node worker 테스트가 포함됐다. | +| API contract | Pass | `InboundGateway.close()`의 "추가 전달 중단" 계약과 구현이 일치한다. | +| code quality | Pass | close guard가 submit, decode await 이후, collect, sinkChain 내부 전달 직전에 적용되어 async 경계가 명확하다. | +| plan deviation | Pass | 후속 계획 범위인 lifecycle 수정과 테스트만 수행했다. | +| verification trust | Pass | 재현 스크립트, TypeScript check/test, 전체 로컬 매트릭스를 리뷰에서 직접 재실행해 PASS를 확인했다. | + +### 발견된 문제 + +- 없음 + +### 리뷰 검증 + +```bash +$ cd typescript && node --import tsx --input-type=module +delivered: [] +``` + +```bash +$ cd typescript && npm run check && npm test -- communicator +> proto-socket@1.0.5 check +> tsc --noEmit + +> proto-socket@1.0.5 test +> vitest run communicator + + RUN v3.2.4 /config/workspace/proto-socket/typescript + ✓ test/communicator.test.ts (24 tests) 417ms + Test Files 1 passed (1) + Tests 24 passed (24) +``` + +```bash +$ bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all +(종료 코드 0) + +Proto 동기화: schema sync PASS +동일언어: Dart / Go / Kotlin / Python / TypeScript 모두 PASS +언어 PASS 매트릭스: 5 서버 x 7 클라이언트 전부 PASS +크로스테스트 상세: 일반 언어간 20방향 16/16, Dart.web/Dart.web(WSS) 10방향 2/2, FAIL lines 0 +결과 기록 파일: agent-test/runs/20260602-074418-proto-socket-full-matrix.md +``` + +### 다음 단계 + +- PASS: active 파일을 archive하고 `complete.log`를 작성한 뒤 task directory를 `agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/`로 이동한다. Runtime은 Milestone PASS completion metadata를 사용해 roadmap 갱신 여부를 판단한다. diff --git a/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/complete.log b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/complete.log new file mode 100644 index 0000000..7379916 --- /dev/null +++ b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/complete.log @@ -0,0 +1,43 @@ +# Complete - m-inbound-queue-ordering/08_typescript_gateway + +## 완료 일시 + +2026-06-02 + +## 요약 + +TypeScript worker gateway 구현은 2회 리뷰 루프에서 close lifecycle Required issue를 보완한 뒤 최종 PASS로 완료됐다. + +## 루프 이력 + +| Plan | Review | Verdict | 메모 | +|------|--------|---------|------| +| `plan_cloud_G08_0.log` | `code_review_cloud_G08_0.log` | FAIL | `WorkerGateway.close()` 이후 pending decode 결과가 sink/onError로 전달되는 Required issue 발견. | +| `plan_cloud_G08_1.log` | `code_review_cloud_G08_1.log` | PASS | close guard, queued/reorder 정리, stale delivery 테스트 추가 후 PASS. | + +## 구현/정리 내용 + +- TypeScript inbound gateway 계약, WorkerGateway reorder/sink 직렬화, Node worker_threads decode gateway, browser-safe fallback surface를 추가했다. +- `WorkerGateway.close()`가 queued/in-flight result와 error를 close 이후 전달하지 않도록 lifecycle guard와 reorder cleanup을 적용했다. +- TypeScript communicator/gateway tests에 Node worker gateway, fallback, out-of-order completion, close stale result/error, attach 경계 검증을 추가했다. + +## 최종 검증 + +- `cd typescript && node --import tsx --input-type=module` - PASS; close 이후 pending decode result가 `delivered: []`로 확인됨. +- `cd typescript && npm run check && npm test -- communicator` - PASS; `test/communicator.test.ts` 24 tests PASS. +- `bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all` - PASS; 결과 기록 파일 `agent-test/runs/20260602-074418-proto-socket-full-matrix.md`. + +## Roadmap Completion + +- Milestone: `agent-roadmap/milestones/inbound-queue-ordering.md` +- Completed task ids: + - `typescript-gateway`: PASS; evidence=`agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_1.log`, `agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_1.log`; verification=`agent-test/runs/20260602-074418-proto-socket-full-matrix.md` +- Not completed task ids: 없음 + +## 잔여 Nit + +- 없음 + +## 후속 작업 + +- 없음 diff --git a/agent-task/m-inbound-queue-ordering/08_typescript_gateway/PLAN-cloud-G08.md b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_0.log similarity index 100% rename from agent-task/m-inbound-queue-ordering/08_typescript_gateway/PLAN-cloud-G08.md rename to agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_0.log diff --git a/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_1.log b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_1.log new file mode 100644 index 0000000..6d271b9 --- /dev/null +++ b/agent-task/archive/2026/06/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_1.log @@ -0,0 +1,118 @@ + + +# TypeScript Worker Gateway Follow-up 계획 + +## 이 파일을 읽는 구현 에이전트에게 + +이 계획은 `code_review_cloud_G08_0.log`의 Required issue만 해소한다. 구현 후 `CODE_REVIEW-cloud-G08.md`의 구현 에이전트 소유 섹션에 실제 변경 내용과 검증 stdout/stderr를 채운 뒤 active 파일을 그대로 두고 리뷰를 요청한다. 최종화는 code-review 전용이다. + +## 배경 + +초기 TypeScript gateway 구현은 Node worker/fallback/reorder 테스트를 통과했지만, `InboundGateway.close()`가 "추가 전달 중단" 계약을 지키지 않는다. pending decode가 close 이후 완료되면 stale result가 sink 또는 onError로 전달될 수 있고, `Communicator.attachInboundGateway()`/`initialize()`/`close()` 경계에서 이전 gateway 결과가 현재 communicator 상태에 섞일 수 있다. + +## 사용자 리뷰 요청 흐름 + +구현 중 user-only blocker가 있으면 active review stub의 `사용자 리뷰 요청` 섹션을 채운다. + +## Roadmap Targets + +- Milestone: `agent-roadmap/milestones/inbound-queue-ordering.md` +- Task ids: + - `typescript-gateway`: Node/browser worker gateway 후보를 적용한다. +- Completion mode: check-on-pass + +## 분석 결과 + +### 읽은 파일 + +- `agent-task/m-inbound-queue-ordering/08_typescript_gateway/code_review_cloud_G08_0.log` +- `agent-task/m-inbound-queue-ordering/08_typescript_gateway/plan_cloud_G08_0.log` +- `typescript/src/inbound_gateway.ts` +- `typescript/src/communicator.ts` +- `typescript/test/communicator.test.ts` + +### 실패 재현 + +```bash +cd typescript && node --import tsx --input-type=module <<'EOF' +import { WorkerGateway } from './src/inbound_gateway.ts'; +let resolveDecode; +const delivered = []; +const gw = new WorkerGateway({ + workers: 1, + decode: () => new Promise((resolve) => { resolveDecode = resolve; }), + sink: (frame) => { delivered.push(frame.seq); }, +}); +gw.submit({ seq: 1, bytes: new Uint8Array([1]) }); +gw.close(); +resolveDecode({ typeName: 't', data: new Uint8Array(), incomingNonce: 0, responseNonce: 0 }); +await new Promise((r) => setTimeout(r, 0)); +console.log(JSON.stringify(delivered)); +EOF +``` + +현재 출력은 `[1]`이며, 기대 출력은 close 이후 전달 없음(`[]`)이다. + +### 범위 결정 근거 + +- `WorkerGateway` lifecycle과 관련 테스트만 수정한다. +- Node worker pool 구조, public export, transport 경로, wire schema는 변경하지 않는다. +- `Roadmap Targets`는 동일 task completion을 계속 향하므로 유지한다. + +### 빌드 등급 + +- build/review: `cloud-G08`. 이유: async worker lifecycle, close/cancel, stale delivery 차단은 gateway API contract와 concurrency 경계에 직접 영향을 준다. + +## 구현 체크리스트 + +- [ ] `WorkerGateway.close()`가 close 이후 queued/in-flight decode 결과와 error를 sink/onError로 전달하지 않도록 lifecycle guard를 추가한다. +- [ ] close 시 queued jobs와 reorder pending 상태를 정리해 stale result가 contiguous release를 유발하지 않게 한다. +- [ ] `Communicator.attachInboundGateway()` 또는 재초기화 경계에서 이전 gateway의 늦은 결과가 현재 communicator로 들어오지 않는 테스트를 추가한다. +- [ ] TypeScript tests에 `WorkerGateway.close()` 이후 pending decode result/error가 전달되지 않는 검증을 추가한다. +- [ ] 중간 검증으로 `cd typescript && npm run check && npm test -- communicator`를 실행한다. +- [ ] 최종 검증으로 `bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all`을 실행한다. +- [ ] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. + +### [REVIEW_API-1] WorkerGateway close lifecycle + +#### 문제 + +- `WorkerGateway.close()`가 `closed = true`만 설정해 pending async decode 완료 후 `collect()`/`sinkChain` 경로가 계속 실행된다. + +#### 해결 방법 + +- close generation 또는 closed guard를 `workerLoop()` 완료 후, `collect()`, `sinkChain` 내부 전달 직전에 적용한다. +- close 시 queued jobs와 reorder pending을 비우거나 교체해 이전 세대 결과가 release되지 않게 한다. +- `onError`도 close 이후에는 전달하지 않는다. + +#### 수정 파일 및 체크리스트 + +- [ ] `typescript/src/inbound_gateway.ts` +- [ ] `typescript/test/communicator.test.ts` + +#### 테스트 작성 + +- pending decode Promise를 close 이후 resolve했을 때 sink가 호출되지 않는지 검증한다. +- pending decode Promise를 close 이후 reject했을 때 onError가 호출되지 않는지 검증한다. +- 기존 communicator gateway reorder/fallback/Node worker tests가 계속 통과하는지 확인한다. + +#### 중간 검증 + +```bash +cd typescript && npm run check && npm test -- communicator +``` + +## 수정 파일 요약 + +| 파일 | 항목 | +|------|------| +| `typescript/src/inbound_gateway.ts` | REVIEW_API-1 | +| `typescript/test/communicator.test.ts` | REVIEW_API-1 | + +## 최종 검증 + +```bash +bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all +``` + +모든 코드 변경 완료 후 반드시 `CODE_REVIEW-*-G??.md`의 구현 에이전트 소유 섹션을 채운다. 이 파일 작성이 구현의 마지막 단계다. diff --git a/agent-task/m-inbound-queue-ordering/08_typescript_gateway/CODE_REVIEW-cloud-G08.md b/agent-task/m-inbound-queue-ordering/08_typescript_gateway/CODE_REVIEW-cloud-G08.md deleted file mode 100644 index b20d81c..0000000 --- a/agent-task/m-inbound-queue-ordering/08_typescript_gateway/CODE_REVIEW-cloud-G08.md +++ /dev/null @@ -1,73 +0,0 @@ - - -# Code Review Reference - API - -> **[IMPLEMENTING AGENT — READ FIRST]** 구현 후 이 파일을 채우고 active 파일을 그대로 둔 채 리뷰를 요청한다. 최종화는 code-review 전용이다. - -## 개요 - -date=2026-06-02 -task=m-inbound-queue-ordering/08_typescript_gateway, plan=0, tag=API - -## Roadmap Targets - -- Milestone: `agent-roadmap/milestones/inbound-queue-ordering.md` -- Task ids: - - `typescript-gateway`: Node/browser worker gateway 후보를 적용한다. -- Completion mode: check-on-pass - -## 구현 항목별 완료 여부 - -| 항목 | 완료 여부 | -|------|---------| -| [API-1] TypeScript worker gateway | [ ] | - -## 구현 체크리스트 - -- [ ] TypeScript gateway input/output에 내부 `seq`를 포함하고 Node worker_threads 후보를 추가한다. -- [ ] browser target은 Web Worker 또는 safe fallback을 사용하며 unsupported 환경에서 기존 coordinator path를 유지한다. -- [ ] transferable `ArrayBuffer` 사용 여부를 검토하고 복사 비용을 줄인다. -- [ ] worker result를 `seq` 기준으로 reorder한 뒤 기존 `enqueueInbound`/coordinator path로 전달한다. -- [ ] TypeScript tests에 Node gateway on, fallback, out-of-order completion 검증을 추가한다. -- [ ] 중간 검증으로 `cd typescript && npm run check && npm test -- communicator`를 실행한다. -- [ ] 최종 검증으로 `bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all`을 실행한다. -- [ ] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. - -## 코드리뷰 전용 체크리스트 - -- [ ] verdict append, log archive, PASS complete.log/archive move를 수행한다. - -## 계획 대비 변경 사항 - -_구현 에이전트가 기록한다._ - -## 주요 설계 결정 - -_구현 에이전트가 기록한다._ - -## 사용자 리뷰 요청 - -- 상태: 없음 -- 사유 유형: 없음 -- 결정 필요: 없음 -- 차단 근거: 없음 -- 실행한 검증/명령: 없음 -- 자동 후속 불가 이유: 없음 -- 재개 조건: 없음 - -## 리뷰어를 위한 체크포인트 - -- Node/browser worker 경계가 package/runtime import를 깨지 않는지 확인한다. -- fallback 모드가 기존 coordinator path를 유지하는지 확인한다. - -## 검증 결과 - -### API-1 중간 검증 -``` -$ cd typescript && npm run check && npm test -- communicator -``` - -### 최종 검증 -``` -$ bash agent-ops/skills/project/run-proto-socket-test-matrix/scripts/run_matrix.sh --all -``` diff --git a/typescript/src/communicator.ts b/typescript/src/communicator.ts index b7c7c48..8a2335b 100644 --- a/typescript/src/communicator.ts +++ b/typescript/src/communicator.ts @@ -8,8 +8,25 @@ import { toBinary, type PacketBase, } from "./packets/message_common_pb.js"; +import { + decodePacketBaseEnvelope, + type FrameDecoder, + type InboundGateway, + WorkerGateway, +} from "./inbound_gateway.js"; export type { Message, MessageType } from "./packets/message_common_pb.js"; +export { + decodePacketBaseEnvelope, + type DecodedEnvelope, + type DecodedFrame, + type FrameDecoder, + FrameReorderBuffer, + type InboundFrame, + type InboundGateway, + WorkerGateway, + type WorkerGatewayOptions, +} from "./inbound_gateway.js"; export class NotConnectedError extends Error { constructor() { @@ -82,11 +99,17 @@ export class Communicator { private closeResolver: (() => void) | null = null; private transport: Transport | null = null; private writeErrorHandler: ((err: Error) => void) | null = null; + private frameErrorHandler: ((err: Error) => void) | null = null; + private gateway: InboundGateway | null = null; + private frameSeq = 0; private closedPromise: Promise = Promise.resolve(); private resolveClosed: (() => void) | null = null; initialize(transport: Transport, parserMap: ParserMap): void { this.transport = transport; + this.closeGateway(); + this.frameSeq = 0; + this.frameErrorHandler = null; this.parserMap = new Map(); this.schemaMap = new Map(); this.handlers = new Map(); @@ -127,6 +150,97 @@ export class Communicator { this.writeErrorHandler = fn; } + /** + * inbound gateway가 raw frame decode에 실패했을 때 호출되는 handler를 등록한다. transport는 이를 통해 + * gateway path에서도 inline path와 동일한 parse-error disconnect semantics를 적용할 수 있다(write path의 + * {@link setWriteErrorHandler}와 대응). gateway가 없으면 {@link onReceivedFrame}의 inline decode 오류도 + * 이 handler로 보고된다. + */ + setFrameErrorHandler(fn: (err: Error) => void): void { + this.frameErrorHandler = fn; + } + + /** + * receive coordinator 앞단에 default worker-pool gateway를 attach한다. {@link onReceivedFrame}으로 제출된 + * raw frame은 pool이 decode하고 내부 `seq`로 reorder해 coordinator의 {@link enqueueInbound}로 전달하므로 + * FIFO dispatch와 stateful 처리 소유권은 coordinator에 유지된다. + * + * opt-in이며 attach하지 않으면 {@link onReceivedFrame}은 호출 컨텍스트에서 inline decode한다. Node는 + * `worker_threads` 기반 `decode`를, browser는 Web Worker 기반 `decode`를 주입할 수 있고, 미주입 시 inline + * decode가 safe fallback이 된다. + */ + enableInboundGateway( + options: { workers?: number; decode?: FrameDecoder } = {}, + ): void { + this.attachInboundGateway( + new WorkerGateway({ + workers: options.workers, + decode: options.decode, + sink: (frame) => + this.enqueueInbound( + frame.typeName, + frame.data, + frame.incomingNonce, + frame.responseNonce, + ), + // gateway decode 오류를 frame error handler로 노출해 inline path와 동일한 parse-error + // disconnect semantics를 적용한다. gateway 자체는 transport를 종료하지 않는다. + onError: (err) => this.frameErrorHandler?.(err), + }), + ); + } + + /** + * custom {@link InboundGateway}를 receive coordinator 앞단에 설치한다. gateway sink는 coordinator 소유권을 + * 보존하기 위해 decoded frame을 {@link enqueueInbound}로 전달해야 한다. default worker-pool gateway는 + * {@link enableInboundGateway}를 사용한다. + */ + attachInboundGateway(gateway: InboundGateway): void { + this.closeGateway(); + this.gateway = gateway; + this.frameSeq = 0; + } + + private closeGateway(): void { + if (this.gateway !== null) { + this.gateway.close(); + this.gateway = null; + } + } + + /** + * transport read loop가 raw `PacketBase` frame을 ingest한다. + * + * gateway가 attach되어 있으면 frame에 내부 `seq`를 부여해 off-coordinator decode + reorder로 제출한다. + * decode에 실패한 frame은 `seq` 순서에 맞춰 {@link setFrameErrorHandler}에 등록된 handler로 보고된다. + * + * gateway가 없으면 호출 컨텍스트에서 inline decode 후 {@link enqueueInbound}로 전달한다(safe fallback). + */ + onReceivedFrame(raw: Uint8Array): void { + if (!this.isAliveFlag || this.isClosing) { + return; + } + const gateway = this.gateway; + if (gateway !== null) { + this.frameSeq += 1; + gateway.submit({ seq: this.frameSeq, bytes: raw }); + return; + } + let envelope; + try { + envelope = decodePacketBaseEnvelope(raw); + } catch (err) { + this.frameErrorHandler?.(toError(err)); + return; + } + void this.enqueueInbound( + envelope.typeName, + envelope.data, + envelope.incomingNonce, + envelope.responseNonce, + ); + } + private nextNonce(): number { if (this.nonce >= Communicator.MAX_NONCE) { this.nonce = 0; @@ -143,6 +257,7 @@ export class Communicator { this.isAliveFlag = false; this.isClosing = false; + this.closeGateway(); const pendingRequests = Array.from(this.pendingRequests.values()); this.pendingRequests.clear(); for (const pending of pendingRequests) { @@ -178,6 +293,7 @@ export class Communicator { return; } this.isClosing = true; + this.closeGateway(); if (this.inboundQueue.length > 0 || this.inFlightCount > 0) { await new Promise((resolve) => { this.closeResolver = resolve; diff --git a/typescript/src/inbound_gateway.ts b/typescript/src/inbound_gateway.ts new file mode 100644 index 0000000..a354678 --- /dev/null +++ b/typescript/src/inbound_gateway.ts @@ -0,0 +1,220 @@ +import { fromBinary, PacketBaseSchema } from "./packets/message_common_pb.js"; + +/** + * 수신 측이 gateway 제출 전에 부여하는 내부 raw 프레임이다. + * + * `seq`는 단조 증가하는 내부 시퀀스로, 병렬 decode 후 입력 순서를 복원하는 용도로만 쓰인다. + * wire format이나 공개 프로토콜 필드에는 절대 노출하지 않는다. + */ +export interface InboundFrame { + seq: number; + bytes: Uint8Array; +} + +/** gateway decode 결과의 envelope 부분이다(순서 복원용 `seq` 제외). */ +export interface DecodedEnvelope { + typeName: string; + data: Uint8Array; + incomingNonce: number; + responseNonce: number; +} + +/** + * gateway가 PacketBase envelope을 decode한 결과다. `Communicator`의 inbound dispatch가 기대하는 + * envelope 필드와, 입력 순서 복원을 위한 원본 `seq`를 함께 싣는다. + */ +export interface DecodedFrame extends DecodedEnvelope { + seq: number; +} + +/** + * 수신 coordinator 앞단에서 순수 raw-bytes 작업(PacketBase envelope decode/경량 전처리)만 병렬 + * 처리하는 worker gateway 계약이다. + * + * gateway는 stateful dispatch, `pendingRequests`, listener, write queue를 소유하지 않는다. decode + * 결과는 입력 `seq` 순서로 sink에 전달되어 coordinator가 결정적으로 dispatch할 수 있게 한다. + */ +export interface InboundGateway { + /** raw 프레임을 gateway에 제출한다. fire-and-forget이며 decode 결과는 이후 sink에 `seq` 순서로 전달된다. */ + submit(frame: InboundFrame): void; + /** gateway 자원을 해제하고 추가 전달을 멈춘다. */ + close(): void; +} + +/** raw 프레임을 PacketBase envelope으로 decode하는 순수 함수다. 공유 상태가 없어 병렬 실행에 안전하다. */ +export type FrameDecoder = (bytes: Uint8Array) => DecodedEnvelope | Promise; + +/** raw bytes를 PacketBase envelope으로 inline decode한다. gateway 미설정 시의 safe fallback 경로다. */ +export function decodePacketBaseEnvelope(bytes: Uint8Array): DecodedEnvelope { + const base = fromBinary(PacketBaseSchema, bytes); + return { + typeName: base.typeName, + data: base.data, + incomingNonce: base.nonce, + responseNonce: base.responseNonce, + }; +} + +interface GatewayResult { + seq: number; + frame: DecodedFrame | null; + error: Error | null; +} + +/** + * out-of-order로 도착할 수 있는 결과에서 입력 순서를 복원한다. `seq`는 1부터 연속이라고 가정한다. + * 결과는 앞선 모든 `seq`가 release될 때까지 버퍼되며, 이후 자신과 연속된 후속 결과까지 순서대로 반환한다. + * concurrent 사용에 안전하지 않으며 collector 단독에서만 호출한다. + */ +export class FrameReorderBuffer { + private nextSeq = 1; + private readonly pending = new Map(); + + /** `res`를 버퍼하고 지금 순서대로 release 가능한 결과들을 `seq` 순으로 반환한다. */ + release(res: GatewayResult): GatewayResult[] { + if (res.seq < this.nextSeq) { + // 이미 release되었거나 중복/지연 도착이므로 두 번 방출하지 않도록 무시한다. + return []; + } + this.pending.set(res.seq, res); + const ready: GatewayResult[] = []; + for (;;) { + const r = this.pending.get(this.nextSeq); + if (r === undefined) { + break; + } + ready.push(r); + this.pending.delete(this.nextSeq); + this.nextSeq += 1; + } + return ready; + } + + /** 버퍼된 결과를 비운다. gateway close 시 stale result가 contiguous release를 유발하지 않게 한다. */ + clear(): void { + this.pending.clear(); + } +} + +export interface WorkerGatewayOptions { + /** decode 결과를 입력 `seq` 순서로 받는 coordinator sink다. 반환 Promise는 직렬화되어 backpressure를 전달한다. */ + sink: (frame: DecodedFrame) => void | Promise; + /** raw bytes decode 함수. 기본은 inline PacketBase decode이며, Web Worker 등 off-thread decode를 주입할 수 있다. */ + decode?: FrameDecoder; + /** decode 오류를 `seq` 순서로 보고한다. transport가 inline path와 동일한 parse-error 정책을 적용할 수 있게 한다. */ + onError?: (err: Error) => void; + /** 동시 in-flight decode 상한. `<= 0`이면 4를 적용한다. */ + workers?: number; +} + +/** + * async worker-pool gateway다. 제출된 프레임을 최대 `workers`개의 async slot이 동시에 decode하고, 단일 + * reorder 버퍼가 결과를 `seq`로 재정렬해 입력 순서로 sink에 전달한다. sink 호출은 단일 chain으로 직렬화되어 + * 순서가 보존되므로, stateful dispatch는 수신 coordinator에 남는다. + * + * 기본 decode는 호출 스레드에서 PacketBase envelope을 inline 처리하므로 worker 미지원 환경의 safe + * fallback이 된다. browser에서는 Web Worker 기반 `decode`를, Node에서는 `worker_threads` 기반 + * `decode`를 주입해 실제 off-thread 처리를 적용할 수 있다. + * + * decode에 실패한 프레임은 sink가 아니라 `onError`로 `seq` 순서에 맞춰 보고되어 한 프레임의 실패가 이후 + * 프레임을 막지 않는다. gateway 자체는 transport를 종료하지 않는다. + */ +export class WorkerGateway implements InboundGateway { + private readonly decode: FrameDecoder; + private readonly sink: (frame: DecodedFrame) => void | Promise; + private readonly onError?: (err: Error) => void; + private readonly maxWorkers: number; + private readonly reorder = new FrameReorderBuffer(); + private readonly jobs: InboundFrame[] = []; + private activeWorkers = 0; + private closed = false; + private sinkChain: Promise = Promise.resolve(); + + constructor(options: WorkerGatewayOptions) { + this.decode = options.decode ?? decodePacketBaseEnvelope; + this.sink = options.sink; + this.onError = options.onError; + const workers = options.workers ?? 0; + this.maxWorkers = workers > 0 ? workers : 4; + } + + submit(frame: InboundFrame): void { + if (this.closed) { + return; + } + this.jobs.push(frame); + if (this.activeWorkers < this.maxWorkers) { + this.activeWorkers += 1; + void this.workerLoop(); + } + } + + close(): void { + if (this.closed) { + return; + } + // 추가 전달을 멈춘다(InboundGateway.close() 계약). queued jobs와 reorder pending을 비워 close 이후 + // 완료되는 in-flight decode가 stale result를 sink/onError로 흘리거나 contiguous release를 유발하지 못하게 + // 한다. closed guard는 collect/sinkChain 전달 직전에도 적용된다. + this.closed = true; + this.jobs.length = 0; + this.reorder.clear(); + } + + private async workerLoop(): Promise { + try { + for (;;) { + const frame = this.jobs.shift(); + if (frame === undefined) { + return; + } + let result: GatewayResult; + try { + const env = await this.decode(frame.bytes); + result = { seq: frame.seq, frame: { seq: frame.seq, ...env }, error: null }; + } catch (err) { + result = { seq: frame.seq, frame: null, error: toError(err) }; + } + // decode가 await 중에 close되었을 수 있다. close 이후에는 stale result/error를 전달하지 않는다. + if (this.closed) { + return; + } + this.collect(result); + } + } finally { + this.activeWorkers -= 1; + } + } + + private collect(result: GatewayResult): void { + if (this.closed) { + return; + } + const ready = this.reorder.release(result); + if (ready.length === 0) { + return; + } + // reorder는 contiguous prefix만 release하므로 batch들은 호출 순서가 곧 seq 순서다. sinkChain으로 + // 직렬화해 async sink가 입력 순서를 깨지 않게 한다. + this.sinkChain = this.sinkChain.then(async () => { + for (const r of ready) { + // 스케줄 이후 close되었으면 남은 전달을 중단한다. + if (this.closed) { + return; + } + if (r.error !== null) { + // decode 오류를 seq 순서로 보고하고 계속 진행한다. transport 종료 여부는 onError가 결정한다. + this.onError?.(r.error); + continue; + } + if (r.frame !== null) { + await this.sink(r.frame); + } + } + }); + } +} + +function toError(err: unknown): Error { + return err instanceof Error ? err : new Error(String(err)); +} diff --git a/typescript/src/node.ts b/typescript/src/node.ts index 07fbf81..94be73b 100644 --- a/typescript/src/node.ts +++ b/typescript/src/node.ts @@ -1,6 +1,7 @@ export * from "./base_client.js"; export * from "./communicator.js"; export * from "./packets/message_common_pb.js"; +export * from "./node_inbound_gateway.js"; export * from "./tcp_client.js"; export * from "./tcp_server.js"; export * from "./node_ws_client.js"; diff --git a/typescript/src/node_inbound_gateway.ts b/typescript/src/node_inbound_gateway.ts new file mode 100644 index 0000000..0afa2f0 --- /dev/null +++ b/typescript/src/node_inbound_gateway.ts @@ -0,0 +1,168 @@ +import { availableParallelism } from "node:os"; +import { fileURLToPath } from "node:url"; +import { Worker } from "node:worker_threads"; + +import { + type DecodedEnvelope, + type FrameDecoder, + type InboundFrame, + type InboundGateway, + WorkerGateway, +} from "./inbound_gateway.js"; + +interface PendingDecode { + resolve: (env: DecodedEnvelope) => void; + reject: (err: Error) => void; +} + +interface WorkerReply { + id: number; + env?: DecodedEnvelope; + error?: string; +} + +/** + * `worker_threads` 기반 PacketBase decode pool이다. N개 worker에 frame을 round-robin으로 분배해 off-thread로 + * decode하고, 결과를 Promise로 돌려준다. pool 자체는 stateful dispatch나 순서 복원을 소유하지 않는다. 순서 + * 복원과 sink 직렬화는 이 pool을 `decode` 백엔드로 사용하는 {@link WorkerGateway}가 담당한다. + * + * decode된 payload는 owned buffer로 transfer되어 결과 path의 clone 복사를 줄인다. 입력 raw frame은 source가 + * buffer 전체를 소유할 때만 transfer하고, pool(공유) buffer view이면 owned buffer로 복사해 detach 부작용을 + * 피한다. + */ +class WorkerThreadDecodePool { + private readonly workers: Worker[]; + private readonly pending = new Map(); + private nextId = 0; + private rr = 0; + private closed = false; + + constructor(workerCount: number) { + this.workers = Array.from({ length: workerCount }, () => { + const worker = spawnDecodeWorker(); + worker.on("message", (reply: WorkerReply) => this.onReply(reply)); + worker.on("error", (err: Error) => this.failAll(err)); + worker.unref(); + return worker; + }); + } + + decode(bytes: Uint8Array): Promise { + if (this.closed) { + return Promise.reject(new Error("inbound gateway worker pool is closed")); + } + const id = this.nextId++; + const worker = this.workers[this.rr % this.workers.length]; + this.rr += 1; + return new Promise((resolve, reject) => { + this.pending.set(id, { resolve, reject }); + const owned = detachableBuffer(bytes); + worker?.postMessage({ id, bytes: new Uint8Array(owned) }, [owned]); + }); + } + + close(): void { + if (this.closed) { + return; + } + this.closed = true; + this.failAll(new Error("inbound gateway worker pool is closed")); + for (const worker of this.workers) { + void worker.terminate(); + } + } + + private onReply(reply: WorkerReply): void { + const pending = this.pending.get(reply.id); + if (pending === undefined) { + return; + } + this.pending.delete(reply.id); + if (reply.error !== undefined) { + pending.reject(new Error(reply.error)); + return; + } + if (reply.env !== undefined) { + pending.resolve(reply.env); + } + } + + private failAll(err: Error): void { + const pending = Array.from(this.pending.values()); + this.pending.clear(); + for (const p of pending) { + p.reject(err); + } + } +} + +/** + * source view가 buffer 전체를 소유하면 그 ArrayBuffer를 그대로(transfer 대상) 반환하고, pool/공유 view이면 + * owned buffer로 복사해 반환한다. 후자에서만 복사가 발생한다. + */ +function detachableBuffer(bytes: Uint8Array): ArrayBuffer { + if ( + bytes.byteOffset === 0 && + bytes.byteLength === bytes.buffer.byteLength && + bytes.buffer instanceof ArrayBuffer + ) { + return bytes.buffer; + } + return bytes.slice().buffer as ArrayBuffer; +} + +/** + * decode worker를 띄운다. 빌드 산출물(.js)에서는 worker 모듈을 직접 로드한다. 빌드 전(.ts) 실행 환경(vitest/ + * tsx 등)에서는 worker `execArgv`가 `--import`를 허용하지 않으므로, worker 안에서 tsx loader를 register한 뒤 + * `.ts` worker를 동적 import하는 부트스트랩을 eval로 실행한다. tsx는 devDependency이며 이 경로는 빌드 전 + * 환경에서만 사용된다. + */ +function spawnDecodeWorker(): Worker { + const here = fileURLToPath(import.meta.url); + if (here.endsWith(".ts")) { + const workerUrl = new URL("./node_inbound_gateway_worker.ts", import.meta.url); + const bootstrap = + `import { register } from 'tsx/esm/api';\n` + + `register();\n` + + `await import(${JSON.stringify(workerUrl.href)});\n`; + return new Worker(bootstrap, { eval: true }); + } + return new Worker(new URL("./node_inbound_gateway_worker.js", import.meta.url)); +} + +/** + * Node `worker_threads`로 PacketBase decode를 off-thread 처리하는 inbound gateway를 만든다. 결과는 내부 + * `seq`로 reorder된 뒤 입력 순서로 sink에 전달되므로, 수신 coordinator가 FIFO dispatch와 stateful 처리를 + * 단독 소유한다는 계약을 유지한다. + * + * `workers <= 0`이면 사용 가능한 병렬도로 기본값을 적용한다. transferable `ArrayBuffer`로 decode 결과 payload + * 복사를 줄인다. browser 등 `worker_threads` 미지원 환경에서는 이 모듈을 import하지 말고, inline decode를 + * 쓰는 {@link WorkerGateway}(또는 `Communicator.enableInboundGateway`)를 safe fallback으로 사용한다. + */ +export function createNodeWorkerGateway(options: { + sink: (frame: { typeName: string; data: Uint8Array; incomingNonce: number; responseNonce: number }) => void | Promise; + onError?: (err: Error) => void; + workers?: number; +}): InboundGateway { + const requested = options.workers ?? 0; + const workerCount = requested > 0 ? requested : Math.max(1, availableParallelism()); + const pool = new WorkerThreadDecodePool(workerCount); + const decode: FrameDecoder = (bytes) => pool.decode(bytes); + + const gateway = new WorkerGateway({ + workers: workerCount, + decode, + sink: options.sink, + onError: options.onError, + }); + + return { + submit(frame: InboundFrame): void { + gateway.submit(frame); + }, + close(): void { + gateway.close(); + pool.close(); + }, + }; +} diff --git a/typescript/src/node_inbound_gateway_worker.ts b/typescript/src/node_inbound_gateway_worker.ts new file mode 100644 index 0000000..3378ac3 --- /dev/null +++ b/typescript/src/node_inbound_gateway_worker.ts @@ -0,0 +1,30 @@ +import { parentPort } from "node:worker_threads"; + +import { decodePacketBaseEnvelope } from "./inbound_gateway.js"; + +/** + * Node `worker_threads` gateway worker의 entry다. main thread가 보낸 raw `PacketBase` frame을 off-thread에서 + * decode하고, 결과 envelope을 다시 main thread로 보낸다. decode는 순수 작업이며 공유 상태를 건드리지 않는다. + * + * decode된 payload(`data`)는 새로 할당된 owned buffer이므로 transferable로 전달해 clone 복사를 피한다. + */ +interface DecodeRequest { + id: number; + bytes: Uint8Array; +} + +const port = parentPort; +if (port === null) { + throw new Error("node_inbound_gateway_worker must run as a worker_threads worker"); +} + +port.on("message", (req: DecodeRequest) => { + try { + const env = decodePacketBaseEnvelope(req.bytes); + // 갓 디코드된 payload buffer는 owned이므로 transfer해 clone 복사를 줄인다. + const transfer = env.data.buffer instanceof ArrayBuffer ? [env.data.buffer] : []; + port.postMessage({ id: req.id, env }, transfer); + } catch (err) { + port.postMessage({ id: req.id, error: err instanceof Error ? err.message : String(err) }); + } +}); diff --git a/typescript/test/communicator.test.ts b/typescript/test/communicator.test.ts index 61b3bc6..ec95aef 100644 --- a/typescript/test/communicator.test.ts +++ b/typescript/test/communicator.test.ts @@ -6,6 +6,13 @@ import { type Transport, parserFromSchema, } from "../src/communicator.js"; +import { + type DecodedEnvelope, + decodePacketBaseEnvelope, + FrameReorderBuffer, + WorkerGateway, +} from "../src/inbound_gateway.js"; +import { createNodeWorkerGateway } from "../src/node_inbound_gateway.js"; import { create, PacketBaseSchema, @@ -72,6 +79,30 @@ async function waitFor(predicate: () => boolean): Promise { throw new Error("condition not met"); } +async function waitUntil(predicate: () => boolean, timeoutMs = 2000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (predicate()) { + return; + } + await new Promise((r) => setTimeout(r, 5)); + } + throw new Error("condition not met"); +} + +function rawFrame(index: number, nonce: number, responseNonce = 0): Uint8Array { + const data = toBinary(TestDataSchema, create(TestDataSchema, { index })); + return toBinary( + PacketBaseSchema, + create(PacketBaseSchema, { + typeName: TestDataSchema.typeName, + nonce, + responseNonce, + data, + }), + ); +} + describe("Communicator", () => { test("addListener/addRequestListener mutual exclusion", () => { const comm = new Communicator(); @@ -482,3 +513,224 @@ describe("Communicator", () => { }); }); +describe("FrameReorderBuffer", () => { + test("out-of-order 결과를 contiguous prefix가 채워질 때만 seq 순서로 release한다", () => { + const buffer = new FrameReorderBuffer(); + + expect(buffer.release({ seq: 3, frame: null, error: null })).toHaveLength(0); + expect(buffer.release({ seq: 2, frame: null, error: null })).toHaveLength(0); + + const flushed = buffer.release({ seq: 1, frame: null, error: null }); + expect(flushed.map((r) => r.seq)).toEqual([1, 2, 3]); + + // nextSeq 아래의 seq는 중복/지연 도착이므로 무시한다. + expect(buffer.release({ seq: 1, frame: null, error: null })).toHaveLength(0); + }); +}); + +describe("Communicator inbound gateway", () => { + test("gateway on: worker가 out-of-order로 완료해도 입력 seq 순서로 dispatch한다", async () => { + const comm = new Communicator(); + comm.initialize(new MockTransport(), parserMap()); + + const received: number[] = []; + comm.addListener(TestDataSchema.typeName, (msg) => { + received.push((msg as unknown as { index: number }).index); + }); + + const total = 6; + // 앞선 seq일수록 더 오래 decode시켜 완료 순서를 역전시킨다. + comm.enableInboundGateway({ + workers: 4, + decode: async (bytes) => { + const env = decodePacketBaseEnvelope(bytes); + await new Promise((r) => setTimeout(r, (total - env.incomingNonce) * 5)); + return env; + }, + }); + + for (let i = 1; i <= total; i += 1) { + comm.onReceivedFrame(rawFrame(i, i)); + } + + await waitUntil(() => received.length === total, 3000); + expect(received).toEqual([1, 2, 3, 4, 5, 6]); + await comm.close(); + }); + + test("gateway 미설정 fallback: onReceivedFrame이 inline decode로 coordinator path를 유지한다", async () => { + const comm = new Communicator(); + comm.initialize(new MockTransport(), parserMap()); + + const received: number[] = []; + comm.addListener(TestDataSchema.typeName, (msg) => { + received.push((msg as unknown as { index: number }).index); + }); + + for (let i = 1; i <= 5; i += 1) { + comm.onReceivedFrame(rawFrame(i, i)); + } + + for (let i = 0; i < 20; i += 1) { + if (received.length === 5) break; + await Promise.resolve(); + } + + expect(received).toEqual([1, 2, 3, 4, 5]); + await comm.close(); + }); + + test("gateway decode 오류는 frame error handler로 보고되고 이후 프레임을 막지 않는다", async () => { + const comm = new Communicator(); + comm.initialize(new MockTransport(), parserMap()); + + const received: number[] = []; + comm.addListener(TestDataSchema.typeName, (msg) => { + received.push((msg as unknown as { index: number }).index); + }); + const errors: Error[] = []; + comm.setFrameErrorHandler((err) => { + errors.push(err); + }); + + comm.enableInboundGateway({ workers: 2 }); + + comm.onReceivedFrame(rawFrame(1, 1)); + comm.onReceivedFrame(new Uint8Array([0xff, 0xff, 0xff])); // 디코드 불가 프레임 + comm.onReceivedFrame(rawFrame(3, 3)); + + await waitUntil(() => received.length === 2, 2000); + expect(received).toEqual([1, 3]); + expect(errors).toHaveLength(1); + await comm.close(); + }); + + test("Node worker_threads gateway: off-thread decode가 입력 순서를 보존한다", async () => { + const comm = new Communicator(); + comm.initialize(new MockTransport(), parserMap()); + + const received: number[] = []; + comm.addListener(TestDataSchema.typeName, (msg) => { + received.push((msg as unknown as { index: number }).index); + }); + + comm.attachInboundGateway( + createNodeWorkerGateway({ + sink: (frame) => + comm.enqueueInbound( + frame.typeName, + frame.data, + frame.incomingNonce, + frame.responseNonce, + ), + workers: 4, + }), + ); + + const total = 8; + for (let i = 1; i <= total; i += 1) { + comm.onReceivedFrame(rawFrame(i, i)); + } + + await waitUntil(() => received.length === total, 8000); + expect(received).toEqual([1, 2, 3, 4, 5, 6, 7, 8]); + await comm.close(); + }, 10000); + + test("attachInboundGateway 경계: 닫힌 이전 gateway의 늦은 결과가 현재 communicator로 들어오지 않는다", async () => { + const comm = new Communicator(); + comm.initialize(new MockTransport(), parserMap()); + + const received: number[] = []; + comm.addListener(TestDataSchema.typeName, (msg) => { + received.push((msg as unknown as { index: number }).index); + }); + + let resolveDecode!: (env: DecodedEnvelope) => void; + comm.enableInboundGateway({ + workers: 1, + decode: () => + new Promise((resolve) => { + resolveDecode = resolve; + }), + }); + comm.onReceivedFrame(rawFrame(1, 1)); + await waitFor(() => resolveDecode !== undefined); + + // 새 gateway attach가 이전 gateway를 close한다. + comm.attachInboundGateway( + new WorkerGateway({ + sink: (frame) => + comm.enqueueInbound( + frame.typeName, + frame.data, + frame.incomingNonce, + frame.responseNonce, + ), + }), + ); + + // 이전 gateway의 pending decode를 close 이후 완료시킨다. + const data = toBinary(TestDataSchema, create(TestDataSchema, { index: 1 })); + resolveDecode({ + typeName: TestDataSchema.typeName, + data, + incomingNonce: 1, + responseNonce: 0, + }); + await new Promise((r) => setTimeout(r, 10)); + + expect(received).toEqual([]); + await comm.close(); + }); +}); + +describe("WorkerGateway lifecycle", () => { + test("close() 이후 pending decode result는 sink로 전달되지 않는다", async () => { + let resolveDecode!: (env: DecodedEnvelope) => void; + const delivered: number[] = []; + const gateway = new WorkerGateway({ + workers: 1, + decode: () => + new Promise((resolve) => { + resolveDecode = resolve; + }), + sink: (frame) => { + delivered.push(frame.seq); + }, + }); + + gateway.submit({ seq: 1, bytes: new Uint8Array([1]) }); + await waitFor(() => resolveDecode !== undefined); + gateway.close(); + resolveDecode({ typeName: "t", data: new Uint8Array(), incomingNonce: 0, responseNonce: 0 }); + await new Promise((r) => setTimeout(r, 0)); + + expect(delivered).toEqual([]); + }); + + test("close() 이후 pending decode reject는 onError로 전달되지 않는다", async () => { + let rejectDecode!: (err: Error) => void; + const errors: Error[] = []; + const gateway = new WorkerGateway({ + workers: 1, + decode: () => + new Promise((_resolve, reject) => { + rejectDecode = reject; + }), + sink: () => {}, + onError: (err) => { + errors.push(err); + }, + }); + + gateway.submit({ seq: 1, bytes: new Uint8Array([1]) }); + await waitFor(() => rejectDecode !== undefined); + gateway.close(); + rejectDecode(new Error("decode failed")); + await new Promise((r) => setTimeout(r, 0)); + + expect(errors).toEqual([]); + }); +}); +