Skip to content

fix(openai): streaming image relay + image edit (sync upstream pair) - #160

Closed
jjcc123312 wants to merge 2 commits into
mainfrom
feat/openai-image-streaming
Closed

fix(openai): streaming image relay + image edit (sync upstream pair)#160
jjcc123312 wants to merge 2 commits into
mainfrom
feat/openai-image-streaming

Conversation

@jjcc123312

Copy link
Copy Markdown

Summary

Syncs the upstream OpenAI Images API streaming relay + image edit feature into the fork as a paired cherry-pick:

  1. d2576ddcd — fix(openai): support streaming image relay and image edit for images API (fix(openai): support streaming image relay and image edit for images API  QuantumNous/new-api#4608)
  2. 59a93cf5c — fix(openai): align image streaming relay governance (splits relay-openai.go into relay_image.go / relay_realtime.go / usage.go, routes image streaming through shared stream handling)

These two commits are a pair and were applied in order (d2576 first).

Conflict resolutions

  • dto/openai_image.go — kept the fork's Stream *bool field with the Rule-5 pointer comment. (Upstream's d2576 used a non-pointer bool; upstream's own follow-up 59a93 independently converts it to *bool + IsStream pointer-deref — matching what the fork already had.)
  • relay/helper/valid_request.go — multipart stream parsing assigns imageRequest.Stream = common.GetPointer(stream) (pointer form for the fork's *bool).
  • relay/channel/openai/adaptor.go — took the upstream streaming branch in DoResponse (if info.IsStream { OpenaiImageStreamHandler } else { OpenaiImageHandler }, incl. the OpenaiHandlerWithUsageOpenaiImageHandler rename). Kept the fork's blockrun-era ConvertImageRequest behavior that strips image streaming fields (stream / partial_images) on the passthrough conversion path — new-api only synthesizes SSE for opt-in channels (e.g. blockrun); forwarding those fields would 400 upstreams that don't support image streaming.
  • relay/channel/openai/relay-openai.go — the 675-line deletion is functions moved to the new files. The fork's codex usage-chunk path stays here (OaiStreamHandler / sendStreamData) and is verified intact; the codex ShouldIncludeUsage wiring in chat_via_responses.go / helper.go is untouched by this sync.
  • Test reconciliationTestConvertImageEditRequestMultipart (from 59a93) asserted upstream's stream-preservation behavior; updated it to assert the fork's deliberate strip behavior on the passthrough path. Other image stream/edit tests adapted to the *bool field.

Two bundled unrelated changes that rode along in 59a93

  • common/init.go — 59a93 bumped GlobalApiRateLimitNum 180→360 and GlobalWebRateLimitNum 60→120. Reverted both back to 180 / 60 (kept code defaults; prod overrides these via env anyway, so this keeps the change minimal).
  • relay/helper/stream_scanner.go — 59a93 bumped DefaultMaxScannerBufferSize 64MB→128MB. Kept the 128MB value (approved) and fixed the stale // 64MB comment to // 128MB.

Verification

  • go build ./relay/... ./common/... ./dto/... — exit 0
  • go vet ./relay/... ./dto/... — no warnings in touched files (only pre-existing warnings in unrelated channel adaptors / common/custom-event.go, confirmed present on main)
  • No leftover conflict markers in relay / common / dto
  • Rate-limit revert confirmed: GlobalApiRateLimitNum=180, GlobalWebRateLimitNum=60
  • SSE buffer confirmed: DefaultMaxScannerBufferSize = 128 << 20
  • go test ./relay/channel/blockrun/... ./relay/channel/openai/... ./relay/helper/... — all pass (blockrun is a hot path built on the openai adaptor; its tests pass)

🤖 Generated with Claude Code

gaoren002 and others added 2 commits June 17, 2026 18:12
…API (QuantumNous#4608)

* fix(openai): support streaming image relay

* fix(openai): keep image edit multipart body reusable

* test(openai): cover image stream usage details

* test(openai): cover image edit fallback stream field

* fix(openai): wrap image json fallback as stream

* fix(relay): support OpenAI image streaming

* fix(openai): record image stream upstream error events

* fix(openai): harden image stream relay

* fix(openai): return image JSON errors

* fix(relay): reset stream status per scanner run

* fix(relay): drop upstream credit passthrough

* fix(openai): keep image errors minimal

* fix(openai): keep image error status from response

---------

Co-authored-by: CaIon <i@caion.me>
Route OpenAI image streaming through shared stream handling, split image/realtime/usage helpers for maintainability, and include the related image request and rate limit updates.

(cherry picked from commit 59a93cf)

Fork adaptations during cherry-pick:
- Keep dto.ImageRequest.Stream as *bool with Rule-5 pointer comment (already pointer on fork)
- Keep blockrun-era ConvertImageRequest stripping of image stream fields on passthrough
- Revert bundled GlobalApiRateLimitNum 360->180 and GlobalWebRateLimitNum 120->60 (prod uses env overrides)
- Keep SSE DefaultMaxScannerBufferSize 128MB bump; fix stale 64MB comment

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
@KingCesc

Copy link
Copy Markdown

🤖 OpenCodeReview · 评审 commit 875adfd9 · 共 10 条

relay/helper/valid_request.go

  • L149: [严重] 这里改用 ParseMultipartFormReusable 会先通过 storage.Bytes() 将整个 multipart 请求体一次性读入内存后再解析;图片编辑请求通常包含较大的文件上传,在并发或接近请求体上限时容易造成明显内存峰值甚至 OOM,相比原来的 c.MultipartForm() 会把大文件落盘风险更高。建议为 multipart 复用解析提供流式实现(直接基于可 Seek 的 body reader + boundary 创建 multipart.Reader,避免 Bytes()),或至少对磁盘存储场景继续使用标准 ParseMultipartForm/临时文件落盘路径。
form, err := common.ParseMultipartFormReusable(c) // TODO: 避免一次性将大文件 multipart body 读入内存,改为流式解析

relay/channel/xai/adaptor.go

  • L117: [严重] 这里没有根据 info.IsStream 分流处理图片响应。客户端请求流式图片生成/编辑时,当前仍走非流式 OpenaiImageHandler,会直接返回 JSON 而不是 SSE(也不会使用 OpenaiImageJSONAsStreamHandler 做 JSON->SSE 兼容转换),导致流式客户端协议解析失败,图片流式能力不可用。建议与 OpenAI adaptor 保持一致,流式请求调用 OpenaiImageStreamHandler,非流式再调用 OpenaiImageHandler
if info.IsStream {
			usage, err = openai.OpenaiImageStreamHandler(c, info, resp)
		} else {
			usage, err = openai.OpenaiImageHandler(c, info, resp)
		}

relay/channel/openai/relay_realtime.go

  • L34-36: [阻塞] usagelocalUsagesumUsage 会被两个 relay goroutine 和主 goroutine 同时读写、累加和重置,但没有任何同步保护。这里会产生数据竞争,直接影响 token 统计和 PreWssConsumeQuota 的扣费结果,可能导致重复扣费或漏扣费。建议将用量统计/扣费集中到单一 goroutine 串行处理,或用 sync.Mutex 保护所有对这些对象的读写,并在两个 goroutine 确认退出后再做最终结算。
usage := &dto.RealtimeUsage{}
	localUsage := &dto.RealtimeUsage{}
	sumUsage := &dto.RealtimeUsage{}
	// TODO: 将 usage/localUsage/sumUsage 的所有读写串行化,或使用 mutex 统一保护,避免并发扣费数据竞争。
  • L126: [阻塞] RealtimeEvent.Response 在 DTO 中是指针,这里直接访问 realtimeEvent.Response.Usage。如果上游返回 response.done 但缺少 response 字段,会触发空指针 panic,导致实时会话异常中断,并可能让未结算用量丢失。建议先校验 Response != nil,缺失时按上游协议错误返回并停止当前会话。
if realtimeEvent.Response == nil {
						errChan <- fmt.Errorf("response.done event missing response")
						return
					}
					realtimeUsage := realtimeEvent.Response.Usage
  • L207-209: [阻塞] 收到转发、解析、token 统计或扣费错误后这里只记录日志,函数最终仍返回 nil, sumUsage;上层会把本次 realtime relay 当作成功处理。对于扣费/转发链路,这会造成状态不一致或漏扣费。建议在 errChan 分支保存并返回明确错误,终止后续成功路径。
case err := <-errChan:
		logger.LogError(c, "realtime error: "+err.Error())
		return types.NewError(err, types.ErrorCodeBadResponse), sumUsage
  • L0: [阻塞] 末尾补扣 usage/localUsage 时忽略了 preConsumeUsage 的错误。PreWssConsumeQuota 会读取用户额度和 token 信息,失败时如果被吞掉,会导致本次请求返回成功但实际未扣费或扣费不完整。建议检查错误并返回给上层,必要时触发补偿/终止会话。
if usage.TotalTokens != 0 {
		if err := preConsumeUsage(c, info, usage, sumUsage); err != nil {
			return types.NewError(err, types.ErrorCodeBadResponse), sumUsage
		}
	}

	if localUsage.TotalTokens != 0 {
		if err := preConsumeUsage(c, info, localUsage, sumUsage); err != nil {
			return types.NewError(err, types.ErrorCodeBadResponse), sumUsage
		}
	}
  • L28-32: [严重] 主协程在任一方向关闭或出错后直接退出,但没有关闭 clientConn/targetConn,也没有取消上下文或等待两个 goroutine 结束。另一侧 ReadMessage 可能长期阻塞,连接和 goroutine 会泄漏;errChan 缓冲满后发送错误也可能阻塞。建议使用可取消 context/errgroup 管理生命周期,在退出路径用 sync.Once 统一关闭两个 websocket,并等待 relay goroutine 结束。
clientClosed := make(chan struct{})
	targetClosed := make(chan struct{})
	sendChan := make(chan []byte, 100)
	receiveChan := make(chan []byte, 100)
	errChan := make(chan error, 2)
	// TODO: 使用 context cancellation/errgroup 管理两个 relay goroutine,并在退出时统一关闭 clientConn 与 targetConn,避免连接和 goroutine 泄漏。

relay/channel/openai/usage.go

  • L72-80: [阻塞] 这里只判断指针是否存在会把 cached_tokens: 0 当成有效结果返回,导致同一响应里的 prompt_cache_hit_tokens 正数不再被后续兜底使用;缓存 token 会被漏记,进而影响缓存命中计费/额度扣减。建议仅在该字段为正数时返回,否则继续检查后备字段。
if payload.Usage.PromptTokensDetails.CachedTokens != nil && *payload.Usage.PromptTokensDetails.CachedTokens > 0 {
		return *payload.Usage.PromptTokensDetails.CachedTokens, true
	}
	if payload.Usage.CachedTokens != nil && *payload.Usage.CachedTokens > 0 {
		return *payload.Usage.CachedTokens, true
	}
	if payload.Usage.PromptCacheHitTokens != nil && *payload.Usage.PromptCacheHitTokens > 0 {
		return *payload.Usage.PromptCacheHitTokens, true
	}
  • L129-132: [阻塞] 从上游响应体提取的 cache_n 未校验非负,异常或兼容接口返回负数时会写入 PromptTokensDetails.CachedTokens,而后续计费会用该值从 prompt tokens 中抵扣,可能产生负抵扣/错误账单。建议对所有从响应体提取的缓存 token 做统一范围校验,至少要求为正数,并可进一步限制不超过 prompt/input tokens。
if payload.Timings.CachedTokens == nil || *payload.Timings.CachedTokens <= 0 {
		return 0, false
	}
	return *payload.Timings.CachedTokens, true

relay/channel/openai/relay_image.go

  • L161-162: [严重] 这里只判断 len(payload.Error) > 0 会把 "error": null 也识别为错误事件(json.RawMessage 内容为 null,长度大于 0)。如果上游正常图片流事件携带空 error 字段,会被标记为软错误,导致 StreamStatus/日志按异常流处理,后续错误统计或告警可能误报。建议仅在 error 非 null 且包含有效错误信息时再判定为错误。
payloadType := strings.ToLower(strings.TrimSpace(payload.Type))
	if payloadType == "error" || payloadType == "upstream_error" {
		return true
	}
	if len(payload.Error) == 0 || strings.EqualFold(strings.TrimSpace(string(payload.Error)), "null") {
		return false
	}
	return true

@jjcc123312

Copy link
Copy Markdown
Author

暂不同步:此 PR 同步的是上游代码,且 OpenCodeReview 指出上游实现存在质量问题(非本次合并引入)。当前无强需求,改为按需单独同步。先关闭,需要时再开。

@jjcc123312 jjcc123312 closed this Jun 17, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants