diff --git a/.env.example b/.env.example index 27d53641..4807d011 100644 --- a/.env.example +++ b/.env.example @@ -116,10 +116,6 @@ LOG_FILE_ENABLED=true # Number of historical days to retain combined logs, plus the current day; 0 disables cleanup; error-only logs keep 30 history days plus the current day. Required: no. Default: 7. LOG_RETENTION_DAYS=7 -# 是否在每日维护中删除过期 usage_events 原始事件;默认不清理,设为 true 后会删除早于 90 个本地自然日前 00:00 的数据。必填:否。默认值:false。 -# Whether daily maintenance deletes expired raw usage_events; disabled by default. When true, rows earlier than local midnight 90 calendar days ago are deleted. Required: no. Default: false. -CLEANUP_USAGE_EVENTS_ENABLED=false - # 是否启用 SQLite 数据库备份。必填:否。默认值:true。 # Whether to enable SQLite database backups. Required: no. Default: true. BACKUP_ENABLED=true diff --git a/README.md b/README.md index dee082ed..ff097a3e 100644 --- a/README.md +++ b/README.md @@ -416,11 +416,12 @@ Scheduled Auth Files quota refresh is configured from the gear button in the Aut | `LOG_LEVEL` | No | `info` | Log level | | `LOG_FILE_ENABLED` | No | `true` | Write persistent log files | | `LOG_RETENTION_DAYS` | No | `7` | Combined-log history days, plus the current day; `0` disables cleanup. Error-only logs keep 30 history days plus the current day | -| `CLEANUP_USAGE_EVENTS_ENABLED` | No | `false` | Delete expired raw `usage_events` during daily maintenance; when enabled, rows earlier than local midnight 90 calendar days ago are deleted | | `BACKUP_ENABLED` | No | `true` | Enable SQLite database backups | | `BACKUP_INTERVAL` | No | `24h` | Database backup interval | | `BACKUP_RETENTION_DAYS` | No | `7` | Backup retention days | +Keeper automatically moves raw `usage_events` older than 90 local calendar days into the permanently retained `usage_events_archive` cold table during the daily 04:30 maintenance window. The archive is reserved for future schema-migration rebuilds and is not queried by normal dashboard APIs. + When file logging is enabled, `cpa-usage-keeper-YYYY-MM-DD.log` contains all emitted levels. Error, fatal, and panic entries are also copied to `cpa-usage-keeper-error-YYYY-MM-DD.log`, which keeps the previous 30 local calendar dates plus the current date. ### Built-In HTTPS diff --git a/README.zh.md b/README.zh.md index 115b2e8f..69f05316 100644 --- a/README.zh.md +++ b/README.zh.md @@ -416,11 +416,12 @@ Auth Files 定时限额刷新在 Auth Files 巡检弹窗的小齿轮中配置。 | `LOG_LEVEL` | 否 | `info` | 日志级别 | | `LOG_FILE_ENABLED` | 否 | `true` | 是否写入持久化日志文件 | | `LOG_RETENTION_DAYS` | 否 | `7` | 综合日志保留历史天数,并额外保留当天;`0` 表示不自动清理。仅错误日志固定保留历史 30 天及当天 | -| `CLEANUP_USAGE_EVENTS_ENABLED` | 否 | `false` | 是否在每日维护中删除过期 `usage_events` 原始事件;启用后会删除早于 90 个本地自然日前 00:00 的数据 | | `BACKUP_ENABLED` | 否 | `true` | 是否启用 SQLite 数据库备份 | | `BACKUP_INTERVAL` | 否 | `24h` | 数据库备份间隔 | | `BACKUP_RETENTION_DAYS` | 否 | `7` | 备份保留天数 | +Keeper 会在每天 04:30 的维护窗口中,把早于 90 个本地自然日的原始 `usage_events` 自动移动到永久保留的 `usage_events_archive` 冷表。该冷表用于未来 schema migration 重建增量数据,正常仪表盘 API 不查询 archive。 + 启用文件日志后,`cpa-usage-keeper-YYYY-MM-DD.log` 会记录所有已输出级别;error、fatal 和 panic 级别还会同时写入 `cpa-usage-keeper-error-YYYY-MM-DD.log`,该文件固定保留历史 30 个本地自然日及当天。 ### 内置 HTTPS diff --git a/go.mod b/go.mod index 9ceaf61e..0a9b8cba 100644 --- a/go.mod +++ b/go.mod @@ -7,6 +7,7 @@ require ( github.com/joho/godotenv v1.5.1 github.com/mattn/go-sqlite3 v1.14.48 github.com/sirupsen/logrus v1.9.3 + golang.org/x/sys v0.20.0 golang.org/x/text v0.20.0 gorm.io/driver/sqlite v1.5.7 gorm.io/gorm v1.26.1 @@ -38,7 +39,6 @@ require ( golang.org/x/arch v0.8.0 // indirect golang.org/x/crypto v0.23.0 // indirect golang.org/x/net v0.25.0 // indirect - golang.org/x/sys v0.20.0 // indirect google.golang.org/protobuf v1.34.1 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/internal/app/app.go b/internal/app/app.go index bc57d03a..117304a8 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -181,9 +181,8 @@ func NewWithConfig(cfg config.Config) (*App, error) { usageAggregationRunner := poller.NewUsageAggregationRunner(db) // syncService 仍然是 metadata 和 usage 处理共享的业务服务入口。 syncService := service.NewSyncServiceWithOptions(db, service.SyncServiceOptions{ - BaseURL: cfg.CPABaseURL, - Client: cpaClient, - CleanupUsageEventsEnabled: cfg.CleanupUsageEventsEnabled, + BaseURL: cfg.CPABaseURL, + Client: cpaClient, // usage_events 事务提交后通过这个缓存做非阻塞增量追加,供 Overview realtime 和右边界补偿复用。 RecentUsageEvents: recentUsageCache, // usage 与 metadata 提交后只唤醒单 writer runner,不在前台链路执行派生聚合。 diff --git a/internal/config/config.go b/internal/config/config.go index 36d5252e..89957501 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -78,8 +78,6 @@ type Config struct { BackupInterval time.Duration // BackupRetentionDays 是备份文件保留天数。 BackupRetentionDays int - // CleanupUsageEventsEnabled 控制每日维护是否删除过期 usage_events 原始事件。 - CleanupUsageEventsEnabled bool // RequestTimeout 是访问 CPA HTTP 和 Redis TCP 的超时时间。 RequestTimeout time.Duration // TLSSkipVerify 控制是否跳过 CPA HTTPS 和 Redis 队列 TLS 的证书验证。 @@ -185,11 +183,6 @@ func Load(options LoadOptions) (*Config, error) { if backupRetentionDays < 0 { return nil, fmt.Errorf("BACKUP_RETENTION_DAYS must be non-negative") } - cleanupUsageEventsEnabled, err := getBool("CLEANUP_USAGE_EVENTS_ENABLED", false) - if err != nil { - return nil, err - } - logFileEnabled, err := getBool("LOG_FILE_ENABLED", true) if err != nil { return nil, err @@ -263,7 +256,6 @@ func Load(options LoadOptions) (*Config, error) { BackupDir: filepath.Join(workDir, workDirBackupsName), BackupInterval: backupInterval, BackupRetentionDays: backupRetentionDays, - CleanupUsageEventsEnabled: cleanupUsageEventsEnabled, RequestTimeout: requestTimeout, TLSSkipVerify: tlsSkipVerify, LogLevel: getString("LOG_LEVEL", "info"), diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 278fa34c..728b505c 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -15,7 +15,7 @@ import ( var configEnvKeys = []string{ "APP_HOST", "APP_PORT", "APP_BASE_PATH", "CPA_PUBLIC_URL", "WORK_DIR", "CPA_BASE_URL", "CPA_MANAGEMENT_KEY", "POLL_INTERVAL", "USAGE_SYNC_MODE", "REDIS_QUEUE_ADDR", "REDIS_QUEUE_TLS", "REDIS_QUEUE_BATCH_SIZE", "REDIS_QUEUE_IDLE_INTERVAL", - "SQLITE_PATH", "BACKUP_ENABLED", "BACKUP_DIR", "BACKUP_INTERVAL", "BACKUP_RETENTION_DAYS", "CLEANUP_USAGE_EVENTS_ENABLED", + "SQLITE_PATH", "BACKUP_ENABLED", "BACKUP_DIR", "BACKUP_INTERVAL", "BACKUP_RETENTION_DAYS", "REQUEST_TIMEOUT", "LOG_LEVEL", "LOG_FILE_ENABLED", "LOG_DIR", "LOG_RETENTION_DAYS", "AUTH_ENABLED", "LOGIN_PASSWORD", "AUTH_SESSION_TTL", "TZ", "TLS_SKIP_VERIFY", "QUOTA_REFRESH_WORKER_LIMIT", } diff --git a/internal/config/test/config_cleanup_test.go b/internal/config/test/config_cleanup_test.go index 741f7cbd..64b709e4 100644 --- a/internal/config/test/config_cleanup_test.go +++ b/internal/config/test/config_cleanup_test.go @@ -13,43 +13,12 @@ var isolatedConfigEnvKeys = []string{ "APP_HOST", "APP_PORT", "APP_BASE_PATH", "CPA_PUBLIC_URL", "WORK_DIR", "CPA_BASE_URL", "CPA_MANAGEMENT_KEY", "CPA_REQUEST_LOG_ACCESS_ENABLED", "REDIS_QUEUE_ADDR", "REDIS_QUEUE_TLS", "REDIS_QUEUE_BATCH_SIZE", "REDIS_QUEUE_IDLE_INTERVAL", - "BACKUP_ENABLED", "BACKUP_INTERVAL", "BACKUP_RETENTION_DAYS", "CLEANUP_USAGE_EVENTS_ENABLED", + "BACKUP_ENABLED", "BACKUP_INTERVAL", "BACKUP_RETENTION_DAYS", "REQUEST_TIMEOUT", "LOG_LEVEL", "LOG_FILE_ENABLED", "LOG_DIR", "LOG_RETENTION_DAYS", "AUTH_ENABLED", "LOGIN_PASSWORD", "AUTH_SESSION_TTL", "TZ", "TLS_ENABLED", "TLS_CERT_FILE", "TLS_KEY_FILE", "TLS_SKIP_VERIFY", "QUOTA_REFRESH_WORKER_LIMIT", } -func TestLoadFromEnvDefaultsUsageEventCleanupDisabled(t *testing.T) { - isolateConfigEnv(t) - t.Setenv("CPA_BASE_URL", "http://127.0.0.1:"+cpa.ManagementRedisDefaultPort) - t.Setenv("CPA_MANAGEMENT_KEY", "secret") - - cfg, err := config.LoadFromEnv() - if err != nil { - t.Fatalf("LoadFromEnv returned error: %v", err) - } - - if cfg.CleanupUsageEventsEnabled { - t.Fatal("expected usage_events cleanup to be disabled by default") - } -} - -func TestLoadFromEnvReadsUsageEventCleanupFlag(t *testing.T) { - isolateConfigEnv(t) - t.Setenv("CPA_BASE_URL", "http://127.0.0.1:"+cpa.ManagementRedisDefaultPort) - t.Setenv("CPA_MANAGEMENT_KEY", "secret") - t.Setenv("CLEANUP_USAGE_EVENTS_ENABLED", "true") - - cfg, err := config.LoadFromEnv() - if err != nil { - t.Fatalf("LoadFromEnv returned error: %v", err) - } - - if !cfg.CleanupUsageEventsEnabled { - t.Fatal("expected usage_events cleanup to be enabled") - } -} - func TestLoadFromEnvDefaultsCPARequestLogAccessDisabled(t *testing.T) { isolateConfigEnv(t) t.Setenv("CPA_BASE_URL", "http://127.0.0.1:"+cpa.ManagementRedisDefaultPort) diff --git a/internal/entities/all.go b/internal/entities/all.go index ddf5466a..1787cec3 100644 --- a/internal/entities/all.go +++ b/internal/entities/all.go @@ -4,6 +4,7 @@ package entities func All() []any { return []any{ &UsageEvent{}, + &UsageEventArchive{}, &RedisUsageInbox{}, &ModelPriceSetting{}, &ModelPriceRule{}, diff --git a/internal/entities/test/entities_test.go b/internal/entities/test/entities_test.go index 9a2ced1e..bc14cf70 100644 --- a/internal/entities/test/entities_test.go +++ b/internal/entities/test/entities_test.go @@ -15,6 +15,7 @@ func TestAllIncludesCoreModels(t *testing.T) { items := All() expected := []any{ &UsageEvent{}, + &UsageEventArchive{}, &RedisUsageInbox{}, &ModelPriceSetting{}, &ModelPriceRule{}, diff --git a/internal/entities/usage_event_archive.go b/internal/entities/usage_event_archive.go new file mode 100644 index 00000000..961171a5 --- /dev/null +++ b/internal/entities/usage_event_archive.go @@ -0,0 +1,46 @@ +package entities + +import "time" + +// UsageEventStorageColumns 是 hot/archive 原始事件复制使用的完整持久化列契约。 +const UsageEventStorageColumns = "id, event_key, api_group_key, provider, endpoint, auth_type, request_id, client_ip, x_forwarded_for, user_agent, model, model_alias, reasoning_effort, service_tier, response_service_tier, executor_type, timestamp, source, auth_index, failed, generate, latency_ms, ttft_ms, input_tokens, output_tokens, reasoning_tokens, cached_tokens, cache_read_tokens, cache_creation_tokens, total_tokens, created_at" + +// UsageEventArchive 永久保存已经离开 hot usage_events 的原始事件。 +// 字段必须与 UsageEvent 的持久化列保持一致,但 archive 不承担在线查询,因此不复制二级索引。 +type UsageEventArchive struct { + ID int64 `gorm:"primaryKey;autoIncrement:false"` + EventKey string + APIGroupKey string + Provider string `gorm:"column:provider"` + Endpoint string `gorm:"column:endpoint"` + AuthType string `gorm:"column:auth_type"` + RequestID string `gorm:"column:request_id"` + ClientIP *string `gorm:"column:client_ip"` + XForwardedFor *string `gorm:"column:x_forwarded_for"` + UserAgent *string `gorm:"column:user_agent"` + Model string + ModelAlias *string `gorm:"column:model_alias"` + ReasoningEffort string `gorm:"column:reasoning_effort;not null;default:''"` + ServiceTier string `gorm:"column:service_tier;not null;default:''"` + ResponseServiceTier string `gorm:"column:response_service_tier;not null;default:''"` + ExecutorType string `gorm:"column:executor_type;not null;default:''"` + Timestamp time.Time `gorm:"serializer:storageTime"` + Source string + AuthIndex string + Failed bool + Generate *bool `gorm:"column:generate;not null;default:true"` + LatencyMS int64 + TTFTMS *int64 `gorm:"column:ttft_ms"` + InputTokens int64 + OutputTokens int64 + ReasoningTokens int64 + CachedTokens int64 + CacheReadTokens int64 `gorm:"not null;default:0"` + CacheCreationTokens int64 `gorm:"not null;default:0"` + TotalTokens int64 + CreatedAt time.Time `gorm:"serializer:storageTime"` +} + +func (UsageEventArchive) TableName() string { + return "usage_events_archive" +} diff --git a/internal/repository/db.go b/internal/repository/db.go index 3bee7fa1..01cfba4e 100644 --- a/internal/repository/db.go +++ b/internal/repository/db.go @@ -335,95 +335,51 @@ func InsertUsageEvents(db *gorm.DB, events []entities.UsageEvent) (int, int, err return inserted, 0, nil } -type CleanupStorageOptions struct { - // CleanupUsageEvents 控制是否删除过期 usage_events 原始事件;默认 false 表示保留原始事件。 - CleanupUsageEvents bool -} - const usageEventsRetentionDays = 90 -// CleanupStorage 是每日维护任务的统一仓储清理入口:先清 inbox/raw events,再清 Activity 与 Latency 限期统计,最后执行 VACUUM。 -// VACUUM 必须在删除完成后单独执行,任何一步失败都会停止后续步骤并把已完成部分的结果返回给上层日志。 -func CleanupStorage(db *gorm.DB, now time.Time, options ...CleanupStorageOptions) (dto.StorageCleanupResult, error) { - opts := CleanupStorageOptions{} - if len(options) > 0 { - opts = options[0] - } +// CleanupStorage 是每日维护任务的统一仓储入口:先清 inbox、归档 raw events,再清限期统计并条件式整理空闲页。 +func CleanupStorage(db *gorm.DB, now time.Time) (dto.StorageCleanupResult, error) { redisResult, err := CleanupRedisUsageInbox(db, now) if err != nil { return dto.StorageCleanupResult{RedisInbox: redisResult}, err } - var usageEventsDeleted int64 - if opts.CleanupUsageEvents { - usageEventsDeleted, err = cleanupUsageEvents(db, now) - if err != nil { - return dto.StorageCleanupResult{RedisInbox: redisResult, UsageEventsDeleted: usageEventsDeleted}, err - } + usageEventsArchive, err := ArchiveExpiredUsageEvents(databaseContext(db), db, now) + result := dto.StorageCleanupResult{ + RedisInbox: redisResult, + UsageEventsArchived: usageEventsArchive.Archived, + UsageEventsArchiveStatus: usageEventsArchive.Status, + } + if err != nil { + return result, err } // Activity 的 short/medium/long 分别按自身 retention 清理,daily 永久保留。 if err := CleanupUsageActivityStats(db, now); err != nil { - return dto.StorageCleanupResult{RedisInbox: redisResult, UsageEventsDeleted: usageEventsDeleted}, err + return result, err } - // Latency 小时保留 3 天、自然日保留 365 天,并复用本轮最后一次 VACUUM。 + // Latency 小时保留 3 天、自然日保留 365 天;空闲页由后面的条件式 VACUUM 统一评估。 if err := CleanupUsageLatencyStats(db, now); err != nil { - return dto.StorageCleanupResult{RedisInbox: redisResult, UsageEventsDeleted: usageEventsDeleted}, err - } - // SQLite 删除不会立即缩小文件,维护窗口最后统一 VACUUM。 - if err := db.Exec("VACUUM").Error; err != nil { - return dto.StorageCleanupResult{RedisInbox: redisResult, UsageEventsDeleted: usageEventsDeleted}, err + return result, err } - return dto.StorageCleanupResult{RedisInbox: redisResult, UsageEventsDeleted: usageEventsDeleted}, nil -} - -// cleanupUsageEvents 严格删除早于“time.Local 当日零点向前 90 个自然日”边界的原始 usage_events。 -func cleanupUsageEvents(db *gorm.DB, now time.Time) (int64, error) { - if db == nil { - return 0, fmt.Errorf("database is nil") - } - // deleted 只在安全水位检查和 DELETE 同一事务提交后返回。 - deleted := int64(0) - // 安全检查与删除共用事务,禁止新 usage event 在两者之间插入并被误删。 - err := db.Transaction(func(tx *gorm.DB) error { - // 三类聚合任一落后时,本轮跳过 raw event 删除并等待下一次维护窗口。 - safe, err := usageEventAggregationsCaughtUp(tx) - if err != nil { - return err - } - // 跳过删除不是维护失败,Activity retention 和 VACUUM 仍可继续执行。 - if !safe { - return nil - } - // 只有安全水位已经覆盖当前最大 ID 时才计算时间保留线。 - cutoff := usageEventsCleanupCutoff(now) - // DELETE 保持严格小于 cutoff,边界时刻本身必须保留。 - result := tx.Unscoped().Where("timestamp < ?", timeutil.FormatStorageTime(cutoff)).Delete(&entities.UsageEvent{}) - if result.Error != nil { - return fmt.Errorf("cleanup usage events: %w", result.Error) - } - // 事务成功返回前记录本次真实删除行数。 - deleted = result.RowsAffected - return nil - }) - // 事务失败时不报告未提交的删除数量。 + vacuumResult, err := maybeVacuumStorage(db) + result.Vacuum = vacuumResult if err != nil { - return 0, err + return result, err } - // 返回已经提交的删除数量;聚合落后时固定为 0。 - return deleted, nil + return result, nil } func usageEventAggregationsCaughtUp(tx *gorm.DB) (bool, error) { // 当前最大 event ID 是三个全局 checkpoint 和每行 Identity cursor 的共同安全目标。 var maxEventID int64 if err := tx.Model(&entities.UsageEvent{}).Select("COALESCE(MAX(id), 0)").Scan(&maxEventID).Error; err != nil { - return false, fmt.Errorf("load usage event cleanup watermark: %w", err) + return false, fmt.Errorf("load usage event archive watermark: %w", err) } - // 空 raw event 表没有待聚合数据,可以直接执行空删除。 + // 空 hot 表没有待聚合数据,也没有需要归档的事件。 if maxEventID == 0 { return true, nil } - // 同一事务内一次读取三行,避免分别查询时观察到不一致水位或继续依赖已删除旧表。 + // 同一事务内一次读取三行,避免分别查询时观察到不一致水位。 var checkpoints []entities.UsageAggregationCheckpoint names := []entities.UsageAggregationCheckpointName{ entities.UsageAggregationCheckpointOverview, @@ -431,7 +387,7 @@ func usageEventAggregationsCaughtUp(tx *gorm.DB) (bool, error) { entities.UsageAggregationCheckpointLatency, } if err := tx.Where("name IN ?", names).Find(&checkpoints).Error; err != nil { - return false, fmt.Errorf("load usage aggregation cleanup watermarks: %w", err) + return false, fmt.Errorf("load usage aggregation archive watermarks: %w", err) } // 缺行、重复异常或任一 cursor 落后都必须保守保留 raw events。 ready := make(map[entities.UsageAggregationCheckpointName]bool, len(names)) @@ -444,16 +400,16 @@ func usageEventAggregationsCaughtUp(tx *gorm.DB) (bool, error) { } } - // Identity 没有全局 checkpoint;复用兼容追赶的同一 EXISTS 判断,但保持在当前 cleanup 事务内读取。 + // Identity 没有全局 checkpoint;复用兼容追赶的同一 EXISTS 判断,但保持在当前 archive 事务内读取。 pendingIdentity, err := hasPendingUsageIdentityAggregation(tx) if err != nil { - return false, fmt.Errorf("check identity cleanup watermark: %w", err) + return false, fmt.Errorf("check identity archive watermark: %w", err) } - // 任一 active/deleted identity 仍有 delta 时都禁止删除;不存在匹配 delta 才算安全。 + // 任一 active/deleted identity 仍有 delta 时都禁止事件离开 hot 表。 return !pendingIdentity, nil } -func usageEventsCleanupCutoff(now time.Time) time.Time { +func usageEventsArchiveCutoff(now time.Time) time.Time { localNow := now.In(time.Local) localDayStart := time.Date(localNow.Year(), localNow.Month(), localNow.Day(), 0, 0, 0, 0, time.Local) return localDayStart.AddDate(0, 0, -usageEventsRetentionDays) diff --git a/internal/repository/disk_space_unix.go b/internal/repository/disk_space_unix.go new file mode 100644 index 00000000..9d29f091 --- /dev/null +++ b/internal/repository/disk_space_unix.go @@ -0,0 +1,20 @@ +//go:build !windows + +package repository + +import ( + "fmt" + + "golang.org/x/sys/unix" +) + +func storageAvailableDiskBytes(path string) (uint64, error) { + var stats unix.Statfs_t + if err := unix.Statfs(path, &stats); err != nil { + return 0, fmt.Errorf("load available disk space for %s: %w", path, err) + } + if stats.Bsize <= 0 { + return 0, fmt.Errorf("load available disk space for %s: invalid block size %d", path, stats.Bsize) + } + return uint64(stats.Bavail) * uint64(stats.Bsize), nil +} diff --git a/internal/repository/disk_space_windows.go b/internal/repository/disk_space_windows.go new file mode 100644 index 00000000..347cea6e --- /dev/null +++ b/internal/repository/disk_space_windows.go @@ -0,0 +1,21 @@ +//go:build windows + +package repository + +import ( + "fmt" + + "golang.org/x/sys/windows" +) + +func storageAvailableDiskBytes(path string) (uint64, error) { + directoryName, err := windows.UTF16PtrFromString(path) + if err != nil { + return 0, fmt.Errorf("encode disk path %s: %w", path, err) + } + var available uint64 + if err := windows.GetDiskFreeSpaceEx(directoryName, &available, nil, nil); err != nil { + return 0, fmt.Errorf("load available disk space for %s: %w", path, err) + } + return available, nil +} diff --git a/internal/repository/dto/storage.go b/internal/repository/dto/storage.go index 4df645f2..11f4175d 100644 --- a/internal/repository/dto/storage.go +++ b/internal/repository/dto/storage.go @@ -1,7 +1,35 @@ package dto +type UsageEventArchiveStatus string + +const ( + UsageEventArchiveStatusArchived UsageEventArchiveStatus = "archived" + UsageEventArchiveStatusEmpty UsageEventArchiveStatus = "empty" + UsageEventArchiveStatusAggregationLagging UsageEventArchiveStatus = "aggregation_lagging" +) + +// UsageEventArchiveResult 记录本轮 raw event 归档数量与停止原因。 +type UsageEventArchiveResult struct { + Archived int64 + Status UsageEventArchiveStatus +} + // StorageCleanupResult 是仓储层每日清理的结果。 type StorageCleanupResult struct { - RedisInbox RedisUsageInboxCleanupResult - UsageEventsDeleted int64 + RedisInbox RedisUsageInboxCleanupResult + UsageEventsArchived int64 + UsageEventsArchiveStatus UsageEventArchiveStatus + Vacuum StorageVacuumResult +} + +// StorageVacuumResult 记录每日维护对 SQLite 空闲页的条件式整理决定。 +type StorageVacuumResult struct { + Performed bool + SkippedReason string + PageSize int64 + PageCount int64 + FreelistCount int64 + FreeBytes uint64 + FreeRatio float64 + AvailableDiskBytes uint64 } diff --git a/internal/repository/migration/20260730_usage_event_archive.go b/internal/repository/migration/20260730_usage_event_archive.go new file mode 100644 index 00000000..409f5173 --- /dev/null +++ b/internal/repository/migration/20260730_usage_event_archive.go @@ -0,0 +1,16 @@ +package migration + +import ( + "fmt" + + "cpa-usage-keeper/internal/entities" + + "gorm.io/gorm" +) + +func createUsageEventArchiveMigration(db *gorm.DB) error { + if err := db.AutoMigrate(&entities.UsageEventArchive{}); err != nil { + return fmt.Errorf("create usage event archive schema: %w", err) + } + return nil +} diff --git a/internal/repository/migration/migration.go b/internal/repository/migration/migration.go index ad512b03..09075792 100644 --- a/internal/repository/migration/migration.go +++ b/internal/repository/migration/migration.go @@ -71,6 +71,8 @@ const ( migrationUsageLatencyStats = "20260726_usage_latency_stats" // migrationAddUsageEventClientMetadata 保存 CPA 新增的客户端请求元数据,历史行保持 NULL。 migrationAddUsageEventClientMetadata = "20260729_add_usage_event_client_metadata" + // migrationCreateUsageEventArchive 创建永久冷表;运行期归档在 schema 完成后才会启动。 + migrationCreateUsageEventArchive = "20260730_create_usage_event_archive" ) type schemaMigration struct { @@ -184,6 +186,7 @@ func orderedMigrations() []databaseMigration { // Latency 回填逐页提交,外层长事务会破坏断点续跑语义。 {version: migrationUsageLatencyStats, run: usageLatencyStatsMigration, disableTransaction: true}, {version: migrationAddUsageEventClientMetadata, run: addUsageEventClientMetadataMigration}, + {version: migrationCreateUsageEventArchive, run: createUsageEventArchiveMigration}, } } diff --git a/internal/repository/migration/migration_test.go b/internal/repository/migration/migration_test.go index b5a3fd98..143eac44 100644 --- a/internal/repository/migration/migration_test.go +++ b/internal/repository/migration/migration_test.go @@ -75,6 +75,7 @@ func TestOrderedMigrationsPreservesExecutionOrder(t *testing.T) { "20260726_usage_aggregation_checkpoints", "20260726_usage_latency_stats", "20260729_add_usage_event_client_metadata", + "20260730_create_usage_event_archive", } assertStringSlicesEqual(t, want, got) } diff --git a/internal/repository/migration/test/usage_event_archive_test.go b/internal/repository/migration/test/usage_event_archive_test.go new file mode 100644 index 00000000..bfc40036 --- /dev/null +++ b/internal/repository/migration/test/usage_event_archive_test.go @@ -0,0 +1,141 @@ +package test + +import ( + "fmt" + "path/filepath" + "strings" + "testing" + "time" + + "cpa-usage-keeper/internal/entities" + "cpa-usage-keeper/internal/repository/migration" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +const usageEventArchiveMigrationVersion = "20260730_create_usage_event_archive" + +func TestUsageEventArchiveMigrationCreatesColdTableOnExistingDatabase(t *testing.T) { + db, err := gorm.Open(sqlite.Open(filepath.Join(t.TempDir(), "existing.db")), &gorm.Config{}) + if err != nil { + t.Fatalf("open existing database: %v", err) + } + closeMigrationTestDatabase(t, db) + if err := db.AutoMigrate(&entities.UsageEvent{}); err != nil { + t.Fatalf("create existing usage_events: %v", err) + } + if err := migration.MarkAllAsApplied(db); err != nil { + t.Fatalf("mark historical migrations applied: %v", err) + } + if err := db.Table("schema_migrations").Where("version = ?", usageEventArchiveMigrationVersion).Delete(nil).Error; err != nil { + t.Fatalf("make archive migration pending: %v", err) + } + + if err := migration.Run(db); err != nil { + t.Fatalf("Run returned error: %v", err) + } + if !db.Migrator().HasTable("usage_events_archive") { + t.Fatal("expected archive migration to create usage_events_archive") + } + var count int64 + if err := db.Table("schema_migrations").Where("version = ?", usageEventArchiveMigrationVersion).Count(&count).Error; err != nil { + t.Fatalf("count archive migration: %v", err) + } + if count != 1 { + t.Fatalf("expected archive migration recorded once, got %d", count) + } +} + +func TestUsageEventReplayPagesMergeArchiveAndHotByGlobalID(t *testing.T) { + db, err := gorm.Open(sqlite.Open(filepath.Join(t.TempDir(), "replay.db")), &gorm.Config{}) + if err != nil { + t.Fatalf("open replay database: %v", err) + } + closeMigrationTestDatabase(t, db) + if err := db.AutoMigrate(&entities.UsageEvent{}, &entities.UsageEventArchive{}); err != nil { + t.Fatalf("create replay schema: %v", err) + } + now := time.Date(2026, 7, 30, 4, 30, 0, 0, time.Local) + archiveRows := []entities.UsageEventArchive{ + {ID: 1, EventKey: "archive-1", AuthType: "oauth", AuthIndex: "auth-file-1", Source: "archive-source", Timestamp: now.AddDate(0, 0, -120), TotalTokens: 1}, + {ID: 4, EventKey: "late-archive-4", Timestamp: now.AddDate(0, 0, -100), TotalTokens: 4}, + } + hotRows := []entities.UsageEvent{ + {ID: 2, EventKey: "hot-2", AuthType: "apikey", AuthIndex: "provider-2", Source: "hot-source", Timestamp: now.AddDate(0, 0, -10), TotalTokens: 2}, + {ID: 3, EventKey: "hot-3", Timestamp: now.AddDate(0, 0, -9), TotalTokens: 3}, + {ID: 5, EventKey: "hot-5", Timestamp: now, TotalTokens: 5}, + } + if err := db.Create(&archiveRows).Error; err != nil { + t.Fatalf("seed archive replay rows: %v", err) + } + if err := db.Create(&hotRows).Error; err != nil { + t.Fatalf("seed hot replay rows: %v", err) + } + + targetID, err := migration.LoadUsageAggregationReplayTargetEventID(db) + if err != nil { + t.Fatalf("load replay target: %v", err) + } + if targetID != 5 { + t.Fatalf("expected replay target 5, got %d", targetID) + } + var gotIDs []int64 + gotEvents := make(map[int64]entities.UsageEvent) + afterID := int64(0) + for { + page, err := migration.LoadUsageAggregationReplayEventPage(db, afterID, targetID, 2) + if err != nil { + t.Fatalf("load replay page after %d: %v", afterID, err) + } + if len(page) == 0 { + break + } + for _, event := range page { + gotIDs = append(gotIDs, event.ID) + gotEvents[event.ID] = event + } + afterID = page[len(page)-1].ID + } + if fmt.Sprint(gotIDs) != fmt.Sprint([]int64{1, 2, 3, 4, 5}) { + t.Fatalf("expected globally ordered replay IDs, got %v", gotIDs) + } + if event := gotEvents[1]; event.AuthType != "oauth" || event.AuthIndex != "auth-file-1" || event.Source != "archive-source" { + t.Fatalf("expected archive replay to preserve identity fields, got %+v", event) + } + if event := gotEvents[2]; event.AuthType != "apikey" || event.AuthIndex != "provider-2" || event.Source != "hot-source" { + t.Fatalf("expected hot replay to preserve identity fields, got %+v", event) + } +} + +func TestUsageEventReplayQueryUsesPrimaryKeyMerge(t *testing.T) { + db, err := gorm.Open(sqlite.Open(filepath.Join(t.TempDir(), "plan.db")), &gorm.Config{}) + if err != nil { + t.Fatalf("open replay plan database: %v", err) + } + closeMigrationTestDatabase(t, db) + if err := db.AutoMigrate(&entities.UsageEvent{}, &entities.UsageEventArchive{}); err != nil { + t.Fatalf("create replay plan schema: %v", err) + } + var rows []struct { + Detail string `gorm:"column:detail"` + } + if err := db.Raw(`EXPLAIN QUERY PLAN + SELECT id FROM ( + SELECT id FROM usage_events_archive WHERE id > ? AND id <= ? + UNION ALL + SELECT id FROM usage_events WHERE id > ? AND id <= ? + ) ORDER BY id ASC LIMIT ?`, 0, 100, 0, 100, 10).Scan(&rows).Error; err != nil { + t.Fatalf("explain replay query: %v", err) + } + details := make([]string, 0, len(rows)) + for _, row := range rows { + details = append(details, row.Detail) + } + plan := strings.Join(details, "\n") + for _, want := range []string{"MERGE (UNION ALL)", "usage_events_archive USING INTEGER PRIMARY KEY", "usage_events USING INTEGER PRIMARY KEY"} { + if !strings.Contains(plan, want) { + t.Fatalf("expected replay query plan to contain %q, got:\n%s", want, plan) + } + } +} diff --git a/internal/repository/migration/test/usage_latency_stats_test.go b/internal/repository/migration/test/usage_latency_stats_test.go index bd63d420..4907b153 100644 --- a/internal/repository/migration/test/usage_latency_stats_test.go +++ b/internal/repository/migration/test/usage_latency_stats_test.go @@ -113,8 +113,8 @@ func TestUsageLatencyStatsMigrationResumesAfterCommittedPage(t *testing.T) { } } -func TestUsageLatencyStatsMigrationEnablesUsageEventCleanup(t *testing.T) { - // Latency migration 只有推进到当前最大事件 ID,才应与另外两类水位共同放行 raw cleanup。 +func TestUsageLatencyStatsMigrationEnablesUsageEventArchive(t *testing.T) { + // Latency migration 只有推进到当前最大事件 ID,才应与另外两类水位共同放行 raw archive。 previousLocal := time.Local time.Local = time.UTC t.Cleanup(func() { time.Local = previousLocal }) @@ -140,19 +140,19 @@ func TestUsageLatencyStatsMigrationEnablesUsageEventCleanup(t *testing.T) { t.Fatalf("seed caught-up overview/activity checkpoints: %v", err) } - result, err := repository.CleanupStorage(db, now, repository.CleanupStorageOptions{CleanupUsageEvents: true}) + result, err := repository.CleanupStorage(db, now) if err != nil { - t.Fatalf("cleanup after latency migration: %v", err) + t.Fatalf("archive after latency migration: %v", err) } - if result.UsageEventsDeleted != 2 { - t.Fatalf("expected latency migration to release two old events, got %+v", result) + if result.UsageEventsArchived != 2 { + t.Fatalf("expected latency migration to allow archiving two old events, got %+v", result) } } func prepareUsageLatencyMigrationDatabase(t *testing.T, db *gorm.DB, events []entities.UsageEvent) { // 只创建 migration 的真实前置表,避免无关历史 migration 参与结果。 t.Helper() - if err := db.AutoMigrate(&entities.UsageEvent{}, &entities.UsageAggregationCheckpoint{}); err != nil { + if err := db.AutoMigrate(&entities.UsageEvent{}, &entities.UsageEventArchive{}, &entities.UsageAggregationCheckpoint{}); err != nil { t.Fatalf("create latency migration schema: %v", err) } if err := db.CreateInBatches(&events, 200).Error; err != nil { diff --git a/internal/repository/migration/usage_event_replay.go b/internal/repository/migration/usage_event_replay.go new file mode 100644 index 00000000..c72403f1 --- /dev/null +++ b/internal/repository/migration/usage_event_replay.go @@ -0,0 +1,62 @@ +package migration + +import ( + "fmt" + + "cpa-usage-keeper/internal/entities" + + "gorm.io/gorm" +) + +const usageAggregationReplayPageLimit = 1000 + +// LoadUsageAggregationReplayTargetEventID 固定 archive 与 hot 在 migration 启动时的全局最大 ID。 +func LoadUsageAggregationReplayTargetEventID(db *gorm.DB) (int64, error) { + if db == nil { + return 0, fmt.Errorf("database is nil") + } + var targetID int64 + if err := db.Raw(` + SELECT MAX( + COALESCE((SELECT MAX(id) FROM usage_events_archive), 0), + COALESCE((SELECT MAX(id) FROM usage_events), 0) + )`).Scan(&targetID).Error; err != nil { + return 0, fmt.Errorf("load usage event replay target: %w", err) + } + return targetID, nil +} + +// LoadUsageAggregationReplayEventPage 只供启动 migration 按全局 ID 顺序重放 archive 与 hot 事件。 +func LoadUsageAggregationReplayEventPage(db *gorm.DB, afterID, targetID int64, limit int) ([]entities.UsageEvent, error) { + if db == nil { + return nil, fmt.Errorf("database is nil") + } + if afterID < 0 || targetID < 0 { + return nil, fmt.Errorf("usage event replay bounds must be non-negative: after=%d target=%d", afterID, targetID) + } + if targetID <= afterID { + return []entities.UsageEvent{}, nil + } + if limit <= 0 { + return nil, fmt.Errorf("usage event replay page limit must be positive") + } + if limit > usageAggregationReplayPageLimit { + limit = usageAggregationReplayPageLimit + } + + // 两个分支都按 INTEGER PRIMARY KEY 范围读取,SQLite 会用 MERGE UNION 恢复全局入库顺序。 + columns := entities.UsageEventStorageColumns + query := fmt.Sprintf(` + SELECT %s FROM ( + SELECT %s FROM usage_events_archive WHERE id > ? AND id <= ? + UNION ALL + SELECT %s FROM usage_events WHERE id > ? AND id <= ? + ) AS usage_event_replay + ORDER BY id ASC + LIMIT ?`, columns, columns, columns) + var events []entities.UsageEvent + if err := db.Raw(query, afterID, targetID, afterID, targetID, limit).Scan(&events).Error; err != nil { + return nil, fmt.Errorf("load usage event replay page: %w", err) + } + return events, nil +} diff --git a/internal/repository/storage_vacuum.go b/internal/repository/storage_vacuum.go new file mode 100644 index 00000000..1cee049f --- /dev/null +++ b/internal/repository/storage_vacuum.go @@ -0,0 +1,129 @@ +package repository + +import ( + "fmt" + "path/filepath" + + "cpa-usage-keeper/internal/repository/dto" + + "gorm.io/gorm" + "gorm.io/plugin/dbresolver" +) + +const ( + storageVacuumMinFreeBytes = uint64(1 << 30) + storageVacuumMinFreeRatio = 0.20 +) + +// StorageVacuumStats 是 SQLite 主库当前页使用情况。 +type StorageVacuumStats struct { + PageSize int64 + PageCount int64 + FreelistCount int64 +} + +func (s StorageVacuumStats) freeBytes() uint64 { + if s.PageSize <= 0 || s.FreelistCount <= 0 { + return 0 + } + return uint64(s.PageSize) * uint64(s.FreelistCount) +} + +func (s StorageVacuumStats) databaseBytes() uint64 { + if s.PageSize <= 0 || s.PageCount <= 0 { + return 0 + } + return uint64(s.PageSize) * uint64(s.PageCount) +} + +func (s StorageVacuumStats) freeRatio() float64 { + if s.PageCount <= 0 || s.FreelistCount <= 0 { + return 0 + } + return float64(s.FreelistCount) / float64(s.PageCount) +} + +// StorageVacuumRequired 只根据当前空闲页和临时磁盘空间决定,不保存或检查上次整理时间。 +func StorageVacuumRequired(stats StorageVacuumStats, availableDiskBytes uint64) bool { + if stats.freeBytes() < storageVacuumMinFreeBytes || stats.freeRatio() < storageVacuumMinFreeRatio { + return false + } + // 完整 VACUUM 可能同时保留临时库与 journal/WAL,按两份原库再加 1 GiB 余量保守判断。 + databaseBytes := stats.databaseBytes() + if databaseBytes > (^uint64(0)-storageVacuumMinFreeBytes)/2 { + return false + } + requiredDiskBytes := databaseBytes*2 + storageVacuumMinFreeBytes + return availableDiskBytes >= requiredDiskBytes +} + +func maybeVacuumStorage(db *gorm.DB) (dto.StorageVacuumResult, error) { + stats, databasePath, err := loadStorageVacuumStats(db) + if err != nil { + return dto.StorageVacuumResult{}, err + } + result := dto.StorageVacuumResult{ + PageSize: stats.PageSize, + PageCount: stats.PageCount, + FreelistCount: stats.FreelistCount, + FreeBytes: stats.freeBytes(), + FreeRatio: stats.freeRatio(), + } + if stats.freeBytes() < storageVacuumMinFreeBytes || stats.freeRatio() < storageVacuumMinFreeRatio { + result.SkippedReason = "threshold_not_met" + return result, nil + } + if databasePath == "" { + result.SkippedReason = "database_path_unavailable" + return result, nil + } + availableDiskBytes, err := storageAvailableDiskBytes(filepath.Dir(databasePath)) + if err != nil { + result.SkippedReason = "disk_space_unavailable" + return result, nil + } + result.AvailableDiskBytes = availableDiskBytes + if !StorageVacuumRequired(stats, availableDiskBytes) { + result.SkippedReason = "disk_space_insufficient" + return result, nil + } + if err := Vacuum(db.Clauses(dbresolver.Write)); err != nil { + return result, err + } + result.Performed = true + return result, nil +} + +func loadStorageVacuumStats(db *gorm.DB) (StorageVacuumStats, string, error) { + if db == nil { + return StorageVacuumStats{}, "", fmt.Errorf("database is nil") + } + writeDB := db.Clauses(dbresolver.Write) + stats := StorageVacuumStats{} + for _, query := range []struct { + sql string + dest *int64 + name string + }{ + {sql: "PRAGMA page_size", dest: &stats.PageSize, name: "page_size"}, + {sql: "PRAGMA page_count", dest: &stats.PageCount, name: "page_count"}, + {sql: "PRAGMA freelist_count", dest: &stats.FreelistCount, name: "freelist_count"}, + } { + if err := writeDB.Raw(query.sql).Scan(query.dest).Error; err != nil { + return StorageVacuumStats{}, "", fmt.Errorf("load sqlite %s: %w", query.name, err) + } + } + var databases []struct { + Name string `gorm:"column:name"` + File string `gorm:"column:file"` + } + if err := writeDB.Raw("PRAGMA database_list").Scan(&databases).Error; err != nil { + return StorageVacuumStats{}, "", fmt.Errorf("load sqlite database path: %w", err) + } + for _, database := range databases { + if database.Name == "main" { + return stats, database.File, nil + } + } + return stats, "", nil +} diff --git a/internal/repository/test/db_test.go b/internal/repository/test/db_test.go index d72b2fc9..63589aa3 100644 --- a/internal/repository/test/db_test.go +++ b/internal/repository/test/db_test.go @@ -709,7 +709,7 @@ func TestDatabaseTimeFieldsUseProjectTimezoneRFC3339Nano(t *testing.T) { } } -func TestCleanupStorageCleansRedisInboxAndVacuums(t *testing.T) { +func TestCleanupStorageCleansRedisInboxAndAppliesVacuumPolicy(t *testing.T) { previousLocal := time.Local location, err := time.LoadLocation("Asia/Shanghai") if err != nil { @@ -783,23 +783,23 @@ func TestCleanupStorageRetainsNinetyLocalDays(t *testing.T) { }); err != nil { t.Fatalf("InsertUsageEvents returned error: %v", err) } - // 准备:只有 Overview 与 Activity 都追平后,旧 raw events 才具备删除安全水位。 + // 准备:只有 Overview 与 Activity 都追平后,旧 raw events 才具备归档安全水位。 if err := repository.AggregateUsageOverviewStats(context.Background(), db, now); err != nil { t.Fatalf("aggregate overview before cleanup: %v", err) } if err := repository.AggregateUsageActivityStats(context.Background(), db, now); err != nil { t.Fatalf("aggregate activity before cleanup: %v", err) } - // Task 1 尚未实现 Latency 回填,测试显式推进第三行后才满足最终 cleanup 门禁。 + // 测试显式推进第三个全局 checkpoint 后,才满足 raw event 归档门禁。 seedCaughtUpLatencyCheckpoint(t, db, now) - // 执行:在两个全局 checkpoint 已追平且没有 identity delta 时运行清理。 - result, err := repository.CleanupStorage(db, now, repository.CleanupStorageOptions{CleanupUsageEvents: true}) + // 执行:三个全局 checkpoint 已追平且没有 identity delta 时运行维护。 + result, err := repository.CleanupStorage(db, now) if err != nil { t.Fatalf("CleanupStorage returned error: %v", err) } - if result.UsageEventsDeleted != 1 { - t.Fatalf("expected one old usage event to be deleted, got %+v", result) + if result.UsageEventsArchived != 1 { + t.Fatalf("expected one old usage event to be archived, got %+v", result) } var remainingKeys []string @@ -841,15 +841,15 @@ func TestCleanupStorageUsesLocalCalendarDaysAcrossDST(t *testing.T) { if err := repository.AggregateUsageActivityStats(context.Background(), db, now); err != nil { t.Fatalf("aggregate activity before cleanup: %v", err) } - // Latency 行必须和另外两类一起覆盖当前最大 ID,raw event 才可删除。 + // Latency 行必须和另外两类一起覆盖当前最大 ID,raw event 才可离开 hot 表。 seedCaughtUpLatencyCheckpoint(t, db, now) - result, err := repository.CleanupStorage(db, now, repository.CleanupStorageOptions{CleanupUsageEvents: true}) + result, err := repository.CleanupStorage(db, now) if err != nil { t.Fatalf("CleanupStorage returned error: %v", err) } - if result.UsageEventsDeleted != 1 { - t.Fatalf("expected only the event before the local calendar cutoff to be deleted, got %+v", result) + if result.UsageEventsArchived != 1 { + t.Fatalf("expected only the event before the local calendar cutoff to be archived, got %+v", result) } var remainingKeys []string @@ -874,15 +874,18 @@ func TestCleanupStorageDefersUsageEventsUntilOverviewAndActivityCatchUp(t *testi t.Fatalf("InsertUsageEvents returned error: %v", err) } - // 执行:全局 checkpoint 尚未创建时尝试清理 raw events。 - result, err := repository.CleanupStorage(db, now, repository.CleanupStorageOptions{CleanupUsageEvents: true}) + // 执行:全局 checkpoint 尚未创建时尝试归档 raw events。 + result, err := repository.CleanupStorage(db, now) if err != nil { t.Fatalf("CleanupStorage returned error: %v", err) } // 断言:清理必须让路,避免异步聚合再也读取不到两条事件。 - if result.UsageEventsDeleted != 0 { + if result.UsageEventsArchived != 0 { t.Fatalf("expected pending overview/activity events to remain, got %+v", result) } + if result.UsageEventsArchiveStatus != dto.UsageEventArchiveStatusAggregationLagging { + t.Fatalf("expected aggregation lagging archive status, got %+v", result) + } var remainingCount int64 if err := db.Model(&entities.UsageEvent{}).Count(&remainingCount).Error; err != nil { @@ -924,11 +927,11 @@ func TestCleanupStorageDefersUsageEventsUntilLatencyCatchUp(t *testing.T) { t.Fatalf("seed lagging latency checkpoint: %v", err) } - result, err := repository.CleanupStorage(db, now, repository.CleanupStorageOptions{CleanupUsageEvents: true}) + result, err := repository.CleanupStorage(db, now) if err != nil { t.Fatalf("CleanupStorage returned error: %v", err) } - if result.UsageEventsDeleted != 0 { + if result.UsageEventsArchived != 0 { t.Fatalf("expected lagging latency aggregation to retain raw events, got %+v", result) } var remainingCount int64 @@ -976,12 +979,12 @@ func TestCleanupStorageDefersUsageEventsUntilIdentityCatchUp(t *testing.T) { } // 执行:identity cursor 尚未越过第二条匹配事件时运行清理。 - result, err := repository.CleanupStorage(db, now, repository.CleanupStorageOptions{CleanupUsageEvents: true}) + result, err := repository.CleanupStorage(db, now) if err != nil { t.Fatalf("CleanupStorage returned error: %v", err) } // 断言:两条 raw events 都必须保留,等待 Identity 下一批安全累计。 - if result.UsageEventsDeleted != 0 { + if result.UsageEventsArchived != 0 { t.Fatalf("expected pending identity events to remain, got %+v", result) } diff --git a/internal/repository/test/storage_vacuum_test.go b/internal/repository/test/storage_vacuum_test.go new file mode 100644 index 00000000..acd5f6d1 --- /dev/null +++ b/internal/repository/test/storage_vacuum_test.go @@ -0,0 +1,48 @@ +package test + +import ( + "testing" + + "cpa-usage-keeper/internal/repository" +) + +func TestStorageVacuumRequiredUsesFreeBytesRatioAndDiskSpace(t *testing.T) { + const ( + pageSize = int64(4096) + oneGiB = uint64(1 << 30) + ) + tests := []struct { + name string + pageCount int64 + freelistCount int64 + available uint64 + want bool + }{ + {name: "all conditions", pageCount: 1_500_000, freelistCount: 300_000, available: 13 * oneGiB, want: true}, + {name: "free bytes below one GiB", pageCount: 1_000_000, freelistCount: 250_000, available: 8 * oneGiB}, + {name: "free ratio below twenty percent", pageCount: 1_500_000, freelistCount: 299_999, available: 13 * oneGiB}, + {name: "temporary disk space insufficient", pageCount: 1_500_000, freelistCount: 300_000, available: 12 * oneGiB}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + stats := repository.StorageVacuumStats{PageSize: pageSize, PageCount: test.pageCount, FreelistCount: test.freelistCount} + if got := repository.StorageVacuumRequired(stats, test.available); got != test.want { + t.Fatalf("StorageVacuumRequired() = %v, want %v for %+v available=%d", got, test.want, stats, test.available) + } + }) + } +} + +func TestStorageVacuumRequiredHasNoTimeIntervalCondition(t *testing.T) { + stats := repository.StorageVacuumStats{PageSize: 4096, PageCount: 1_500_000, FreelistCount: 300_000} + if !repository.StorageVacuumRequired(stats, 13<<30) { + t.Fatal("expected current page and disk conditions alone to allow vacuum") + } +} + +func TestStorageVacuumRequiredRejectsOverflowingSpaceEstimate(t *testing.T) { + stats := repository.StorageVacuumStats{PageSize: 1 << 32, PageCount: 1 << 31, FreelistCount: 1 << 30} + if repository.StorageVacuumRequired(stats, ^uint64(0)) { + t.Fatal("expected overflowing temporary space estimate to skip vacuum") + } +} diff --git a/internal/repository/test/usage_event_archive_test.go b/internal/repository/test/usage_event_archive_test.go new file mode 100644 index 00000000..120796c1 --- /dev/null +++ b/internal/repository/test/usage_event_archive_test.go @@ -0,0 +1,226 @@ +package test + +import ( + "context" + "database/sql" + "fmt" + "strings" + "testing" + "time" + + "cpa-usage-keeper/internal/entities" + "cpa-usage-keeper/internal/repository" + + "gorm.io/gorm" +) + +type sqliteTableColumn struct { + Name string `gorm:"column:name"` + Type string `gorm:"column:type"` + NotNull int `gorm:"column:notnull"` + DefaultSQL sql.NullString `gorm:"column:dflt_value"` + PrimaryKey int `gorm:"column:pk"` +} + +func TestUsageEventArchiveSchemaMatchesHotColumnsWithoutSecondaryIndexes(t *testing.T) { + db := openTestDatabase(t) + + if !db.Migrator().HasTable("usage_events_archive") { + t.Fatal("expected fresh database to create usage_events_archive") + } + hotColumns := loadSQLiteTableColumns(t, db, "usage_events") + archiveColumns := loadSQLiteTableColumns(t, db, "usage_events_archive") + if fmt.Sprint(hotColumns) != fmt.Sprint(archiveColumns) { + t.Fatalf("usage event archive schema mismatch:\n hot=%+v\n archive=%+v", hotColumns, archiveColumns) + } + storageColumnNames := strings.Split(strings.ReplaceAll(entities.UsageEventStorageColumns, " ", ""), ",") + hotColumnNames := make([]string, 0, len(hotColumns)) + for _, column := range hotColumns { + hotColumnNames = append(hotColumnNames, column.Name) + } + if fmt.Sprint(storageColumnNames) != fmt.Sprint(hotColumnNames) { + t.Fatalf("usage event archive copy columns mismatch:\n copy=%v\n schema=%v", storageColumnNames, hotColumnNames) + } + + var archiveSQL string + if err := db.Raw("SELECT sql FROM sqlite_master WHERE type = 'table' AND name = ?", "usage_events_archive").Scan(&archiveSQL).Error; err != nil { + t.Fatalf("load usage_events_archive schema SQL: %v", err) + } + if strings.Contains(strings.ToUpper(archiveSQL), "AUTOINCREMENT") { + t.Fatalf("expected archive primary key not to use AUTOINCREMENT, got %s", archiveSQL) + } + + var archiveIndexes []string + if err := db.Raw("SELECT name FROM sqlite_master WHERE type = 'index' AND tbl_name = ? ORDER BY name", "usage_events_archive").Scan(&archiveIndexes).Error; err != nil { + t.Fatalf("load usage_events_archive indexes: %v", err) + } + if len(archiveIndexes) != 0 { + t.Fatalf("expected archive to have no secondary indexes, got %v", archiveIndexes) + } +} + +func TestArchiveExpiredUsageEventsPreservesOriginalRowAndHotSequence(t *testing.T) { + db := openTestDatabase(t) + now := time.Date(2026, 7, 30, 4, 30, 0, 0, time.Local) + generate := false + clientIP := "203.0.113.10" + ttft := int64(321) + events := []entities.UsageEvent{ + { + EventKey: "archive-me", APIGroupKey: "group-a", Provider: "openai", Endpoint: "/v1/responses", + AuthType: "oauth", RequestID: "request-a", ClientIP: &clientIP, Model: "gpt-5", ReasoningEffort: "high", + ServiceTier: "priority", ResponseServiceTier: "priority", ExecutorType: "codex", Timestamp: now.AddDate(0, 0, -91), + Source: "auth-a", AuthIndex: "auth-a", Failed: true, Generate: &generate, LatencyMS: 999, TTFTMS: &ttft, + InputTokens: 10, OutputTokens: 20, ReasoningTokens: 5, CachedTokens: 4, CacheReadTokens: 3, CacheCreationTokens: 2, TotalTokens: 35, + }, + {EventKey: "recent", Model: "gpt-5", Timestamp: now.Add(-time.Hour), TotalTokens: 1}, + } + if _, _, err := repository.InsertUsageEvents(db, events); err != nil { + t.Fatalf("seed usage events: %v", err) + } + var original entities.UsageEvent + if err := db.Where("event_key = ?", "archive-me").Take(&original).Error; err != nil { + t.Fatalf("load original usage event: %v", err) + } + previousMaxID := eventsMaxID(t, db) + seedArchiveCaughtUpCheckpoints(t, db, previousMaxID, now) + + result, err := repository.CleanupStorage(db, now) + if err != nil { + t.Fatalf("CleanupStorage returned error: %v", err) + } + if result.UsageEventsArchived != 1 { + t.Fatalf("expected one archived usage event, got %+v", result) + } + + var archived entities.UsageEventArchive + if err := db.Where("id = ?", original.ID).Take(&archived).Error; err != nil { + t.Fatalf("load archived usage event: %v", err) + } + if archived.ID != original.ID || archived.EventKey != original.EventKey || archived.RequestID != original.RequestID || archived.TotalTokens != original.TotalTokens { + t.Fatalf("archive row did not preserve original values: original=%+v archive=%+v", original, archived) + } + var oldHotCount int64 + if err := db.Model(&entities.UsageEvent{}).Where("id = ?", original.ID).Count(&oldHotCount).Error; err != nil { + t.Fatalf("count archived hot event: %v", err) + } + if oldHotCount != 0 { + t.Fatalf("expected archived event to leave hot table, got %d", oldHotCount) + } + + if _, _, err := repository.InsertUsageEvents(db, []entities.UsageEvent{{EventKey: "after-archive", Timestamp: now, TotalTokens: 2}}); err != nil { + t.Fatalf("insert event after archive: %v", err) + } + var insertedAfter entities.UsageEvent + if err := db.Where("event_key = ?", "after-archive").Take(&insertedAfter).Error; err != nil { + t.Fatalf("load event inserted after archive: %v", err) + } + if insertedAfter.ID <= previousMaxID { + t.Fatalf("expected hot AUTOINCREMENT sequence to continue after archive, got id %d", insertedAfter.ID) + } +} + +func TestArchiveExpiredUsageEventsRollsBackOnArchivePrimaryKeyConflict(t *testing.T) { + db := openTestDatabase(t) + now := time.Date(2026, 7, 30, 4, 30, 0, 0, time.Local) + if _, _, err := repository.InsertUsageEvents(db, []entities.UsageEvent{{EventKey: "hot-old", Timestamp: now.AddDate(0, 0, -91), TotalTokens: 9}}); err != nil { + t.Fatalf("seed old usage event: %v", err) + } + var hot entities.UsageEvent + if err := db.Where("event_key = ?", "hot-old").Take(&hot).Error; err != nil { + t.Fatalf("load hot usage event: %v", err) + } + seedArchiveCaughtUpCheckpoints(t, db, hot.ID, now) + if err := db.Create(&entities.UsageEventArchive{ID: hot.ID, EventKey: "existing-archive", Timestamp: hot.Timestamp, TotalTokens: 1}).Error; err != nil { + t.Fatalf("seed conflicting archive row: %v", err) + } + + if _, err := repository.ArchiveExpiredUsageEvents(context.Background(), db, now); err == nil { + t.Fatal("expected archive primary key conflict to fail") + } + var hotCount int64 + if err := db.Model(&entities.UsageEvent{}).Where("id = ?", hot.ID).Count(&hotCount).Error; err != nil { + t.Fatalf("count hot row after conflict: %v", err) + } + if hotCount != 1 { + t.Fatalf("expected hot row to remain after archive conflict, got %d", hotCount) + } + var archiveKey string + if err := db.Model(&entities.UsageEventArchive{}).Where("id = ?", hot.ID).Pluck("event_key", &archiveKey).Error; err != nil { + t.Fatalf("load archive row after conflict: %v", err) + } + if archiveKey != "existing-archive" { + t.Fatalf("expected existing archive row to remain unchanged, got %q", archiveKey) + } +} + +func TestArchiveExpiredUsageEventsProcessesMoreThanOneBatch(t *testing.T) { + db := openTestDatabase(t) + now := time.Date(2026, 7, 30, 4, 30, 0, 0, time.Local) + events := make([]entities.UsageEvent, 5001) + for index := range events { + events[index] = entities.UsageEvent{ + EventKey: fmt.Sprintf("old-%05d", index), + Timestamp: now.AddDate(0, 0, -91).Add(time.Duration(index) * time.Nanosecond), + TotalTokens: int64(index + 1), + } + } + if _, _, err := repository.InsertUsageEvents(db, events); err != nil { + t.Fatalf("seed multi-batch usage events: %v", err) + } + seedArchiveCaughtUpCheckpoints(t, db, eventsMaxID(t, db), now) + + archiveResult, err := repository.ArchiveExpiredUsageEvents(context.Background(), db, now) + if err != nil { + t.Fatalf("ArchiveExpiredUsageEvents returned error: %v", err) + } + if archiveResult.Archived != int64(len(events)) { + t.Fatalf("expected %d archived rows, got %+v", len(events), archiveResult) + } + var hotCount int64 + if err := db.Model(&entities.UsageEvent{}).Count(&hotCount).Error; err != nil { + t.Fatalf("count hot rows after multi-batch archive: %v", err) + } + if hotCount != 0 { + t.Fatalf("expected no hot rows after multi-batch archive, got %d", hotCount) + } + var archiveCount int64 + if err := db.Model(&entities.UsageEventArchive{}).Count(&archiveCount).Error; err != nil { + t.Fatalf("count archive rows after multi-batch archive: %v", err) + } + if archiveCount != int64(len(events)) { + t.Fatalf("expected %d archive rows, got %d", len(events), archiveCount) + } +} + +func loadSQLiteTableColumns(t *testing.T, db interface { + Raw(string, ...any) *gorm.DB +}, table string) []sqliteTableColumn { + t.Helper() + var columns []sqliteTableColumn + if err := db.Raw("PRAGMA table_info(" + table + ")").Scan(&columns).Error; err != nil { + t.Fatalf("load %s columns: %v", table, err) + } + return columns +} + +func eventsMaxID(t *testing.T, db *gorm.DB) int64 { + t.Helper() + var maxID int64 + if err := db.Model(&entities.UsageEvent{}).Select("COALESCE(MAX(id), 0)").Scan(&maxID).Error; err != nil { + t.Fatalf("load max usage event id: %v", err) + } + return maxID +} + +func seedArchiveCaughtUpCheckpoints(t *testing.T, db *gorm.DB, maxID int64, now time.Time) { + t.Helper() + rows := []entities.UsageAggregationCheckpoint{ + {Name: entities.UsageAggregationCheckpointOverview, LastAggregatedUsageEventID: maxID, CreatedAt: now, UpdatedAt: now}, + {Name: entities.UsageAggregationCheckpointActivity, LastAggregatedUsageEventID: maxID, CreatedAt: now, UpdatedAt: now}, + {Name: entities.UsageAggregationCheckpointLatency, LastAggregatedUsageEventID: maxID, CreatedAt: now, UpdatedAt: now}, + } + if err := db.Create(&rows).Error; err != nil { + t.Fatalf("seed archive cleanup checkpoints: %v", err) + } +} diff --git a/internal/repository/test/usage_latency_stats_test.go b/internal/repository/test/usage_latency_stats_test.go index 5e298b8a..7a8826dd 100644 --- a/internal/repository/test/usage_latency_stats_test.go +++ b/internal/repository/test/usage_latency_stats_test.go @@ -56,7 +56,7 @@ func TestAggregateUsageLatencyStatsKeepsHoursForThreeDaysAndDaysForThreeHundredS } func TestCleanupStorageRemovesExpiredLatencyHourAndDayRows(t *testing.T) { - // 每日维护复用现有 cleanup/VACUUM,只删除三天前小时行和365天窗口外的自然日行。 + // 每日维护保留原有 Latency retention,并由统一条件式 VACUUM 评估空闲页。 db := openTestDatabase(t) now := time.Date(2026, 7, 26, 12, 30, 0, 0, time.Local) templates := buildLatencyRowsForTest(t, []entities.UsageEvent{validLatencyEvent(1, "template", now, 100, 900)}) diff --git a/internal/repository/usage_event_archive.go b/internal/repository/usage_event_archive.go new file mode 100644 index 00000000..eb36781c --- /dev/null +++ b/internal/repository/usage_event_archive.go @@ -0,0 +1,116 @@ +package repository + +import ( + "context" + "fmt" + "time" + + "cpa-usage-keeper/internal/entities" + "cpa-usage-keeper/internal/repository/dto" + "cpa-usage-keeper/internal/timeutil" + + "gorm.io/gorm" +) + +// 每个候选 ID 在单条 IN 查询中占一个绑定变量,复用 repository 的 SQLite 保守变量预算。 +const usageEventArchiveBatchSize = sqliteVariableLimit + +// ArchiveExpiredUsageEvents 分批把超过 hot 保留线且已完成派生聚合的事件原子移动到冷表。 +func ArchiveExpiredUsageEvents(ctx context.Context, db *gorm.DB, now time.Time) (dto.UsageEventArchiveResult, error) { + if db == nil { + return dto.UsageEventArchiveResult{}, fmt.Errorf("database is nil") + } + if ctx == nil { + ctx = context.Background() + } + result := dto.UsageEventArchiveResult{Status: dto.UsageEventArchiveStatusEmpty} + for { + if err := ctx.Err(); err != nil { + return result, err + } + archived, status, err := archiveExpiredUsageEventsBatch(ctx, db, now) + if err != nil { + return result, err + } + result.Archived += archived + if status == dto.UsageEventArchiveStatusAggregationLagging { + result.Status = status + return result, nil + } + if archived == 0 { + return result, nil + } + result.Status = dto.UsageEventArchiveStatusArchived + } +} + +func archiveExpiredUsageEventsBatch(ctx context.Context, db *gorm.DB, now time.Time) (int64, dto.UsageEventArchiveStatus, error) { + archived := int64(0) + status := dto.UsageEventArchiveStatusEmpty + err := db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + cutoff := usageEventsArchiveCutoff(now) + var ids []int64 + if err := tx.Model(&entities.UsageEvent{}). + Select("id"). + Where("timestamp < ?", timeutil.FormatStorageTime(cutoff)). + Order("timestamp asc, id asc"). + Limit(usageEventArchiveBatchSize). + Pluck("id", &ids).Error; err != nil { + return fmt.Errorf("load usage events archive batch: %w", err) + } + if len(ids) == 0 { + return nil + } + + // 只有确有过期候选时才检查保守门禁,避免没有归档积压时误报水位阻塞。 + safe, err := usageEventAggregationsCaughtUp(tx) + if err != nil { + return err + } + if !safe { + status = dto.UsageEventArchiveStatusAggregationLagging + return nil + } + + // INSERT SELECT 避免千万级归档在 Go 内存中反序列化完整事件;列清单是 hot/archive 共享契约。 + columns := entities.UsageEventStorageColumns + insertSQL := fmt.Sprintf("INSERT INTO usage_events_archive (%s) SELECT %s FROM usage_events WHERE id IN ?", columns, columns) + insertResult := tx.Exec(insertSQL, ids) + if insertResult.Error != nil { + return fmt.Errorf("archive usage events: %w", insertResult.Error) + } + if insertResult.RowsAffected != int64(len(ids)) { + return fmt.Errorf("archive usage events: expected %d inserted rows, got %d", len(ids), insertResult.RowsAffected) + } + + var archivedCount int64 + if err := tx.Model(&entities.UsageEventArchive{}).Where("id IN ?", ids).Count(&archivedCount).Error; err != nil { + return fmt.Errorf("verify archived usage events: %w", err) + } + if archivedCount != int64(len(ids)) { + return fmt.Errorf("verify archived usage events: expected %d rows, got %d", len(ids), archivedCount) + } + + deleteResult := tx.Unscoped().Where("id IN ?", ids).Delete(&entities.UsageEvent{}) + if deleteResult.Error != nil { + return fmt.Errorf("delete archived usage events from hot table: %w", deleteResult.Error) + } + if deleteResult.RowsAffected != int64(len(ids)) { + return fmt.Errorf("delete archived usage events from hot table: expected %d rows, got %d", len(ids), deleteResult.RowsAffected) + } + archived = int64(len(ids)) + status = dto.UsageEventArchiveStatusArchived + return nil + }) + if err != nil { + return 0, status, err + } + return archived, status, nil +} + +func databaseContext(db *gorm.DB) context.Context { + if db != nil && db.Statement != nil && db.Statement.Context != nil { + return db.Statement.Context + } + return context.Background() +} diff --git a/internal/service/sync.go b/internal/service/sync.go index 5d7e95d2..208ab446 100644 --- a/internal/service/sync.go +++ b/internal/service/sync.go @@ -11,6 +11,7 @@ import ( "cpa-usage-keeper/internal/entities" "cpa-usage-keeper/internal/quota" "cpa-usage-keeper/internal/repository" + repositorydto "cpa-usage-keeper/internal/repository/dto" "cpa-usage-keeper/internal/service/tokenprocessor" "cpa-usage-keeper/internal/timeutil" @@ -59,16 +60,14 @@ type SyncService struct { // usageAggregation 只接收提交后通知,不允许热路径同步调用聚合仓储函数。 usageAggregation UsageAggregationNotifier // usageHeaderQuota 与聚合 runner 解耦,在 Quota worker 内按一分钟窗口自行合并。 - usageHeaderQuota UsageHeaderSnapshotAppender - cleanupUsageEventsEnabled bool + usageHeaderQuota UsageHeaderSnapshotAppender } // NewSyncService 按生产配置组装 CPA metadata client;远端 usage 拉取由 poller 独立负责。 func NewSyncService(db *gorm.DB, cfg config.Config) *SyncService { return NewSyncServiceWithOptions(db, SyncServiceOptions{ - BaseURL: cfg.CPABaseURL, - Client: cpa.NewClient(cfg.CPABaseURL, cfg.CPAManagementKey, cfg.RequestTimeout, cfg.TLSSkipVerify), - CleanupUsageEventsEnabled: cfg.CleanupUsageEventsEnabled, + BaseURL: cfg.CPABaseURL, + Client: cpa.NewClient(cfg.CPABaseURL, cfg.CPAManagementKey, cfg.RequestTimeout, cfg.TLSSkipVerify), }) } @@ -82,8 +81,7 @@ type SyncServiceOptions struct { // UsageAggregationNotifier 注入 App 唯一的单 writer runner。 UsageAggregationNotifier UsageAggregationNotifier // UsageHeaderQuota 独立接收原始 Header;是否配置聚合 notifier 不影响它。 - UsageHeaderQuota UsageHeaderSnapshotAppender - CleanupUsageEventsEnabled bool + UsageHeaderQuota UsageHeaderSnapshotAppender } // NewSyncServiceWithOptions 是统一构造入口,负责填充默认时钟和 metadata fetcher。 @@ -106,8 +104,7 @@ func NewSyncServiceWithOptions(db *gorm.DB, opts SyncServiceOptions) *SyncServic // 构造时只保存 notifier 接口,不启动额外 goroutine。 usageAggregation: opts.UsageAggregationNotifier, // Header appender 始终独立于聚合 notifier,生产 App 会同时注入两个接收方。 - usageHeaderQuota: opts.UsageHeaderQuota, - cleanupUsageEventsEnabled: opts.CleanupUsageEventsEnabled, + usageHeaderQuota: opts.UsageHeaderQuota, } } @@ -188,19 +185,27 @@ func (s *SyncService) CleanupRedisUsageInbox(ctx context.Context) error { return err } -// CleanupStorage 是每日 04:30 维护任务调用的统一入口:先清 Redis inbox,按配置清 usage_events,最后 VACUUM 收缩 SQLite。 +// CleanupStorage 是每日 04:30 维护任务入口:归档过期 usage_events 后清理限期统计并按空闲页条件整理 SQLite。 func (s *SyncService) CleanupStorage(ctx context.Context) error { if err := s.validate(syncMetadataOptional); err != nil { return err } - result, err := repository.CleanupStorage(s.db, s.now(), repository.CleanupStorageOptions{ - CleanupUsageEvents: s.cleanupUsageEventsEnabled, + result, err := repository.CleanupStorage(s.db.WithContext(ctx), s.now()) + entry := logrus.WithFields(logrus.Fields{ + "redis_processed_deleted": result.RedisInbox.ProcessedDeleted, + "redis_failed_deleted": result.RedisInbox.FailedDeleted, + "usage_events_archived": result.UsageEventsArchived, + "usage_events_archive_status": result.UsageEventsArchiveStatus, + "vacuum_performed": result.Vacuum.Performed, + "vacuum_skipped_reason": result.Vacuum.SkippedReason, + "sqlite_free_bytes": result.Vacuum.FreeBytes, + "sqlite_free_ratio": result.Vacuum.FreeRatio, }) - logrus.WithFields(logrus.Fields{ - "redis_processed_deleted": result.RedisInbox.ProcessedDeleted, - "redis_failed_deleted": result.RedisInbox.FailedDeleted, - "usage_events_deleted": result.UsageEventsDeleted, - }).Debug("storage cleanup finished") + if result.UsageEventsArchiveStatus == repositorydto.UsageEventArchiveStatusAggregationLagging { + entry.Warn("usage event archive deferred because aggregations are lagging") + } else { + entry.Debug("storage cleanup finished") + } return err } diff --git a/internal/service/test/sync_cleanup_test.go b/internal/service/test/sync_cleanup_test.go index a27cee00..e7999cfa 100644 --- a/internal/service/test/sync_cleanup_test.go +++ b/internal/service/test/sync_cleanup_test.go @@ -1,8 +1,10 @@ package test import ( + "bytes" "context" "path/filepath" + "strings" "testing" "time" @@ -10,13 +12,15 @@ import ( "cpa-usage-keeper/internal/entities" "cpa-usage-keeper/internal/repository" "cpa-usage-keeper/internal/service" + "github.com/sirupsen/logrus" "gorm.io/gorm" ) -func TestSyncServiceCleanupStorageSkipsUsageEventsByDefault(t *testing.T) { +func TestSyncServiceCleanupStorageDefersArchiveUntilAggregationsCatchUp(t *testing.T) { db := openSyncCleanupTestDatabase(t) now := time.Date(2026, 6, 16, 9, 0, 0, 0, time.Local) seedSyncCleanupUsageEventsAt(t, db, now.AddDate(0, 0, -91), now.Add(-time.Hour)) + logs := captureSyncCleanupLogs(t, logrus.WarnLevel) syncer := service.NewSyncServiceWithOptions(db, service.SyncServiceOptions{ Now: func() time.Time { return now }, }) @@ -32,25 +36,26 @@ func TestSyncServiceCleanupStorageSkipsUsageEventsByDefault(t *testing.T) { if count != 2 { t.Fatalf("expected usage_events to be retained by default, got %d rows", count) } + content := logs.String() + if !strings.Contains(content, "level=warning") || !strings.Contains(content, "usage event archive deferred because aggregations are lagging") { + t.Fatalf("expected aggregation lag warning, got %q", content) + } } -func TestSyncServiceCleanupStorageDeletesUsageEventsWhenEnabled(t *testing.T) { +func TestSyncServiceCleanupStorageArchivesUsageEventsAfterAggregationsCatchUp(t *testing.T) { // 准备:写入一条过期和一条近期事件,并先追平两个全局聚合 checkpoint。 db := openSyncCleanupTestDatabase(t) now := time.Date(2026, 6, 16, 9, 0, 0, 0, time.Local) seedSyncCleanupUsageEventsAt(t, db, now.AddDate(0, 0, -91), now.Add(-time.Hour)) catchUpSyncCleanupAggregations(t, db, now) - syncer := service.NewSyncServiceWithOptions(db, service.SyncServiceOptions{ - Now: func() time.Time { return now }, - CleanupUsageEventsEnabled: true, - }) + syncer := service.NewSyncServiceWithOptions(db, service.SyncServiceOptions{Now: func() time.Time { return now }}) - // 执行:显式启用 raw usage event 清理。 + // 执行:固定维护策略把超过 90 天且已完成聚合的 raw event 移入 archive。 if err := syncer.CleanupStorage(context.Background()); err != nil { t.Fatalf("CleanupStorage returned error: %v", err) } - // 断言:安全水位满足后只删除早于 90 天边界的过期事件。 + // 断言:安全水位满足后 hot 只保留近期事件,旧事件完整进入 archive。 var remainingKeys []string if err := db.Model(&entities.UsageEvent{}).Order("event_key asc").Pluck("event_key", &remainingKeys).Error; err != nil { t.Fatalf("load remaining usage events: %v", err) @@ -58,19 +63,25 @@ func TestSyncServiceCleanupStorageDeletesUsageEventsWhenEnabled(t *testing.T) { if len(remainingKeys) != 1 || remainingKeys[0] != "recent" { t.Fatalf("expected only recent usage event to remain, got %v", remainingKeys) } + var archivedKeys []string + if err := db.Model(&entities.UsageEventArchive{}).Order("event_key asc").Pluck("event_key", &archivedKeys).Error; err != nil { + t.Fatalf("load archived usage events: %v", err) + } + if len(archivedKeys) != 1 || archivedKeys[0] != "old" { + t.Fatalf("expected old usage event in archive, got %v", archivedKeys) + } } -func TestNewSyncServiceCleanupStorageReadsCleanupFlagFromConfig(t *testing.T) { - // 准备:通过生产构造器配置清理开关,并让两个全局聚合 checkpoint 先追平。 +func TestNewSyncServiceCleanupStorageUsesFixedArchivePolicy(t *testing.T) { + // 准备:通过生产构造器创建 service,并让三个全局聚合 checkpoint 先追平。 db := openSyncCleanupTestDatabase(t) now := time.Now().In(time.Local) seedSyncCleanupUsageEventsAt(t, db, now.AddDate(0, 0, -91), now) catchUpSyncCleanupAggregations(t, db, now) syncer := service.NewSyncService(db, config.Config{ - CPABaseURL: "https://cpa.example.com", - CPAManagementKey: "secret", - RequestTimeout: time.Second, - CleanupUsageEventsEnabled: true, + CPABaseURL: "https://cpa.example.com", + CPAManagementKey: "secret", + RequestTimeout: time.Second, }) // 执行:调用生产 SyncService 的统一维护入口。 @@ -78,13 +89,13 @@ func TestNewSyncServiceCleanupStorageReadsCleanupFlagFromConfig(t *testing.T) { t.Fatalf("CleanupStorage returned error: %v", err) } - // 断言:配置开关生效且只删除安全水位内的过期事件。 + // 断言:无需配置开关,固定策略只把安全水位内的过期事件移出 hot。 var remainingKeys []string if err := db.Model(&entities.UsageEvent{}).Order("event_key asc").Pluck("event_key", &remainingKeys).Error; err != nil { t.Fatalf("load remaining usage events: %v", err) } if len(remainingKeys) != 1 || remainingKeys[0] != "recent" { - t.Fatalf("expected production config cleanup flag to retain only recent usage event, got %v", remainingKeys) + t.Fatalf("expected fixed archive policy to retain only recent usage event, got %v", remainingKeys) } } @@ -129,5 +140,22 @@ func catchUpSyncCleanupAggregations(t *testing.T, db *gorm.DB, now time.Time) { if err := repository.AggregateUsageLatencyStats(context.Background(), db, now); err != nil { t.Fatalf("aggregate latency before cleanup: %v", err) } - // 断言由调用用例通过最终删除结果完成,helper 不额外读取数据库。 + // 断言由调用用例通过最终归档结果完成,helper 不额外读取数据库。 +} + +func captureSyncCleanupLogs(t *testing.T, level logrus.Level) *bytes.Buffer { + t.Helper() + logs := &bytes.Buffer{} + previousOutput := logrus.StandardLogger().Out + previousFormatter := logrus.StandardLogger().Formatter + previousLevel := logrus.GetLevel() + logrus.SetOutput(logs) + logrus.SetFormatter(&logrus.TextFormatter{DisableTimestamp: true}) + logrus.SetLevel(level) + t.Cleanup(func() { + logrus.SetOutput(previousOutput) + logrus.SetFormatter(previousFormatter) + logrus.SetLevel(previousLevel) + }) + return logs }