Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 11 additions & 2 deletions docs/ACCOUNT_POOL_LEASES_CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,8 +42,8 @@ account-pools:
- Key 绑定实例后只使用被授权的租约池;未绑定实例的普通 Key 只能使用非租约池,包括 scope=all 的 Key。
- `X-CPA-Instance-ID` 必须与 Key 绑定实例匹配,`X-CPA-User-ID` 必须是正整数。身份仅在客户端鉴权后读取,专用头随即从请求中移除,不转发给模型上游。
- 归属按“实例 + 用户 ID”的摘要确定,与 CPA Key 和下游令牌无关。同一用户换渠道或 Key 不会再占一个池;若新 Key 无权使用原池,返回权限错误。
- 固定 60 分钟,从首次分配开始;普通调用不续期。首次从授权组内选择一个可用账号,此后只使用该账号。账号暂时不可用或不支持后续请求模型时返回错误,不再切换到组内其他账号
- 无空闲且可用的授权账号返回 503 `pool_busy`,不排无界队列、不抢占。已绑定账号失败时不自动换号
- 固定 60 分钟,从首次分配开始;普通调用不续期。首次从授权组内选择一个可用账号,正常情况下继续使用该账号。明确额度耗尽、刷新后仍未授权或账号失效时,可安全切换同组空闲账号,保留原到期时间。临时限流、服务过载、请求错误或模型不支持不会触发换号
- 无空闲且可用的授权账号返回 503 `pool_busy`,不排无界队列、不抢占。没有同组空闲账号或仍有其他在途请求时,返回 503 `pool_lease_failover_pending`,后续请求再尝试分配
- 到期但仍有请求执行时,租约进入待释放状态;旧请求可以结束,新请求不能继续使用该租约,其他用户也不能接管。
- 流式取消后仍等待上游生产通道结束才释放在途计数,不以断开下游连接为理由提前把池交给别人。若上游未能正常结束,需要先排查/结束旧请求,不能直接清零计数。
- 带 `previous_response_id` 的请求在已无有效租约时返回 409,要求开始新会话,避免续接到另一账号。
Expand All @@ -67,3 +67,12 @@ account-pools:
包括:并发首请求单次分配、同组多账号分配给不同用户、跨 Key 复用、请求不续期、到期在途保护、流式取消、上游失败不换池、状态恢复/缺失/损坏、文件锁、落盘失败、分组变更与失效会话,以及 New API 身份上下文传递和专用头移除。

本阶段为本地代码接入,不自动发布或修改线上配置。New API 与 CPA 两边部署并完成实例/Key 配置后才生效。

## 额度耗尽时的账号替换

- 普通请求、计数请求及流式启动失败可在安全边界自动换号并重试一次;入口也检查此前已记录的账号故障。每次执行循环最多换号一次,避免无限重试。
- 替换必须原子落盘成功才生效,更新租约 ID 以隔离旧的会话路由缓存;用户和原到期时间保持不变。其他用户占用的账号和其他分组均不可选。
- 旧账号仍有并发请求或被丢弃的上游流尚未结束时禁止替换,防止旧账号过早释放。已经向客户端输出的流不会自动重放;下一次完整请求可换号。
- 携带 previous_response_id 的会话无法安全跨账号,返回 409 `pool_lease_session_expired`,需使用完整历史新建会话。当前采用保守保护:发生过换号的租约在剩余有效期内不接受 previous_response_id,避免误把旧账号响应 ID 交给新账号;完整历史请求不受影响。
- 短暂的通用 429 仍按现有退避处理;只有明确 usage_limit_reached、insufficient_quota、usage limit 等账号额度错误才会换号。裸 403 或内容策略错误不会触发换号。
- 此版本不改变前端 API 契约,可继续使用 v1.22.2-cpa.21。线上升级须由管理员明确执行,不随发布自动部署。
47 changes: 47 additions & 0 deletions internal/poollease/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ type Candidate struct {
}

