docs: record distributed admission delivery
Some checks are pending
ci / test (ubuntu-latest) (push) Waiting to run
ci / test (windows-latest) (push) Waiting to run
ci / race (push) Waiting to run
ci / integration (push) Waiting to run

This commit is contained in:
youfak 2026-07-30 21:04:19 +08:00
parent e54fc84a81
commit c3b5b25597
11 changed files with 63 additions and 13 deletions

View File

@ -154,7 +154,7 @@ Redis Adapter 应通过单个 Lua 脚本、Redis Function 或等价的原子原
- 两个并发成功响应的 Proxy ID 集合交集为空。 - 两个并发成功响应的 Proxy ID 集合交集为空。
- 原子操作失败时整个批次不返回,也不得留下部分移除结果。 - 原子操作失败时整个批次不返回,也不得留下部分移除结果。
- `allOrNothing` 不足时零个条目退出活动池。 - `allOrNothing` 不足时零个条目退出活动池。
- Redis 活动池不可用时返回 503不以内存副本冒充成功。 - Redis 活动池或分布式限流不可用时返回 503不以内存副本冒充成功。
- PostgreSQL 不可用不阻断提取;需要 PostgreSQL 的 Admin 管理写入单独降级。 - PostgreSQL 不可用不阻断提取;需要 PostgreSQL 的 Admin 管理写入单独降级。
## 7. 幂等 ## 7. 幂等
@ -175,7 +175,9 @@ TTL 到期或 Redis 数据丢失后不再保证旧 Key 去重,系统不回退
- 直连请求使用来源 IP 形成匿名 Client。 - 直连请求使用来源 IP 形成匿名 Client。
- 只有来源属于 `trustedProxies` 时才接受转发头。 - 只有来源属于 `trustedProxies` 时才接受转发头。
- 全局和每 Client 限流在查询库存前执行。 - 全局和每 Client 限流在查询库存前执行本地监听器先做早期拒绝Redis 使用
服务端时间和单个 Lua 原子操作执行跨 Controller 副本的权威额度。
- Client 身份进入 Redis 前转换为定长摘要,窗口切换时原子删除上一窗口字段。
- 过滤条件、数量、请求体和 Header 都有长度/数量上限。 - 过滤条件、数量、请求体和 Header 都有长度/数量上限。
## 9. 错误模型 ## 9. 错误模型
@ -190,7 +192,7 @@ TTL 到期或 Redis 数据丢失后不再保证旧 Key 去重,系统不回退
- `415`:请求体不是 `application/json` - `415`:请求体不是 `application/json`
- `422`:数量、枚举或过滤组合违反业务约束。 - `422`:数量、枚举或过滤组合违反业务约束。
- `429`:全局或 Client 速率限制,响应 `Retry-After` - `429`:全局或 Client 速率限制,响应 `Retry-After`
- `503`Redis 活动池不可用、原子提取不可执行或服务正在排空。 - `503`Redis 活动池/分布式限流不可用、原子提取不可执行或服务正在排空。
- `500`:未分类的内部错误;响应不包含底层错误文本。 - `500`:未分类的内部错误;响应不包含底层错误文本。
错误响应不得包含 Provider Secret、Proxy 凭据、SQL 或内部拓扑。 错误响应不得包含 Provider Secret、Proxy 凭据、SQL 或内部拓扑。

View File

@ -104,6 +104,10 @@ limits:
`trustedProxies` 只决定何时接受 `Forwarded``X-Forwarded-For`,不能替代 `trustedProxies` 只决定何时接受 `Forwarded``X-Forwarded-For`,不能替代
`allowCIDRs`。来自非可信代理的转发头必须忽略。 `allowCIDRs`。来自非可信代理的转发头必须忽略。
`requestsPerMinute``requestsPerMinutePerClient` 为非负整数,最大值为
`2^53-1`Distribution 的非零额度由 Redis Lua 计数,因此配置校验统一限制在
Lua 可精确表示的整数范围内。
### 3.1 认证模式 ### 3.1 认证模式
- `none`:无身份认证,访问控制与限流仍生效。 - `none`:无身份认证,访问控制与限流仍生效。

View File

