update core services and move progress files to archive

This commit is contained in:
toki 2026-06-20 19:25:18 +09:00
parent 6d4553bdc9
commit 8157ba060b
15 changed files with 993 additions and 37 deletions

View file

@ -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 판별을 좁게 수정한다.

View file

@ -0,0 +1,210 @@
<!-- task=m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress plan=1 tag=REVIEW_AUTHORING_RUNTIME -->
# 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-<milestone-slug>`이면 런타임이 읽을 완료 이벤트 메타데이터를 보고하고, 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로 이동한다.

View file

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

View file

@ -0,0 +1,171 @@
<!-- task=m-plane-origin-authoring-roundtrip-sync/02+01_iop_progress plan=1 tag=REVIEW_AUTHORING_RUNTIME -->
# 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`의 구현 에이전트 소유 섹션을 채운다. 이 파일 작성이 구현의 마지막 단계다.

View file

@ -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`이 함께 기록됩니다.

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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

View file

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