type Lease struct {
Reassigned bool `json:"reassigned,omitempty"`
Credential string `json:"credential,omitempty"`
LegacyGroup bool `json:"legacy_group,omitempty"`
ID string `json:"id"`
Expand Down Expand Up @@ -276,3 +277,49 @@ func (s *Store) Snapshot(now time.Time) []Lease {
})
return rows
}

// Replace moves a lease only when its sole in-flight request has stopped using
// the old credential. Rotating the ID invalidates old routing namespaces.
func (s *Store) Replace(expected Lease, candidates []Candidate, now time.Time) (Lease, func(), error) {
s.mu.Lock()
defer s.mu.Unlock()
if err := s.check(expected.Policy, now); err != nil {
return Lease{}, nil, err
}
old := s.leases[expected.ID]
if old == nil || old.Owner != expected.Owner || old.Credential != expected.Credential || old.LegacyGroup || old.Active != 1 || !now.Before(old.Expires) {
return Lease{}, nil, ErrBusy
}
for _, c := range candidates {
if c.Group != old.Group || c.Credential == "" || c.Credential == old.Credential || s.occupied(c, "") {
continue
}
id := make([]byte, 16)
if _, err := rand.Read(id); err != nil {
return Lease{}, nil, err
}
next := *old
next.ID, next.Credential, next.Reassigned = hex.EncodeToString(id), c.Credential, true
delete(s.leases, old.ID)
s.leases[next.ID] = &next
if err := s.save(); err != nil {
delete(s.leases, next.ID)
s.leases[old.ID] = old
return Lease{}, nil, err
}
return next, s.releaser(next.ID), nil
}
return Lease{}, nil, ErrBusy
}

// Hold retains a live lease while an abandoned upstream producer is drained.
func (s *Store) Hold(id string) (func(), bool) {
s.mu.Lock()
defer s.mu.Unlock()
l := s.leases[id]
if l == nil || l.Active == 0 {
return nil, false
}
l.Active++
return s.releaser(id), true
}
55 changes: 55 additions & 0 deletions internal/poollease/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -306,3 +306,58 @@ func TestStateRejectsDuplicateAccountReservations(t *testing.T) {
t.Fatal("duplicate account accepted")
}
}

func TestReplacementPreservesExpiryAndExclusiveOwnership(t *testing.T) {
s := testStore(t)
now := time.Now()
allowed := map[string]bool{"a": true, "b": true}
candidates := []Candidate{{"a", "one"}, {"a", "two"}, {"a", "three"}, {"b", "other"}}
old, done, err := s.Acquire("alice", "v1", allowed, candidates, now)
if err != nil {
t.Fatal(err)
}
defer done()
_, done2, _ := s.Acquire("bob", "v1", allowed, candidates, now)
defer done2()
_, parallel, _ := s.Acquire("alice", "v1", allowed, candidates, now)
if _, _, err = s.Replace(old, candidates, now); !errors.Is(err, ErrBusy) {
t.Fatal("replaced during concurrent execution", err)
}
parallel()
next, release, err := s.Replace(old, candidates, now.Add(time.Minute))
if err != nil {
t.Fatal(err)
}
if next.Credential != "three" || next.Group != old.Group || !next.Expires.Equal(old.Expires) || next.ID == old.ID {
t.Fatalf("bad replacement %+v", next)
}
done()
if rows := s.Snapshot(now); len(rows) != 2 {
t.Fatal(rows)
}
if _, _, err = s.Replace(old, candidates, now); !errors.Is(err, ErrBusy) {
t.Fatal("stale replacement accepted")
}
if _, _, err = s.Replace(next, []Candidate{{"a", "two"}, {"b", "other"}}, now); !errors.Is(err, ErrBusy) {
t.Fatal("cross group or occupied replacement")
}
release()
}

