From 77ce9db722bf20821d1d41cdd3d1ce1c53d458bf Mon Sep 17 00:00:00 2001 From: break Date: Mon, 10 Aug 2026 16:58:33 +0800 Subject: [PATCH] =?UTF-8?q?=E5=B9=B6=E5=8F=91=E7=9B=91=E6=8E=A7=E4=BF=AE?= =?UTF-8?q?=E5=A4=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../service/polling/concurrency_service.go | 15 ++++++++-- .../.openspec.yaml | 2 ++ .../design.md | 30 +++++++++++++++++++ .../proposal.md | 24 +++++++++++++++ .../specs/polling-operations/spec.md | 13 ++++++++ .../tasks.md | 9 ++++++ openspec/specs/polling-operations/spec.md | 12 ++++++++ 7 files changed, 102 insertions(+), 3 deletions(-) create mode 100644 openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/.openspec.yaml create mode 100644 openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/design.md create mode 100644 openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/proposal.md create mode 100644 openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/specs/polling-operations/spec.md create mode 100644 openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/tasks.md diff --git a/internal/service/polling/concurrency_service.go b/internal/service/polling/concurrency_service.go index f80a9d9..6c438e1 100644 --- a/internal/service/polling/concurrency_service.go +++ b/internal/service/polling/concurrency_service.go @@ -2,6 +2,7 @@ package polling import ( "context" + "strings" "time" "github.com/redis/go-redis/v9" @@ -63,7 +64,7 @@ func (s *ConcurrencyService) List(ctx context.Context) ([]*ConcurrencyStatus, er } // 从 Redis 获取当前并发数 - currentKey := constants.RedisPollingConcurrencyCurrentKey(cfg.TaskType) + currentKey := pollingConcurrencyCurrentKey(cfg.TaskType) current, err := s.redis.Get(ctx, currentKey).Int64() if err != nil && err != redis.Nil { current = 0 @@ -98,7 +99,7 @@ func (s *ConcurrencyService) GetByTaskType(ctx context.Context, taskType string) } // 从 Redis 获取当前并发数 - currentKey := constants.RedisPollingConcurrencyCurrentKey(cfg.TaskType) + currentKey := pollingConcurrencyCurrentKey(cfg.TaskType) current, err := s.redis.Get(ctx, currentKey).Int64() if err != nil && err != redis.Nil { current = 0 @@ -176,7 +177,7 @@ func (s *ConcurrencyService) ResetConcurrency(ctx context.Context, taskType stri } // 重置 Redis 当前计数为 0 - currentKey := constants.RedisPollingConcurrencyCurrentKey(taskType) + currentKey := pollingConcurrencyCurrentKey(config.TaskType) before, getErr := s.redis.Get(ctx, currentKey).Int64() beforeExists := getErr == nil if getErr != nil && getErr != redis.Nil { @@ -257,6 +258,14 @@ func (s *ConcurrencyService) SyncConfigToRedis(ctx context.Context, config *mode return s.redis.Set(ctx, configKey, config.MaxConcurrency, 24*time.Hour).Err() } +// pollingConcurrencyCurrentKey 将配置中的短任务类型转换为 Worker 使用的完整计数键。 +func pollingConcurrencyCurrentKey(taskType string) string { + if !strings.HasPrefix(taskType, "polling:") { + taskType = "polling:" + taskType + } + return constants.RedisPollingConcurrencyCurrentKey(taskType) +} + // getTaskTypeName 获取任务类型的中文名称 func (s *ConcurrencyService) getTaskTypeName(taskType string) string { switch taskType { diff --git a/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/.openspec.yaml b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/.openspec.yaml new file mode 100644 index 0000000..d7bc011 --- /dev/null +++ b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-08-10 diff --git a/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/design.md b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/design.md new file mode 100644 index 0000000..4ae9739 --- /dev/null +++ b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/design.md @@ -0,0 +1,30 @@ +## Context + +轮询执行器以完整任务类型维护当前并发计数,而管理服务以短任务类型读取和重置计数;最大并发配置仍按短任务类型保存。 + +## Goals / Non-Goals + +**Goals:** +- 管理服务与执行器使用同一当前计数键。 +- 保留短任务类型的数据库与 Redis 最大并发配置格式。 + +**Non-Goals:** +- 不调整最大并发值、Redis Lua 限流算法或任务调度。 +- 不新增数据库记录或 Schema。 + +## Decisions + +管理服务在读取和重置当前计数前,将数据库短任务类型规范化为轮询完整任务类型。这样只修正观测与重置目标,执行器和既有配置保持不变。 + +不修改执行器改用短类型计数,避免改变已运行 Worker 的共享信号量键并造成发布期间的计数分裂。 + +## Risks / Trade-offs + +- [旧错误短类型计数键残留] → 新代码忽略该键;其无 TTL 的残留值不再影响限流或展示。 +- [重置运行中计数] → 保持现有管理语义,后续任务完成时仍会递减同一完整类型键。 + +## Migration Plan + +1. 发布 API 与 Worker 代码。 +2. 在存在轮询任务时核对列表、详情和重置读取同一完整类型计数键。 +3. 回滚时恢复上一版本;不涉及数据迁移。 diff --git a/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/proposal.md b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/proposal.md new file mode 100644 index 0000000..78dbcd2 --- /dev/null +++ b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/proposal.md @@ -0,0 +1,24 @@ +## Why + +轮询并发管理接口按数据库中的短任务类型读取 Redis 当前计数,而轮询执行器按完整任务类型写入计数,导致使用率与重置操作面向错误键并持续显示为零。 + +## What Changes + +- 统一轮询并发管理读写的当前计数键与执行器键格式。 +- 使列表、详情和重置接口展示并操作实际生效的全局并发计数。 +- 不改变最大并发配置键、限流算法或现有任务类型。 + +## Capabilities + +### New Capabilities + +- 无。 + +### Modified Capabilities + +- `polling-operations`: 轮询并发控制接口返回实际任务计数并重置同一计数。 + +## Impact + +- 涉及 `internal/service/polling/concurrency_service.go` 的并发状态读取与重置。 +- 不涉及数据库 Schema、路由或新增依赖。 diff --git a/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/specs/polling-operations/spec.md b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/specs/polling-operations/spec.md new file mode 100644 index 0000000..9fc4242 --- /dev/null +++ b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/specs/polling-operations/spec.md @@ -0,0 +1,13 @@ +## ADDED Requirements + +### Requirement: 轮询并发计数可观测 + +系统 SHALL 在轮询并发配置列表和详情中返回与实际限流器相同任务类型的当前计数、可用并发和使用率;重置操作 MUST 重置该同一计数。 + +#### Scenario: 查询运行中的任务计数 +- **WHEN** 某轮询任务正在占用并发配额 +- **THEN** 查询该任务类型的并发状态返回非零当前计数,并据此计算可用并发和使用率 + +#### Scenario: 重置任务计数 +- **WHEN** 授权操作者重置某轮询任务类型的并发计数 +- **THEN** 后续状态查询返回该任务类型的当前计数为零,且不影响其他任务类型的计数 diff --git a/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/tasks.md b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/tasks.md new file mode 100644 index 0000000..f131c75 --- /dev/null +++ b/openspec/changes/archive/2026-08-10-fix-polling-concurrency-usage-metrics/tasks.md @@ -0,0 +1,9 @@ +## 1. 并发计数键修复 + +- [x] 1.1 统一轮询并发列表、详情与重置操作使用完整任务类型的当前计数键。 +- [x] 1.2 保持最大并发配置键使用短任务类型,确认不改变限流器行为。 + +## 2. 验证 + +- [x] 2.1 使用 Redis 计数键验证列表、详情和重置读取同一任务计数。 +- [x] 2.2 运行 gofmt、API/Worker 构建、OpenAPI、OpenSpec 和上下文健康检查。 diff --git a/openspec/specs/polling-operations/spec.md b/openspec/specs/polling-operations/spec.md index 707bfaa..bbc05db 100644 --- a/openspec/specs/polling-operations/spec.md +++ b/openspec/specs/polling-operations/spec.md @@ -26,6 +26,18 @@ - **WHEN** 授权操作者取消任务 - **THEN** 系统先记录已取消;若后台批量处理随后结束,当前实现可能把同一任务覆盖为已完成 +### Requirement: 轮询并发计数可观测 + +系统 SHALL 在轮询并发配置列表和详情中返回与实际限流器相同任务类型的当前计数、可用并发和使用率;重置操作 MUST 重置该同一计数。 + +#### Scenario: 查询运行中的任务计数 +- **WHEN** 某轮询任务正在占用并发配额 +- **THEN** 查询该任务类型的并发状态返回非零当前计数,并据此计算可用并发和使用率 + +#### Scenario: 重置任务计数 +- **WHEN** 授权操作者重置某轮询任务类型的并发计数 +- **THEN** 后续状态查询返回该任务类型的当前计数为零,且不影响其他任务类型的计数 + ## 可达操作索引 本节只用于入口导航,不是行为 Requirement;业务义务以上述 Requirements 为准。