fix: 接入内嵌 worker 并完成集成验收 (#13, #15)

This commit is contained in:
ila
2026-08-21 12:36:51 +08:00
parent 919e11f664
commit 1ad8dca9cf
14 changed files with 595 additions and 46 deletions
+8 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Project-Profile
wiki_url: https://git.ilapage.cn/OPC/chorus/wiki/Project-Profile.-
wiki_revision: 9e3b50db59ed7ff77bc900d2dfa61c66ce97fd83
synchronized_at: 2026-08-21T03:22:44Z
wiki_revision: 79ad76521eb76e279cf1d99188835d43dadd7ca0
synchronized_at: 2026-08-21T04:33:29Z
<!-- gitea-wiki-mirror:end -->
# 项目档案
@@ -158,3 +158,9 @@ synchronized_at: 2026-08-21T03:22:44Z
- MVP-0 落地后:`go build ./...`、`go vet ./...`、`go test ./...` 必须通过。
- 涉及 `retryable` 判定、SSRF 校验、密钥加解密和迁移的修改必须有针对性单元测试。
- 未执行或无法覆盖的验证必须记录到工单。
## MVP-0 集成候选状态(2026-08-21)
- #5~#12 和 #14 已由用户验收;#13 集成验收执行中发现 #10 的内嵌 worker 未接入 `portal/main.go`,已用缺陷 #15 恢复既有确认行为。
- 当前候选版本已在本机 MySQL 8.4.8 隔离库完成空库 `up/down/up`、重复种子、Go build/vet/race、真实 portal 进程认领、Playwright 四视口和终态验证。
- #13 与 #15 均停在待验收前的证据整理阶段;在用户明确验收前,MVP-0、MVP #4 和 Epic #3 仍未完成。
- 验证只使用构造用户、构造图片、mock 上游与测试 fixture;没有连接生产/共享库、调用真实 Provider 额度或发布生产。
+8 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Architecture-and-Code-Map
wiki_url: https://git.ilapage.cn/OPC/chorus/wiki/Architecture-and-Code-Map.-
wiki_revision: 08b7352c995a45a8254943958da8ada72c210b13
synchronized_at: 2026-08-21T03:22:57Z
wiki_revision: b61f12999ae1b9fdb3c8a49949fbf96697c14145
synchronized_at: 2026-08-21T04:33:38Z
<!-- gitea-wiki-mirror:end -->
# 架构与代码地图
@@ -96,6 +96,12 @@ MVP-0 第一张上传图自动作为 `primary`,其余为 `reference`;不提
- 429/5xx/超时/连接错误以及 400/401/内容策略拒绝均映射为固定脱敏错误;MVP-0 记录 retryable 分类但不换 Provider、不自动重试。
- `internal/platform/mockprovider` 和 `cmd/chorus-mock-provider` 提供本地协议 fixture;自动测试通过注入受控 DNS/dial 连接它,不放宽生产对回环/私网地址的 SSRF 拦截。
- MySQL 8 集成测试分别完成一次 chat 和 images_edits,断言 rendered_prompt、attempt ProviderModel/latency、输出记录、原图和缩略图;全程不访问真实 Provider。
### #15 portal 组合根与运行闭环
- `portal/main.go` 现在实际组装 MySQL queue controller、AES-GCM key ring、SSRF 安全 HTTP client、Provider factory/catalog、本地存储和内嵌 worker;同步 handler 仍只写 pending,不持有 Provider 依赖。
- worker 的 claim 必须经过 `queue.Controller`;收到退出信号后先阻止新 claim,HTTP server 停止接收请求,并等待在途 worker 在租约上下文内完成。HTTP 或 worker 任一组件意外结束都会让进程失败关闭。
- 非生产默认 lease 60 秒、Provider HTTP 超时 45 秒、轮询 250ms、响应上限 32MiB;生产必须显式配置,且 lease 必须比 HTTP 超时多 5 秒以上。
- `CHORUS_TEST_DISABLE_WORKER=true` 只允许 `CHORUS_ENV=test`,用于浏览器稳定观察状态的测试 fixture;开发和生产环境启用会直接报错。真实 Provider/worker 成功仍通过注入受控 DNS/DialContext 的 MySQL 集成测试验证,生产 SSRF 逻辑没有旁路。
### #11 已落地的 portal 认证与用户 API
- `portal/session` 使用服务端内存会话和 HMAC 不透明 Cookie;登录成功与退出都会轮换 session ID 和 CSRF token。Cookie 为 HttpOnly/SameSite=Lax,生产启用 Secure;写请求只接受 `X-CSRF-Token`,避免 multipart 在正文限流前被隐式解析。
+21 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Local-Development-and-Verification
wiki_url: https://git.ilapage.cn/OPC/chorus/wiki/Local-Development-and-Verification.-
wiki_revision: 357def01026f5c12859461d54739c52a0896da4e
synchronized_at: 2026-08-21T03:27:24Z
wiki_revision: 31d2d641a46dbf9bf0b9275542831a8149ea603b
synchronized_at: 2026-08-21T04:33:43Z
<!-- gitea-wiki-mirror:end -->
# 本地开发与验证
@@ -271,3 +271,22 @@ git diff --check
```
只运行受当前范围影响且实际存在的产品命令;不存在或因工具版本不能执行的项写入工单。涉及迁移、安全、队列、权限或 UI 时,还必须执行上面的专项验证并记录结果。
## MVP-0 集成验收记录(2026-08-21)
本次 #13 使用专用可丢弃库验证,不能把示例中的 down 或 fixture 指向现有开发库、共享库或生产库。结论如下:
- MySQL 8.4.8:全量 `up → seed 两次 → down -all → up → seed 两次` 通过;种子稳定为 1 用户、1 Provider、2 ProviderModel、2 Prompt,down 后业务表为 0。
- Go:`go build ./...`、`go vet ./...`、`go test -race -count=1 -p=1 ./...` 通过;MySQL 测试覆盖幂等、跨用户、租约恢复、旧 token CAS、attempts、文本/图片 mock 成功和鉴权文件。
- 浏览器:`pnpm --dir portal/web test:e2e` 在 375/768/1024/1440 通过;1024 额外覆盖 pending、running、文本成功、图片成功、失败终态及终态移除轮询属性。
- 进程:真实内嵌 worker 从隔离库认领 6 条 pending 并形成终态;回环 Provider 被 SSRF 在连接前拒绝;两次前台 `Ctrl+C` 均干净退出。
内嵌 worker 的运行参数:
| 变量 | 非生产默认 | 生产要求 |
|---|---:|---|
| `CHORUS_WORKER_LEASE_SECONDS` | 60 | 必须显式设置,且大于 HTTP 超时 + 5 秒 |
| `CHORUS_WORKER_POLL_MILLISECONDS` | 250 | 必须显式设置 |
| `CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS` | 45 | 必须显式设置 |
| `CHORUS_PROVIDER_MAX_RESPONSE_BYTES` | 33554432 | 必须显式设置 |
`CHORUS_TEST_DISABLE_WORKER=true` 仅供 `CHORUS_ENV=test` 的 Playwright fixture 使用。测试 helper 位于 `portal/web/e2e/fixture`,会按构造用户归属写入受控状态和测试图片;不得打包部署,也不得用于开发或生产数据。mock 上游的自动测试通过依赖注入连接本地 fixture,不允许把回环地址加入生产 SSRF 白名单。
+6 -6
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Product-Requirements-Overview
wiki_url: https://git.ilapage.cn/OPC/chorus/wiki/Product-Requirements-Overview.-
wiki_revision: 939520c952298b440bd1a92e07c18c4305c4a77e
synchronized_at: 2026-08-21T03:58:18Z
wiki_revision: fa500788e6c0ecc276dcfcb0b43154c30b773949
synchronized_at: 2026-08-21T04:33:58Z
<!-- gitea-wiki-mirror:end -->
# 产品需求总览
@@ -29,13 +29,13 @@ synchronized_at: 2026-08-21T03:58:18Z
## 当前需求索引
工单 [#1](https://git.ilapage.cn/OPC/chorus/issues/1) 已完成长期技术基线整理和验收。项目由 [Epic #3](https://git.ilapage.cn/OPC/chorus/issues/3) 统一跟踪,当前阶段由 [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) 汇总。MVP-0 生产闭环已拆分为环境 [#5](https://git.ilapage.cn/OPC/chorus/issues/5)、实现 [#6–#12](https://git.ilapage.cn/OPC/chorus/issues/6) 和集成验收 [#13](https://git.ilapage.cn/OPC/chorus/issues/13)。截至 2026-08-21,#5~#12 已全部验收;#13 的独立集成验收前置条件已满足,MVP-0 尚未完成。
工单 [#1](https://git.ilapage.cn/OPC/chorus/issues/1) 已完成长期技术基线整理和验收。项目由 [Epic #3](https://git.ilapage.cn/OPC/chorus/issues/3) 统一跟踪,当前阶段由 [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) 汇总。MVP-0 生产闭环已拆分为环境 [#5](https://git.ilapage.cn/OPC/chorus/issues/5)、实现 [#6–#12](https://git.ilapage.cn/OPC/chorus/issues/6) 和集成验收 [#13](https://git.ilapage.cn/OPC/chorus/issues/13)。截至 2026-08-21,#5~#12 与 #14 已全部验收;#13 集成执行中发现内嵌 worker 组合根遗漏并建立修复 #15。#13/#15 已完成实现与自动验证证据整理,均仍需用户验收,因此 MVP-0 尚未完成。
| 需求领域 | 用户与场景 | 需求状态 | MVP | 详细说明 | 实施工单 | 设计证据 |
|---|---|---|---|---|---|---|
| 核心生成域与单上游 | 用户提交提示词/原图得到结果 | 已实现(#6~#10 已验收,待 #13 集成验收) | [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) | [架构](Architecture-and-Code-Map.-)、[业务规则](Business-Rules-and-Glossary.-) | [#6 迁移](https://git.ilapage.cn/OPC/chorus/issues/6)、[#7 core](https://git.ilapage.cn/OPC/chorus/issues/7)、[#8 platform](https://git.ilapage.cn/OPC/chorus/issues/8)、[#9 queue](https://git.ilapage.cn/OPC/chorus/issues/9)、[#10 worker](https://git.ilapage.cn/OPC/chorus/issues/10) | 架构/数据/状态设计已确认 |
| 用户端生成页 | 种子用户登录并完成一次异步生成 | 已实现(#11/#12 已验收,待 #13 集成验收) | [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) | 本页“用户端布局与状态” | [#11 portal 后端](https://git.ilapage.cn/OPC/chorus/issues/11)、[#12 用户界面](https://git.ilapage.cn/OPC/chorus/issues/12) | [原型设计 #2](https://git.ilapage.cn/OPC/chorus/issues/2);`prototypes/2/v1/index.html`,2026-08-20 用户已确认 |
| 默认 prompt template | 系统以数据配置而非硬编码合成提示词 | 已实现(#7 已验收,待 #13 集成验收) | [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) | [业务规则](Business-Rules-and-Glossary.-) | [#7](https://git.ilapage.cn/OPC/chorus/issues/7) | 2026-08-20 已确认中性模板与 `{{.UserPrompt}}` |
| 核心生成域与单上游 | 用户提交提示词/原图得到结果 | 已实现(#6~#10 已验收;#13/#15 待验收) | [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) | [架构](Architecture-and-Code-Map.-)、[业务规则](Business-Rules-and-Glossary.-) | [#6 迁移](https://git.ilapage.cn/OPC/chorus/issues/6)、[#7 core](https://git.ilapage.cn/OPC/chorus/issues/7)、[#8 platform](https://git.ilapage.cn/OPC/chorus/issues/8)、[#9 queue](https://git.ilapage.cn/OPC/chorus/issues/9)、[#10 worker](https://git.ilapage.cn/OPC/chorus/issues/10) | 架构/数据/状态设计已确认 |
| 用户端生成页 | 种子用户登录并完成一次异步生成 | 已实现(#11/#12 已验收;#13 待验收) | [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) | 本页“用户端布局与状态” | [#11 portal 后端](https://git.ilapage.cn/OPC/chorus/issues/11)、[#12 用户界面](https://git.ilapage.cn/OPC/chorus/issues/12) | [原型设计 #2](https://git.ilapage.cn/OPC/chorus/issues/2);`prototypes/2/v1/index.html`,2026-08-20 用户已确认 |
| 默认 prompt template | 系统以数据配置而非硬编码合成提示词 | 已实现(#7 已验收;#13 待验收) | [MVP-0 #4](https://git.ilapage.cn/OPC/chorus/issues/4) | [业务规则](Business-Rules-and-Glossary.-) | [#7](https://git.ilapage.cn/OPC/chorus/issues/7) | 2026-08-20 已确认中性模板与 `{{.UserPrompt}}` |
| 多 Provider 路由与故障转移 | 单上游故障时继续服务 | 已确认 | MVP-1 | 业务规则“选路与上游” | 待建 | 架构/状态设计 |
| 管理端配置与记录 | 运营配置模型、路由池并排障 | 已确认 | MVP-1 | 本页“管理端页面” | 待建 | 复用型/定制页原型待确认 |
| 图片角色编辑 | 用户编辑 role_rule、角色、备注与顺序 | 已确认 | MVP-1 | 业务规则“提示词与上传” | 待建 | 组件状态与键盘交互原型 |
+18 -2
View File
@@ -2,8 +2,8 @@
generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件)
wiki_page: Deployment-and-Operations
wiki_url: https://git.ilapage.cn/OPC/chorus/wiki/Deployment-and-Operations.-
wiki_revision: 23cabbc250817324c8f00dacc8e3d15c512141d4
synchronized_at: 2026-08-21T03:23:42Z
wiki_revision: 80dac27043770615620a5d346e5b0e6f72451ab4
synchronized_at: 2026-08-21T04:34:00Z
<!-- gitea-wiki-mirror:end -->
# 部署与运维
@@ -118,6 +118,22 @@ MVP-0 的终端用户会话保存在单个 portal 进程内存中,Cookie 只
生产 Cookie 必须为 Secure/HttpOnly/SameSite=Lax。反向代理不得记录 Cookie 或 CSRF header;请求日志和 GORM SQL 使用参数化输出。JSON 正文、multipart 总量、单图字节/像素/数量和历史查询都有显式生产配置,缺少任一阈值时启动失败。
上传先在应用层校验,再在 generation 事务回调中保存带 owner/generation metadata 的文件。事务失败只清理本请求保存的 key;存储根目录不能作为 nginx/static 目录暴露。
## MVP-0 worker 已实现配置与退出行为
portal 单二进制会同时启动 HTTP server 和内嵌 worker。生产除现有 portal 限制外,还必须显式提供:
| 变量 | 约束 |
|---|---|
| `CHORUS_WORKER_LEASE_SECONDS` | 正整数;必须大于 `CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS + 5` |
| `CHORUS_WORKER_POLL_MILLISECONDS` | 正整数;控制空队列轮询间隔 |
| `CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS` | 正整数;安全 HTTP client 总超时和响应头超时 |
| `CHORUS_PROVIDER_MAX_RESPONSE_BYTES` | 正整数;Provider JSON/Base64/下载响应读取上限 |
`CHORUS_MASTER_KEY` 的原始字节长度必须是 AES 支持的 16、24 或 32 字节;MVP-0 写入信封使用 `key_id=primary`。生产缺失或长度无效时 portal 启动失败且不输出密钥。开发环境未配置主密钥时只生成进程内临时密钥,适用于默认 `auth_type=none` mock;要测试持久化 Bearer 凭据必须显式提供稳定的仓库外密钥。
`CHORUS_TEST_DISABLE_WORKER` 是测试专用开关,只在 `CHORUS_ENV=test` 接受;开发和生产设置为 true 会启动失败。部署制品和服务配置不得设置该变量,`portal/web/e2e/fixture` 也不得进入生产运行命令。
#13 已在 Windows 前台进程两次用 `Ctrl+C` 验证停止认领、关闭 HTTP 并等待 worker 返回。真实 Linux/systemd 的 TERM、超时和重启恢复仍属于首次发布前门禁,不能把本机验证写成已上线。
## worker 优雅退出
- 收到 TERM 后先停止认领新任务;
+47 -17
View File
@@ -20,21 +20,26 @@ const (
)
type Config struct {
Environment Environment
DBDSN string
ListenAddress string
StorageRoot string
SessionKey string
MasterKey string
SessionTTL time.Duration
LoginAttempts int
LoginWindow time.Duration
MaxPromptBytes int
MaxImages int
MaxImageBytes int64
MaxUploadBytes int64
MaxImagePixels uint64
HistoryLimit int
Environment Environment
DBDSN string
ListenAddress string
StorageRoot string
SessionKey string
MasterKey string
SessionTTL time.Duration
LoginAttempts int
LoginWindow time.Duration
MaxPromptBytes int
MaxImages int
MaxImageBytes int64
MaxUploadBytes int64
MaxImagePixels uint64
HistoryLimit int
WorkerLeaseDuration time.Duration
WorkerPollInterval time.Duration
ProviderHTTPTimeout time.Duration
ProviderMaxResponseBytes int64
TestDisableWorker bool
}
func Load() (Config, error) {
@@ -93,7 +98,7 @@ func LoadFromLookup(lookup func(string) (string, bool)) (Config, error) {
missing = append(missing, name)
}
}
for _, name := range []string{"CHORUS_SESSION_TTL_MINUTES", "CHORUS_LOGIN_ATTEMPTS", "CHORUS_LOGIN_WINDOW_SECONDS", "CHORUS_MAX_PROMPT_BYTES", "CHORUS_MAX_IMAGES", "CHORUS_MAX_IMAGE_BYTES", "CHORUS_MAX_UPLOAD_BYTES", "CHORUS_MAX_IMAGE_PIXELS", "CHORUS_HISTORY_LIMIT"} {
for _, name := range []string{"CHORUS_SESSION_TTL_MINUTES", "CHORUS_LOGIN_ATTEMPTS", "CHORUS_LOGIN_WINDOW_SECONDS", "CHORUS_MAX_PROMPT_BYTES", "CHORUS_MAX_IMAGES", "CHORUS_MAX_IMAGE_BYTES", "CHORUS_MAX_UPLOAD_BYTES", "CHORUS_MAX_IMAGE_PIXELS", "CHORUS_HISTORY_LIMIT", "CHORUS_WORKER_LEASE_SECONDS", "CHORUS_WORKER_POLL_MILLISECONDS", "CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS", "CHORUS_PROVIDER_MAX_RESPONSE_BYTES"} {
if read(name) == "" {
missing = append(missing, name)
}
@@ -136,9 +141,34 @@ func LoadFromLookup(lookup func(string) (string, bool)) (Config, error) {
if cfg.HistoryLimit, err = intConfig(read("CHORUS_HISTORY_LIMIT"), 50); err != nil {
return Config{}, fmt.Errorf("CHORUS_HISTORY_LIMIT is invalid")
}
if cfg.WorkerLeaseDuration, err = durationConfig(read("CHORUS_WORKER_LEASE_SECONDS"), 60, time.Second); err != nil {
return Config{}, fmt.Errorf("CHORUS_WORKER_LEASE_SECONDS is invalid")
}
if cfg.WorkerPollInterval, err = durationConfig(read("CHORUS_WORKER_POLL_MILLISECONDS"), 250, time.Millisecond); err != nil {
return Config{}, fmt.Errorf("CHORUS_WORKER_POLL_MILLISECONDS is invalid")
}
if cfg.ProviderHTTPTimeout, err = durationConfig(read("CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS"), 45, time.Second); err != nil {
return Config{}, fmt.Errorf("CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS is invalid")
}
providerMaxResponseBytes, err := intConfig(read("CHORUS_PROVIDER_MAX_RESPONSE_BYTES"), 32<<20)
if err != nil {
return Config{}, fmt.Errorf("CHORUS_PROVIDER_MAX_RESPONSE_BYTES is invalid")
}
cfg.ProviderMaxResponseBytes = int64(providerMaxResponseBytes)
if value := read("CHORUS_TEST_DISABLE_WORKER"); value != "" {
if cfg.TestDisableWorker, err = strconv.ParseBool(value); err != nil {
return Config{}, fmt.Errorf("CHORUS_TEST_DISABLE_WORKER is invalid")
}
if cfg.TestDisableWorker && environment != Test {
return Config{}, fmt.Errorf("CHORUS_TEST_DISABLE_WORKER is only allowed in test")
}
}
if cfg.MaxUploadBytes < cfg.MaxImageBytes || cfg.MaxImages > 32 || cfg.HistoryLimit > 200 {
return Config{}, fmt.Errorf("portal limits are inconsistent")
}
if cfg.WorkerLeaseDuration <= cfg.ProviderHTTPTimeout+5*time.Second {
return Config{}, fmt.Errorf("worker lease must exceed provider HTTP timeout by more than 5 seconds")
}
return cfg, nil
}
@@ -146,7 +176,7 @@ func (c Config) ValidateProduction() error {
if c.Environment != Production {
return nil
}
if c.DBDSN == "" || c.MasterKey == "" || len(c.SessionKey) < 32 || c.StorageRoot == "" || c.SessionTTL <= 0 || c.MaxImageBytes <= 0 || c.MaxUploadBytes < c.MaxImageBytes {
if c.DBDSN == "" || c.MasterKey == "" || len(c.SessionKey) < 32 || c.StorageRoot == "" || c.SessionTTL <= 0 || c.MaxImageBytes <= 0 || c.MaxUploadBytes < c.MaxImageBytes || c.WorkerLeaseDuration <= c.ProviderHTTPTimeout+5*time.Second || c.ProviderMaxResponseBytes <= 0 {
return errors.New("production configuration is incomplete")
}
return nil
+25
View File
@@ -3,6 +3,7 @@ package config
import (
"strings"
"testing"
"time"
)
func lookup(values map[string]string) func(string) (string, bool) {
@@ -24,6 +25,9 @@ func TestLoadDevelopmentRequiresDSNWithoutLeakingIt(t *testing.T) {
if cfg.DBDSN != secret || cfg.ListenAddress != "127.0.0.1:8080" {
t.Fatalf("unexpected config: %#v", cfg)
}
if cfg.WorkerLeaseDuration != 60*time.Second || cfg.WorkerPollInterval != 250*time.Millisecond || cfg.ProviderHTTPTimeout != 45*time.Second || cfg.ProviderMaxResponseBytes != 32<<20 {
t.Fatalf("unexpected worker defaults: %#v", cfg)
}
}
func TestLoadProductionFailsClosed(t *testing.T) {
@@ -58,6 +62,8 @@ func TestLoadProductionPortalLimits(t *testing.T) {
"CHORUS_SESSION_TTL_MINUTES": "480", "CHORUS_LOGIN_ATTEMPTS": "5", "CHORUS_LOGIN_WINDOW_SECONDS": "60",
"CHORUS_MAX_PROMPT_BYTES": "8000", "CHORUS_MAX_IMAGES": "8", "CHORUS_MAX_IMAGE_BYTES": "10485760",
"CHORUS_MAX_UPLOAD_BYTES": "33554432", "CHORUS_MAX_IMAGE_PIXELS": "40000000", "CHORUS_HISTORY_LIMIT": "50",
"CHORUS_WORKER_LEASE_SECONDS": "60", "CHORUS_WORKER_POLL_MILLISECONDS": "250",
"CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS": "45", "CHORUS_PROVIDER_MAX_RESPONSE_BYTES": "33554432",
}
cfg, err := LoadFromLookup(lookup(values))
if err != nil {
@@ -71,3 +77,22 @@ func TestLoadProductionPortalLimits(t *testing.T) {
t.Fatal("inconsistent upload limits accepted")
}
}
func TestLoadRejectsUnsafeWorkerLease(t *testing.T) {
_, err := LoadFromLookup(lookup(map[string]string{
"CHORUS_DSN": "test-dsn", "CHORUS_WORKER_LEASE_SECONDS": "50", "CHORUS_PROVIDER_HTTP_TIMEOUT_SECONDS": "45",
}))
if err == nil || !strings.Contains(err.Error(), "worker lease") {
t.Fatalf("expected worker lease error, got %v", err)
}
}
func TestWorkerCanOnlyBeDisabledInTest(t *testing.T) {
if _, err := LoadFromLookup(lookup(map[string]string{"CHORUS_DSN": "test-dsn", "CHORUS_TEST_DISABLE_WORKER": "true"})); err == nil {
t.Fatal("development accepted disabled worker")
}
cfg, err := LoadFromLookup(lookup(map[string]string{"CHORUS_ENV": "test", "CHORUS_DSN": "test-dsn", "CHORUS_TEST_DISABLE_WORKER": "true"}))
if err != nil || !cfg.TestDisableWorker {
t.Fatalf("test worker fixture was rejected: %#v %v", cfg, err)
}
}
+18
View File
@@ -5,6 +5,8 @@ import (
"errors"
"sync"
"time"
"git.ilapage.cn/OPC/chorus/internal/core/model"
)
var ErrClaimsStopped = errors.New("queue claims are stopped")
@@ -33,3 +35,19 @@ func (c *Controller) ClaimNext(ctx context.Context, owner string, leaseDuration
}
return c.repository.ClaimNext(ctx, owner, leaseDuration)
}
func (c *Controller) AssignProvider(ctx context.Context, generationID uint64, leaseToken string, providerModelID uint64) (bool, error) {
return c.repository.AssignProvider(ctx, generationID, leaseToken, providerModelID)
}
func (c *Controller) Succeed(ctx context.Context, generationID uint64, leaseToken string, outputs []model.GenerationOutput, attempt model.Attempt) (bool, error) {
return c.repository.Succeed(ctx, generationID, leaseToken, outputs, attempt)
}
func (c *Controller) Fail(ctx context.Context, generationID uint64, leaseToken, code, message string, attempt model.Attempt) (bool, error) {
return c.repository.Fail(ctx, generationID, leaseToken, code, message, attempt)
}
func (c *Controller) Inputs(ctx context.Context, generationID uint64) ([]model.GenerationInput, error) {
return c.repository.Inputs(ctx, generationID)
}
+3
View File
@@ -28,6 +28,9 @@ func (f *fakeRepository) Succeed(context.Context, uint64, string, []model.Genera
func (f *fakeRepository) Fail(context.Context, uint64, string, string, string, model.Attempt) (bool, error) {
return false, nil
}
func (f *fakeRepository) Inputs(context.Context, uint64) ([]model.GenerationInput, error) {
return nil, nil
}
func TestControllerStopsNewClaims(t *testing.T) {
repository := &fakeRepository{}
+1
View File
@@ -18,4 +18,5 @@ type Repository interface {
AssignProvider(ctx context.Context, generationID uint64, leaseToken string, providerModelID uint64) (bool, error)
Succeed(ctx context.Context, generationID uint64, leaseToken string, outputs []model.GenerationOutput, attempt model.Attempt) (bool, error)
Fail(ctx context.Context, generationID uint64, leaseToken string, code, message string, attempt model.Attempt) (bool, error)
Inputs(ctx context.Context, generationID uint64) ([]model.GenerationInput, error)
}
+120 -14
View File
@@ -2,7 +2,9 @@ package main
import (
"context"
"crypto/rand"
"errors"
"fmt"
"log"
"net/http"
"os"
@@ -12,11 +14,14 @@ import (
"git.ilapage.cn/OPC/chorus/internal/config"
"git.ilapage.cn/OPC/chorus/internal/core/queue"
platformcrypto "git.ilapage.cn/OPC/chorus/internal/platform/crypto"
safehttp "git.ilapage.cn/OPC/chorus/internal/platform/http"
platformstorage "git.ilapage.cn/OPC/chorus/internal/platform/storage"
"git.ilapage.cn/OPC/chorus/portal/auth"
"git.ilapage.cn/OPC/chorus/portal/handler"
"git.ilapage.cn/OPC/chorus/portal/service"
"git.ilapage.cn/OPC/chorus/portal/session"
"git.ilapage.cn/OPC/chorus/portal/worker"
"github.com/gin-gonic/gin"
"gorm.io/driver/mysql"
"gorm.io/gorm"
@@ -59,6 +64,52 @@ func run() error {
if err != nil {
return err
}
queueController := queue.NewController(queueRepository)
keyMaterial := []byte(cfg.MasterKey)
if len(keyMaterial) == 0 {
keyMaterial = make([]byte, 32)
if _, err := rand.Read(keyMaterial); err != nil {
return errors.New("generate development provider key")
}
}
keyRing, err := platformcrypto.NewKeyRing("primary", map[string][]byte{"primary": keyMaterial})
for index := range keyMaterial {
keyMaterial[index] = 0
}
if err != nil {
return errors.New("configure provider key ring")
}
httpClient, err := safehttp.New(safehttp.Config{Timeout: cfg.ProviderHTTPTimeout, MaxRedirects: 3})
if err != nil {
return err
}
providerFactory, err := worker.NewOpenAIFactory(safehttp.NewProviderClient(httpClient), cfg.ProviderMaxResponseBytes)
if err != nil {
return err
}
catalog, err := worker.NewGORMCatalog(db, keyRing)
if err != nil {
return err
}
workerStorage, err := worker.NewLocalStorage(storage)
if err != nil {
return err
}
backgroundWorker, err := worker.New(worker.Config{
Owner: fmt.Sprintf("portal-%d", os.Getpid()),
LeaseDuration: cfg.WorkerLeaseDuration,
PollInterval: cfg.WorkerPollInterval,
}, queueController, queueController, catalog, providerFactory, workerStorage)
if err != nil {
return err
}
runWorker := backgroundWorker.Run
if cfg.TestDisableWorker {
runWorker = func(ctx context.Context) error {
<-ctx.Done()
return nil
}
}
sessions, err := session.New([]byte(cfg.SessionKey), cfg.SessionTTL, cfg.Environment == config.Production)
if err != nil {
return err
@@ -78,20 +129,75 @@ func run() error {
server := &http.Server{Addr: cfg.ListenAddress, Handler: router, ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 30 * time.Second, WriteTimeout: 30 * time.Second, IdleTimeout: 60 * time.Second}
shutdown, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
errorsChannel := make(chan error, 1)
go func() {
return runComponents(shutdown, func() error {
log.Printf("chorus portal listening on %s", cfg.ListenAddress)
errorsChannel <- server.ListenAndServe()
}()
select {
case err := <-errorsChannel:
if !errors.Is(err, http.ErrServerClosed) {
return err
}
return nil
case <-shutdown.Done():
ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
defer cancel()
return server.ListenAndServe()
}, func(ctx context.Context) error {
return server.Shutdown(ctx)
}
}, runWorker)
}
type componentResult struct {
name string
err error
}
func runComponents(ctx context.Context, runHTTP func() error, shutdownHTTP func(context.Context) error, runWorker func(context.Context) error) error {
workCtx, cancel := context.WithCancel(context.Background())
defer cancel()
results := make(chan componentResult, 2)
go func() { results <- componentResult{name: "http", err: runHTTP()} }()
go func() { results <- componentResult{name: "worker", err: runWorker(workCtx)} }()
var first componentResult
hasFirst := false
select {
case <-ctx.Done():
case first = <-results:
hasFirst = true
}
cancel()
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 15*time.Second)
shutdownErr := shutdownHTTP(shutdownCtx)
shutdownCancel()
remaining := 2
if hasFirst {
remaining--
}
var componentErr error
if hasFirst {
componentErr = unexpectedComponentError(first)
}
for remaining > 0 {
result := <-results
remaining--
if err := completedComponentError(result); componentErr == nil && err != nil {
componentErr = err
}
}
if shutdownErr != nil {
return fmt.Errorf("shutdown portal HTTP server: %w", shutdownErr)
}
return componentErr
}
func unexpectedComponentError(result componentResult) error {
if result.name == "http" && errors.Is(result.err, http.ErrServerClosed) {
return errors.New("portal HTTP server stopped unexpectedly")
}
if result.err == nil {
return fmt.Errorf("portal %s stopped unexpectedly", result.name)
}
return fmt.Errorf("portal %s failed: %w", result.name, result.err)
}
func completedComponentError(result componentResult) error {
if result.name == "http" && errors.Is(result.err, http.ErrServerClosed) {
return nil
}
if result.err != nil {
return fmt.Errorf("portal %s failed during shutdown: %w", result.name, result.err)
}
return nil
}
+65
View File
@@ -0,0 +1,65 @@
package main
import (
"context"
"errors"
"net/http"
"strings"
"testing"
"time"
)
func TestRunComponentsWaitsForInflightWorker(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
httpStopped := make(chan struct{})
workerStarted := make(chan struct{})
workerCanceled := make(chan struct{})
releaseWorker := make(chan struct{})
done := make(chan error, 1)
go func() {
done <- runComponents(ctx, func() error {
<-httpStopped
return http.ErrServerClosed
}, func(context.Context) error {
close(httpStopped)
return nil
}, func(workerCtx context.Context) error {
close(workerStarted)
<-workerCtx.Done()
close(workerCanceled)
<-releaseWorker
return nil
})
}()
<-workerStarted
cancel()
<-workerCanceled
select {
case err := <-done:
t.Fatalf("runComponents returned before inflight worker completed: %v", err)
case <-time.After(20 * time.Millisecond):
}
close(releaseWorker)
if err := <-done; err != nil {
t.Fatal(err)
}
}
func TestRunComponentsFailsClosedWhenWorkerStops(t *testing.T) {
httpStopped := make(chan struct{})
workerErr := errors.New("worker unavailable")
err := runComponents(context.Background(), func() error {
<-httpStopped
return http.ErrServerClosed
}, func(context.Context) error {
close(httpStopped)
return nil
}, func(context.Context) error {
return workerErr
})
if err == nil || !strings.Contains(err.Error(), workerErr.Error()) {
t.Fatalf("expected worker failure, got %v", err)
}
}
+174
View File
@@ -0,0 +1,174 @@
package main
import (
"bytes"
"context"
"fmt"
"image"
"image/color"
"image/png"
"os"
"strconv"
"time"
"git.ilapage.cn/OPC/chorus/internal/core/model"
platformstorage "git.ilapage.cn/OPC/chorus/internal/platform/storage"
"gorm.io/driver/mysql"
"gorm.io/gorm"
)
const fixtureLeaseToken = "e2e-fixture-token"
func main() {
if err := run(); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
}
func run() error {
if len(os.Args) != 4 {
return fmt.Errorf("usage: fixture <running|text-success|image-success|failed> <generation-id> <user-email>")
}
dsn := os.Getenv("CHORUS_E2E_DSN")
if dsn == "" {
return fmt.Errorf("CHORUS_E2E_DSN is required")
}
generationID, err := strconv.ParseUint(os.Args[2], 10, 64)
if err != nil || generationID == 0 {
return fmt.Errorf("generation id is invalid")
}
db, err := gorm.Open(mysql.Open(dsn), &gorm.Config{})
if err != nil {
return fmt.Errorf("connect fixture database")
}
sqlDB, err := db.DB()
if err != nil {
return fmt.Errorf("configure fixture database")
}
defer sqlDB.Close()
switch os.Args[1] {
case "running":
return setRunning(db, generationID, os.Args[3])
case "text-success":
return setTextSuccess(db, generationID, os.Args[3])
case "image-success":
return setImageSuccess(db, generationID, os.Args[3])
case "failed":
return setFailed(db, generationID, os.Args[3])
default:
return fmt.Errorf("fixture action is invalid")
}
}
func ownedGeneration(db *gorm.DB, generationID uint64, email string) (model.Generation, error) {
var generation model.Generation
err := db.Table("generations g").Select("g.*").Joins("JOIN users u ON u.id=g.user_id").Where("g.id=? AND u.email=?", generationID, email).Take(&generation).Error
if err != nil {
return model.Generation{}, fmt.Errorf("find owned generation")
}
return generation, nil
}
func setRunning(db *gorm.DB, generationID uint64, email string) error {
generation, err := ownedGeneration(db, generationID, email)
if err != nil || generation.Status != model.StatusPending {
return fmt.Errorf("generation is not owned pending")
}
result := db.Model(&model.Generation{}).Where("id=? AND status=?", generationID, model.StatusPending).Updates(map[string]any{
"status": model.StatusRunning, "lease_owner": "e2e-fixture", "lease_token": fixtureLeaseToken,
"lease_until": time.Now().Add(5 * time.Minute), "attempt_count": 1,
"attempts": gorm.Expr("JSON_ARRAY(JSON_OBJECT('attempt',1,'started_at',NOW(6),'latency_ms',1))"),
})
if result.Error != nil || result.RowsAffected != 1 {
return fmt.Errorf("set fixture generation running")
}
return nil
}
func setTextSuccess(db *gorm.DB, generationID uint64, email string) error {
generation, err := ownedGeneration(db, generationID, email)
if err != nil || generation.Kind != model.KindText || generation.Status != model.StatusRunning {
return fmt.Errorf("generation is not owned running text")
}
text := "MVP-0 browser text result"
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.Create(&model.GenerationOutput{GenerationID: generationID, Kind: model.KindText, TextContent: &text}).Error; err != nil {
return err
}
return finish(tx, generationID, model.StatusSucceeded, "", "")
})
}
func setImageSuccess(db *gorm.DB, generationID uint64, email string) error {
generation, err := ownedGeneration(db, generationID, email)
if err != nil || generation.Kind != model.KindImage || generation.Status != model.StatusRunning {
return fmt.Errorf("generation is not owned running image")
}
root := os.Getenv("CHORUS_E2E_STORAGE_ROOT")
if root == "" {
return fmt.Errorf("CHORUS_E2E_STORAGE_ROOT is required")
}
local, err := platformstorage.NewLocal(platformstorage.Config{Root: root, MaxObjectBytes: 1 << 20, MaxImagePixels: 10000, ThumbnailMaxSide: 32, AllowedImageMIME: map[string]bool{"image/png": true}})
if err != nil {
return err
}
key := fmt.Sprintf("e2e/%d/output", generationID)
objects, err := local.PutImage(context.Background(), platformstorage.ImageRequest{Key: key, ThumbnailKey: key + "-thumbnail", OwnerID: generation.UserID, GenerationID: generationID, ContentType: "image/png", Source: bytes.NewReader(testPNG())})
if err != nil {
return err
}
committed := false
defer func() {
if !committed {
_ = local.Delete(context.Background(), objects.Original.Key)
_ = local.Delete(context.Background(), objects.Thumbnail.Key)
}
}()
mimeType := objects.Original.ContentType
size := uint64(objects.Original.Size)
originalKey, thumbnailKey := objects.Original.Key, objects.Thumbnail.Key
err = db.Transaction(func(tx *gorm.DB) error {
output := model.GenerationOutput{GenerationID: generationID, Kind: model.KindImage, StorageKey: &originalKey, ThumbnailStorageKey: &thumbnailKey, MIMEType: &mimeType, SizeBytes: &size}
if err := tx.Create(&output).Error; err != nil {
return err
}
return finish(tx, generationID, model.StatusSucceeded, "", "")
})
committed = err == nil
return err
}
func setFailed(db *gorm.DB, generationID uint64, email string) error {
generation, err := ownedGeneration(db, generationID, email)
if err != nil || generation.Status != model.StatusRunning {
return fmt.Errorf("generation is not owned running")
}
return finish(db, generationID, model.StatusFailed, "bad_request", "upstream request failed")
}
func finish(db *gorm.DB, generationID uint64, status model.GenerationStatus, code, message string) error {
updates := map[string]any{"status": status, "completed_at": time.Now(), "lease_owner": nil, "lease_token": nil, "lease_until": nil}
if status == model.StatusFailed {
updates["error_code"] = code
updates["error_message"] = message
}
result := db.Model(&model.Generation{}).Where("id=? AND status=? AND lease_token=?", generationID, model.StatusRunning, fixtureLeaseToken).Updates(updates)
if result.Error != nil || result.RowsAffected != 1 {
return fmt.Errorf("finish fixture generation")
}
return nil
}
func testPNG() []byte {
img := image.NewRGBA(image.Rect(0, 0, 4, 4))
for y := 0; y < 4; y++ {
for x := 0; x < 4; x++ {
img.Set(x, y, color.RGBA{R: 44, G: 111, B: 173, A: 255})
}
}
var output bytes.Buffer
_ = png.Encode(&output, img)
return output.Bytes()
}
+81 -1
View File
@@ -1,4 +1,6 @@
const { test, expect } = require("@playwright/test");
const { execFileSync } = require("node:child_process");
const path = require("node:path");
const email = process.env.CHORUS_E2E_EMAIL;
const password = process.env.CHORUS_E2E_PASSWORD;
@@ -43,7 +45,7 @@ test.describe("portal responsive workflow", () => {
await expect(page.locator("#composer-form")).toBeVisible();
await expect(page.locator("#result-section")).toBeVisible();
await expect(page.locator(".history-end")).toBeVisible();
await expect(page.locator(".history-end, .history-empty")).toBeVisible();
await expectNoHorizontalOverflow(page);
await expectMinimumTargetSize(page, ".button, .icon-button, .mode-button, .history-item");
const width = testInfo.project.use.viewport.width;
@@ -88,6 +90,84 @@ test.describe("portal responsive workflow", () => {
});
});
test("browser observes text, image, and failed terminal states", async ({ page }, testInfo) => {
test.skip(testInfo.project.name !== "desktop-small", "terminal state flow runs once");
test.skip(!email || !password || !process.env.CHORUS_E2E_DSN || !process.env.CHORUS_E2E_STORAGE_ROOT, "E2E fixture environment is required");
test.setTimeout(60_000);
await login(page);
await expect(page.locator("#composer-form")).toBeVisible();
await page.locator('[data-mode="text"]').click();
await expect(page.locator('[data-mode="text"]')).toHaveAttribute("aria-pressed", "true");
await page.locator("#prompt").fill("browser text success");
await page.locator("#generate-button").click();
await expect(page).toHaveURL(/\/generations\/\d+$/);
const textID = generationID(page.url());
fixture("running", textID);
await page.reload();
await expect(page.getByText("正在生成内容")).toBeVisible();
fixture("text-success", textID);
await expect(page.getByText("MVP-0 browser text result")).toBeVisible({ timeout: 5_000 });
await expect(page.locator("#result-section")).not.toHaveAttribute("hx-get");
await page.goto("/");
await page.locator('[data-mode="image"]').click();
await expect(page.locator('[data-mode="image"]')).toHaveAttribute("aria-pressed", "true");
await page.locator("#image-upload").setInputFiles({
name: "source.png",
mimeType: "image/png",
buffer: Buffer.from("iVBORw0KGgoAAAANSUhEUgAAAAQAAAAECAIAAAAmkwkpAAAAFElEQVR4nGPkWfSfAQkwMaABFQAAZQ4BAagN7JwAAAAASUVORK5CYII=", "base64")
});
await page.locator("#prompt").fill("browser image success");
await page.locator("#generate-button").click();
await expect(page).toHaveURL(/\/generations\/\d+$/);
const imageID = generationID(page.url());
fixture("running", imageID);
await page.reload();
await expect(page.getByText("正在生成内容")).toBeVisible();
fixture("image-success", imageID);
await expect(page.locator(".result-gallery img")).toBeVisible({ timeout: 5_000 });
await expect(page.locator("#result-section")).not.toHaveAttribute("hx-get");
await page.goto("/");
await page.locator('[data-mode="text"]').click();
await expect(page.locator('[data-mode="text"]')).toHaveAttribute("aria-pressed", "true");
await page.locator("#prompt").fill("browser classified failure");
await page.locator("#generate-button").click();
await expect(page).toHaveURL(/\/generations\/\d+$/);
const failedID = generationID(page.url());
fixture("running", failedID);
await page.reload();
fixture("failed", failedID);
await expect(page.getByText("这次生成没有完成")).toBeVisible({ timeout: 5_000 });
await expect(page.getByText("upstream request failed")).toBeVisible();
await expect(page.locator("#result-section")).not.toHaveAttribute("hx-get");
});
async function login(page) {
await page.goto("/login");
await page.locator("#email").fill(email);
await page.locator("#password").fill(password);
await page.locator("#login-button").click();
await expect(page).toHaveURL(/\/$/);
}
function generationID(url) {
const match = url.match(/\/generations\/(\d+)$/);
if (!match) throw new Error(`generation id missing from ${url}`);
return match[1];
}
function fixture(action, id) {
const go = process.env.CHORUS_GO || "go";
const repoRoot = path.resolve(__dirname, "../../..");
execFileSync(go, ["run", "./portal/web/e2e/fixture", action, String(id), email], {
cwd: repoRoot,
env: process.env,
stdio: "pipe"
});
}
async function expectNoHorizontalOverflow(page) {
const hasNoOverflow = await page.evaluate(() => document.documentElement.scrollWidth <= document.documentElement.clientWidth);
expect(hasNoOverflow).toBe(true);