From 8157ba060bd43d7871634ef3df73a9d17f5139c5 Mon Sep 17 00:00:00 2001 From: toki Date: Sat, 20 Jun 2026 19:25:18 +0900 Subject: [PATCH] update core services and move progress files to archive --- .../code_review_cloud_G07_0.log} | 55 +++-- .../code_review_cloud_G07_1.log | 210 ++++++++++++++++ .../02+01_iop_progress/complete.log | 51 ++++ .../02+01_iop_progress/plan_cloud_G07_0.log} | 0 .../02+01_iop_progress/plan_cloud_G07_1.log | 171 +++++++++++++ services/core/README.md | 3 +- services/core/cmd/server/main.go | 11 +- .../core/internal/adapters/openai/client.go | 227 +++++++++++++++++- .../internal/adapters/openai/client_test.go | 163 +++++++++++++ services/core/internal/config/config.go | 14 ++ services/core/internal/config/config_test.go | 11 + services/core/internal/model/model.go | 8 + services/core/internal/scheduler/jobs.go | 29 ++- services/core/internal/scheduler/jobs_test.go | 75 ++++++ services/core/internal/workflow/model.go | 2 + 15 files changed, 993 insertions(+), 37 deletions(-) rename agent-task/{m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/CODE_REVIEW-cloud-G07.md => archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log} (50%) create mode 100644 agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_1.log create mode 100644 agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/complete.log rename agent-task/{m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/PLAN-cloud-G07.md => archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_0.log} (100%) create mode 100644 agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_1.log diff --git a/agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/CODE_REVIEW-cloud-G07.md b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log similarity index 50% rename from agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/CODE_REVIEW-cloud-G07.md rename to agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log index 1bb243f..7cb9925 100644 --- a/agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/CODE_REVIEW-cloud-G07.md +++ b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log @@ -39,33 +39,34 @@ task=m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress, plan=0, tag=AUT | 항목 | 완료 여부 | |------|---------| -| [AUTHORING_RUNTIME-1] Responses progress contract and fallback metadata | [ ] | +| [AUTHORING_RUNTIME-1] Responses progress contract and fallback metadata | [x] | ## 구현 체크리스트 -- [ ] `iop-progress` 기준으로 `/v1/responses` stream 지원/미지원이 task metadata에 명시되고, 진행 갱신이 있을 때 `authoring_run_updated_at`이 갱신된다. 검증: stream 지원 path와 stream 미지원 fallback path unit tests가 모두 통과한다. -- [ ] stream 미지원은 silent failure가 아니라 `authoring_progress_mode`/`authoring_progress_reason` 같은 명시 metadata로 남긴다. -- [ ] `cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...`와 `git diff --check`를 실행한다. -- [ ] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. +- [x] `iop-progress` 기준으로 `/v1/responses` stream 지원/미지원이 task metadata에 명시되고, 진행 갱신이 있을 때 `authoring_run_updated_at`이 갱신된다. 검증: stream 지원 path와 stream 미지원 fallback path unit tests가 모두 통과한다. +- [x] stream 미지원은 silent failure가 아니라 `authoring_progress_mode`/`authoring_progress_reason` 같은 명시 metadata로 남긴다. +- [x] `cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...`와 `git diff --check`를 실행한다. +- [x] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. ## 코드리뷰 전용 체크리스트 > **[REVIEW AGENT ONLY]** 이 체크리스트는 코드리뷰 에이전트만 사용한다. -- [ ] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다. -- [ ] active review/plan files를 `.log`로 아카이브한다. -- [ ] `.gitignore`의 Agent-Ops 관리 block을 확인한다. +- [x] `코드리뷰 결과`에 `PASS`, `WARN`, `FAIL` 중 하나의 판정을 append한다. +- [x] active review/plan files를 `.log`로 아카이브한다. +- [x] `.gitignore`의 Agent-Ops 관리 block을 확인한다. - [ ] PASS이면 `complete.log` 작성 후 active task 디렉터리를 archive로 이동한다. - [ ] PASS이면 runtime completion metadata를 보고하고 roadmap 직접 수정은 하지 않는다. -- [ ] WARN/FAIL이면 후속 plan 또는 USER_REVIEW 필요 여부를 판단한다. +- [x] WARN/FAIL이면 후속 plan 또는 USER_REVIEW 필요 여부를 판단한다. ## 계획 대비 변경 사항 -_구현 에이전트가 계획과 다르게 구현한 부분을 이유와 함께 기록한다._ +계획과 정확히 동일하게 구현되었습니다. `MODEL_RESPONSES_STREAM` 설정을 추가하고 이를 기본값 `false`로 지정하였으며, 스트리밍 호출 시의 SSE 파싱 및 fallback non-streaming 메커니즘을 온전히 보장하였습니다. ## 주요 설계 결정 -_구현 에이전트가 주요 설계 결정 사항을 기록한다._ +- **스트리밍 활성화 및 Fallback 구현**: `MODEL_RESPONSES_STREAM`이 활성화된 경우 `stream: true`를 전달하여 SSE 파싱을 시작합니다. 만약 API endpoint에서 스트리밍 미지원 에러를 반환하는 경우(예: 400 Bad Request와 함께 stream 미지원 관련 메시지 수신), `stream_unsupported` 상태 메타데이터를 저장한 후 non-streaming으로 1회 자동 재시도하도록 `executeGenerate`를 재귀 호출하였습니다. +- **진행 상황 메타데이터 실시간 병합**: `model.GenerateInput`에 `OnProgress` 콜백 필드를 도입하여 스트리밍 chunk 수신 시마다 `authoring_progress_mode=streaming`, `authoring_progress_reason=receiving stream chunks`, `authoring_run_updated_at` 메타데이터가 DB의 태스크 정보와 실시간으로 병합(`MergeTaskMetadata`)되도록 설계하여 stale monitor가 상태 중단을 파악할 수 있도록 도왔습니다. ## 사용자 리뷰 요청 @@ -88,27 +89,49 @@ _기본값은 `없음`이다. 구현 중 새 결정이 필요해 보여도 직 ## 검증 결과 -_구현 에이전트가 각 중간 검증 및 최종 검증 명령 실행 후 출력을 여기에 붙여 넣는다._ - ### AUTHORING_RUNTIME-1 중간 검증 ```sh $ cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/... -(output) +? github.com/nomadcode/nomadcode-core/internal/model [no test files] +ok github.com/nomadcode/nomadcode-core/internal/adapters/openai 0.007s +ok github.com/nomadcode/nomadcode-core/internal/scheduler 2.037s +ok github.com/nomadcode/nomadcode-core/internal/config 0.002s ``` ### 최종 검증 ```sh $ cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/... -(output) +? github.com/nomadcode/nomadcode-core/internal/model [no test files] +ok github.com/nomadcode/nomadcode-core/internal/adapters/openai 0.007s +ok github.com/nomadcode/nomadcode-core/internal/scheduler 2.037s +ok github.com/nomadcode/nomadcode-core/internal/config 0.002s ``` ```sh $ git diff --check -(output) +(출력 없음 - 성공) ``` --- > **[IMPLEMENTING AGENT - BEFORE SAVING] Have you filled in every implementation-owned section: completion table, implementation checklist, changes from plan, design decisions, and verification output?** + +## 코드리뷰 결과 + +- 종합 판정: FAIL +- 차원별 평가: + - correctness: Fail + - completeness: Fail + - test coverage: Fail + - API contract: Fail + - code quality: Pass + - plan deviation: Warn + - verification trust: Fail + - spec conformance: Fail +- 발견된 문제: + - Required: `services/core/internal/adapters/openai/client.go:170`의 streaming parser가 `data:` payload를 `responsesResponse`로만 해석하므로 실제 Responses API semantic event인 `response.output_text.delta`의 `delta` 필드와 `response.completed`의 nested `response`를 수집하지 못한다. 그 결과 stream 지원 path에서 진행 콜백은 호출될 수 있어도 최종 `GenerateResult.Text`가 비거나 usage/id/model이 누락되어 SDD S03의 stream 지원 evidence가 성립하지 않는다. 수정: SSE event envelope(`type`, `delta`, `response`, `error`)를 별도 타입으로 파싱하고, `response.output_text.delta`는 text builder에 append하며 `response.completed`는 final response/usage를 반영하는 테스트를 추가한다. + - Required: `services/core/internal/adapters/openai/client.go:133`이 stream 요청의 모든 `400 Bad Request`를 `stream_unsupported`로 분류해 non-streaming으로 재시도한다. invalid model, invalid input, auth/validation성 400도 stream 미지원 metadata로 오분류될 수 있어 silent fallback 방지 요구와 맞지 않는다. 수정: 명시적인 stream unsupported error type/code/message일 때만 fallback하고, 일반 400은 원래 error를 반환하는 negative test를 추가한다. + - Required: `services/core/internal/adapters/openai/client_test.go:228`의 streaming success test가 실제 semantic event shape가 아니라 `output_text`가 있는 pseudo full-response chunk만 검증한다. 수정: `event: response.output_text.delta` / `data: {"type":"response.output_text.delta","delta":"..."}` 및 `response.completed` fixture를 사용해 텍스트/usage/progress를 검증한다. +- 다음 단계: FAIL 후속 plan/review를 작성해 streaming event parser와 fallback 판별을 좁게 수정한다. diff --git a/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_1.log b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_1.log new file mode 100644 index 0000000..e8ffba4 --- /dev/null +++ b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_1.log @@ -0,0 +1,210 @@ + + +# Code Review Reference - REVIEW_AUTHORING_RUNTIME + +> **[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 selected SDD decision or selected Milestone `구현 잠금 > 결정 필요` item, fill `사용자 리뷰 요청` with evidence and stop with active files in place; code-review decides whether to write `USER_REVIEW.md`. +> Do not ask the user directly, present choices in chat, or call `request_user_input` during implementation. +> Finalization (`코드리뷰 결과`, log rename, `complete.log`, archive moves, `코드리뷰 전용 체크리스트`) is review-agent-only. + +## 개요 + +date=2026-06-20 +task=m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress, plan=1, tag=REVIEW_AUTHORING_RUNTIME + +## Roadmap Targets + +- Milestone: `agent-roadmap/phase/agent-ops-mcp-control-plane/milestones/plane-origin-authoring-roundtrip-sync.md` +- Task ids: + - `iop-progress`: `/v1/responses` streaming 지원 여부를 dev IOP 계약으로 확인하고, stream 미지원 환경에서는 `authoring_run_updated_at`과 task metadata 기반 stale monitor로 대체한다. +- Completion mode: check-on-pass + +## Spec Targets + +- SDD: `agent-roadmap/sdd/agent-ops-mcp-control-plane/plane-origin-authoring-roundtrip-sync/SDD.md` +- Acceptance scenarios: + - `S03`: task=`iop-progress`; evidence=`IOP streaming capability check and fallback stale-monitor evidence` +- Completion mode: spec-check-on-pass + +## Archive Evidence Snapshot + +- Current archived plan: `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_0.log` +- Current archived review: `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log` +- Verdict: FAIL +- Required summary: + - `services/core/internal/adapters/openai/client.go:170` parses streaming `data:` payloads only as `responsesResponse`, so real Responses semantic events such as `response.output_text.delta` do not contribute text and `response.completed` does not contribute final response/usage. + - `services/core/internal/adapters/openai/client.go:133` treats any streaming `400 Bad Request` as `stream_unsupported`, which can misclassify ordinary validation/model errors and retry incorrectly. + - `services/core/internal/adapters/openai/client_test.go:228` uses pseudo full-response chunks instead of semantic Responses SSE event fixtures, so the broken stream path still passes. +- Affected files: + - `services/core/internal/adapters/openai/client.go` + - `services/core/internal/adapters/openai/client_test.go` + - `services/core/README.md` if wording changes are needed after behavior is narrowed. +- Verification evidence from failed loop: + - `cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...` passed locally. + - `cd services/core && go test -count=1 ./...` passed locally. + - `git diff --check` passed locally. +- Roadmap/spec carryover: + - Roadmap Task `iop-progress` remains claimed only when this follow-up passes. + - SDD scenario `S03` remains claimed only when semantic stream support and explicit unsupported fallback evidence are both present. +- Narrow reread allowed: + - `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_0.log` + - `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log` + +## 이 파일을 읽는 리뷰 에이전트에게 + +> **[REVIEW AGENT ONLY]** 아래 종결 절차는 코드리뷰 에이전트 전용이다. 구현 에이전트는 이 섹션을 실행하지 않는다. + +각 항목의 구현을 실제 소스 파일과 대조하고, `검증 결과` 섹션의 출력이 코드와 일치하는지 확인하세요. + +--- + +## 구현 항목별 완료 여부 + +| 항목 | 완료 여부 | +|------|---------| +| [REVIEW_AUTHORING_RUNTIME-1] Semantic Responses SSE parser | [x] | +| [REVIEW_AUTHORING_RUNTIME-2] Explicit unsupported fallback only | [x] | +| [REVIEW_AUTHORING_RUNTIME-3] Regression evidence and docs alignment | [x] | + +## 구현 체크리스트 + +- [x] `REVIEW_AUTHORING_RUNTIME-1` Responses streaming parser가 semantic SSE event of `response.output_text.delta`, `response.completed`, `error`를 처리하고 final text/usage/id/model을 보존한다. +- [x] `REVIEW_AUTHORING_RUNTIME-2` stream fallback은 명시적인 stream unsupported error에서만 발생하고, 일반 400/validation error는 `stream_unsupported` metadata 없이 원래 error로 반환된다. +- [x] `REVIEW_AUTHORING_RUNTIME-3` streaming tests가 pseudo `output_text` chunk 대신 semantic SSE fixture와 일반 400 negative case를 검증한다. +- [x] `cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...`, `cd services/core && go test -count=1 ./...`, `git diff --check`를 실행한다. +- [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_cloud_G07_N.log`로 아카이브한다. +- [x] active `PLAN-*-G??.md`를 `plan_cloud_G07_M.log`로 아카이브한다. +- [x] `.gitignore`의 Agent-Ops 관리 block이 `agent-task/**/*.md`와 `agent-task/**/*.log`를 unignore하고 `agent-roadmap/current.md`를 ignore하는지 확인한다. +- [x] PASS이면 `complete.log`를 작성하고 active task 디렉터리를 archive로 이동한다. +- [x] PASS이고 task group이 `m-`이면 런타임이 읽을 완료 이벤트 메타데이터를 보고하고, roadmap 수정이나 `update-roadmap` 직접 호출을 하지 않는다. +- [ ] WARN/FAIL이고 user-review gate가 트리거되지 않았으면 다음 active `PLAN-*-G??.md`와 `CODE_REVIEW-*-G??.md`를 작성하고 `complete.log`를 작성하지 않는다. +- [ ] USER_REVIEW이면 `USER_REVIEW.md`를 작성하고 active `PLAN-*.md`, `CODE_REVIEW-*.md`, `complete.log`를 남기지 않는다. + +## 계획 대비 변경 사항 + +계획과 정확히 동일하게 구현되었습니다. + +## 주요 설계 결정 + +- **Responses SSE Event Parser 구현**: `responsesStreamEvent` DTO를 정의하여 `response.output_text.delta` 이벤트(텍스트 델타), `response.completed` 이벤트(최종 결과 및 Usage 병합), `error` 이벤트(명시적 에러 반환)를 SSE 규격에 맞게 파싱하고, 기존 비정규 chunk 호환을 위한 fallback 로직을 함께 구현하였습니다. +- **스트리밍 미지원 Fallback 범위 좁히기**: 단순 400 Bad Request에 대해 무작정 fallback하지 않고, `isStreamUnsupportedError` 헬퍼를 통해 에러 메시지나 유형 등에 `stream`, `unsupported`와 같은 미지원 시그널이 포함된 경우에만 `stream_unsupported` 메타데이터와 함께 non-streaming으로 fallback하도록 제어했습니다. 일반 validation/model error는 fallback 없이 즉시 에러로 반환됩니다. + +## 사용자 리뷰 요청 + +_기본값은 `없음`이다. 구현 중 새 결정이 필요해 보여도 직접 질문하거나 선택지를 제시하거나 `request_user_input`을 호출하지 않는다. 이 섹션은 선택된 SDD 결정 또는 선택된 Milestone `구현 잠금 > 결정 필요` 항목이 실구현을 차단할 때만 채운다. 외부 환경/secret/서비스 준비, 검증 증거 공백, 반복 실패, 일반 범위 조정은 사용자 리뷰 요청이 아니며 `검증 결과`, `계획 대비 변경 사항`, 또는 code-review의 일반 follow-up plan으로 처리한다._ + +- 상태: 없음 +- 사유 유형: 없음 +- 연결 대상: 없음 +- 결정 필요: 없음 +- 차단 근거: 없음 +- 실행한 검증/명령: 없음 +- 자동 후속 불가 이유: 없음 +- 재개 조건: 없음 + +## 리뷰어를 위한 체크포인트 + +- Semantic Responses SSE event fixture에서 `response.output_text.delta`의 `delta`가 최종 텍스트에 누적되는지 확인한다. +- `response.completed`에서 id/model/usage가 보존되는지 확인한다. +- 일반 400 validation error가 `stream_unsupported`로 오분류되어 non-streaming fallback되지 않는지 확인한다. +- README fallback 조건이 구현과 일치하는지 확인한다. + +## 검증 결과 + +### REVIEW_AUTHORING_RUNTIME-1 중간 검증 + +```sh +$ cd services/core && go test -count=1 ./internal/adapters/openai/... +ok github.com/nomadcode/nomadcode-core/internal/adapters/openai 0.008s +``` + +### REVIEW_AUTHORING_RUNTIME-2 중간 검증 + +```sh +$ cd services/core && go test -count=1 ./internal/adapters/openai/... +ok github.com/nomadcode/nomadcode-core/internal/adapters/openai 0.008s +``` + +### REVIEW_AUTHORING_RUNTIME-3 중간 검증 + +```sh +$ cd services/core && go test -count=1 ./internal/adapters/openai/... +ok github.com/nomadcode/nomadcode-core/internal/adapters/openai 0.008s + +$ git diff --check +(출력 없음 - 성공) +``` + +### 최종 검증 + +```sh +$ cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/... +? github.com/nomadcode/nomadcode-core/internal/model [no test files] +ok github.com/nomadcode/nomadcode-core/internal/adapters/openai 0.008s +ok github.com/nomadcode/nomadcode-core/internal/scheduler 2.034s +ok github.com/nomadcode/nomadcode-core/internal/config 0.003s +``` + +```sh +$ cd services/core && go test -count=1 ./... +ok github.com/nomadcode/nomadcode-core/cmd/plane-smoke 0.003s +ok github.com/nomadcode/nomadcode-core/cmd/server 0.009s +ok github.com/nomadcode/nomadcode-core/internal/adapters/a2a 0.011s +ok github.com/nomadcode/nomadcode-core/internal/adapters/jira 0.014s +ok github.com/nomadcode/nomadcode-core/internal/adapters/mattermost 0.010s +ok github.com/nomadcode/nomadcode-core/internal/adapters/openai 0.015s +ok github.com/nomadcode/nomadcode-core/internal/adapters/plane 0.012s +? github.com/nomadcode/nomadcode-core/internal/agent [no test files] +ok github.com/nomadcode/nomadcode-core/internal/authoring 0.004s +ok github.com/nomadcode/nomadcode-core/internal/config 0.004s +? github.com/nomadcode/nomadcode-core/internal/db [no test files] +ok github.com/nomadcode/nomadcode-core/internal/gitoevents 0.008s +ok github.com/nomadcode/nomadcode-core/internal/gitosync 1.386s +ok github.com/nomadcode/nomadcode-core/internal/http 0.017s +? github.com/nomadcode/nomadcode-core/internal/model [no test files] +ok github.com/nomadcode/nomadcode-core/internal/notification 0.005s +ok github.com/nomadcode/nomadcode-core/internal/projectsync 0.007s +ok github.com/nomadcode/nomadcode-core/internal/protosocket 0.021s +ok github.com/nomadcode/nomadcode-core/internal/roadmapsync 0.006s +ok github.com/nomadcode/nomadcode-core/internal/roadmapsyncpipeline 0.006s +ok github.com/nomadcode/nomadcode-core/internal/scheduler 2.033s +ok github.com/nomadcode/nomadcode-core/internal/storage 0.004s +ok github.com/nomadcode/nomadcode-core/internal/workflow 0.004s +ok github.com/nomadcode/nomadcode-core/internal/workitem 0.006s +ok github.com/nomadcode/nomadcode-core/internal/workitempipeline 0.005s +``` + +```sh +$ git diff --check +(출력 없음 - 성공) +``` + +--- + +> **[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 + - spec conformance: Pass +- 발견된 문제: 없음 +- 다음 단계: PASS 완료 처리. active plan/review를 로그로 아카이브하고 `complete.log` 작성 후 task directory를 archive로 이동한다. diff --git a/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/complete.log b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/complete.log new file mode 100644 index 0000000..0c0355c --- /dev/null +++ b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/complete.log @@ -0,0 +1,51 @@ +# Complete - m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress + +## 완료 일시 + +2026-06-20 + +## 요약 + +Plane-Origin Authoring Roundtrip Sync `iop-progress` follow-up review completed after 2 loops; final verdict PASS. + +## 루프 이력 + +| Plan | Review | Verdict | 메모 | +|------|--------|---------|------| +| `plan_cloud_G07_0.log` | `code_review_cloud_G07_0.log` | FAIL | Real Responses semantic SSE events and explicit unsupported fallback classification were incomplete. | +| `plan_cloud_G07_1.log` | `code_review_cloud_G07_1.log` | PASS | Semantic stream delta/completed/error handling, explicit unsupported fallback, regression tests, and docs alignment verified. | + +## 구현/정리 내용 + +- Added Responses streaming progress support with semantic SSE parsing for `response.output_text.delta`, `response.completed`, and `error`. +- Narrowed stream fallback so non-streaming retry happens only for explicit stream unsupported errors, while ordinary validation/model errors return immediately. +- Added regression tests for semantic SSE success, stream unsupported fallback, and no fallback on validation errors. +- Documented `MODEL_RESPONSES_STREAM` progress metadata and fallback behavior in `services/core/README.md`. + +## 최종 검증 + +- `cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...` - PASS; targeted model/openai/scheduler/config packages passed. +- `cd services/core && go test -count=1 ./...` - PASS; full core package suite passed. +- `git diff --check` - PASS; no whitespace errors. + +## Roadmap Completion + +- Milestone: `agent-roadmap/phase/agent-ops-mcp-control-plane/milestones/plane-origin-authoring-roundtrip-sync.md` +- Completed task ids: + - `iop-progress`: PASS; evidence=`plan_cloud_G07_1.log`, `code_review_cloud_G07_1.log`; verification=`cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...`, `cd services/core && go test -count=1 ./...`, `git diff --check` +- Not completed task ids: 없음 + +## Spec Completion + +- SDD: `agent-roadmap/sdd/agent-ops-mcp-control-plane/plane-origin-authoring-roundtrip-sync/SDD.md` +- Completed scenario ids: + - `S03`: PASS; task=`iop-progress`; evidence=`plan_cloud_G07_1.log`, `code_review_cloud_G07_1.log`; verification=`cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...`, `cd services/core && go test -count=1 ./...`, `git diff --check` +- Not completed scenario ids: 없음 + +## 잔여 Nit + +- 없음 + +## 후속 작업 + +- 없음 diff --git a/agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/PLAN-cloud-G07.md b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_0.log similarity index 100% rename from agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/PLAN-cloud-G07.md rename to agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_0.log diff --git a/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_1.log b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_1.log new file mode 100644 index 0000000..53141b6 --- /dev/null +++ b/agent-task/archive/2026/06/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_1.log @@ -0,0 +1,171 @@ + + +# Implementation Plan - REVIEW_AUTHORING_RUNTIME + +## 이 파일을 읽는 구현 에이전트에게 + +이 plan은 code-review FAIL 후속 루프다. 아래 `Archive Evidence Snapshot`에 명시된 이전 plan/review log만 좁게 참고하고, `agent-task/archive/**`는 읽지 않는다. 구현 중 직접 사용자에게 질문하지 않는다. 선택된 SDD 결정 또는 Milestone lock 결정이 실제 구현을 차단할 때만 active `CODE_REVIEW-cloud-G07.md`의 `사용자 리뷰 요청` 섹션에 근거와 재개 조건을 기록한다. + +## Roadmap Targets + +- Milestone: `agent-roadmap/phase/agent-ops-mcp-control-plane/milestones/plane-origin-authoring-roundtrip-sync.md` +- Task ids: + - `iop-progress`: `/v1/responses` streaming 지원 여부를 dev IOP 계약으로 확인하고, stream 미지원 환경에서는 `authoring_run_updated_at`과 task metadata 기반 stale monitor로 대체한다. +- Completion mode: check-on-pass + +## Spec Targets + +- SDD: `agent-roadmap/sdd/agent-ops-mcp-control-plane/plane-origin-authoring-roundtrip-sync/SDD.md` +- Acceptance scenarios: + - `S03`: task=`iop-progress`; evidence=`IOP streaming capability check and fallback stale-monitor evidence` +- Completion mode: spec-check-on-pass + +## Archive Evidence Snapshot + +- Current archived plan: `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_0.log` +- Current archived review: `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log` +- Verdict: FAIL +- Required summary: + - `services/core/internal/adapters/openai/client.go:170` parses streaming `data:` payloads only as `responsesResponse`, so real Responses semantic events such as `response.output_text.delta` do not contribute text and `response.completed` does not contribute final response/usage. + - `services/core/internal/adapters/openai/client.go:133` treats any streaming `400 Bad Request` as `stream_unsupported`, which can misclassify ordinary validation/model errors and retry incorrectly. + - `services/core/internal/adapters/openai/client_test.go:228` uses pseudo full-response chunks instead of semantic Responses SSE event fixtures, so the broken stream path still passes. +- Affected files: + - `services/core/internal/adapters/openai/client.go` + - `services/core/internal/adapters/openai/client_test.go` + - `services/core/README.md` if wording changes are needed after behavior is narrowed. +- Verification evidence from failed loop: + - `cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...` passed locally. + - `cd services/core && go test -count=1 ./...` passed locally. + - `git diff --check` passed locally. +- Roadmap/spec carryover: + - Roadmap Task `iop-progress` remains claimed only when this follow-up passes. + - SDD scenario `S03` remains claimed only when semantic stream support and explicit unsupported fallback evidence are both present. +- Narrow reread allowed: + - `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/plan_cloud_G07_0.log` + - `agent-task/m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress/code_review_cloud_G07_0.log` + +## 범위 결정 근거 + +- 포함: OpenAI-compatible Responses SSE semantic event parsing, explicit stream unsupported fallback classification, regression tests and minimal docs alignment. +- 제외: scheduler stale threshold policy, live Plane full-cycle evidence, push-only flow, identity/projection sync. +- Lane/grade: `cloud-G07`. 후속 범위는 좁지만 외부 Responses streaming semantic event 계약을 정확히 반영해야 하고, 이전 루프의 unit tests가 실제 계약을 대표하지 못한 evidence-trust 실패가 있었다. + +## 구현 체크리스트 + +- [ ] `REVIEW_AUTHORING_RUNTIME-1` Responses streaming parser가 semantic SSE event의 `response.output_text.delta`, `response.completed`, `error`를 처리하고 final text/usage/id/model을 보존한다. +- [ ] `REVIEW_AUTHORING_RUNTIME-2` stream fallback은 명시적인 stream unsupported error에서만 발생하고, 일반 400/validation error는 `stream_unsupported` metadata 없이 원래 error로 반환된다. +- [ ] `REVIEW_AUTHORING_RUNTIME-3` streaming tests가 pseudo `output_text` chunk 대신 semantic SSE fixture와 일반 400 negative case를 검증한다. +- [ ] `cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/...`, `cd services/core && go test -count=1 ./...`, `git diff --check`를 실행한다. +- [ ] CODE_REVIEW-*-G??.md의 구현 에이전트 소유 섹션을 실제 구현 내용과 검증 출력으로 채운다. 이 항목이 완료되기 전에는 구현이 완료된 것이 아니다. + +## 의존 관계 및 구현 순서 + +- Runtime dependency from directory name remains `02+01_iop_progress`. +- 추가 hidden dependency를 만들지 않는다. +- 이전 loop log는 위 `Archive Evidence Snapshot`에 명시된 두 파일만 필요한 경우 좁게 읽는다. + +### [REVIEW_AUTHORING_RUNTIME-1] Semantic Responses SSE parser + +문제: + +현재 streaming parser는 `data:` payload를 `responsesResponse`로 unmarshal한 뒤 `output_text`/`output`만 읽는다. 실제 Responses stream은 semantic event이며 text delta는 `type=response.output_text.delta`와 `delta` 필드로 온다. final usage와 response id/model은 `response.completed`의 nested `response`에서 올 수 있다. + +해결 방법: + +- streaming 전용 event DTO를 추가한다. + - fields: `Type`, `Delta`, `Response`, `Error`. + - `response.output_text.delta`: `Delta`를 text builder에 append하고 progress callback을 호출한다. + - `response.completed`: nested `Response`를 final response에 merge하고 usage/id/model을 반영한다. + - `error`: 명시 error를 반환한다. +- 기존 pseudo full-response chunk 호환이 필요하면 fallback branch로 유지하되, semantic event 처리가 우선이어야 한다. +- JSON parse 실패를 무조건 삼키지 말고 malformed data가 final output을 만들 수 없으면 error evidence가 남도록 한다. + +수정 파일 및 체크리스트: + +- [ ] `services/core/internal/adapters/openai/client.go`: semantic stream event DTO와 merge 로직 추가. +- [ ] `services/core/internal/adapters/openai/client_test.go`: semantic delta/completed fixture 기반 테스트로 검증. + +중간 검증: + +```sh +cd services/core && go test -count=1 ./internal/adapters/openai/... +``` + +### [REVIEW_AUTHORING_RUNTIME-2] Explicit unsupported fallback only + +문제: + +stream 요청에서 모든 `400 Bad Request`를 `stream_unsupported`로 분류한다. 이러면 invalid model/input 같은 일반 validation error도 stream 미지원으로 metadata가 남고 non-streaming 재시도가 발생할 수 있다. + +해결 방법: + +- fallback predicate를 별도 helper로 만든다. +- status code만으로 fallback하지 않는다. +- error message/type/code에 `stream`, `streaming`, `unsupported`, `not supported`, `unsupported_parameter` 같은 stream unsupported 신호가 있을 때만 fallback한다. +- 일반 400은 원래 `responseError`를 반환하고 progress callback을 호출하지 않는다. + +수정 파일 및 체크리스트: + +- [ ] `services/core/internal/adapters/openai/client.go`: fallback 판별 helper 추가. +- [ ] `services/core/internal/adapters/openai/client_test.go`: stream unsupported는 fallback, 일반 400은 no fallback/no progress로 검증. + +중간 검증: + +```sh +cd services/core && go test -count=1 ./internal/adapters/openai/... +``` + +### [REVIEW_AUTHORING_RUNTIME-3] Regression evidence and docs alignment + +문제: + +기존 streaming test는 `output_text`를 가진 pseudo chunk만 검증해서 실제 semantic event parser 결함을 잡지 못했다. + +해결 방법: + +- `TestGenerateStreamingSuccess`를 semantic SSE fixture로 바꾸거나 새 테스트를 추가한다. +- fixture는 `event: response.output_text.delta` line과 `data: {"type":"response.output_text.delta","delta":"..."}` line을 함께 포함한다. +- final usage는 `data: {"type":"response.completed","response":{...}}`에서 검증한다. +- README 표현이 “400이면 fallback”처럼 너무 넓으면 명시 stream unsupported error일 때만 fallback한다고 좁힌다. + +수정 파일 및 체크리스트: + +- [ ] `services/core/internal/adapters/openai/client_test.go`: semantic SSE fixture와 negative 400 test 추가/갱신. +- [ ] `services/core/README.md`: fallback 조건 문구를 구현과 일치시킨다. + +중간 검증: + +```sh +cd services/core && go test -count=1 ./internal/adapters/openai/... +git diff --check +``` + +## 수정 파일 요약 + +| 파일 | 항목 | +|------|------| +| `services/core/internal/adapters/openai/client.go` | REVIEW_AUTHORING_RUNTIME-1, REVIEW_AUTHORING_RUNTIME-2 | +| `services/core/internal/adapters/openai/client_test.go` | REVIEW_AUTHORING_RUNTIME-1, REVIEW_AUTHORING_RUNTIME-2, REVIEW_AUTHORING_RUNTIME-3 | +| `services/core/README.md` | REVIEW_AUTHORING_RUNTIME-3 | + +## 최종 검증 + +```sh +cd services/core && go test -count=1 ./internal/model/... ./internal/adapters/openai/... ./internal/scheduler/... ./internal/config/... +``` + +Expected: pass with semantic stream and explicit unsupported fallback tests. + +```sh +cd services/core && go test -count=1 ./... +``` + +Expected: pass for full core regression coverage. + +```sh +git diff --check +``` + +Expected: no output. + +모든 코드 변경 완료 후 반드시 `CODE_REVIEW-*-G??.md`의 구현 에이전트 소유 섹션을 채운다. 이 파일 작성이 구현의 마지막 단계다. diff --git a/services/core/README.md b/services/core/README.md index 58caf95..1d821e2 100644 --- a/services/core/README.md +++ b/services/core/README.md @@ -90,7 +90,7 @@ MODEL_TIMEOUT_SEC="900" \ ./bin/run ``` -모델 호출은 OpenAI-compatible Responses API의 non-streaming `POST /v1/responses` 형식을 사용합니다. `MODEL_BASE_URL`이 `/v1`까지만 가리키면 Core가 `/responses`를 붙여 호출합니다. NomadCode의 task/workspace/session 문맥은 OpenAI-compatible 표면을 깨는 별도 top-level wrapper가 아니라 `metadata` 확장으로 전달합니다. workspace-bound authoring에서는 slot checkout path가 `metadata.workspace` flat string으로 전달됩니다. direct Ollama 호환 경로에서는 `MODEL_CONTEXT_SIZE`를 Ollama 전용 option인 `options.num_ctx`로 전달할 수 있지만, 기본 dev IOP Edge 호출에서는 `0`으로 둡니다. +모델 호출은 OpenAI-compatible Responses API의 `POST /v1/responses` 형식을 사용합니다. 기본적으로 non-streaming 형식을 사용하지만, `MODEL_RESPONSES_STREAM` 설정이 `true`로 설정된 경우 `stream: true`로 스트리밍 호출(SSE)을 시도합니다. 스트리밍 호출이 실행되는 동안 클라이언트는 수신된 chunk를 파싱하면서 실시간으로 progress 콜백(`OnProgress`)을 호출하여 `authoring_progress_mode=streaming`, `authoring_progress_reason=receiving stream chunks` 및 최신 타임스탬프(`authoring_run_updated_at`)를 태스크 메타데이터에 머지합니다. 만약 엔드포인트(IOP)가 스트리밍을 지원하지 않아 명시적인 stream unsupported 에러(에러 메시지나 유형 등에 `stream`, `unsupported` 등의 신호 포함)를 응답하는 경우, silent failure 없이 자동으로 1회 non-streaming 호출로 fallback하며 `authoring_progress_mode=stream_unsupported`, `authoring_progress_reason=<원래 에러 메시지>` 메타데이터를 저장하고 재시도합니다. 일반적인 유효성 검증이나 모델 이름 오류 같은 일반 에러는 fallback 없이 즉시 에러로 반환됩니다. `MODEL_RESPONSES_STREAM`이 `false`인 경우 처음부터 non-streaming으로 동작하며 `authoring_progress_mode=non_streaming` 메타데이터가 기록됩니다. `MODEL_BASE_URL`이 `/v1`까지만 가리키면 Core가 `/responses`를 붙여 호출합니다. NomadCode의 task/workspace/session 문맥은 OpenAI-compatible 표면을 깨는 별도 top-level wrapper가 아니라 `metadata` 확장으로 전달합니다. workspace-bound authoring에서는 slot checkout path가 `metadata.workspace` flat string으로 전달됩니다. direct Ollama 호환 경로에서는 `MODEL_CONTEXT_SIZE`를 Ollama 전용 option인 `options.num_ctx`로 전달할 수 있지만, 기본 dev IOP Edge 호출에서는 `0`으로 둡니다. A2A agent 호출 인터페이스는 JSON-RPC 2.0 `message/send`, `tasks/get`, `tasks/cancel`을 우선 지원합니다. `A2A_EDGE_URL`을 설정하면 worker는 A2A `message/send`를 blocking 호출하고, 완료된 task/message 응답만 local task completion으로 반영합니다. 기존 `A2A_AGENT_URL`도 fallback alias로 받습니다. 현재 NomadCode의 기본 실행 경로는 OpenAI-compatible Responses 호출이며, A2A는 후속 agent delegation 작업에서 기본화 여부를 다시 결정합니다. @@ -254,6 +254,7 @@ NomadCode Core의 비동기 작업 재시도 및 타임아웃 처리는 다음 - 실행 중인 authoring 작업은 `authoring_run_state=in_progress`와 `authoring_run_updated_at`이 함께 기록됩니다. 이 타임스탬프가 `AUTHORING_STALE_AFTER_SEC`를 초과하면 stale로 관찰됩니다. - queue 대기 중인 authoring 작업은 `task.UpdatedAt`을 기준 타임스탬프로 사용하여 동일한 `AUTHORING_STALE_AFTER_SEC` 임계값에 비교할 수 있습니다. queue 대기 중에는 `WORKFLOW_TASK_TIMEOUT_SEC`가 적용되지 않으므로 조기 실패가 발생하지 않습니다. - stale로 판단된 작업에 대한 조치(알림, 재-Enqueue 등)는 후속 모니터링 컴포넌트의 책임이며, 이 설정은 판단 기준만 제공합니다. + - `MODEL_RESPONSES_STREAM` 설정으로 스트리밍이 활성화되면 실시간 진행 상황을 `authoring_progress_mode`, `authoring_progress_reason`, `authoring_run_updated_at` 메타데이터로 갱신하여 stale monitor가 진행 상황 중단을 파악할 수 있도록 돕습니다. 명시적인 스트리밍 미지원 에러가 반환되는 환경에서는 fallback을 통해 `authoring_progress_mode=stream_unsupported`와 fallback 사유를 남깁니다. 처음부터 비스트리밍일 경우 `authoring_progress_mode=non_streaming`으로 기록됩니다. - **retry 메타데이터 (Retry Metadata)**: - 모든 실패 작업(generic 및 authoring)은 `FailTaskWithMetadata` 경로를 통해 `retryable` 메타데이터를 기록합니다. `retryable=true`이면 `attempt < DefaultTaskMaxAttempts`이므로 River가 다음 시도를 스케줄링합니다. - authoring 태스크 실패 시 `authoring_run_state=failed`, `authoring_failure_type`, `authoring_failure_category`, `authoring_run_updated_at`, `retryable`이 함께 기록됩니다. diff --git a/services/core/cmd/server/main.go b/services/core/cmd/server/main.go index 53ff39c..1223f9f 100644 --- a/services/core/cmd/server/main.go +++ b/services/core/cmd/server/main.go @@ -78,11 +78,12 @@ func run(logger *slog.Logger) error { protosocket.NewTaskEventBroadcaster(protoSocketServer), ) modelClient := modelopenai.NewClient(modelopenai.Config{ - BaseURL: cfg.ModelBaseURL, - APIKey: cfg.ModelAPIKey, - Model: cfg.ModelName, - ContextSize: cfg.ModelContextSize, - TimeoutSec: cfg.ModelTimeoutSec, + BaseURL: cfg.ModelBaseURL, + APIKey: cfg.ModelAPIKey, + Model: cfg.ModelName, + ContextSize: cfg.ModelContextSize, + TimeoutSec: cfg.ModelTimeoutSec, + ModelResponsesStream: cfg.ModelResponsesStream, }, logger) var agentClient agent.Client diff --git a/services/core/internal/adapters/openai/client.go b/services/core/internal/adapters/openai/client.go index 48e45bc..0a5f870 100644 --- a/services/core/internal/adapters/openai/client.go +++ b/services/core/internal/adapters/openai/client.go @@ -1,6 +1,7 @@ package openai import ( + "bufio" "bytes" "context" "encoding/json" @@ -22,11 +23,12 @@ const ( ) type Config struct { - BaseURL string - APIKey string - Model string - ContextSize int - TimeoutSec int + BaseURL string + APIKey string + Model string + ContextSize int + TimeoutSec int + ModelResponsesStream bool } type Client struct { @@ -63,12 +65,18 @@ func (c *Client) Generate(ctx context.Context, input model.GenerateInput) (model return model.GenerateResult{}, fmt.Errorf("model input is required") } + streamRequested := c.cfg.ModelResponsesStream + respResult, err := c.executeGenerate(ctx, input, endpoint, modelName, streamRequested) + return respResult, err +} + +func (c *Client) executeGenerate(ctx context.Context, input model.GenerateInput, endpoint, modelName string, stream bool) (model.GenerateResult, error) { reqBody := responsesRequest{ Model: modelName, Input: input.Input, Instructions: input.Instructions, Metadata: buildRequestMetadata(input), - Stream: false, + Stream: stream, Temperature: input.Temperature, TopP: input.TopP, MaxOutputTokens: input.MaxOutputTokens, @@ -89,6 +97,9 @@ func (c *Client) Generate(ctx context.Context, input model.GenerateInput) (model if c.cfg.APIKey != "" { req.Header.Set("Authorization", "Bearer "+c.cfg.APIKey) } + if stream { + req.Header.Set("Accept", "text/event-stream") + } if c.logger != nil { c.logger.Info( @@ -96,22 +107,185 @@ func (c *Client) Generate(ctx context.Context, input model.GenerateInput) (model "endpoint", endpoint, "model", modelName, "num_ctx", c.cfg.ContextSize, + "stream", stream, ) } resp, err := c.httpClient.Do(req) if err != nil { + if stream && isStreamUnsupportedError(err) { + if input.OnProgress != nil { + input.OnProgress(model.GenerateProgress{ + Mode: "stream_unsupported", + Reason: err.Error(), + LastEventTime: time.Now().UTC(), + }) + } + return c.executeGenerate(ctx, input, endpoint, modelName, false) + } return model.GenerateResult{}, err } defer resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + raw, _ := io.ReadAll(io.LimitReader(resp.Body, 8<<20)) + errRes := responseError(resp.StatusCode, raw) + if stream && isStreamUnsupportedError(errRes) { + if input.OnProgress != nil { + input.OnProgress(model.GenerateProgress{ + Mode: "stream_unsupported", + Reason: errRes.Error(), + LastEventTime: time.Now().UTC(), + }) + } + return c.executeGenerate(ctx, input, endpoint, modelName, false) + } + return model.GenerateResult{}, errRes + } + + if stream { + reader := bufio.NewReader(resp.Body) + var finalResponse responsesResponse + var textBuilder strings.Builder + var currentEvent string + for { + line, err := reader.ReadString('\n') + if err != nil { + if err == io.EOF { + break + } + return model.GenerateResult{}, err + } + line = strings.TrimSpace(line) + if line == "" { + continue + } + if strings.HasPrefix(line, "event: ") { + currentEvent = strings.TrimPrefix(line, "event: ") + continue + } + if !strings.HasPrefix(line, "data: ") { + continue + } + dataStr := strings.TrimPrefix(line, "data: ") + if dataStr == "[DONE]" { + break + } + + var ev responsesStreamEvent + errEv := json.Unmarshal([]byte(dataStr), &ev) + if errEv == nil && (ev.Type != "" || ev.Error != nil || ev.Response != nil) { + eventType := ev.Type + if eventType == "" { + eventType = currentEvent + } + + switch eventType { + case "response.output_text.delta": + textBuilder.WriteString(ev.Delta) + if input.OnProgress != nil { + input.OnProgress(model.GenerateProgress{ + Mode: "streaming", + Reason: "receiving stream chunks", + LastEventTime: time.Now().UTC(), + }) + } + case "response.completed": + if ev.Response != nil { + finalResponse.ID = ev.Response.ID + finalResponse.Model = ev.Response.Model + if ev.Response.Usage != nil { + finalResponse.Usage = ev.Response.Usage + } + } + case "error": + if ev.Error != nil { + return model.GenerateResult{}, fmt.Errorf("responses API stream error: %s", ev.Error.Message) + } + return model.GenerateResult{}, fmt.Errorf("responses API stream error: unknown error") + } + } else { + var chunk responsesResponse + errChunk := json.Unmarshal([]byte(dataStr), &chunk) + if errChunk == nil { + if chunk.ID != "" { + finalResponse.ID = chunk.ID + } + if chunk.Model != "" { + finalResponse.Model = chunk.Model + } + if chunk.Usage != nil { + if finalResponse.Usage == nil { + finalResponse.Usage = &responsesUsage{} + } + if chunk.Usage.InputTokens != 0 { + finalResponse.Usage.InputTokens = chunk.Usage.InputTokens + } + if chunk.Usage.OutputTokens != 0 { + finalResponse.Usage.OutputTokens = chunk.Usage.OutputTokens + } + if chunk.Usage.TotalTokens != 0 { + finalResponse.Usage.TotalTokens = chunk.Usage.TotalTokens + } + if chunk.Usage.PromptTokens != 0 { + finalResponse.Usage.PromptTokens = chunk.Usage.PromptTokens + } + if chunk.Usage.CompletionTokens != 0 { + finalResponse.Usage.CompletionTokens = chunk.Usage.CompletionTokens + } + } + chunkText := chunk.text() + if chunkText != "" { + textBuilder.WriteString(chunkText) + } + if input.OnProgress != nil { + input.OnProgress(model.GenerateProgress{ + Mode: "streaming", + Reason: "receiving stream chunks", + LastEventTime: time.Now().UTC(), + }) + } + } else { + return model.GenerateResult{}, fmt.Errorf("failed to parse SSE data: %s (errEv: %v, errChunk: %v)", dataStr, errEv, errChunk) + } + } + currentEvent = "" + } + + finalResponse.OutputText = textBuilder.String() + raw, _ := json.Marshal(finalResponse) + + usage := model.Usage{} + if finalResponse.Usage != nil { + usage.InputTokens = firstNonZero(finalResponse.Usage.InputTokens, finalResponse.Usage.PromptTokens) + usage.OutputTokens = firstNonZero(finalResponse.Usage.OutputTokens, finalResponse.Usage.CompletionTokens) + usage.TotalTokens = finalResponse.Usage.TotalTokens + if usage.TotalTokens == 0 { + usage.TotalTokens = usage.InputTokens + usage.OutputTokens + } + } + + return model.GenerateResult{ + ID: finalResponse.ID, + Model: firstNonEmpty(finalResponse.Model, modelName), + Text: finalResponse.text(), + Usage: usage, + Raw: json.RawMessage(raw), + }, nil + } + + if input.OnProgress != nil { + input.OnProgress(model.GenerateProgress{ + Mode: "non_streaming", + Reason: "non-streaming mode", + LastEventTime: time.Now().UTC(), + }) + } + raw, err := io.ReadAll(io.LimitReader(resp.Body, 8<<20)) if err != nil { return model.GenerateResult{}, err } - if resp.StatusCode < 200 || resp.StatusCode >= 300 { - return model.GenerateResult{}, responseError(resp.StatusCode, raw) - } var parsed responsesResponse if err := json.Unmarshal(raw, &parsed); err != nil { @@ -280,3 +454,38 @@ func firstNonEmpty(values ...string) string { } return "" } + +type responsesStreamEvent struct { + Type string `json:"type"` + Delta string `json:"delta,omitempty"` + Response *struct { + ID string `json:"id,omitempty"` + Model string `json:"model,omitempty"` + Usage *responsesUsage `json:"usage,omitempty"` + } `json:"response,omitempty"` + Error *struct { + Message string `json:"message,omitempty"` + Type string `json:"type,omitempty"` + Code string `json:"code,omitempty"` + } `json:"error,omitempty"` +} + +func isStreamUnsupportedError(err error) bool { + if err == nil { + return false + } + errStr := strings.ToLower(err.Error()) + signals := []string{ + "stream", + "streaming", + "unsupported", + "not supported", + "unsupported_parameter", + } + for _, sig := range signals { + if strings.Contains(errStr, sig) { + return true + } + } + return false +} diff --git a/services/core/internal/adapters/openai/client_test.go b/services/core/internal/adapters/openai/client_test.go index 02eabe2..ff71a63 100644 --- a/services/core/internal/adapters/openai/client_test.go +++ b/services/core/internal/adapters/openai/client_test.go @@ -5,6 +5,7 @@ import ( "encoding/json" "net/http" "net/http/httptest" + "strings" "testing" "time" @@ -223,3 +224,165 @@ func TestResponsesURL(t *testing.T) { } } } + +func TestGenerateStreamingSuccess(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + _, _ = w.Write([]byte("event: response.output_text.delta\n")) + _, _ = w.Write([]byte("data: {\"type\":\"response.output_text.delta\",\"delta\":\"hello\"}\n\n")) + _, _ = w.Write([]byte("event: response.output_text.delta\n")) + _, _ = w.Write([]byte("data: {\"type\":\"response.output_text.delta\",\"delta\":\" world\"}\n\n")) + _, _ = w.Write([]byte("event: response.completed\n")) + _, _ = w.Write([]byte("data: {\"type\":\"response.completed\",\"response\":{\"id\":\"chunk-1\",\"model\":\"m-stream\",\"usage\":{\"prompt_tokens\":3,\"completion_tokens\":4,\"total_tokens\":7}}}\n\n")) + _, _ = w.Write([]byte("data: [DONE]\n\n")) + })) + defer server.Close() + + client := NewClient(Config{ + BaseURL: server.URL, + Model: "m-stream", + ModelResponsesStream: true, + }, nil) + + var progressEvents []model.GenerateProgress + result, err := client.Generate(context.Background(), model.GenerateInput{ + Input: "say hello", + OnProgress: func(p model.GenerateProgress) { + progressEvents = append(progressEvents, p) + }, + }) + if err != nil { + t.Fatalf("Generate returned error: %v", err) + } + + if result.Text != "hello world" { + t.Fatalf("expected 'hello world', got %q", result.Text) + } + if result.ID != "chunk-1" { + t.Fatalf("expected ID 'chunk-1', got %q", result.ID) + } + if result.Usage.TotalTokens != 7 { + t.Fatalf("expected TotalTokens 7, got %d", result.Usage.TotalTokens) + } + if len(progressEvents) != 2 { + t.Fatalf("expected 2 progress events, got %d", len(progressEvents)) + } + for _, p := range progressEvents { + if p.Mode != "streaming" { + t.Errorf("expected progress mode 'streaming', got %q", p.Mode) + } + } +} + +func TestGenerateStreamingFallbackToNonStreaming(t *testing.T) { + var requestCount int + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + requestCount++ + var reqBody map[string]any + if err := json.NewDecoder(r.Body).Decode(&reqBody); err != nil { + t.Fatalf("decode request: %v", err) + } + + if requestCount == 1 { + if streamVal, ok := reqBody["stream"].(bool); !ok || !streamVal { + t.Fatalf("expected first request to have stream: true, got %v", reqBody["stream"]) + } + w.WriteHeader(http.StatusBadRequest) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"error": {"message": "Streaming is not supported for this model", "type": "unsupported_parameter"}}`)) + return + } + + if streamVal, ok := reqBody["stream"].(bool); !ok || streamVal { + t.Fatalf("expected second request to have stream: false, got %v", reqBody["stream"]) + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{ + "id": "resp-fallback", + "model": "m-fallback", + "output_text": "hello from fallback", + "usage": { + "total_tokens": 10 + } + }`)) + })) + defer server.Close() + + client := NewClient(Config{ + BaseURL: server.URL, + Model: "m-stream", + ModelResponsesStream: true, + }, nil) + + var progressEvents []model.GenerateProgress + result, err := client.Generate(context.Background(), model.GenerateInput{ + Input: "say hello", + OnProgress: func(p model.GenerateProgress) { + progressEvents = append(progressEvents, p) + }, + }) + if err != nil { + t.Fatalf("Generate returned error: %v", err) + } + + if result.Text != "hello from fallback" { + t.Fatalf("expected 'hello from fallback', got %q", result.Text) + } + if result.ID != "resp-fallback" { + t.Fatalf("expected ID 'resp-fallback', got %q", result.ID) + } + if requestCount != 2 { + t.Fatalf("expected 2 requests, got %d", requestCount) + } + + if len(progressEvents) != 2 { + t.Fatalf("expected 2 progress events, got %d", len(progressEvents)) + } + if progressEvents[0].Mode != "stream_unsupported" { + t.Errorf("expected first progress mode 'stream_unsupported', got %q", progressEvents[0].Mode) + } + if !strings.Contains(progressEvents[0].Reason, "Streaming is not supported") { + t.Errorf("expected reason to contain 'Streaming is not supported', got %q", progressEvents[0].Reason) + } + if progressEvents[1].Mode != "non_streaming" { + t.Errorf("expected second progress mode 'non_streaming', got %q", progressEvents[1].Mode) + } +} + +func TestGenerateStreamingNoFallbackOnValidationError(t *testing.T) { + var requestCount int + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + requestCount++ + w.WriteHeader(http.StatusBadRequest) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"error": {"message": "Invalid model name: codex-xyz", "type": "invalid_request_error"}}`)) + })) + defer server.Close() + + client := NewClient(Config{ + BaseURL: server.URL, + Model: "codex-xyz", + ModelResponsesStream: true, + }, nil) + + var progressEvents []model.GenerateProgress + _, err := client.Generate(context.Background(), model.GenerateInput{ + Input: "say hello", + OnProgress: func(p model.GenerateProgress) { + progressEvents = append(progressEvents, p) + }, + }) + if err == nil { + t.Fatal("expected Generate to fail") + } + + if !strings.Contains(err.Error(), "Invalid model name") { + t.Fatalf("unexpected error message: %v", err) + } + if requestCount != 1 { + t.Fatalf("expected only 1 request (no fallback), got %d", requestCount) + } + if len(progressEvents) != 0 { + t.Fatalf("expected no progress events, got %d", len(progressEvents)) + } +} diff --git a/services/core/internal/config/config.go b/services/core/internal/config/config.go index 1c0b526..1d542de 100644 --- a/services/core/internal/config/config.go +++ b/services/core/internal/config/config.go @@ -48,6 +48,7 @@ type Config struct { GitoDevelopRepoPath string GitoRemoteName string RoadmapCreationTodoStateID string + ModelResponsesStream bool } // GitoBranchEventsEnabled reports whether the Gito branch event consumer should @@ -115,6 +116,7 @@ func Load() Config { GitoDevelopRepoPath: os.Getenv("GITO_DEVELOP_REPO_PATH"), GitoRemoteName: getEnv("GITO_REMOTE_NAME", "origin"), RoadmapCreationTodoStateID: firstEnv("ROADMAP_CREATION_TODO_STATE_ID", "PLANE_TODO_STATE_ID"), + ModelResponsesStream: getEnvBool("MODEL_RESPONSES_STREAM", false), } } @@ -146,3 +148,15 @@ func getEnvInt(key string, fallback int) int { } return parsed } + +func getEnvBool(key string, fallback bool) bool { + value := os.Getenv(key) + if value == "" { + return fallback + } + parsed, err := strconv.ParseBool(value) + if err != nil { + return fallback + } + return parsed +} diff --git a/services/core/internal/config/config_test.go b/services/core/internal/config/config_test.go index 8141c72..d43f7e3 100644 --- a/services/core/internal/config/config_test.go +++ b/services/core/internal/config/config_test.go @@ -47,11 +47,22 @@ func TestLoadModelAndA2ADefaults(t *testing.T) { if cfg.ModelTimeoutSec != 900 { t.Fatalf("ModelTimeoutSec: got %d", cfg.ModelTimeoutSec) } + if cfg.ModelResponsesStream != false { + t.Fatalf("ModelResponsesStream: got %t, want false", cfg.ModelResponsesStream) + } if cfg.A2ATimeoutSec != 300 { t.Fatalf("A2ATimeoutSec: got %d", cfg.A2ATimeoutSec) } } +func TestConfigLoadsModelResponsesStreamOverride(t *testing.T) { + t.Setenv("MODEL_RESPONSES_STREAM", "true") + cfg := Load() + if cfg.ModelResponsesStream != true { + t.Fatalf("ModelResponsesStream override: got %t, want true", cfg.ModelResponsesStream) + } +} + func TestConfigLoadsAuthoringStaleAfterSecDefault(t *testing.T) { cfg := Load() if cfg.AuthoringStaleAfterSec != 1200 { diff --git a/services/core/internal/model/model.go b/services/core/internal/model/model.go index b059728..acb5593 100644 --- a/services/core/internal/model/model.go +++ b/services/core/internal/model/model.go @@ -3,12 +3,19 @@ package model import ( "context" "encoding/json" + "time" ) type Client interface { Generate(ctx context.Context, input GenerateInput) (GenerateResult, error) } +type GenerateProgress struct { + Mode string // "streaming", "stream_unsupported", "non_streaming" + Reason string + LastEventTime time.Time +} + // WorkspaceMetadata carries workspace-bound authoring context. On the wire, // only the Path is serialized as a flat string to metadata.workspace in the // /v1/responses payload. Other fields are used on the client-side/authoring @@ -33,6 +40,7 @@ type GenerateInput struct { MaxOutputTokens int Temperature *float64 TopP *float64 + OnProgress func(GenerateProgress) } type GenerateResult struct { diff --git a/services/core/internal/scheduler/jobs.go b/services/core/internal/scheduler/jobs.go index 4d0a656..48d1ef4 100644 --- a/services/core/internal/scheduler/jobs.go +++ b/services/core/internal/scheduler/jobs.go @@ -169,6 +169,21 @@ func (w *TaskWorker) runTask(ctx context.Context, task storage.Task) (json.RawMe if w.Model == nil { return nil, "", fmt.Errorf("model client is required for authoring tasks") } + var progressMode string + var progressReason string + generateInput.OnProgress = func(p model.GenerateProgress) { + progressMode = p.Mode + progressReason = p.Reason + updates := map[string]any{ + workflow.MetadataKeyAuthoringProgressMode: p.Mode, + workflow.MetadataKeyAuthoringProgressReason: p.Reason, + workflow.MetadataKeyAuthoringRunUpdatedAt: p.LastEventTime.Format(time.RFC3339), + } + if _, err := w.Lifecycle.MergeTaskMetadata(ctx, task.ID, updates); err != nil && w.Logger != nil { + w.Logger.Warn("authoring progress metadata merge failed", "task_id", task.ID, "error", err) + } + } + generated, err := w.Model.Generate(ctx, generateInput) if err != nil { return nil, "", err @@ -177,12 +192,14 @@ func (w *TaskWorker) runTask(ctx context.Context, task storage.Task) (json.RawMe BridgeSuccess: true, }) result := map[string]any{ - "message": generated.Text, - "mode": "authoring_run", - "model": generated.Model, - "response_id": generated.ID, - "authoring_run_state": decision.State, - "authoring_run_updated_at": time.Now().UTC().Format(time.RFC3339), + "message": generated.Text, + "mode": "authoring_run", + "model": generated.Model, + "response_id": generated.ID, + "authoring_run_state": decision.State, + "authoring_run_updated_at": time.Now().UTC().Format(time.RFC3339), + workflow.MetadataKeyAuthoringProgressMode: progressMode, + workflow.MetadataKeyAuthoringProgressReason: progressReason, "usage": map[string]int{ "input_tokens": generated.Usage.InputTokens, "output_tokens": generated.Usage.OutputTokens, diff --git a/services/core/internal/scheduler/jobs_test.go b/services/core/internal/scheduler/jobs_test.go index a531488..a464366 100644 --- a/services/core/internal/scheduler/jobs_test.go +++ b/services/core/internal/scheduler/jobs_test.go @@ -1067,3 +1067,78 @@ func TestWorkKeepsAuthoringTaskRunningUntilDevelopMatch(t *testing.T) { t.Errorf("expected MergeTaskMetadata to be called with authoring_run_state=in_progress, got %v", fakeLifecycle.mergedMetadata) } } + +func TestRunTaskAuthoringProgressCallback(t *testing.T) { + fakeLifecycle := &fakeTaskLifecycle{ + task: storage.Task{ + ID: "task-auth-progress", + Source: "plane", + Status: "running", + Metadata: checkoutTaskMeta("/home/user/workspace/nomadcode/slots/000", "develop"), + }, + } + + progressTime := time.Date(2026, 6, 20, 12, 0, 0, 0, time.UTC) + fakeModel := fakeModelClient{ + generate: func(_ context.Context, input model.GenerateInput) (model.GenerateResult, error) { + if input.OnProgress != nil { + input.OnProgress(model.GenerateProgress{ + Mode: "streaming", + Reason: "receiving stream chunks", + LastEventTime: progressTime, + }) + } + return model.GenerateResult{ + ID: "resp-progress", + Model: "m", + Text: "authoring progress completed", + }, nil + }, + } + + worker := &TaskWorker{ + Lifecycle: fakeLifecycle, + Model: fakeModel, + } + + task := storage.Task{ + ID: "task-auth-progress", + Source: "plane", + Metadata: checkoutTaskMeta("/home/user/workspace/nomadcode/slots/000", "develop"), + } + + raw, msg, err := worker.runTask(context.Background(), task) + if err != nil { + t.Fatalf("runTask returned error: %v", err) + } + if msg == "" { + t.Fatal("expected non-empty message") + } + + foundProgress := false + for _, m := range fakeLifecycle.mergedMetadata { + if m[workflow.MetadataKeyAuthoringProgressMode] == "streaming" { + foundProgress = true + if m[workflow.MetadataKeyAuthoringProgressReason] != "receiving stream chunks" { + t.Errorf("unexpected progress reason: %v", m[workflow.MetadataKeyAuthoringProgressReason]) + } + if m[workflow.MetadataKeyAuthoringRunUpdatedAt] != progressTime.Format(time.RFC3339) { + t.Errorf("unexpected updated at: %v", m[workflow.MetadataKeyAuthoringRunUpdatedAt]) + } + } + } + if !foundProgress { + t.Errorf("expected MergeTaskMetadata to be called with streaming progress, got %v", fakeLifecycle.mergedMetadata) + } + + var result map[string]any + if err := json.Unmarshal(raw, &result); err != nil { + t.Fatalf("unmarshal result: %v", err) + } + if result[workflow.MetadataKeyAuthoringProgressMode] != "streaming" { + t.Errorf("expected progress mode streaming in raw result, got %v", result[workflow.MetadataKeyAuthoringProgressMode]) + } + if result[workflow.MetadataKeyAuthoringProgressReason] != "receiving stream chunks" { + t.Errorf("expected progress reason in raw result, got %v", result[workflow.MetadataKeyAuthoringProgressReason]) + } +} diff --git a/services/core/internal/workflow/model.go b/services/core/internal/workflow/model.go index 4a078be..582d257 100644 --- a/services/core/internal/workflow/model.go +++ b/services/core/internal/workflow/model.go @@ -58,4 +58,6 @@ const ( MetadataKeyAuthoringRunUpdatedAt = "authoring_run_updated_at" MetadataKeyAuthoringFailureType = "authoring_failure_type" MetadataKeyAuthoringFailureCategory = "authoring_failure_category" + MetadataKeyAuthoringProgressMode = "authoring_progress_mode" + MetadataKeyAuthoringProgressReason = "authoring_progress_reason" )