将非流式 Gemini generateContent 改为 River Worker 执行,API 通过批量状态查询等待完成,避免高并发同步请求耗尽 API 数据库连接池。\n\n生成图片在持久化时保存带哈希的资产引用,响应恢复时下载并校验后重建 Base64,数据库不保存媒体原文。请求体物化增加前置内存门禁,并扩展双 API/Worker PostgreSQL 压力测试。\n\n验证:go test ./... -count=1;go vet ./...;pnpm openapi;迁移安全检查;64 与 256 请求双角色异步 Gemini 压力测试。
51 lines
1.3 KiB
Go
51 lines
1.3 KiB
Go
package runner
|
|
|
|
import (
|
|
"testing"
|
|
"time"
|
|
)
|
|
|
|
func TestTaskCompletionWaitersAreBatchedAndSignaled(t *testing.T) {
|
|
service := &Service{
|
|
taskCompletionWaiters: map[string]map[chan struct{}]struct{}{},
|
|
taskCompletionPollWake: make(chan struct{}, 1),
|
|
}
|
|
first, unregisterFirst := service.registerTaskCompletionWaiter("task-1")
|
|
defer unregisterFirst()
|
|
second, unregisterSecond := service.registerTaskCompletionWaiter("task-2")
|
|
|
|
taskIDs := service.taskCompletionWaiterIDs(10)
|
|
if len(taskIDs) != 2 {
|
|
t.Fatalf("batched task IDs=%v", taskIDs)
|
|
}
|
|
service.signalTaskCompletion("task-1")
|
|
select {
|
|
case <-first:
|
|
case <-time.After(time.Second):
|
|
t.Fatal("task completion waiter was not signaled")
|
|
}
|
|
select {
|
|
case <-second:
|
|
t.Fatal("unrelated task waiter was signaled")
|
|
default:
|
|
}
|
|
|
|
unregisterSecond()
|
|
if got := service.taskCompletionWaiterIDs(10); len(got) != 1 || got[0] != "task-1" {
|
|
t.Fatalf("waiter unregister left unexpected IDs=%v", got)
|
|
}
|
|
}
|
|
|
|
func TestTerminalTaskStatus(t *testing.T) {
|
|
for _, status := range []string{"succeeded", "failed", "cancelled", "manual_review"} {
|
|
if !terminalTaskStatus(status) {
|
|
t.Fatalf("status %q should be terminal", status)
|
|
}
|
|
}
|
|
for _, status := range []string{"queued", "running", "pending"} {
|
|
if terminalTaskStatus(status) {
|
|
t.Fatalf("status %q should not be terminal", status)
|
|
}
|
|
}
|
|
}
|