@ -174,8 +174,9 @@ Upstream、`endBehavior` 默认 `stop`并覆盖列表末端停止disabled
candidate eligibility, Gateway reserve, ownership, removal, and short-lived idempotency. candidate eligibility, Gateway reserve, ownership, removal, and short-lived idempotency.
- [x] Implement Redis Worker ownership, drain/ACK, expiry reclaim, inventory and bounded - [x] Implement Redis Worker ownership, drain/ACK, expiry reclaim, inventory and bounded
sweep primitives with a monotonic global epoch. sweep primitives with a monotonic global epoch.
- [ ] Implement Redis Provider leader, distributed rate, Client limit and Worker - [x] Implement Redis Provider leader, distributed request quota, Client limit and
heartbeat; wire automatic Provider inventory rebuild after Redis loss. automatic Provider inventory rebuild after Redis state loss.
- [ ] Implement the Worker heartbeat receiving path and session lifecycle.
- [x] Keep Provider output in Redis TTL activity state and node memory only; keep the - [x] Keep Provider output in Redis TTL activity state and node memory only; keep the
Gateway request path on immutable local snapshots with no Redis/PostgreSQL calls. Gateway request path on immutable local snapshots with no Redis/PostgreSQL calls.
- [x] Expose Distribution extraction/status and Admin status/enable/disable/switch/reload - [x] Expose Distribution extraction/status and Admin status/enable/disable/switch/reload
@ -236,9 +237,13 @@ Redis Provider Permit 现已把 requestInterval、maxInFlight 与 maxTotal 放
启动后动态扩容。Supervisor 以 PostgreSQL 权威 HMAC 指纹和 revision 栅栏协调 启动后动态扩容。Supervisor 以 PostgreSQL 权威 HMAC 指纹和 revision 栅栏协调
多副本 reload管理库瞬断沿用 last-known 状态,本地共享源落后时停止旧 Provider 多副本 reload管理库瞬断沿用 last-known 状态,本地共享源落后时停止旧 Provider
源匹配并预检后自动替换,迟到旧 revision 不覆盖新配置。 源匹配并预检后自动替换,迟到旧 revision 不覆盖新配置。
Distribution 现通过公用 `admission.Admitter` 接入独立 `redisadmission` Adapter
全局和单 Client 分钟额度使用 Redis 服务端时间并在单个 Lua 原子边界内检查、递增,
Controller 多副本共享同一计数。Client 身份只以 SHA-256 摘要进入 Redis窗口切换
原子回收历史字段Redis 异常 fail-closed 并返回 503真实额度耗尽返回 429。
Gateway 请求热路径仍只使用本地准入,不增加 Redis/PostgreSQL 调用。
WorkerControlPlane gRPC 接收端、session 签发/心跳、Snapshot ACK WorkerControlPlane gRPC 接收端、session 签发/心跳、Snapshot ACK
账本、Client 分布式限流和健康执行链仍待完成,因此本轮不勾选 Task 10 的组合 账本和健康执行链仍待完成,因此 Task 10 尚未全部完成。
验收项。
## Task 11: Checker and Health Reducer ## Task 11: Checker and Health Reducer

View File

@ -240,7 +240,7 @@ Outbox 发布器必须以稳定 consumer ID 有界领取;发布成功后原子
回查 PostgreSQL 或 Provider。 回查 PostgreSQL 或 Provider。
2. Distribution 立即失败关闭并返回 503禁止本地内存提取或 PostgreSQL 兜底。 2. Distribution 立即失败关闭并返回 503禁止本地内存提取或 PostgreSQL 兜底。
3. Controller 停止活动池写入和需要分布式互斥的工作,防止多个 Fetch Leader 3. Controller 停止活动池写入和需要分布式互斥的工作,防止多个 Fetch Leader
本地限流不能声称满足全局额度。 Distribution 的 Redis 权威限流同时失败关闭,本地早期限流不冒充跨副本额度。
4. 恢复后确认 Leader 唯一和租约 epoch 单调,由 Provider 重新 Fetch 并构建 TTL 4. 恢复后确认 Leader 唯一和租约 epoch 单调,由 Provider 重新 Fetch 并构建 TTL
活动池,再恢复 Distribution。Redis 整体丢失会终止原活动池代次的排他状态和 活动池,再恢复 Distribution。Redis 整体丢失会终止原活动池代次的排他状态和
短期幂等窗口;高可用、持久化、监控和告警必须明确并降低该风险。 短期幂等窗口;高可用、持久化、监控和告警必须明确并降低该风险。

