Files
ModelRouter/internal/scheduler/scheduler.go
JianFeeeee c19b8e6394 fix(scheduler): AUTO 遇空内容响应降级到下一个 slot(客户端曾报 "no content")
生产故障:pi 客户端报 `model "AUTO" returned a completed response with no content`,
重试延迟 8 秒。实测 AUTO 20 次有 2 次返回空 content,全部是 claude-opus-4-8。

根因(直连上游抓包确认):思考型模型先吐 reasoning_content,max_tokens 小到
思考阶段就把预算用完时,上游返回 200 / finish_reason=length,28 个 chunk 全是
reasoning_content、content 一片空白。runTier 只看 err == nil 就当成功返回,
客户端拿到一个空响应。

修复:
- resultIsEmpty:非流式路径把「无 content、无 tool_calls、无 image」的响应当作
  slot 失败继续降级。注意 ReasoningContent 不算内容——客户端要的是文本,为
  另一个模型的思考阶段扣住请求比降级更糟。
- peekStream:流式路径在出现首个真实内容前缓冲 reasoning 前导,流结束仍无内容
  则回报空结果,让 chainDrive 换 slot。缓冲只覆盖思考前导,拿到内容后立即
  转发。tool_calls delta 算内容,agent 回合不会被误判。
- 3 条判据 + 3 个变异(恒 false / 恒 true / reasoning 算内容)全部被捕获。

同时修三个 WebUI 布局缺陷(都靠截图而非 DOM 断言发现):
- 插件侧栏项只渲染图标没有标题:btn.innerHTML 只塞 pluginIconHTML(pg.icon),
  与原生页的「图标 + <span>标题</span>」不一致,侧栏是一排无名图标。
- 计费维度表 8 列挤在 465px 卡片里:table-layout:fixed 把每列压到 62px,
  23/80 个单元格溢出、数字互相重叠。改为 6 列(token 细分合并为
  「输入(新鲜+缓存)」,细分进 title)+ table-layout:auto,实测 0/60 溢出。
- 数字列 word-break:break-all 让每个字符独占一行(USD 0.56 竖排成 U/S/D),
  改 nowrap + 容器横向滚动。
2026-10-02 15:10:01 +08:00

785 lines
28 KiB
Go

