mirror of
https://gitcode.com/JianFeeeee/ModelRouter.git
synced 2026-10-05 07:02:29 +00:00
被问"还有 auto 调度相关 stage 呢?"问出来的真实缺口。
## 问题
chainDrive 只返回 (resp, src, model, err),调用方只知道**最终哪个槽位赢了**。
遍历过程中算出来又丢掉的东西——哪些档被跳过、为什么跳过、哪些槽位硬失败、
哪档全忙——一律不可见。ChainErr 里其实有这些,但**只在全部失败时**才填,
而它是 error 返回值不是记录。于是:
"tier 1 冷却所以降级到 tier 3" == "tier 1 正常接单"
对插件而言 tier 只是个常量 -2("resolved by the chain"),信息量为零。而这
恰恰是优先级链存在的全部理由,也是"我那个贵模型为什么没被用"的答案。
## 做法(scheduler 侧零新依赖)
新增 TraceEvent / TraceSink,chainDrive 多一个可选 sink 参数:
- TraceEvent 是本包的普通 struct,sink 是 func 参数 ⇒ **不新增 import**,
scheduler 仍然可独立测试
- sink 为 nil 时每次 emit 只多一次 nil 判断;没有插件的网关在 AUTO 热路径上
零开销(gateway 的 chainTraceSink 直接返回 nil)
- 事件是纯观测:scheduler 不基于它做任何分支,gateway 也不把它喂回路由/
冷却/配额
四种 kind:tier_skip / slot_fail / tier_busy / selected,selected 每次成功
遍历恰好一次且是最后一步。顺序保证所有 step 在 routed 之前。
## 暴露给插件
新增 chain_step stage(逐个步骤),并在 request_end 载荷里加三个便于做报表的
字段:chain_walk(上限 12 步,防审计记录膨胀)、degraded、tier_served。
## ★ 计费口径(我按推荐的做,已写进文档,需要你确认)
**按实际服务的模型计费**:降级到 tier 3 仍按 tier 3 的价算,轨迹只作观测。
理由与 §7.5 的边界一致——插件只报表不执法,两套口径混在一起会引出"降级该不该
多收钱"这种无法从代码判断的争议。若要改成"按本该用的档计价",需要在 models
价目里允许按 tier 定价,这我没做,因为那是个产品决策。
## 计费插件同步消费
by_tier_served / skip_reasons / degraded_reqs 三个新维度。skip_reasons 的等待
时长做了归一(`no free slot within <wait>`),否则 busy-wait 文案一变就多一行。
降级次数在 request_end 里计而不是在 chain_step 里计:一次降级的请求要走多步,
按步计会重复计数。
## 判据(346 个测试全绿,新增 15 个)
scheduler 6 个:正常路径只发一个 selected / 跳档+降级可见 / 硬失败与跳档
严格区分(不可混为一谈,否则抖动上游看起来像空闲上游)/
nil sink 安全 / 全失败时轨迹与 ChainErr 并存且不互相破坏 /
空链不发事件
gateway 1 个端到端:tier 1 全 500 → 插件收到 slot_fail(tier 1) +
selected(tier 2),request_end 的 tier_served=2 且 degraded=true
lua 2 个:降级计数与按实际模型计价 / 跳过原因归一聚合
lua 1 个:chain_step 是真 stage 且顺序正确
3 个变异都红:去掉 slot_fail(3 个判据红)/ 去掉 tier_skip(1 个)/
去掉 degraded 字段(1 个)。
192 lines
6.4 KiB
Go
192 lines
6.4 KiB
Go
package scheduler
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"strings"
|
|
"testing"
|
|
)
|
|
|
|
// The chain trace is the only way a caller learns that a request was DEGRADED
|
|
// — served by a lower tier than the one that should have taken it. chainDrive's
|
|
// return value carries only the winner, so without these events "tier 1 was
|
|
// cooling and we dropped to tier 2" is indistinguishable from "tier 1 served
|
|
// it", which is the exact question a priority chain exists to answer.
|
|
|
|
// recorder collects trace events for assertions.
|
|
type recorder struct{ events []TraceEvent }
|
|
|
|
func (r *recorder) sink(ev TraceEvent) { r.events = append(r.events, ev) }
|
|
|
|
func (r *recorder) kinds() []TraceKind {
|
|
out := make([]TraceKind, 0, len(r.events))
|
|
for _, e := range r.events {
|
|
out = append(out, e.Kind)
|
|
}
|
|
return out
|
|
}
|
|
|
|
func (r *recorder) find(k TraceKind) *TraceEvent {
|
|
for i := range r.events {
|
|
if r.events[i].Kind == k {
|
|
return &r.events[i]
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// TestTraceSelectedOnlyOnHappyPath: a clean walk emits exactly one event.
|
|
func TestTraceSelectedOnlyOnHappyPath(t *testing.T) {
|
|
p1 := fakeProv("p1", "m1")
|
|
ch := BuildChain([]Rule{{Model: "m1", Source: "p1", Tier: 1}},
|
|
func(model, source string) Provider { return p1 })
|
|
var rec recorder
|
|
_, _, _, err := New(3).ChainChat(context.Background(), ch, chatReq(), nil, rec.sink)
|
|
if err != nil {
|
|
t.Fatalf("chain: %v", err)
|
|
}
|
|
if got := rec.kinds(); len(got) != 1 || got[0] != TraceSelected {
|
|
t.Errorf("events = %v, want a single selected", got)
|
|
}
|
|
e := rec.find(TraceSelected)
|
|
if e.Tier != 1 || e.Source != "p1" || e.Model != "m1" {
|
|
t.Errorf("selected event = %+v, want tier 1 / p1 / m1", e)
|
|
}
|
|
if e.Attempt != 1 {
|
|
t.Errorf("Attempt = %d, want 1", e.Attempt)
|
|
}
|
|
}
|
|
|
|
// TestTraceRecordsTierSkipAndDegradation is the core case: tier 1 is
|
|
// unschedulable, tier 2 answers. The trace must show the skip AND the eventual
|
|
// selection, so a consumer can see the request was served one tier down.
|
|
func TestTraceRecordsTierSkipAndDegradation(t *testing.T) {
|
|
// p1 is unavailable (not probeable), so tier 1 yields no candidates.
|
|
p1 := fakeProv("p1", "m1")
|
|
p1.available.Store(false)
|
|
p2 := fakeProv("p2", "m2")
|
|
ch := BuildChain([]Rule{
|
|
{Model: "m1", Source: "p1", Tier: 1},
|
|
{Model: "m2", Source: "p2", Tier: 2},
|
|
}, bySource(p1, p2))
|
|
var rec recorder
|
|
_, src, model, err := New(3).ChainChat(context.Background(), ch, chatReq(), nil, rec.sink)
|
|
if err != nil {
|
|
t.Fatalf("chain: %v", err)
|
|
}
|
|
if src != "p2" || model != "m2" {
|
|
t.Fatalf("served by %s/%s, want p2/m2", src, model)
|
|
}
|
|
skip := rec.find(TraceTierSkip)
|
|
if skip == nil {
|
|
t.Fatalf("no tier_skip event; events = %v", rec.kinds())
|
|
}
|
|
if skip.Tier != 1 {
|
|
t.Errorf("skip tier = %d, want 1", skip.Tier)
|
|
}
|
|
if !strings.Contains(skip.Reason, "cooling") {
|
|
t.Errorf("skip reason = %q, want it to mention cooling", skip.Reason)
|
|
}
|
|
sel := rec.find(TraceSelected)
|
|
if sel == nil || sel.Tier != 2 {
|
|
t.Errorf("selected = %+v, want tier 2", sel)
|
|
}
|
|
// The order matters: the skip must be observable BEFORE the selection.
|
|
if rec.events[0].Kind != TraceTierSkip || rec.events[len(rec.events)-1].Kind != TraceSelected {
|
|
t.Errorf("event order = %v, want skip first and selected last", rec.kinds())
|
|
}
|
|
}
|
|
|
|
// TestTraceRecordsHardSlotFailures: a slot that returns an upstream error is a
|
|
// different event from a skip — the request tried it and it failed. Losing that
|
|
// distinction makes a flaky upstream look like an idle one.
|
|
func TestTraceRecordsHardSlotFailures(t *testing.T) {
|
|
p1 := fakeProv("p1", "m1")
|
|
p1.fail.Store(true) // Chat returns "upstream error"
|
|
p2 := fakeProv("p2", "m2")
|
|
ch := BuildChain([]Rule{
|
|
{Model: "m1", Source: "p1", Tier: 1},
|
|
{Model: "m2", Source: "p2", Tier: 2},
|
|
}, bySource(p1, p2))
|
|
var rec recorder
|
|
_, src, _, err := New(3).ChainChat(context.Background(), ch, chatReq(), nil, rec.sink)
|
|
if err != nil {
|
|
t.Fatalf("chain: %v", err)
|
|
}
|
|
if src != "p2" {
|
|
t.Fatalf("served by %s, want p2", src)
|
|
}
|
|
fail := rec.find(TraceSlotFail)
|
|
if fail == nil {
|
|
t.Fatalf("no slot_fail event; events = %v", rec.kinds())
|
|
}
|
|
if fail.Tier != 1 || fail.Source != "p1" || fail.Model != "m1" {
|
|
t.Errorf("slot_fail = %+v, want tier 1 / p1 / m1", fail)
|
|
}
|
|
if !strings.Contains(fail.Err, "upstream error") {
|
|
t.Errorf("slot_fail error = %q, want the upstream text", fail.Err)
|
|
}
|
|
// A hard failure must NOT be reported as a skip.
|
|
if rec.find(TraceTierSkip) != nil {
|
|
t.Error("a hard failure was also reported as a tier_skip")
|
|
}
|
|
}
|
|
|
|
// TestTraceNilSinkIsSafe: the gateway passes nil when no plugin is loaded, so
|
|
// every emit path must tolerate it. This is the "plugins are optional" property
|
|
// on the scheduler side.
|
|
func TestTraceNilSinkIsSafe(t *testing.T) {
|
|
p1 := fakeProv("p1", "m1")
|
|
p1.fail.Store(true)
|
|
p2 := fakeProv("p2", "m2")
|
|
ch := BuildChain([]Rule{
|
|
{Model: "m1", Source: "p1", Tier: 1},
|
|
{Model: "m2", Source: "p2", Tier: 2},
|
|
}, bySource(p1, p2))
|
|
if _, _, _, err := New(3).ChainChat(context.Background(), ch, chatReq(), nil, nil); err != nil {
|
|
t.Fatalf("a nil trace sink broke the walk: %v", err)
|
|
}
|
|
}
|
|
|
|
// TestTraceOnTotalFailure: when every tier fails, the walk still emits its
|
|
// per-step events AND returns the ChainErr. The trace is additive — it must not
|
|
// replace or disturb the error contract callers depend on for the 503.
|
|
func TestTraceOnTotalFailure(t *testing.T) {
|
|
p1 := fakeProv("p1", "m1")
|
|
p1.fail.Store(true)
|
|
p2 := fakeProv("p2", "m2")
|
|
p2.fail.Store(true)
|
|
ch := BuildChain([]Rule{
|
|
{Model: "m1", Source: "p1", Tier: 1},
|
|
{Model: "m2", Source: "p2", Tier: 2},
|
|
}, bySource(p1, p2))
|
|
var rec recorder
|
|
_, _, _, err := New(3).ChainChat(context.Background(), ch, chatReq(), nil, rec.sink)
|
|
var ce *ChainErr
|
|
if !errors.As(err, &ce) {
|
|
t.Fatalf("err = %v, want a *ChainErr so the gateway can answer 503", err)
|
|
}
|
|
if len(ce.Tiers) != 2 {
|
|
t.Errorf("ChainErr.Tiers = %d, want 2 (the error contract must be unchanged)", len(ce.Tiers))
|
|
}
|
|
if n := len(rec.kinds()); n != 2 {
|
|
t.Errorf("events = %v, want two slot_fail and no selection", rec.kinds())
|
|
}
|
|
if rec.find(TraceSelected) != nil {
|
|
t.Error("a selected event was emitted for a walk that served nothing")
|
|
}
|
|
}
|
|
|
|
// TestTraceSkipsEmptyChain: no chain configured must not emit anything; the
|
|
// gateway answers 503 before scheduling in that case anyway.
|
|
func TestTraceSkipsEmptyChain(t *testing.T) {
|
|
var rec recorder
|
|
_, _, _, err := New(3).ChainChat(context.Background(), &Chain{}, chatReq(), nil, rec.sink)
|
|
if err == nil {
|
|
t.Fatal("expected an error for an empty chain")
|
|
}
|
|
if len(rec.events) != 0 {
|
|
t.Errorf("events = %v, want none", rec.kinds())
|
|
}
|
|
}
|