func TestReplacementPersistenceFailureRestoresOldLease(t *testing.T) {
s := testStore(t)
now := time.Now()
c := []Candidate{{"a", "one"}, {"a", "two"}}
old, done, _ := s.Acquire("alice", "v1", map[string]bool{"a": true}, c, now)
defer done()
original := s.path
s.path = filepath.Join(t.TempDir(), "missing", "state")
if _, _, err := s.Replace(old, c, now); err == nil {
t.Fatal("expected save failure")
}
s.path = original
rows := s.Snapshot(now)
if len(rows) != 1 || rows[0].ID != old.ID || rows[0].Credential != old.Credential {
t.Fatal(rows)
}
}
4 changes: 4 additions & 0 deletions sdk/cliproxy/auth/account_pools_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,13 +19,17 @@ type poolCaptureExecutor struct {
mu sync.Mutex
ids []string
failures map[string]bool
errors map[string]error
}

func (*poolCaptureExecutor) Identifier() string { return "pool-test" }
func (e *poolCaptureExecutor) Execute(_ context.Context, a *Auth, _ coreexecutor.Request, _ coreexecutor.Options) (coreexecutor.Response, error) {
e.mu.Lock()
defer e.mu.Unlock()
e.ids = append(e.ids, a.ID)
if err := e.errors[a.ID]; err != nil {
return coreexecutor.Response{}, err
}
if e.failures[a.ID] {
return coreexecutor.Response{}, &Error{HTTPStatus: 503, Message: "upstream unavailable"}
}
Expand Down
40 changes: 38 additions & 2 deletions sdk/cliproxy/auth/conductor_execution.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ func (m *Manager) Execute(ctx context.Context, providers []string, req cliproxye
if errLease != nil {
return cliproxyexecutor.Response{}, errLease
}
defer releaseLease()
defer func() { releaseLease() }()

req, opts = cliproxysession.Enrich(req, opts)
normalized := m.normalizeProviders(providers)
Expand All @@ -59,6 +59,7 @@ func (m *Manager) Execute(ctx context.Context, providers []string, req cliproxye
_, maxRetryCredentials, maxWait := m.retrySettings()

var lastErr error
leaseReplaced := false
retryModel := authSelectionModelFromOptions(opts, req.Model)
for attempt := 0; ; attempt++ {
resp, errExec := m.executeMixedOnce(ctx, normalized, req, opts, maxRetryCredentials)
Expand All @@ -68,6 +69,17 @@ func (m *Manager) Execute(ctx context.Context, providers []string, req cliproxye
if isRequestTerminatedError(errExec) || isRequestStopError(errExec) {
return cliproxyexecutor.Response{}, unwrapRequestStopError(errExec)
}
if !leaseReplaced {
next, done, changed, err := m.replaceFailedPoolLease(ctx, normalized, req, opts, errExec)
if err != nil {
return cliproxyexecutor.Response{}, err
}
if changed {
releaseLease()
ctx, releaseLease, leaseReplaced = next, done, true
continue
}
}
lastErr = errExec
wait, shouldRetry := m.shouldRetryAfterError(errExec, attempt, normalized, retryModel, maxWait)
if !shouldRetry {
Expand Down Expand Up @@ -101,7 +113,7 @@ func (m *Manager) ExecuteCount(ctx context.Context, providers []string, req clip
if errLease != nil {
return cliproxyexecutor.Response{}, errLease
}
defer releaseLease()
defer func() { releaseLease() }()

req, opts = cliproxysession.Enrich(req, opts)
normalized := m.normalizeProviders(providers)
Expand All @@ -116,6 +128,7 @@ func (m *Manager) ExecuteCount(ctx context.Context, providers []string, req clip
_, maxRetryCredentials, maxWait := m.retrySettings()

var lastErr error
leaseReplaced := false
retryModel := authSelectionModelFromOptions(opts, req.Model)
for attempt := 0; ; attempt++ {
resp, errExec := m.executeCountMixedOnce(ctx, normalized, req, opts, maxRetryCredentials)
Expand All @@ -125,6 +138,17 @@ func (m *Manager) ExecuteCount(ctx context.Context, providers []string, req clip
if isRequestTerminatedError(errExec) || isRequestStopError(errExec) {
return cliproxyexecutor.Response{}, unwrapRequestStopError(errExec)
}
if !leaseReplaced {
next, done, changed, err := m.replaceFailedPoolLease(ctx, normalized, req, opts, errExec)
if err != nil {
return cliproxyexecutor.Response{}, err
}
if changed {
releaseLease()
ctx, releaseLease, leaseReplaced = next, done, true
continue
}
}
lastErr = errExec
wait, shouldRetry := m.shouldRetryAfterError(errExec, attempt, normalized, retryModel, maxWait)
if !shouldRetry {
Expand Down Expand Up @@ -172,6 +196,7 @@ func (m *Manager) ExecuteStream(ctx context.Context, providers []string, req cli
_, maxRetryCredentials, maxWait := m.retrySettings()

var lastErr error
leaseReplaced := false
retryModel := authSelectionModelFromOptions(opts, req.Model)
for attempt := 0; ; attempt++ {
result, errStream := m.executeStreamMixedOnce(ctx, normalized, req, opts, maxRetryCredentials)
Expand All @@ -182,6 +207,17 @@ func (m *Manager) ExecuteStream(ctx context.Context, providers []string, req cli
if isRequestTerminatedError(errStream) || isRequestStopError(errStream) {
return nil, unwrapRequestStopError(errStream)
}
if !leaseReplaced {
next, done, changed, err := m.replaceFailedPoolLease(ctx, normalized, req, opts, errStream)
if err != nil {
return nil, err
}
if changed {
releaseLease()
ctx, releaseLease, leaseReplaced = next, done, true
continue
}
}
lastErr = errStream
wait, shouldRetry := m.shouldRetryAfterError(errStream, attempt, normalized, retryModel, maxWait)
if !shouldRetry {
Expand Down
20 changes: 10 additions & 10 deletions sdk/cliproxy/auth/conductor_stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -172,13 +172,13 @@ func (m *Manager) wrapStreamResult(ctx context.Context, auth *Auth, provider, re
}
for _, chunk := range buffered {
if ok := emit(chunk); !ok {
discardStreamChunks(remaining)
m.discardLeasedStream(ctx, remaining)
return
}
}
for chunk := range remaining {
if ok := emit(chunk); !ok {
discardStreamChunks(remaining)
m.discardLeasedStream(ctx, remaining)
return
}
}
Expand Down Expand Up @@ -302,7 +302,7 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
buffered, closed, bootstrapErr := readStreamBootstrap(ctx, streamResult.Chunks)
if bootstrapErr != nil {
if errCtx := ctx.Err(); errCtx != nil {
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
return nil, errCtx
}
if allowRetry {
Expand All @@ -316,11 +316,11 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
}
}
if errRefresh != nil {
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
bootstrapErr = errRefresh
streamResult = &cliproxyexecutor.StreamResult{}
} else if okRefresh {
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
auth = refreshed
m.replaceHomeExecutionLifecycleAuth(execOpts.ExecutionLifecycle, auth)
publishSelectedAuthMetadata(execOpts.Metadata, auth)
Expand All @@ -342,7 +342,7 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
}
if !ephemeralResult {
if errCancel := claudeOAuthRequestCancellation(ctx, auth, bootstrapErr); errCancel != nil {
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
return nil, errCancel
}
}
Expand All @@ -357,7 +357,7 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
}
applyRequestScopedActionToResult(action, okAction, &result)
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
if isRequestScopedStop(action, okAction) {
return nil, wrapRequestStopError(bootstrapErr)
}
Expand All @@ -375,7 +375,7 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
result.CredentialScope = true
}
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
return nil, bootstrapErr
}
if idx < len(execModels)-1 {
Expand All @@ -386,7 +386,7 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
result.CredentialScope = true
}
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
lastErr = bootstrapErr
if result.CredentialScope {
return nil, newStreamBootstrapError(bootstrapErr, streamResult.Headers)
Expand All @@ -400,7 +400,7 @@ func (m *Manager) executeStreamWithModelPool(ctx context.Context, executor Provi
result.CredentialScope = true
}
m.recordExecutionResult(ctx, result, auth, ephemeralResult)
discardStreamChunks(streamResult.Chunks)
m.discardLeasedStream(ctx, streamResult.Chunks)
return nil, newStreamBootstrapError(bootstrapErr, streamResult.Headers)
}

Expand Down
Loading
Loading