View File

@ -84,11 +84,11 @@ CI 已配置 Linux race job。PostgreSQL 18 和 Redis 8.2 的隔离 Adapter fixt
4. Controller 的 PostgreSQL 连接池、迁移和 pgx Adapter 启动装配已完成; 4. Controller 的 PostgreSQL 连接池、迁移和 pgx Adapter 启动装配已完成;
公用 bootstrap 已通过 PostgreSQL 18 + Redis 8.2 双存储集成Controller 公用 bootstrap 已通过 PostgreSQL 18 + Redis 8.2 双存储集成Controller
三监听器与探针集成已完成;可选聚合指标和完整容器进程部署验证仍待实现。 三监听器与探针集成已完成;可选聚合指标和完整容器进程部署验证仍待实现。
5. Redis Provider Leader、分布式速率与 Client 限制、Worker 心跳和自动重建; 5. Worker heartbeat gRPC 接收路径Redis Provider Leader、分布式请求额度、
TTL 活动池、原子提取和 Worker ownership 已完成。 Distribution Client 限制和 Provider 状态丢失重建已完成。
6. Worker 网络快照流Redis ownership drain/ACK/过期回收已完成。 6. Worker 网络快照流Redis ownership drain/ACK/过期回收已完成。
7. Checker 调度、探测器和健康 reducer。 7. Checker 调度、探测器和健康 reducer。
8. Admin/Distribution 细粒度授权、分布式限流和审计查询。 8. Admin/Distribution 细粒度授权和审计查询Distribution 分布式限流已完成
9. 真实 Compose/Kubernetes 集成、故障演练和代表性集群负载测试。 9. 真实 Compose/Kubernetes 集成、故障演练和代表性集群负载测试。
10. 将五种 Routing 策略和 `onUnavailable` 接入 Gateway/Distribution 运行链, 10. 将五种 Routing 策略和 `onUnavailable` 接入 Gateway/Distribution 运行链,
补齐 Sequential 持久化恢复、跨实例 CAS 和 disabled candidate 语义。 补齐 Sequential 持久化恢复、跨实例 CAS 和 disabled candidate 语义。

View File

@ -72,7 +72,7 @@
| DIST-005 | 返回 expiresAt 与 remainingTtlSeconds | 9334-9360 | `extraction/service_test.go` | | DIST-005 | 返回 expiresAt 与 remainingTtlSeconds | 9334-9360 | `extraction/service_test.go` |
| DIST-006 | 提取前校验 minRemainingTTL 与 maxHealthCheckAge | 9334-9369 | 过滤测试 | | DIST-006 | 提取前校验 minRemainingTTL 与 maxHealthCheckAge | 9334-9369 | 过滤测试 |
| DIST-007 | reserveForGateway 防止 Extract 清空共享池 | 9281-9333 | 共享池测试 | | DIST-007 | reserveForGateway 防止 Extract 清空共享池 | 9281-9333 | 共享池测试 |
| DIST-008 | 提取认证可关闭,关闭后仍有来源识别与全局限制 | 8112-8441 | 来源身份准入`FixedWindow` 并发测试 | | DIST-008 | 提取认证可关闭,关闭后仍有来源识别与全局限制 | 8112-8441 | 来源身份准入、`FixedWindow` 单元测试及 `redisadmission` 双实例/并发集成测试 |
## 健康、安全、运维与测试 ## 健康、安全、运维与测试

View File

@ -30,7 +30,8 @@
原子操作;并发与主从切换下不得部分提交。 原子操作;并发与主从切换下不得部分提交。
- PostgreSQL 配置版本、Upstream/Routing 管理状态、Admin 审计与 outbox 的事务 - PostgreSQL 配置版本、Upstream/Routing 管理状态、Admin 审计与 outbox 的事务
更新及幂等重放;测试库断言不包含 Proxy 明细或逐次提取记录。 更新及幂等重放;测试库断言不包含 Proxy 明细或逐次提取记录。
- Redis TTL 活动池、Leader 租约、限流、短期幂等窗口和失联恢复。 - Redis TTL 活动池、Leader 租约、Provider 请求额度、Distribution Client
跨副本限流、短期幂等窗口和失联恢复。
- Snapshot/Delta/ACK/Report 的版本与校验和兼容性。 - Snapshot/Delta/ACK/Report 的版本与校验和兼容性。
- OpenAPI 错误模型、认证矩阵、批量 fulfillment。 - OpenAPI 错误模型、认证矩阵、批量 fulfillment。
- OpenAPI 本地引用闭合、operationId 唯一、响应集合和 security scheme 引用。 - OpenAPI 本地引用闭合、operationId 唯一、响应集合和 security scheme 引用。

