From 953617b12df0838e38d1e48e9cd0f499f70fa3c2 Mon Sep 17 00:00:00 2001 From: toki Date: Mon, 15 Jun 2026 07:38:37 +0900 Subject: [PATCH] =?UTF-8?q?feat(node):=20FIFO=20admission=20queue=20?= =?UTF-8?q?=EA=B5=AC=ED=98=84=20=EB=B0=8F=20=EA=B4=80=EB=A0=A8=20=ED=8C=8C?= =?UTF-8?q?=EC=9D=BC=20=EC=97=85=EB=8D=B0=EC=9D=B4=ED=8A=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - edge connector 및 service 테스트 개선 - node FIFO admission queue 구현 - 로드맵 및 마일스톤 문서 업데이트 - queue observe snapshot 아카이브 이동 --- .../inference-provider-extension/PHASE.md | 2 +- ...-availability-capacity-queue-foundation.md | 23 +++-- .../code_review_local_G05_0.log} | 93 ++++++++++++++----- .../05+04_queue_observe_snapshot/complete.log | 46 +++++++++ .../plan_local_G05_0.log} | 0 .../internal/controlplane/connector_test.go | 4 +- apps/edge/internal/service/service_test.go | 4 +- apps/node/internal/node/node.go | 11 ++- apps/node/internal/node/node_test.go | 60 ++++++++++++ 9 files changed, 202 insertions(+), 41 deletions(-) rename agent-task/{m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/CODE_REVIEW-local-G05.md => archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/code_review_local_G05_0.log} (53%) create mode 100644 agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/complete.log rename agent-task/{m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/PLAN-local-G05.md => archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/plan_local_G05_0.log} (100%) diff --git a/agent-roadmap/phase/inference-provider-extension/PHASE.md b/agent-roadmap/phase/inference-provider-extension/PHASE.md index 6996ed4..3ef316f 100644 --- a/agent-roadmap/phase/inference-provider-extension/PHASE.md +++ b/agent-roadmap/phase/inference-provider-extension/PHASE.md @@ -20,7 +20,7 @@ Ollama 경로가 안정화된 뒤, 그 결과를 기준선으로 삼아 Lemonade - 경로: `agent-roadmap/archive/phase/inference-provider-extension/milestones/node-multi-target-serving-foundation.md` - 요약: 하나의 Node 연결이 여러 CLI profile, terminal gateway, Ollama/vLLM/SGLang 같은 추론 엔진 연결, 여러 model target을 동시에 제공할 수 있도록 config/proto/routing/runtime 기준선을 정리했고, 최종 코드 리뷰와 검증까지 통과해 archive했다. -- [진행중] Node provider 상태와 Capacity Queue 기반 +- [검토중] Node provider 상태와 Capacity Queue 기반 - 경로: `agent-roadmap/phase/inference-provider-extension/milestones/provider-availability-capacity-queue-foundation.md` - 요약: NomadCode workspace 실행 계약이 닫힌 뒤 provider별 health/model probe 차이를 Node adapter 내부로 숨기고, Edge가 공통으로 볼 수 있는 최소 상태와 capacity/in-flight/queued snapshot, FIFO admission queue 기준선을 만든다. diff --git a/agent-roadmap/phase/inference-provider-extension/milestones/provider-availability-capacity-queue-foundation.md b/agent-roadmap/phase/inference-provider-extension/milestones/provider-availability-capacity-queue-foundation.md index a88532e..2ee66d9 100644 --- a/agent-roadmap/phase/inference-provider-extension/milestones/provider-availability-capacity-queue-foundation.md +++ b/agent-roadmap/phase/inference-provider-extension/milestones/provider-availability-capacity-queue-foundation.md @@ -13,7 +13,7 @@ Ollama, Lemonade, vLLM, SGLang 같은 provider별 상태 확인 방식 차이를 ## 상태 -[진행중] +[검토중] ## 승격 조건 @@ -46,21 +46,21 @@ Provider별 probe 차이를 Node 내부로 감추고 Edge가 라우팅 입력으 Provider별 동시 처리 한도를 IOP가 소유하고, 한도를 넘는 요청을 예측 가능한 FIFO queue로 관리한다. -- [ ] [capacity-config] provider target별 `capacity`, `max_queue`, `queue_timeout`, `request_timeout` 설정 기준을 정리하고 기본 config 예시에 반영한다. -- [ ] [admission-gate] Node provider executor가 `in_flight < capacity`일 때만 provider 호출을 시작하고, capacity가 찬 요청은 provider별 FIFO queue에 넣는다. -- [ ] [queue-release] 실행 중 요청이 응답, 실패, 취소, timeout으로 종료되면 `in_flight`를 줄이고 queue의 첫 요청을 실행 슬롯으로 승격한다. -- [ ] [queue-reject] `max_queue` 초과 또는 `queue_timeout` 초과 요청은 명확한 error/rejection reason으로 종료한다. -- [ ] [queue-observe] Edge가 보는 snapshot에 `capacity`, `in_flight`, `queued`가 포함되어 이후 capacity-aware routing의 입력으로 사용할 수 있다. +- [x] [capacity-config] provider target별 `capacity`, `max_queue`, `queue_timeout`, `request_timeout` 설정 기준을 정리하고 기본 config 예시에 반영한다. +- [x] [admission-gate] Node provider executor가 `in_flight < capacity`일 때만 provider 호출을 시작하고, capacity가 찬 요청은 provider별 FIFO queue에 넣는다. +- [x] [queue-release] 실행 중 요청이 응답, 실패, 취소, timeout으로 종료되면 `in_flight`를 줄이고 queue의 첫 요청을 실행 슬롯으로 승격한다. +- [x] [queue-reject] `max_queue` 초과 또는 `queue_timeout` 초과 요청은 명확한 error/rejection reason으로 종료한다. +- [x] [queue-observe] Edge가 보는 snapshot에 `capacity`, `in_flight`, `queued`가 포함되어 이후 capacity-aware routing의 입력으로 사용할 수 있다. ## 완료 리뷰 -- 상태: 없음 -- 요청일: 없음 -- 완료 근거: `status-model`, `probe-contract`, `edge-snapshot`은 완료되었으나 capacity queue Task가 아직 충족되지 않았다. +- 상태: 요청됨 +- 요청일: 2026-06-15 +- 완료 근거: 모든 기능 Task가 PASS 완료되었고, 마지막 `queue-observe`는 `agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/complete.log`에서 완료 근거와 검증을 남겼다. - 리뷰 필요: - [ ] 사용자가 완료 결과를 확인했다 - [ ] archive 이동을 승인했다 -- 리뷰 코멘트: 없음 +- 리뷰 코멘트: 사용자 최종 확인과 archive 승인 대기 ## 범위 제외 @@ -80,6 +80,9 @@ Provider별 동시 처리 한도를 IOP가 소유하고, 한도를 넘는 요청 - 진행 근거: `status-model`은 `apps/node/internal/runtime/types.go`와 Node `CAPABILITIES` 응답에 제한 상태 모델을 추가했고 `go test ./apps/node/...`로 검증했다. - 완료 근거: `probe-contract`는 `agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/01_probe_contract/complete.log`에서 PASS 완료되었고 `go test -count=1 ./apps/node/...`, `./scripts/e2e-smoke.sh`, `git diff --check`로 검증했다. - 완료 근거: `edge-snapshot`은 `agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/02+01_edge_snapshot/complete.log`에서 PASS 완료되었고 `make proto`, `go test -count=1 ./apps/node/... ./apps/edge/... ./packages/go/... ./proto/gen/...`, `make test-control-plane-edge-wire`, `./scripts/e2e-smoke.sh`, `git diff --check`로 검증했다. +- 완료 근거: `capacity-config`는 `agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/03_queue_capability_contract/complete.log`에서 PASS 완료되었고 `go test -count=1 ./apps/node/internal/runtime ./apps/node/internal/adapters/ollama ./apps/node/internal/adapters/vllm ./apps/node/internal/adapters`로 검증했다. +- 완료 근거: `admission-gate`, `queue-release`, `queue-reject`는 `agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/04+03_fifo_admission_queue/complete.log`에서 PASS 완료되었고 `go test -count=1 ./apps/node/internal/node`, `go test -race -count=1 ./apps/node/internal/node`, `go test -count=1 ./apps/node/...`, `./scripts/e2e-smoke.sh`, `git diff --check`로 검증했다. +- 완료 근거: `queue-observe`는 `agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/complete.log`에서 PASS 완료되었고 `go test -count=1 ./apps/node/internal/node -run 'TestOnCommandRequest_Capabilities'`, `go test -count=1 ./apps/edge/internal/service ./apps/edge/internal/controlplane`, `go test -count=1 ./apps/node/internal/node`, `go test -count=1 ./apps/node/... ./apps/edge/...`, `./scripts/e2e-smoke.sh`, `git diff --check`로 검증했다. - 선행 작업: Node 단일 통로 멀티 타겟 서빙 기반, OpenAI Workspace Agent Execution Contract - 후속 작업: Lemonade provider 서빙 경로 추가, vLLM provider 서빙 경로 추가, SGLang provider 서빙 경로 추가, model group alias와 capacity-aware routing Milestone - 확인 필요: 없음 diff --git a/agent-task/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/CODE_REVIEW-local-G05.md b/agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/code_review_local_G05_0.log similarity index 53% rename from agent-task/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/CODE_REVIEW-local-G05.md rename to agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/code_review_local_G05_0.log index 06d6114..241c263 100644 --- a/agent-task/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/CODE_REVIEW-local-G05.md +++ b/agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/code_review_local_G05_0.log @@ -31,36 +31,39 @@ task=m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snaps | 항목 | 완료 여부 | |------|---------| -| [API-1] Node Capabilities Queued Count | [ ] | -| [API-2] Edge Snapshot Preservation | [ ] | +| [API-1] Node Capabilities Queued Count | [x] | +| [API-2] Edge Snapshot Preservation | [x] | ## 구현 체크리스트 -- [ ] predecessor `04+03_fifo_admission_queue`의 `complete.log`가 active 또는 archive에 있는지 확인한다. -- [ ] Node capabilities command가 admission manager의 actual `in_flight`와 `queued` 값을 읽는다. -- [ ] `NodeCommandResponse.Result["queued"]`와 `ProviderSnapshot.Queued`에 같은 queued count를 넣는다. -- [ ] Edge service/control-plane preservation tests에 queued non-zero fixture를 추가한다. -- [ ] 최종 검증 명령을 실행한다. -- [ ] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. +- [x] predecessor `04+03_fifo_admission_queue`의 `complete.log`가 active 또는 archive에 있는지 확인한다. +- [x] Node capabilities command가 admission manager의 actual `in_flight`와 `queued` 값을 읽는다. +- [x] `NodeCommandResponse.Result["queued"]`와 `ProviderSnapshot.Queued`에 같은 queued count를 넣는다. +- [x] Edge service/control-plane preservation tests에 queued non-zero fixture를 추가한다. +- [x] 최종 검증 명령을 실행한다. +- [x] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. ## 코드리뷰 전용 체크리스트 > **[REVIEW AGENT ONLY]** 이 체크리스트는 코드리뷰 에이전트만 사용한다. -- [ ] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다. -- [ ] active `CODE_REVIEW-*-G??.md`를 `code_review_local_G05_N.log`로 아카이브한다. -- [ ] active `PLAN-*-G??.md`를 `plan_local_G05_M.log`로 아카이브한다. -- [ ] PASS이면 `complete.log`를 작성하고 active `.md` 파일을 남기지 않는다. -- [ ] PASS이면 active task 디렉터리를 archive로 이동한다. -- [ ] PASS이고 task group이 `m-provider-availability-capacity-queue-foundation`이면 완료 이벤트 메타데이터를 보고하고 roadmap을 직접 수정하지 않는다. +- [x] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다. +- [x] active `CODE_REVIEW-*-G??.md`를 `code_review_local_G05_N.log`로 아카이브한다. +- [x] active `PLAN-*-G??.md`를 `plan_local_G05_M.log`로 아카이브한다. +- [x] PASS이면 `complete.log`를 작성하고 active `.md` 파일을 남기지 않는다. +- [x] PASS이면 active task 디렉터리를 archive로 이동한다. +- [x] PASS이고 task group이 `m-provider-availability-capacity-queue-foundation`이면 완료 이벤트 메타데이터를 보고하고 roadmap을 직접 수정하지 않는다. ## 계획 대비 변경 사항 -_구현 에이전트가 계획과 다르게 구현한 부분을 이유와 함께 기록한다._ +- 계획 범위와 동일하게 구현했다. +- 추가로 `./scripts/e2e-smoke.sh`를 실행해 임시 config 기반 `scripts/dev/edge.sh` + `scripts/dev/node.sh` 보조 진단 흐름에서 `/capabilities`와 메시지 왕복이 깨지지 않는지 확인했다. ## 주요 설계 결정 -_구현 에이전트가 주요 설계 결정 사항을 기록한다._ +- Node `CAPABILITIES`의 `queued` 값은 새 카운터를 만들지 않고 기존 per-adapter `fifoGate`의 `queuedCount()`에서 읽도록 했다. +- `NodeCommandResponse.Result["queued"]`와 `ProviderSnapshot.Queued`는 같은 local variable에서 채워 snapshot과 map 응답이 갈라지지 않게 했다. +- Edge production mapper는 이미 typed provider snapshot을 보존하므로, service/control-plane 테스트 fixture를 non-zero queued 값으로 바꿔 보존 회귀만 강화했다. ## 사용자 리뷰 요청 @@ -87,27 +90,75 @@ _구현 에이전트가 각 중간 검증 및 최종 검증 명령 실행 후 ### API-1 중간 검증 ```bash $ go test -count=1 ./apps/node/internal/node -run 'TestOnCommandRequest_Capabilities' -(output) +ok iop/apps/node/internal/node 0.014s ``` ### API-2 중간 검증 ```bash $ go test -count=1 ./apps/edge/internal/service ./apps/edge/internal/controlplane -(output) +ok iop/apps/edge/internal/service 0.004s +ok iop/apps/edge/internal/controlplane 4.448s ``` ### 최종 검증 ```bash $ go test -count=1 ./apps/node/internal/node -(output) +ok iop/apps/node/internal/node 2.693s $ go test -count=1 ./apps/edge/internal/service ./apps/edge/internal/controlplane -(output) +ok iop/apps/edge/internal/service 0.004s +ok iop/apps/edge/internal/controlplane 4.448s $ go test -count=1 ./apps/node/... ./apps/edge/... -(output) +ok iop/apps/node/cmd/node 0.012s +ok iop/apps/node/internal/adapters 0.010s +ok iop/apps/node/internal/adapters/cli 46.674s +? iop/apps/node/internal/adapters/cli/internal/testutil [no test files] +ok iop/apps/node/internal/adapters/cli/status 39.836s +? iop/apps/node/internal/adapters/mock [no test files] +ok iop/apps/node/internal/adapters/ollama 0.008s +ok iop/apps/node/internal/adapters/vllm 0.007s +ok iop/apps/node/internal/bootstrap 0.265s +ok iop/apps/node/internal/node 2.693s +ok iop/apps/node/internal/router 0.006s +? iop/apps/node/internal/runtime [no test files] +ok iop/apps/node/internal/store 0.080s +ok iop/apps/node/internal/terminal 0.544s +ok iop/apps/node/internal/transport 5.144s +ok iop/apps/edge/cmd/edge 0.043s +ok iop/apps/edge/internal/bootstrap 0.014s +ok iop/apps/edge/internal/controlplane 4.451s +ok iop/apps/edge/internal/edgecmd 0.007s +ok iop/apps/edge/internal/events 0.004s +ok iop/apps/edge/internal/input 0.004s +ok iop/apps/edge/internal/input/a2a 0.005s +ok iop/apps/edge/internal/node 0.005s +ok iop/apps/edge/internal/openai 1.509s +ok iop/apps/edge/internal/opsconsole 0.006s +ok iop/apps/edge/internal/service 0.005s +ok iop/apps/edge/internal/transport 2.011s +``` + +### 추가 검증 +```bash +$ ./scripts/e2e-smoke.sh +[e2e] Auxiliary smoke test PASSED. ``` --- > **[IMPLEMENTING AGENT — BEFORE SAVING] Have you filled in every implementation-owned section: completion table, implementation checklist, changes from plan, design decisions, and verification output?** + +## 코드리뷰 결과 + +- 종합 판정: PASS +- 차원별 평가: + - Correctness: Pass + - Completeness: Pass + - Test coverage: Pass + - API contract: Pass + - Code quality: Pass + - Plan deviation: Pass + - Verification trust: Pass +- 발견된 문제: 없음 +- 다음 단계: PASS 종결. active plan/review를 log로 보존하고 `complete.log` 작성 후 task directory를 archive로 이동한다. diff --git a/agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/complete.log b/agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/complete.log new file mode 100644 index 0000000..2cc8e6f --- /dev/null +++ b/agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/complete.log @@ -0,0 +1,46 @@ +# Complete - m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot + +## 완료 일시 + +2026-06-15 + +## 요약 + +Queue observation snapshot 구현을 1회 리뷰했고 최종 판정은 PASS다. + +## 루프 이력 + +| Plan | Review | Verdict | 메모 | +|------|--------|---------|------| +| `plan_local_G05_0.log` | `code_review_local_G05_0.log` | PASS | Node capabilities queued count와 Edge provider snapshot 보존 테스트 구현 및 검증 통과 | + +## 구현/정리 내용 + +- Node `CAPABILITIES` 응답의 `queued` 값을 hardcoded `0` 대신 per-adapter FIFO gate의 actual queued count에서 읽도록 변경했다. +- `NodeCommandResponse.Result["queued"]`와 typed `ProviderSnapshot.Queued`가 같은 queued count를 쓰도록 정리했다. +- capacity가 찬 상태에서 background run이 대기 중일 때 capabilities result와 provider snapshot 모두 `queued=1`을 반환하는 Node regression test를 추가했다. +- Edge service와 Control Plane status snapshot preservation tests의 provider snapshot fixture를 non-zero queued 값으로 바꿔 typed snapshot 보존 경로를 강화했다. + +## 최종 검증 + +- `go test -count=1 ./apps/node/internal/node -run 'TestOnCommandRequest_Capabilities'` - PASS; `ok iop/apps/node/internal/node 0.014s` +- `go test -count=1 ./apps/edge/internal/service ./apps/edge/internal/controlplane` - PASS; service/controlplane 패키지 통과 +- `go test -count=1 ./apps/node/internal/node` - PASS; `ok iop/apps/node/internal/node 2.693s` +- `go test -count=1 ./apps/node/... ./apps/edge/...` - PASS; Node/Edge 하위 패키지 전체 통과 +- `./scripts/e2e-smoke.sh` - PASS; 임시 config 기반 dev edge/node 보조 smoke에서 `/capabilities`, `/transport`, 메시지 2회, background run, `/sessions`, `/terminate-session` 확인 +- `git diff --check` - PASS + +## Roadmap Completion + +- Milestone: `agent-roadmap/phase/inference-provider-extension/milestones/provider-availability-capacity-queue-foundation.md` +- Completed task ids: + - `queue-observe`: PASS; evidence=`agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/plan_local_G05_0.log`, `agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/code_review_local_G05_0.log`; verification=`go test -count=1 ./apps/node/internal/node -run 'TestOnCommandRequest_Capabilities'`, `go test -count=1 ./apps/edge/internal/service ./apps/edge/internal/controlplane`, `go test -count=1 ./apps/node/internal/node`, `go test -count=1 ./apps/node/... ./apps/edge/...`, `./scripts/e2e-smoke.sh`, `git diff --check` +- Not completed task ids: 없음 + +## 잔여 Nit + +- 없음 + +## 후속 작업 + +- 없음 diff --git a/agent-task/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/PLAN-local-G05.md b/agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/plan_local_G05_0.log similarity index 100% rename from agent-task/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/PLAN-local-G05.md rename to agent-task/archive/2026/06/m-provider-availability-capacity-queue-foundation/05+04_queue_observe_snapshot/plan_local_G05_0.log diff --git a/apps/edge/internal/controlplane/connector_test.go b/apps/edge/internal/controlplane/connector_test.go index 93351c1..37eb55b 100644 --- a/apps/edge/internal/controlplane/connector_test.go +++ b/apps/edge/internal/controlplane/connector_test.go @@ -230,7 +230,7 @@ func TestConnectorRespondsToStatusRequestFromProvider(t *testing.T) { Status: "available", Capacity: 4, InFlight: 2, - Queued: 0, + Queued: 3, }, }, }, @@ -297,7 +297,7 @@ func TestConnectorRespondsToStatusRequestFromProvider(t *testing.T) { t.Fatalf("expected 1 provider snapshot, got %d", len(n.GetProviderSnapshots())) } snap := n.GetProviderSnapshots()[0] - if snap.GetAdapter() != "cli" || snap.GetStatus() != "available" || snap.GetCapacity() != 4 || snap.GetInFlight() != 2 || snap.GetQueued() != 0 { + if snap.GetAdapter() != "cli" || snap.GetStatus() != "available" || snap.GetCapacity() != 4 || snap.GetInFlight() != 2 || snap.GetQueued() != 3 { t.Errorf("unexpected snapshot contents: %+v", snap) } } diff --git a/apps/edge/internal/service/service_test.go b/apps/edge/internal/service/service_test.go index daa011e..18e8176 100644 --- a/apps/edge/internal/service/service_test.go +++ b/apps/edge/internal/service/service_test.go @@ -907,7 +907,7 @@ func TestServiceCapabilitiesPreservesProviderSnapshots(t *testing.T) { Status: "available", Capacity: 4, InFlight: 2, - Queued: 0, + Queued: 3, }, }, }, nil @@ -936,7 +936,7 @@ func TestServiceCapabilitiesPreservesProviderSnapshots(t *testing.T) { t.Fatalf("expected 1 provider snapshot, got %d", len(view.ProviderSnapshots)) } snap := view.ProviderSnapshots[0] - if snap.Adapter != "cli" || snap.Status != "available" || snap.Capacity != 4 || snap.InFlight != 2 || snap.Queued != 0 { + if snap.Adapter != "cli" || snap.Status != "available" || snap.Capacity != 4 || snap.InFlight != 2 || snap.Queued != 3 { t.Fatalf("unexpected snapshot contents: %+v", snap) } } diff --git a/apps/node/internal/node/node.go b/apps/node/internal/node/node.go index c84a7e6..239f095 100644 --- a/apps/node/internal/node/node.go +++ b/apps/node/internal/node/node.go @@ -204,9 +204,8 @@ func (n *Node) OnRunRequest(ctx context.Context, sess *transport.Session, req *i return run() } - // A queued background run returns to the caller after enqueue; promotion and - // execution proceed in a goroutine. All other runs (immediately-admitted - // background, or any foreground run) take the synchronous path. + // Background runs return to the caller after enqueue; promotion and + // execution proceed in a goroutine. Foreground runs take the synchronous path. if spec.Background { go func() { _ = launch() }() return nil @@ -387,8 +386,10 @@ func (n *Node) handleCapabilitiesCommand(ctx context.Context, req *iop.NodeComma gate, ok := n.adapterGates[req.GetAdapter()] n.adapterGatesMu.Unlock() inFlight := 0 + queued := 0 if ok { inFlight = gate.activeCount() + queued = gate.queuedCount() } result := map[string]string{ @@ -399,7 +400,7 @@ func (n *Node) handleCapabilitiesCommand(ctx context.Context, req *iop.NodeComma "provider_status": string(runtime.NormalizeProviderStatus(providerStatus)), "capacity": strconv.Itoa(caps.MaxConcurrency), "in_flight": strconv.Itoa(inFlight), - "queued": "0", + "queued": strconv.Itoa(queued), } if providerDetail != "" { result["provider_detail"] = providerDetail @@ -410,7 +411,7 @@ func (n *Node) handleCapabilitiesCommand(ctx context.Context, req *iop.NodeComma Status: string(runtime.NormalizeProviderStatus(providerStatus)), Capacity: int32(caps.MaxConcurrency), InFlight: int32(inFlight), - Queued: 0, + Queued: int32(queued), } return &iop.NodeCommandResponse{ diff --git a/apps/node/internal/node/node_test.go b/apps/node/internal/node/node_test.go index 8c52b26..ba81f2b 100644 --- a/apps/node/internal/node/node_test.go +++ b/apps/node/internal/node/node_test.go @@ -718,6 +718,66 @@ func TestOnCommandRequest_Capabilities_InFlight(t *testing.T) { <-ba.done } +func TestOnCommandRequest_Capabilities_Queued(t *testing.T) { + sa := newQueuedSlowAdapter("slow", 1, 4, 0) + router := &fixedRouter{adapterName: "slow", adapters: map[string]runtime.Adapter{"slow": sa}} + n, st := makeNodeWithConcurrency(t, router, 0) + + if err := n.OnRunRequest(context.Background(), &transport.Session{}, &iop.RunRequest{ + RunId: "cap-hold", Adapter: "slow", Target: "v1", Background: true, + }); err != nil { + t.Fatalf("holding run: %v", err) + } + waitStarted(t, sa, "cap-hold") + + if err := n.OnRunRequest(context.Background(), &transport.Session{}, &iop.RunRequest{ + RunId: "cap-queued", Adapter: "slow", Target: "v1", Background: true, + }); err != nil { + t.Fatalf("queued run: %v", err) + } + requireQueued(t, st, "cap-queued") + + resp, err := n.OnCommandRequest(context.Background(), &transport.Session{}, &iop.NodeCommandRequest{ + RequestId: "req-cap-queued", + Type: iop.NodeCommandType_NODE_COMMAND_TYPE_CAPABILITIES, + Adapter: "slow", + Target: "v1", + }) + if err != nil { + t.Fatalf("OnCommandRequest: %v", err) + } + if resp.GetError() != "" { + t.Fatalf("expected no error, got %q", resp.GetError()) + } + if got := resp.GetResult()["capacity"]; got != "1" { + t.Fatalf("result[capacity]: got %q want %q", got, "1") + } + if got := resp.GetResult()["in_flight"]; got != "1" { + t.Fatalf("result[in_flight]: got %q want %q", got, "1") + } + if got := resp.GetResult()["queued"]; got != "1" { + t.Fatalf("result[queued]: got %q want %q", got, "1") + } + if len(resp.GetProviderSnapshots()) != 1 { + t.Fatalf("expected 1 provider snapshot, got %d", len(resp.GetProviderSnapshots())) + } + snap := resp.GetProviderSnapshots()[0] + if snap.GetCapacity() != 1 { + t.Fatalf("snap.Capacity: got %d want %d", snap.GetCapacity(), 1) + } + if snap.GetInFlight() != 1 { + t.Fatalf("snap.InFlight: got %d want %d", snap.GetInFlight(), 1) + } + if snap.GetQueued() != 1 { + t.Fatalf("snap.Queued: got %d want %d", snap.GetQueued(), 1) + } + + sa.releaseRun("cap-hold") + waitStarted(t, sa, "cap-queued") + sa.releaseRun("cap-queued") + requireStatusEventually(t, st, "cap-queued", "completed") +} + func TestOnCommandRequest_CapabilitiesProviderStatusModel(t *testing.T) { cases := []struct { name string