// Package scheduler implements request scheduling across providers: direct
// fallback scheduling over candidate lists, and the AUTO chain (tiers with
// per-tier round-robin cursors, preference ordering, token-quota windows and
// per-(source,model) cooldown awareness) per the target architecture in
// plan.md.
package scheduler
import (
"context"
"errors"
"fmt"
"sort"
"strings"
"sync"
"sync/atomic"
"time"
"llmsproxy/internal/types"
)
// busyWait is how long a fully-busy tier is polled for a free slot before the
// request falls through to the next tier (bounded wait, plan 2.3).
var busyWait = 2 * time.Second
// busyPoll is the polling interval while waiting for a busy tier.
var busyPoll = 100 * time.Millisecond
// Scheduler drives one chat tool call across the candidate provider chain.
type Scheduler struct {
// MaxRetries how many fallback providers to try before failing.
MaxRetries int
}
func New(maxRetries int) *Scheduler {
if maxRetries < 0 {
maxRetries = 0
}
return &Scheduler{MaxRetries: maxRetries}
}
// Provider is the minimal interface the scheduler needs to schedule over.
type Provider interface {
Name() string
ModelFor(reqModel string) string
ModelAvailable(model string) bool
// ModelSchedulable is the probe-aware availability gate: ok reports
// whether the model may take a request now, isProbe marks that it is only
// allowed as a cooldown probe (the caller must release the permit with
// ProbeDone once the attempt finished).
ModelSchedulable(model string) (ok bool, isProbe bool)
ProbeDone(model string)
Pref(model string) int64
Chat(ctx context.Context, req *types.ChatRequest) (*types.UnifiedResponse, error)
ChatStream(ctx context.Context, req *types.ChatRequest) (<-chan types.UnifiedChunk, error)
Image(ctx context.Context, req *types.ImageGenRequest) (*types.UnifiedResponse, error)
}
// ---- AUTO chain ----
// Rule is one persisted AUTO chain slot (mirror of config.ModelScope).
type Rule struct {
Model string
Source string
Tier int
Quota int64
Period string
Hours int64
}
// Slot is one schedulable chain position: a model pinned to its provider,
// with an optional token-quota window. Slots are immutable after build.
type Slot struct {
Model string
Source string
Quota int64
Period string
Hours int64
Prov Provider
}
// TierNode is one priority tier. Slots keep their configured order (the
// stable base for preference ordering). next is the round-robin cursor: it
// holds the last used slot index (-1 = none yet), so the very first request
// starts at the configured order and later ones rotate.
type TierNode struct {
Tier int
Slots []*Slot
next atomic.Int64
}
// NextStart advances the tier cursor and returns the start index for the next
// scheduling run (first run: index 0).
func (tn *TierNode) NextStart() int64 {
return tn.next.Add(1)
}
// Chain is the immutable AUTO scheduling plan. A rebuilt chain is swapped in
// atomically; per-tier cursors live inside the chain and are shared across
// requests (rotation state resets when the chain is rebuilt, e.g. after
// editing the rules — acceptable, the swap also resets cooldowns).
type Chain struct {
Tiers []*TierNode // ascending tier order (tier 1 = highest priority, tried first)
}
// TierErrors is the per-tier failure summary carried by ChainErr. Errors
// (TierError or skipped-tier reasons) are collected in tier order.
type TierError struct {
Tier int
Source string
Model string
Err error
}
// ChainErr is returned by chain scheduling when every AUTO tier failed. Its
// message summarizes each failed tier (which source/model and why) so a 503
// names the culprits instead of the bare "no provider available".
type ChainErr struct {
Tiers []TierError
Skipped []string // whole-tier reasons (cooling / quota / all busy)
}
func (e *ChainErr) Error() string {
var b strings.Builder
b.WriteString("all auto tiers failed: ")
first := true
for _, t := range e.Tiers {
if !first {
b.WriteString("; ")
}
first = false
fmt.Fprintf(&b, "tier %d %s/%s: %v", t.Tier, t.Source, t.Model, t.Err)
}
for _, s := range e.Skipped {
if !first {
b.WriteString("; ")
}
first = false
b.WriteString(s)
}
return b.String()
}
// BuildChain groups rules into ascending tiers (tier 1 = highest priority,
// tried first) and resolves each slot's provider via prov. Rules whose
// provider resolves to nil are dropped (the source no longer serves the
// model). Slot order within a tier follows the configured rule order.
func BuildChain(rules []Rule, prov func(model, source string) Provider) *Chain {
byTier := map[int][]*Slot{}
var tiers []int
for _, r := range rules {
p := prov(r.Model, r.Source)
if p == nil {
continue
}
if _, ok := byTier[r.Tier]; !ok {
tiers = append(tiers, r.Tier)
}
byTier[r.Tier] = append(byTier[r.Tier], &Slot{
Model: r.Model,
Source: r.Source,
Quota: r.Quota,
Period: r.Period,
Hours: r.Hours,
Prov: p,
})
}
sort.Slice(tiers, func(i, j int) bool { return tiers[i] < tiers[j] })
ch := &Chain{}
for _, t := range tiers {
tn := &TierNode{Tier: t, Slots: byTier[t]}
tn.next.Store(-1)
ch.Tiers = append(ch.Tiers, tn)
}
return ch
}
// tierResult is the outcome of one scheduling run over one tier.
type tierResult struct {
resp *types.UnifiedResponse
chunks <-chan types.UnifiedChunk
src string
model string
hard []TierError // hard failures seen in this pass (nil = none)
}
// candidate is one schedulable slot in a tier pass. probe marks a slot that is
// still cooling but past its half-cooldown mark and holding the probe permit:
// it is tried only after every normal slot, so probe traffic is what the tier
// falls back to instead of what it prefers.
type candidate struct {
slot *Slot
probe bool
}
// collectCands partitions a tier's slots into normal and probe candidates,
// normal first. Quota-exhausted slots are dropped outright. Every probe
// candidate returned holds a probe permit, so the caller MUST call
// releaseProbes on the result exactly once.
func collectCands(slots []*Slot, exhausted func(*Slot) bool) []candidate {
var normal, probes []candidate
for _, sl := range slots {
if exhausted != nil && exhausted(sl) {
continue
}
ok, isProbe := sl.Prov.ModelSchedulable(sl.Model)
if !ok {
continue
}
if isProbe {
probes = append(probes, candidate{slot: sl, probe: true})
continue
}
normal = append(normal, candidate{slot: sl})
}
return append(normal, probes...)
}
// releaseProbes hands every claimed probe permit back, whether or not the probe
// slot was actually used.
func releaseProbes(cands []candidate) {
for _, c := range cands {
if c.probe {
c.slot.Prov.ProbeDone(c.slot.Model)
}
}
}
// normalCount is how many leading candidates are normal (non-probe). The
// round-robin cursor rotates only over those: probe slots are a strictly
// ordered tail, never a rotation target.
func normalCount(cands []candidate) int {
for i, c := range cands {
if c.probe {
return i
}
}
return len(cands)
}
// runTier executes one tier pass starting at the round-robin base index.
// Cooldown is the only hard skip (re-verified per slot); a busy slot is
// skipped without any penalty; a hard failure is recorded and the pass moves
// on to the next slot (plan 2.3: "单请求内不重试已失败槽" — the failed slot is
// not retried, the others still are). hard == nil and no success means every
// candidate was merely busy/cooling, so the caller may wait a bounded time.
//
// Normal candidates rotate by base; probe candidates form a fixed tail tried
// only after every normal slot failed or was busy.
// emptyResultReason describes why a 200-with-no-content response counts as a
// slot failure for AUTO.
//
// WHY: reasoning models (claude-opus-*, codebuddy_glm-*, …) emit
// `reasoning_content` first and only then `content`. When the caller's
// max_tokens is small enough that the thinking phase consumes the whole budget,
// upstream returns 200 / finish_reason=length with 28 chunks of reasoning and
// ZERO content. runTier treated "err == nil" as success and handed that to the
// client, which then failed with "returned a completed response with no
// content" — a client-side error message for what is really a bad slot choice.
//
// So an empty result is a SLOT failure, not a request failure: the gateway
// degrades to the next slot and the user still gets an answer. Measured on
// production AUTO: 2 of 20 requests returned empty content, all of them
// claude-opus-4-8.
//
// A response carrying tool_calls or image data is NOT empty: an agent turn
// legitimately produces tool calls with no text. ReasoningContent does NOT
// rescue it either — see resultIsEmpty.
const emptyResultReason = "upstream returned no content (reasoning-only response, or the token budget was consumed before any text)"
// resultIsEmpty reports whether a successful-but-useless response should be
// treated as a slot failure.
//
// Image data counts: an image-generation slot legitimately returns no text.
// A usage-only response is NOT empty either — the upstream answered, it just
// said nothing, and that is exactly the case worth degrading away from.
func resultIsEmpty(resp *types.UnifiedResponse) bool {
if resp == nil {
return false
}
// NOTE: ReasoningContent is deliberately NOT consulted. My first version
// excluded it ("the model was thinking, that is an answer"), and the test
// built from the real production capture failed immediately: the captured
// response is exactly reasoning_content-with-usage and zero text. The
// client asked for text and there is none; holding a request hostage to
// another model's thinking phase is strictly worse than degrading.
return strings.TrimSpace(resp.Content) == "" &&
len(resp.ToolCalls) == 0 &&
len(resp.ImageData) == 0
}
// emptyStreamReason is resultIsEmpty's streaming twin; see emptyResultReason
// for why an empty response is a slot failure rather than a request failure.
const emptyStreamReason = emptyResultReason
// peekStream wraps a chunk channel so the caller learns whether the stream
// produced real content BEFORE the chunks are forwarded.
//
// Why this is necessary: reasoning models emit reasoning_content first. With a
// small max_tokens the whole budget is spent thinking, the stream ends with
// finish_reason=length and zero content. If the gateway forwarded those chunks
// as they arrived, the client would already have seen a 200 SSE stream and
// could not be given a different slot — its only recourse is the useless
// "returned a completed response with no content" error. Buffering until the
// first real content (or the end of the stream) keeps the degrade path
// available at the cost of holding back the first few chunks.
//
// What is NOT buffered: the wrapper starts forwarding as soon as a chunk with
// non-empty Content or ToolCalls arrives, and keeps forwarding everything from
// then on, so only the reasoning preamble is held. Reasoning-only responses
// are dropped in full and reported as empty, which lets chainDrive try the
// next slot.
func peekStream(in <-chan types.UnifiedChunk) (<-chan types.UnifiedChunk, func() bool) {
out := make(chan types.UnifiedChunk, 16)
var (
mu sync.Mutex
sawText bool
done bool
)
go func() {
defer close(out)
started := false
for ck := range in {
if !started {
// Hold back the reasoning / usage-only preamble. A tool-call
// delta counts as content: an agent turn legitimately emits
// tool_calls with no text.
if strings.TrimSpace(ck.Content) == "" && len(ck.ToolCalls) == 0 {
continue
}
started = true
mu.Lock()
sawText = true
mu.Unlock()
}
out <- ck
}
mu.Lock()
done = true
mu.Unlock()
}()
// peek blocks until the stream either produces content or ends, then
// reports whether any content was seen. Polling a 2ms tick rather than
// using a second channel keeps peekStream single-goroutine and leak-free.
peek := func() bool {
for {
mu.Lock()
seen, finished := sawText, done
mu.Unlock()
if seen || finished {
return seen
}
time.Sleep(2 * time.Millisecond)
}
}
return out, peek
}
func runTier(ctx context.Context, tn *TierNode, cands []candidate, base int64, req *types.ChatRequest, stream bool) tierResult {
n := len(cands)
norm := normalCount(cands)
var hard []TierError
for i := 0; i < n; i++ {
var c candidate
if i < norm {
c = cands[(int(base)+i)%norm] // rotate within the normal head
} else {
c = cands[i] // probe tail keeps its order
}
sl := c.slot
// A probe candidate is intentionally NOT re-checked here: it is cooling
// by definition, and its permit was already claimed.
if !c.probe && !sl.Prov.ModelAvailable(sl.Model) {
continue
}
r := *req
r.Model = sl.Model
if stream {
chunks, err := sl.Prov.ChatStream(ctx, &r)
if err == nil {
guarded, peek := peekStream(chunks)
if peek() {
return tierResult{chunks: guarded, src: sl.Source, model: sl.Model}
}
// The stream finished with no content at all: a
// reasoning-only response. Drain and move on to the next
// slot instead of pinning the client to a useless stream.
hard = append(hard, TierError{
Tier: tn.Tier, Source: sl.Source, Model: sl.Model,
Err: errors.New(emptyStreamReason),
})
continue
}
if ctx.Err() != nil {
return tierResult{}
}
if errors.Is(err, types.ErrBusy) {
continue
}
hard = append(hard, TierError{Tier: tn.Tier, Source: sl.Source, Model: sl.Model, Err: err})
continue
}
resp, err := sl.Prov.Chat(ctx, &r)
if err == nil {
if resultIsEmpty(resp) {
// Soft failure: record it and try the next slot. Deliberately
// NOT a hard TierError — a hard error is reported to the client
// verbatim when the whole chain fails, and "this one model was
// unhelpful" is not the client's problem to debug.
hard = append(hard, TierError{
Tier: tn.Tier, Source: sl.Source, Model: sl.Model,
Err: errors.New(emptyResultReason),
})
continue
}
return tierResult{resp: resp, src: sl.Source, model: sl.Model}
}
if ctx.Err() != nil {
return tierResult{}
}
if errors.Is(err, types.ErrBusy) {
continue
}
hard = append(hard, TierError{Tier: tn.Tier, Source: sl.Source, Model: sl.Model, Err: err})
}
return tierResult{hard: hard}
}
// TraceKind classifies one step of an AUTO chain walk.
type TraceKind string
const (
// TraceTierSkip: the whole tier was skipped — every slot was cooling,
// quota-exhausted, or none was schedulable. Reason says which.
TraceTierSkip TraceKind = "tier_skip"
// TraceSlotFail: one slot failed hard (upstream error / bad adapter). The
// walk continues to the next slot or tier.
TraceSlotFail TraceKind = "slot_fail"
// TraceTierBusy: the tier was fully busy and the bounded wait expired.
TraceTierBusy TraceKind = "tier_busy"
// TraceSelected: this slot served the request. Exactly one per successful
// chain walk, and the last event emitted.
TraceSelected TraceKind = "selected"
)
// TraceEvent is one observable step of an AUTO chain walk.
//
// WHY THIS EXISTS: chainDrive's return value is (resp, src, model, err), so a
// caller learns only which slot finally served the request. Everything the
// scheduler decided on the way there — which tiers it skipped and WHY, which
// slots hard-failed, whether a tier was merely busy — was computed and then
// discarded. That is invisible to operators and to plugins: "tier 1 was cooling
// so we degraded to tier 3" looked exactly like "tier 1 served it".
//
// The walk already accumulates this in ChainErr, but ONLY on total failure, and
// ChainErr is an error return, not a record. Emitting a trace as it happens
// covers the far more common case: a request that SUCCEEDED after degrading.
//
// Design constraints:
// - scheduler stays dependency-free and independently testable. A TraceEvent
// is a plain struct in this package and the sink is a func parameter, so no
// import is added and no test has to change to observe a walk.
// - The sink is optional (nil = emit nothing). The overhead on the hot path
// is one nil check per event.
// - Events are OBSERVATION ONLY. Nothing in the scheduler branches on them,
// and the gateway does not feed them back into routing, cooldown or quota —
// see docs/plugins.md for why accounting and enforcement are kept apart.
type TraceEvent struct {
Kind TraceKind
Tier int
Source string
Model string
Reason string // human-readable, for TraceTierSkip / TraceSlotFail
Err string // the underlying error text, for TraceSlotFail
// Attempt counts the 1-based slot attempt within the whole walk.
Attempt int
}
// TraceSink receives chain-walk events. It must not block: it is called from the
// request path, and a slow sink slows the request.
type TraceSink func(TraceEvent)
// chainDrive runs a request down the chain (plan 2.3): tiers ascending (tier
// 1, the highest priority, first), per-tier round-robin starting at the tier
// cursor, same-tier runs ordered by preference (negative prefs sink but stay
// reachable). Quota-exhausted and cooling slots are filtered up front; a
// fully busy tier is polled for a bounded time before falling through.
// Failures are summarized in *ChainErr for the caller to map to HTTP 503.
//
// trace may be nil; when set it receives one event per observable step.
func (s *Scheduler) chainDrive(ctx context.Context, chain *Chain, req *types.ChatRequest, exhausted func(*Slot) bool, stream bool, trace TraceSink) (*types.UnifiedResponse, <-chan types.UnifiedChunk, string, string, error) {
emit := func(ev TraceEvent) {
if trace != nil {
trace(ev)
}
}
attempt := 0
if chain == nil || len(chain.Tiers) == 0 {
return nil, nil, "", "", fmt.Errorf("no auto slot configured")
}
var ce ChainErr
for _, tn := range chain.Tiers {
// initial filter: quota-exhausted slots are dropped, cooling slots are
// dropped unless they qualify as half-cooldown probes (appended last).
cands := collectCands(tn.Slots, exhausted)
if len(cands) == 0 {
reason := "no schedulable slot (cooling or quota exhausted)"
ce.Skipped = append(ce.Skipped, fmt.Sprintf("tier %d: %s", tn.Tier, reason))
emit(TraceEvent{Kind: TraceTierSkip, Tier: tn.Tier, Reason: reason})
continue
}
// No Pref sort: load balancing is done by round-robin cursor.
// Persistently failing slots are excluded by ModelSchedulable
// (which checks Pref > prefMin).
base := tn.NextStart()
res := runTier(ctx, tn, cands, base, req, stream)
if res.resp != nil || res.chunks != nil {
attempt++
emit(TraceEvent{Kind: TraceSelected, Tier: tn.Tier, Source: res.src, Model: res.model, Attempt: attempt})
releaseProbes(cands)
return res.resp, res.chunks, res.src, res.model, nil
}
if ctx.Err() != nil {
releaseProbes(cands)
return nil, nil, "", "", ctx.Err()
}
if len(res.hard) > 0 {
ce.Tiers = append(ce.Tiers, res.hard...)
for _, h := range res.hard {
attempt++
emit(TraceEvent{
Kind: TraceSlotFail, Tier: tn.Tier, Source: h.Source, Model: h.Model,
Err: types.OneLine(h.Err.Error(), 200), Attempt: attempt,
})
}
releaseProbes(cands)
continue // hard failures: fall through to the next tier, no waiting
}
// every candidate was busy or cooling: bounded poll before downgrading
if err := s.pollBusyTier(ctx, tn, cands, base, req, stream, &ce); err != nil {
releaseProbes(cands)
if r, ok := err.(*tierSuccess); ok {
attempt++
emit(TraceEvent{Kind: TraceSelected, Tier: tn.Tier, Source: r.res.src, Model: r.res.model, Attempt: attempt})
return r.res.resp, r.res.chunks, r.res.src, r.res.model, nil
}
return nil, nil, "", "", err
}
emit(TraceEvent{Kind: TraceTierBusy, Tier: tn.Tier, Reason: fmt.Sprintf("no free slot within %v", busyWait)})
releaseProbes(cands)
}
if len(ce.Tiers) == 0 && len(ce.Skipped) == 0 {
return nil, nil, "", "", fmt.Errorf("no auto slot configured")
}
return nil, nil, "", "", &ce
}
// tierSuccess carries a successful result out of pollBusyTier through the error
// return. It is never surfaced to callers of chainDrive.
type tierSuccess struct{ res tierResult }
func (t *tierSuccess) Error() string { return "tier success" }
// pollBusyTier waits a bounded time for a fully-busy tier to free a slot,
// retrying the pass while cooldowns expire. It returns nil when the tier should
// be abandoned (caller falls through to the next tier), a *tierSuccess when a
// retry succeeded, or a context error.
func (s *Scheduler) pollBusyTier(ctx context.Context, tn *TierNode, cands []candidate, base int64, req *types.ChatRequest, stream bool, ce *ChainErr) error {
deadline := time.Now().Add(busyWait)
timer := time.NewTimer(busyPoll)
defer timer.Stop()
for {
select {
case <-ctx.Done():
return ctx.Err()
case <-timer.C:
}
if time.Now().After(deadline) {
ce.Skipped = append(ce.Skipped, fmt.Sprintf("tier %d: no free slot within %v", tn.Tier, busyWait))
return nil
}
// refresh candidates: cooldowns may have expired meanwhile. Probe
// candidates keep their already-claimed permit and stay eligible.
var again []candidate
for _, c := range cands {
if c.probe || c.slot.Prov.ModelAvailable(c.slot.Model) {
again = append(again, c)
}
}
if len(again) == 0 {
ce.Skipped = append(ce.Skipped, fmt.Sprintf("tier %d: no free slot within %v", tn.Tier, busyWait))
return nil
}
res := runTier(ctx, tn, again, base, req, stream)
if res.resp != nil || res.chunks != nil {
return &tierSuccess{res: res}
}
if ctx.Err() != nil {
return ctx.Err()
}
if len(res.hard) > 0 {
ce.Tiers = append(ce.Tiers, res.hard...)
return nil // hard failure while waiting: stop waiting, fall through
}
timer.Reset(busyPoll)
}
}
// ChainChat runs a non-streaming AUTO request down the chain. exhausted, when
// non-nil, decides slot token-quota exhaustion. Returns the response, the
// serving source and the exact model id used; on total failure a *ChainErr
// summarizing every tier.
func (s *Scheduler) ChainChat(ctx context.Context, chain *Chain, req *types.ChatRequest, exhausted func(*Slot) bool, trace TraceSink) (*types.UnifiedResponse, string, string, error) {
resp, _, src, model, err := s.chainDrive(ctx, chain, req, exhausted, false, trace)
return resp, src, model, err
}
// ChainChatStream runs a streaming AUTO request down the chain. A slot is
// abandoned only on connect failures / busy (before its first chunk); after a
// stream starts it is pinned. Same return contract as ChainChat.
func (s *Scheduler) ChainChatStream(ctx context.Context, chain *Chain, req *types.ChatRequest, exhausted func(*Slot) bool, trace TraceSink) (<-chan types.UnifiedChunk, string, string, error) {
_, chunks, src, model, err := s.chainDrive(ctx, chain, req, exhausted, true, trace)
return chunks, src, model, err
}
// ChainImage runs an image-generation AUTO request down the chain: tiers
// ascending (tier 1 highest priority), per-tier round-robin, same-tier order
// by preference. Each slot's model is pinned to its own image id (ModelFor),
// so a fallback switches per source. Cooling-down slots are skipped. Returns
// the response, serving source and the exact model id used; on total failure
// a *ChainErr summarizing every tier.
func (s *Scheduler) ChainImage(ctx context.Context, chain *Chain, req *types.ImageGenRequest) (*types.UnifiedResponse, string, string, error) {
if chain == nil || len(chain.Tiers) == 0 {
return nil, "", "", fmt.Errorf("no image auto slot configured")
}
var ce ChainErr
for _, tn := range chain.Tiers {
cands := collectCands(tn.Slots, nil)
if len(cands) == 0 {
ce.Skipped = append(ce.Skipped, fmt.Sprintf("tier %d: no schedulable slot (cooling)", tn.Tier))
continue
}
// No Pref sort: load balancing is done by round-robin cursor.
base := tn.NextStart()
norm := normalCount(cands)
var hard []TierError
for i := 0; i < len(cands); i++ {
var c candidate
if i < norm {
c = cands[(int(base)+i)%norm]
} else {
c = cands[i]
}
sl := c.slot
if !c.probe && !sl.Prov.ModelAvailable(sl.Model) {
continue
}
r := *req
r.Model = sl.Model
resp, err := sl.Prov.Image(ctx, &r)
if ctx.Err() != nil {
releaseProbes(cands)
return nil, "", "", ctx.Err()
}
if err == nil {
releaseProbes(cands)
return resp, sl.Source, sl.Model, nil
}
if errors.Is(err, types.ErrBusy) {
continue
}
hard = append(hard, TierError{Tier: tn.Tier, Source: sl.Source, Model: sl.Model, Err: err})
}
releaseProbes(cands)
if len(hard) > 0 {
ce.Tiers = append(ce.Tiers, hard...)
}
}
if len(ce.Tiers) == 0 && len(ce.Skipped) == 0 {
return nil, "", "", fmt.Errorf("no image auto slot configured")
}
return nil, "", "", &ce
}
// ---- direct scheduling ----
// Chat runs a chat request across cands, falling back on failure. Each
// candidate receives a request pinned to its own model (ModelFor), so a
// fallback switches the model id per provider instead of reusing the first
// candidate's model name. On success it returns the response together with
// the name of the provider and the exact model id that served the request.
// A candidate whose (source, model) pair is cooling down is skipped like a
// busy one — direct paths share the "cooldown is the only hard skip"
// semantics of the AUTO chain (plan 2.3/2.5); otherwise persistent direct
// traffic would keep renewing a capped auth cooldown forever.
func (s *Scheduler) Chat(ctx context.Context, cands []Provider, req *types.ChatRequest) (*types.UnifiedResponse, string, string, error) {
attempts := s.MaxRetries + 1
var lastErr error
for i := 0; i < attempts && i < len(cands); i++ {
p := cands[i]
r := *req
r.Model = p.ModelFor(req.Model)
ok, isProbe := p.ModelSchedulable(r.Model)
if !ok {
lastErr = fmt.Errorf("provider %s: model %q cooling down", p.Name(), r.Model)
continue
}
resp, err := p.Chat(ctx, &r)
if isProbe {
p.ProbeDone(r.Model)
}
if ctx.Err() != nil {
return nil, "", "", ctx.Err()
}
if err == nil {
return resp, p.Name(), r.Model, nil
}
lastErr = fmt.Errorf("provider %s: %w", p.Name(), err)
}
if lastErr == nil {
if len(cands) == 0 {
lastErr = fmt.Errorf("no provider available")
}
}
return nil, "", "", lastErr
}
// ChatStream runs a streaming chat across cands, falling back early on
// connect errors. The request model is pinned per candidate like Chat. On
// success it returns the chunk channel plus the serving provider name and
// model id.
func (s *Scheduler) ChatStream(ctx context.Context, cands []Provider, req *types.ChatRequest) (<-chan types.UnifiedChunk, string, string, error) {
attempts := s.MaxRetries + 1
var lastErr error
for i := 0; i < attempts && i < len(cands); i++ {
p := cands[i]
r := *req
r.Model = p.ModelFor(req.Model)
ok, isProbe := p.ModelSchedulable(r.Model)
if !ok {
lastErr = fmt.Errorf("provider %s: model %q cooling down", p.Name(), r.Model)
continue
}
resp, err := p.ChatStream(ctx, &r)
if isProbe {
p.ProbeDone(r.Model)
}
if err == nil {
return resp, p.Name(), r.Model, nil
}
lastErr = fmt.Errorf("provider %s: %w", p.Name(), err)
}
if lastErr == nil && len(cands) == 0 {
lastErr = fmt.Errorf("no provider available")
}
return nil, "", "", lastErr
}
// Image runs an image-generation request across cands; returns the used
// provider name on success.
func (s *Scheduler) Image(ctx context.Context, cands []Provider, req *types.ImageGenRequest) (*types.UnifiedResponse, string, error) {
attempts := s.MaxRetries + 1
var lastErr error
for i := 0; i < attempts && i < len(cands); i++ {
p := cands[i]
im := p.ModelFor(req.Model)
ok, isProbe := p.ModelSchedulable(im)
if !ok {
lastErr = fmt.Errorf("provider %s: model %q cooling down", p.Name(), im)
continue
}
resp, err := p.Image(ctx, req)
if isProbe {
p.ProbeDone(im)
}
if err == nil {
return resp, p.Name(), nil
}
lastErr = fmt.Errorf("provider %s: %w", p.Name(), err)
}
if lastErr == nil && len(cands) == 0 {
lastErr = fmt.Errorf("no provider available")
}
return nil, "", lastErr
}