View File

@ -428,6 +428,13 @@ func TestValidateRejectsInvalidConfigurationMatrix(t *testing.T) {
}, },
want: "requestsPerMinute", want: "requestsPerMinute",
}, },
{
name: "listener request limit exceeds exact counter range",
mutate: func(cfg *Config) {
cfg.Distribution.Limits.RequestsPerMinutePerClient = int(MaximumExactCounter) + 1
},
want: "requestsPerMinutePerClient",
},
{ {
name: "invalid client identification mode", name: "invalid client identification mode",
mutate: func(cfg *Config) { mutate: func(cfg *Config) {

View File

@ -122,6 +122,9 @@ func validateListener(name string, listener Listener, security Security) error {
if limit.value < 0 { if limit.value < 0 {
return fmt.Errorf("validate %s limits.%s: must be non-negative", name, limit.name) return fmt.Errorf("validate %s limits.%s: must be non-negative", name, limit.name)
} }
if int64(limit.value) > MaximumExactCounter {
return fmt.Errorf("validate %s limits.%s: exceeds exact counter range", name, limit.name)
}
} }
host, err := validateListenAddress(name, listener.Listen) host, err := validateListenAddress(name, listener.Listen)
if err != nil { if err != nil {

View File

@ -111,6 +111,9 @@ func (s *Service) Extract(ctx context.Context, request Request) (Response, error
return response, ErrInvalidFulfillment return response, ErrInvalidFulfillment
} }
if err := s.admission.Admit(ctx, admissionKey(request)); err != nil { if err := s.admission.Admit(ctx, admissionKey(request)); err != nil {
if errors.Is(err, admission.ErrUnavailable) {
return response, errors.Join(ErrUnavailable, err)
}
return response, errors.Join(ErrAdmissionRejected, err) return response, errors.Join(ErrAdmissionRejected, err)
} }

View File

@ -9,6 +9,7 @@ import (
"proxy-pool/internal/domain/activitypool" "proxy-pool/internal/domain/activitypool"
domain "proxy-pool/internal/domain/extraction" domain "proxy-pool/internal/domain/extraction"
proxyDomain "proxy-pool/internal/domain/proxy" proxyDomain "proxy-pool/internal/domain/proxy"
platformAdmission "proxy-pool/internal/platform/admission"
) )
func TestServiceAppliesPolicyAndBuildsResponse(t *testing.T) { func TestServiceAppliesPolicyAndBuildsResponse(t *testing.T) {
@ -122,6 +123,30 @@ func TestServiceAppliesAdmissionBeforeStoreUsingStableIdentity(t *testing.T) {
} }
} }
func TestServiceMapsUnavailableAdmissionToServiceUnavailable(t *testing.T) {
store := &recordingStore{}
admitter := &recordingAdmission{err: platformAdmission.ErrUnavailable}
service, err := NewService(store, Policy{
MaxCountPerRequest: 1,
DefaultFulfillment: domain.Partial,
}, admitter, time.Now)
if err != nil {
t.Fatalf("NewService(): %v", err)
}
_, err = service.Extract(context.Background(), Request{
RequestID: "req-1",
ClientID: "client-1",
Count: 1,
})
if !errors.Is(err, ErrUnavailable) || errors.Is(err, ErrAdmissionRejected) {
t.Fatalf("Extract() error = %v, want only ErrUnavailable", err)
}
if store.calls != 0 {
t.Fatalf("store calls = %d, want 0", store.calls)
}
}
func TestServiceUsesSourceIdentityForEphemeralIdempotency(t *testing.T) { func TestServiceUsesSourceIdentityForEphemeralIdempotency(t *testing.T) {
store := &recordingStore{} store := &recordingStore{}
service, err := NewService(store, Policy{ service, err := NewService(store, Policy{