Documentation
¶
Overview ¶
归属说明:retention worker 放 usage 包——原 Recorder.janitorLoop(逐行 DELETE)已删除,保留策略整体移交本 worker;保留天数与 Recorder 配置同源 (config usage.log_retention_days),故与 Recorder 同包管理。
离线聚合 worker(spec 2026-08-14 使用量统计离线聚合化 + spec 2026-08-23 v2 双表扩展):请求路径零统计计算/投递——usage_stats(cube)与 usage_entity_ stats(实体卷积)只由本 worker 每周期从 DB(usage_logs + err_logs)重建 (SQL 侧聚合,DELETE+INSERT 覆盖语义)。放行行(含 abort)从 usage_logs 重建(全字段);纯错误行(4xx/5xx/network)从 err_logs 重建(count 语义); 拒绝行随 err_logs 采样队列丢样(口径注释见 proxy/forward.go recordRejected)。 统计分钟级陈旧(接受,不加 hot counter);quota 回写在线保留(Recorder flushQuota,独立于本 worker)。
Package usage 承载请求明细的异步落库(规格 §7.2/§10.5)、key 额度增量回写 (quota 在线保留——独立于统计,不随离线聚合搬移)、usagelog 保留策略 (retention worker,Phase 5 T4.5:按日分区 DROP 清理)与离线聚合 worker (stats_agg.go:spec 2026-08-14 使用量统计离线聚合化)。 请求路径**零统计计算**(用户裁决 2026-08-14):统计内存桶机制整体删除, Record 锁内仅明细 append + quota 原子累加;usage_stats 由离线聚合 worker 每周期从 DB 重建(DELETE+INSERT 覆盖语义,见 stats_agg.go)。明细经无界 pending 批量落库(O1 管道化:Record 永不阻塞——此前有界 channel cap 16384 饱和阻塞发送是压测 off 路径 16.4k goroutine 卡 chan send、healthz inflight 31-33k @10k 幽灵根因,O3 复测定位 2026-08-09;崩溃丢 ≤1 flush 窗口的崩溃 等价语义不变,pending 内存即唯一积压面,由水线 Warn 观测)。
Index ¶
- type ErrLogConfig
- type ErrLogInserter
- type ErrLogWorker
- func (w *ErrLogWorker) Close(ctx context.Context) error
- func (w *ErrLogWorker) DroppedExempt() int64
- func (w *ErrLogWorker) DroppedReject() int64
- func (w *ErrLogWorker) EnqueueError(l *domain.UsageLog)
- func (w *ErrLogWorker) EnqueueRejected(l *domain.UsageLog)
- func (w *ErrLogWorker) Inserted() int64
- func (w *ErrLogWorker) Name() string
- func (w *ErrLogWorker) Queued() int
- func (w *ErrLogWorker) Start(ctx context.Context) error
- func (w *ErrLogWorker) Stats() any
- type ErrLogWorkerStats
- type LogInserter
- type PartitionManager
- type QuotaWriter
- type Recorder
- func (r *Recorder) AddQuota(keyID int64, delta int64)
- func (r *Recorder) Close(ctx context.Context) error
- func (r *Recorder) Name() string
- func (r *Recorder) Pending() int
- func (r *Recorder) Record(l *domain.UsageLog)
- func (r *Recorder) SetQuotaWriter(q QuotaWriter)
- func (r *Recorder) Start(ctx context.Context) error
- func (r *Recorder) Stats() any
- type RecorderStats
- type RetentionConfig
- type RetentionWorker
- type RetentionWorkerStats
- type StatsAggConfig
- type StatsAggStore
- type StatsAggWorker
- type StatsAggWorkerStats
- type UsageConfig
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type ErrLogConfig ¶
type ErrLogConfig struct {
QueueSize int // 拒绝队列容量(背压面:满 → 丢弃;默认 4096 ≈ 1MB 量级内存)
ExemptQueueSize int // 豁免队列容量(双轨行;默认 1024——恒落盘语义,仅自身溢出才丢)
BatchSize int // 每批落盘行数(默认 500;单条 CreateBulk 有界)
FlushInterval time.Duration // 批间隔(默认 500ms;DB 写速率上界 = BatchSize/FlushInterval)
}
ErrLogConfig err_logs 落盘 worker 节奏。
type ErrLogInserter ¶
ErrLogInserter err_logs 批量插入面(repository.ErrLogRepo 实现)。
type ErrLogWorker ¶
type ErrLogWorker struct {
// contains filtered or unexported fields
}
ErrLogWorker 错误明细落盘 worker(worker.Worker 契约,Name="errlog")。
func NewErrLogWorker ¶
func NewErrLogWorker(cfg ErrLogConfig, writer ErrLogInserter, log *logx.Logger) *ErrLogWorker
func (*ErrLogWorker) Close ¶
func (w *ErrLogWorker) Close(ctx context.Context) error
Close 幂等排空(优雅停机核心):置位 closed(此后 Enqueue 丢弃计数——无尾 窗口静默丢,S4)→ 等 worker goroutine 退出(受预算约束;loop 退出前在途 flush 已同步收尾)→ 两队列剩余条目分批落库(每批 ≤ BatchSize,预算内完整 排空——正常停机双轨行 + 未丢弃拒绝行不丢);排空失败同 flush 语义:豁免行 回灌重试一次(排空循环为紧凑 while、无 ticker 节奏,回灌无限重试即紧循环打 爆 DB——仅重试一次,保留"不无限重试"的既有预算语义)、拒绝行按采样语义 丢弃;ctx 到期 → Warn(含已排空/剩余条数)+ 截断退出(剩余丢弃计数),不 阻塞停机。Close 结束打印 inserted/dropped 终值(S3 对账观测)。未 Start 也 安全(跳过 loop 等待直接排空)。
func (*ErrLogWorker) DroppedExempt ¶
func (w *ErrLogWorker) DroppedExempt() int64
DroppedExempt 双轨行丢弃计数(恒 0 为正常——豁免队列溢出/落库失败止损重试 耗尽;>0 即异常态观测)。
func (*ErrLogWorker) DroppedReject ¶
func (w *ErrLogWorker) DroppedReject() int64
DroppedReject 拒绝行丢弃计数(队列满采样;统计面对账指标)。
func (*ErrLogWorker) EnqueueError ¶
func (w *ErrLogWorker) EnqueueError(l *domain.UsageLog)
EnqueueError 投递一条**双轨行**(finish/recordLog 的已计费错误:abort/failover/ 4xx/5xx/network——usage_logs 错误行):豁免队列——**不参与拒绝风暴采样丢弃, 恒落盘**(架构审查 B2;仅本队列自身溢出才丢——异常态,Warn 恰好一次)。
func (*ErrLogWorker) EnqueueRejected ¶
func (w *ErrLogWorker) EnqueueRejected(l *domain.UsageLog)
EnqueueRejected 投递一条**拒绝行**(recordRejected:401/429/402/400/404 本地 拒绝 + 组限流):普通队列——风暴采样丢弃面(非阻塞 select-default;丢弃 计数 DroppedReject 原子累加,供指标/日志对账)。
func (*ErrLogWorker) Name ¶
func (w *ErrLogWorker) Name() string
Name worker.Worker 契约(wm 反向排空:errlog 在 rec 之后注册 → 先于 rec 排空 错误明细;与计费 flusher 无共享状态,排空顺序互不依赖)。
func (*ErrLogWorker) Queued ¶
func (w *ErrLogWorker) Queued() int
Queued 当前两队列积压总条数(背压观测:恒 ≤ QueueSize + ExemptQueueSize)。
func (*ErrLogWorker) Stats ¶
func (w *ErrLogWorker) Stats() any
Stats 满足 handler.StatsProvider(独立于 worker.Worker 契约;装配链路见 internal/handler/ops.go 文件头)。
type ErrLogWorkerStats ¶
type ErrLogWorkerStats struct {
Queued int `json:"queued"` // 两队列积压总条数(恒 ≤ 容量和)
QueueCap int `json:"queue_cap"` // 拒绝队列容量(风暴采样阈值)
ExemptQueueCap int `json:"exempt_queue_cap"` // 豁免队列容量(双轨行)
DroppedReject int64 `json:"dropped_reject"` // 拒绝行采样丢弃累计
DroppedExempt int64 `json:"dropped_exempt"` // 双轨行丢弃累计(>0 即异常态)
Inserted int64 `json:"inserted"` // 成功落盘累计
WarnedReject bool `json:"warned_reject"` // 拒绝丢弃告警边沿
WarnedExempt bool `json:"warned_exempt"` // 双轨丢弃告警边沿
}
ErrLogWorkerStats err_logs 落盘 worker 状态(队列占用 + 丢弃/落盘计数)。
type LogInserter ¶
type PartitionManager ¶
type PartitionManager interface {
EnsureUsageLogPartitions(ctx context.Context, now, until time.Time) error
DropUsageLogPartitionsBefore(ctx context.Context, cutoff time.Time) (int, error)
EnsureErrLogPartitions(ctx context.Context, now, until time.Time) error
DropErrLogPartitionsBefore(ctx context.Context, cutoff time.Time) (int, error)
EnsureUsageStatsPartitions(ctx context.Context, now, until time.Time) error
DropUsageStatsPartitionsBefore(ctx context.Context, cutoff time.Time) (int, error)
EnsureUsageEntityStatsPartitions(ctx context.Context, now, until time.Time) error
DropUsageEntityStatsPartitionsBefore(ctx context.Context, cutoff time.Time) (int, error)
DeleteRedemptionUsesBefore(ctx context.Context, cutoff time.Time) (int, error)
}
PartitionManager 分区管理面(repository.Repository 实现):保留策略只需 DROP 过期分区 + 预建未来分区,不感知分区表内部 DDL。now/until 由调用方 传入(start 边界由 now 推导,测试可注入时钟)。四表各自独立调度(保留期 独立:LogRetentionDays / ErrLogRetentionDays / StatsRetentionDays——后者由 usage_stats 与 usage_entity_stats 共用,同一循环 DROP+预建)。 DeleteRedemptionUsesBefore 是普通表(redemption_uses 无分区可 DROP)的 有界批删路径——同为保留策略的周期清理手段,归口本接口(F3-2)。
type QuotaWriter ¶
QuotaWriter 批量回写 key 额度消耗(增量;内存权威,DB 滞后 ≤ flush 间隔)。 由 proxy 的 gate 计数 + 本 Recorder 的 flush 节奏落库(Phase 3a:额度后扣)。
type Recorder ¶
type Recorder struct {
// contains filtered or unexported fields
}
func New ¶
func New(cfg UsageConfig, logs LogInserter, log *logx.Logger) *Recorder
func (*Recorder) AddQuota ¶
AddQuota 累加 key 额度增量并入同一 quotaUsed map、同一回写路径(不回桶、 不落统计)。Record 对 KeyID>0 行的累加是唯一生产增量入口(计费 worker 只动 余额不动配额);本方法是测试/手工注入面,汇入同一 map 与回写路径。
func (*Recorder) Close ¶
Close 幂等排空(优雅停机核心):等聚合 goroutine 退出(受预算约束)→ 以 flushMu 获取等待在途批次(SIGTERM 时 ticker 批次可能已在途占住 flushMu 且 pending 已 swap;Close 必须先等其结束,否则 drain 循环见 pendingN==0 会 静默提前返回,在途批次无界运行——O1 复测根因 1)→ 受 shutdown ctx 预算 约束的排空循环(此时无在途批次、flushMu 无竞争)。正常情形完整排空语义 不变(无 deadline ctx = 全部落库);ctx 到期 → Cancel baseCtx(在途落库 快速失败回灌,不丢)+ Warn(flushed/remaining 条数单位一致)+ 截断退出, 不阻塞停机;在途批次收尾超时(A-P2-8-2)→ 放弃排空、Warn 截断退出(在途 由已取消 baseCtx 收尾回灌不丢)。额度面由 flushQuota 以本 ctx 预算收尾 (到期截断 + Warn)。未 Start 也安全(跳过聚合等待;在途 flush 与 pending 残留同样等待/排空)。
func (*Recorder) Record ¶
Record 记录一次放行路径明细(非 billed 行):短锁归并 pending(无界 slice append,O(1) 摊还)+ quota 原子累加——**永不阻塞**(无 channel:此前有界 channel cap 16384 饱和阻塞发送是 off 路径 16.4k goroutine 卡 chan send 幽灵 根因;HTTP 层过载保护由 max_inflight 兜底,pending 内存由水线 Warn 观测, 崩溃丢 ≤1 flush 窗口语义不变)。热路径零额外开销:closed 检查为 1 次 atomic.Load(I-4)。**零统计计算**(spec 2026-08-14):统计桶机制整体删除, 锁内仅剩明细 append + quotaUsed 累加两个 O(1) 操作。
func (*Recorder) SetQuotaWriter ¶
func (r *Recorder) SetQuotaWriter(q QuotaWriter)
SetQuotaWriter 注入额度回写器(装配期调用;nil = 关闭回写)。
type RecorderStats ¶
type RecorderStats struct {
PendingLogs int64 `json:"pending_logs"` // 尚未落库的明细条数
PendingWaterline int64 `json:"pending_waterline"` // 水线(包级 var 直读)
Warned bool `json:"warned"` // 水线告警边沿是否置位
}
RecorderStats usage 明细/额度 worker 状态(spec 2026-08-14:统计桶机制整体 删除——stat_buckets_created 观测随之消失)。
type RetentionConfig ¶
type RetentionConfig struct {
LogRetentionDays int // usage_logs 分区保留天数(config usage.log_retention_days;<= 0 = 不删除)
ErrLogRetentionDays int // err_logs 分区保留天数(config usage.errlog_retention_days,默认 7 天短保留——错误审计;<= 0 = 不删除)
StatsRetentionDays int // usage_stats 分区保留天数(config usage.stats_retention_days,默认 180 天——聚合统计长保留;<= 0 = 不删除)
TickerInterval time.Duration // 巡检周期(生产 1h;测试注入短周期;<= 0 兜底 1h)
}
RetentionConfig retention worker 配置。
type RetentionWorker ¶
type RetentionWorker struct {
// contains filtered or unexported fields
}
RetentionWorker 按日分区保留 worker(worker.Worker 契约,Name="retention"): 每小时巡检一次——四表各自独立调度(usage_logs 按 LogRetentionDays、err_logs 按 ErrLogRetentionDays、usage_stats/usage_entity_stats 共用 StatsRetentionDays (同一循环 DROP+预建,180d),同一循环无新增 goroutine):
- DROP 分区下界 < now - 保留天数的分区(DROP TABLE O(1),比逐行 DELETE 快 5~6 个量级;按分区名日期判定,无需查元数据——usage_stats 保留清理 用户裁决 2026-08-11:PG DELETE 不释放空间,必须分区 DROP)
- 预建 当日 + 未来 1 天 分区(PG 无自动建分区,防日界跨区插入失败)
- redemption_uses 有界批删(F3-2):普通表无分区可 DROP,同一循环内每轮 DELETE 至多 5000 行超窗行(TTL 定死 90 天,见 redemptionUseRetentionDays) ——低频表单轮即清,超大批多轮收敛(每轮上限防长事务持锁)
DROP × 在途插入竞态(评审 I-3):DROP TABLE 需 ACCESS EXCLUSIVE 锁,与 在途插入事务串行;能落进被 DROP 分区(保留期前)的行只有回放/陈旧 created_at 的延迟日志——该分区数据本就在保留语义内(要清理)。万一插入 恰好失败 → 走落库失败路径(Warn + 丢弃,与普通批量落库失败同语义,不自愈 不重试),可接受。
与 Recorder/ErrLogWorker 解耦(不依赖明细管道):DROP 幂等,无排空需求, Close 直接返回 nil。
func NewRetention ¶
func NewRetention(cfg RetentionConfig, parts PartitionManager, log *logx.Logger) *RetentionWorker
func (*RetentionWorker) Close ¶
func (w *RetentionWorker) Close(ctx context.Context) error
Close 幂等(worker.Worker 契约):DROP/预建均幂等,无排空需求。
func (*RetentionWorker) Name ¶
func (w *RetentionWorker) Name() string
Name worker.Worker 契约(注册顺序无依赖——DROP/预建均幂等)。
func (*RetentionWorker) Stats ¶
func (w *RetentionWorker) Stats() any
Stats 满足 handler.StatsProvider(独立于 worker.Worker 契约;装配链路见 internal/handler/ops.go 文件头)。
type RetentionWorkerStats ¶
type RetentionWorkerStats struct {
LastPatrolUnixMs int64 `json:"last_patrol_unix_ms"` // 最近一次巡检完成时刻(0 = 尚未巡检)
LastDroppedLogPartitions int64 `json:"last_dropped_log_partitions"` // 最近成功轮 usage_logs DROP 分区数(失败轮保留上轮值)
LastDroppedErrLogPartitions int64 `json:"last_dropped_errlog_partitions"` // 最近成功轮 err_logs DROP 分区数(失败轮保留上轮值)
LastDroppedStatsPartitions int64 `json:"last_dropped_stats_partitions"` // 最近成功轮 usage_stats DROP 分区数(失败轮保留上轮值)
LastDroppedEntityStatsPartitions int64 `json:"last_dropped_entity_stats_partitions"` // 最近成功轮 usage_entity_stats DROP 分区数(与 stats 同 StatsRetentionDays,失败轮保留上轮值)
LogRetentionDays int `json:"log_retention_days"`
ErrLogRetentionDays int `json:"errlog_retention_days"`
StatsRetentionDays int `json:"stats_retention_days"`
}
RetentionWorkerStats 分区保留 worker 状态(runOnce 收尾原子写,零新增 DB)。
type StatsAggConfig ¶
type StatsAggConfig struct {
Interval time.Duration // 聚合周期(config usage.stats_agg_interval,默认 5m;0 = 禁用——Start 直接返回)
Lag time.Duration // 读窗口滞后(读窗口 [W, T),T = now − Lag;watermark 只推进到 T)
}
StatsAggConfig 离线聚合 worker 配置。
type StatsAggStore ¶
type StatsAggStore interface {
// AcquireStatsAggLock 抢占会话级 advisory lock(专用连接持有到 release;
// 抢锁失败 ok=false——其他实例在聚合,本轮跳过)。
AcquireStatsAggLock(ctx context.Context) (release func(), ok bool, err error)
// LoadAggRange 重算范围 [from,to) 双结果集全量重建:cube 两查询 +
// entity 六查询(含消费明细行数——cube 两查询 count(*) 合计)。
LoadAggRange(ctx context.Context, from, to time.Time) ([]*domain.StatBucket, []*domain.EntityStatBucket, int64, error)
// AggregateRange 单事务 DELETE cube [delFrom,delTo) + INSERT cube +
// DELETE entity [同范围] + INSERT entity + watermark 推进 wmTo(= 读窗口 T,
// ≠ 重算范围上界——P1-A 两范围分离;双表同一事务原子回滚)。
AggregateRange(ctx context.Context, delFrom, delTo, wmTo time.Time, cube []*domain.StatBucket, entity []*domain.EntityStatBucket) error
LoadStatsAggWatermark(ctx context.Context) (time.Time, error)
InitStatsAggWatermark(ctx context.Context, t time.Time) error
}
StatsAggStore 离线聚合存储面(repository.StatRepo 实现):两范围分离 + 单事务覆盖落盘 + watermark + 会话级 advisory lock 全部在 repo 侧 (stat_repo.go/stat_entity_agg.go——pgx 直查直写,ent 无数组列类型)。
type StatsAggWorker ¶
type StatsAggWorker struct {
// contains filtered or unexported fields
}
StatsAggWorker 离线聚合 worker(worker.Worker 契约,Name="stats-agg"): 每周期一个两范围 + 双结果集(cube 两查询 + entity 六查询)+ 单事务流程 (spec §3):
读窗口 [W, T):W = watermark,T = now − 滞后——只推进 watermark,不直接
用于 DELETE(部分小时桶边界问题,见下)
重算范围 [R0, R1) = [trunc_hour(W), trunc_hour(T) + 1h)——DELETE + SELECT
共同边界(cube 与 entity 两表同范围)
LoadAggRange(R0, R1) → AggregateRange(R0, R1, T, cube, entity) 单事务落盘
**两范围分离(评审 P1-A,核心正确性)**:小时桶是部分完成的桶(跨多周期 累积),直接按读窗口 DELETE 会截断当前小时桶([小时起点, W) 的行丢失)→ 每周期欠计。小时对齐扩展后:SELECT 覆盖已消费行无害(DELETE 先清、INSERT 全量覆盖,幂等仍成立——重放同范围结果一致)。watermark 只推进到 T(原始读 位置),**不推进到 R1**——推进到 R1 会永久跳过 [T, R1) 的行(正是 P1-A 要防 的错误形态;签名显式分离见 AggregateRange 的 wmTo 参数)。
幂等/重放(issue #8 教训):DELETE+INSERT+watermark 推进同一事务——崩溃 回滚 → 游标不动 → 重算恢复不双计;重复执行同范围(手动重算修正)结果一致 (覆盖语义)。
并发防护:pg_try_advisory_lock(会话级,专用连接持有整个周期——池连接复用 即丢锁,P3);抢锁失败 → 本轮跳过(其他实例在聚合)。单写者语义由此钉死, 事务内串行无 40P01 重试需求(P2-5 取舍:advisory lock 串行下单写者)。
watermark 存储/初始化/追赶(评审 P2-1):单行 watermark 表(stats_agg_ watermark,bootstrap 建表见 partition.go);**全新库初始化 = now − 滞后** (防首跑扫全史 + DELETE 撞 retention 已 DROP 分区);ON CONFLICT DO NOTHING 容忍多实例并发初始化(败者重读既有值);**追赶上限**:停摆恢复后单周期 窗口 ≤ 1h 分批收敛(防单次超大窗口)。
**手动重建运维口径(Momus B1 勘误)**:worker 对缺失 watermark 行的初始化 硬编码 now−lag——"清空 watermark 行"不会触发历史回算。手动重建统计必须: (1) 清空 usage_stats / usage_entity_stats 数据;(2) **手工种子单行 watermark** 至最早保留小时边界(INSERT INTO stats_agg_watermark (id, watermark) VALUES (1, '<最早保留小时>')),否则历史窗口永远聚合不出数据(worker 只从 watermark 起追赶,且受 1h/周期上限分批收敛)。
func NewStatsAgg ¶
func NewStatsAgg(cfg StatsAggConfig, store StatsAggStore, log *logx.Logger) *StatsAggWorker
func (*StatsAggWorker) Close ¶
func (w *StatsAggWorker) Close(ctx context.Context) error
Close 幂等(worker.Worker 契约):无排空需求(usage_stats 由 DB 侧覆盖语义 收敛,worker 停摆期间的落后窗口由重启后追赶上限分批收敛)。
func (*StatsAggWorker) Start ¶
func (w *StatsAggWorker) Start(ctx context.Context) error
Start worker.Worker 契约:Interval <= 0 = 禁用聚合(config 0 语义——不启动 循环,Close 直接返回;等价于不装配本 worker)。
func (*StatsAggWorker) Stats ¶
func (w *StatsAggWorker) Stats() any
Stats 满足 handler.StatsProvider(观测面:watermark/上轮桶数/上轮行数/上轮 耗时;失败轮保留上轮值——对齐 RetentionWorker 观测纪律)。
type StatsAggWorkerStats ¶
type StatsAggWorkerStats struct {
// WatermarkUnixMs watermark 位置(毫秒;0 = 尚未推进——未初始化/首轮未完成)
WatermarkUnixMs int64 `json:"watermark_unix_ms"`
// LastBuckets 上轮写入桶数(cube + entity 合计;失败轮保留上轮值)
LastBuckets int64 `json:"last_buckets"`
// LastRows 上轮消费明细行数(cube 两查询 count(*) 合计——实体六查询扫的
// 是同批源行不重复计数;失败轮保留上轮值)
LastRows int64 `json:"last_rows"`
// LastDurationMs 上轮耗时(毫秒;失败轮保留上轮值)
LastDurationMs int64 `json:"last_duration_ms"`
}
StatsAggWorkerStats 离线聚合 worker 观测(/ops/workers;runOnce 收尾原子写, 零新增 DB)。
type UsageConfig ¶
type UsageConfig struct {
BatchSize int
FlushInterval time.Duration
QuotaFlushInterval time.Duration // quota 增量批量回写 cadence
Workers int // flush 并行 worker 数(0 = 单 worker;O1 模式分片并行)
// StatsAggInterval 离线聚合周期(spec 2026-08-14;config usage.stats_agg_
// interval,默认 5m;0 = 禁用聚合——不装配聚合 worker 的等价语义)。
StatsAggInterval time.Duration
}