feat(node): FIFO admission queue 구현 및 관련 파일 업데이트

- edge connector 및 service 테스트 개선
- node FIFO admission queue 구현
- 로드맵 및 마일스톤 문서 업데이트
- queue observe snapshot 아카이브 이동
This commit is contained in:
toki 2026-06-15 07:38:37 +09:00
parent 39d1d08f33
commit 953617b12d
9 changed files with 202 additions and 41 deletions

View file

@ -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 기준선을 만든다.

View file

@ -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
- 확인 필요: 없음

View file

@ -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로 이동한다.

View file

@ -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
- 없음
## 후속 작업
- 없음

View file

@ -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)
}
}

View file

@ -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)
}
}

View file

@ -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{

View file

@ -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