From 15bbb953db8a0240c059e501979a04e977b53d18 Mon Sep 17 00:00:00 2001 From: break Date: Wed, 16 Sep 2026 16:41:25 +0800 Subject: [PATCH] =?UTF-8?q?fix(=E9=80=9A=E9=81=93=E6=B5=81=E9=87=8F?= =?UTF-8?q?=E9=98=88=E5=80=BC):=20AUG26-011=20=E4=BF=AE=E5=A4=8D=E5=A4=B1?= =?UTF-8?q?=E8=B4=A5/=E6=9C=AA=E7=9F=A5=E7=BB=93=E6=9E=9C=E6=94=B6?= =?UTF-8?q?=E6=95=9B=E3=80=81=E9=94=81=E5=AE=9A=20carrier=20=E7=BC=BA?= =?UTF-8?q?=E5=A4=B1=E5=87=BA=E8=B7=AF=E4=B8=8E=E5=AE=A1=E8=AE=A1=E5=9B=9E?= =?UTF-8?q?=E5=BD=92?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../application/carrierthreshold/cycle.go | 79 +++++++++++++++---- .../application/carrierthreshold/execute.go | 4 +- .../carrierthreshold/lock_store.go | 72 ++++++++++++----- internal/domain/carrierthreshold/lock.go | 19 +++++ .../service/iot_card/stop_resume_service.go | 8 ++ ...227_add_carrier_traffic_threshold.down.sql | 2 +- ...00227_add_carrier_traffic_threshold.up.sql | 10 ++- .../design.md | 45 ++++++++--- .../proposal.md | 5 +- .../specs/package-lifecycle/spec.md | 55 +++++++++++++ .../tasks.md | 9 +++ 11 files changed, 250 insertions(+), 58 deletions(-) create mode 100644 openspec/changes/add-carrier-channel-traffic-thresholds/specs/package-lifecycle/spec.md diff --git a/internal/application/carrierthreshold/cycle.go b/internal/application/carrierthreshold/cycle.go index 1035783..b8b110b 100644 --- a/internal/application/carrierthreshold/cycle.go +++ b/internal/application/carrierthreshold/cycle.go @@ -18,38 +18,42 @@ const cycleBatchSize = 200 // CycleStats 是一次周期处理的可观察结果。 // -// Scanned 为扫到的跨期持锁数;Unlocked 为本次认领解锁成功的锁数;Resumed 为同时写出复机事件的锁数; -// Skipped 为本地事实缺失、被并发推进或判定失败而未改动的锁数。 -type CycleStats struct{ Scanned, Unlocked, Resumed, Skipped int } +// Scanned 为扫到的待处理持锁数;Unlocked 为本次认领解锁成功的锁数;Resumed 为同时写出复机事件的锁数; +// Anomaly 为因运营商配置缺失或重置日非法而解锁并转人工的锁数;Skipped 为本地事实缺失、被并发推进 +// 或判定失败而未改动的锁数。 +type CycleStats struct{ Scanned, Unlocked, Resumed, Anomaly, Skipped int } -// ProcessDueLocks 扫描已跨期的持锁锁行:认领解锁,并在新周期条件满足时同事务写复机事件。 +// ProcessDueLocks 扫描待处理的持锁锁行:认领解锁,并在新周期条件满足时同事务写复机事件。 // // 过期判断按锁行自身运营商的 data_reset_day(支持换运营商后旧锁仍按其归属处理)。解锁与复机事件 // 在同一事务提交:事务失败则解锁一起回滚,下一分钟重新处理,不会出现「已解锁但没有复机任务」的中间态。 // 任一复机条件不满足时只解锁、不调运营商,并把可观察原因写入锁行。 +// +// 锁行引用的运营商已不存在或重置日非法时无法计算周期归属:这类行按「新周期对仍持锁卡解除通道锁」 +// 的语义解锁并标记异常转人工,绝不写复机事件、不调运营商,也绝不静默跳过(否则持锁卡会永久禁止复机)。 func (s *Service) ProcessDueLocks(ctx context.Context, now time.Time) (CycleStats, error) { stats := CycleStats{} if s == nil || s.db == nil || s.repository == nil { return stats, errors.New(errors.CodeServiceUnavailable, "通道流量阈值周期处理能力未配置") } - locks, err := s.ScanExpiredLocks(ctx, now, cycleBatchSize) + due, err := s.ScanDueLocks(ctx, now, cycleBatchSize) if err != nil { return stats, err } - stats.Scanned = len(locks) - if len(locks) == 0 { + stats.Scanned = len(due) + if len(due) == 0 { return stats, nil } if s.commander == nil { return stats, errors.New(errors.CodeServiceUnavailable, "通道流量阈值停复机执行端口未配置") } var firstErr error - for index := range locks { - lock := locks[index] - if err := s.processDueLock(ctx, &lock, &stats); err != nil { + for index := range due { + item := due[index] + if err := s.processDueLock(ctx, &item, &stats); err != nil { stats.Skipped++ s.logger.Warn("通道阈值跨期处理单条失败", - zap.Uint("lock_id", lock.ID), zap.Uint("card_id", lock.CardID), zap.Error(err)) + zap.Uint("lock_id", item.Lock.ID), zap.Uint("card_id", item.Lock.CardID), zap.Error(err)) if firstErr == nil { firstErr = err } @@ -58,8 +62,12 @@ func (s *Service) ProcessDueLocks(ctx context.Context, now time.Time) (CycleStat return stats, firstErr } -// processDueLock 处理单条跨期锁:认领解锁并在条件满足时同事务写复机事件。 -func (s *Service) processDueLock(ctx context.Context, lock *model.CarrierTrafficThresholdLock, stats *CycleStats) error { +// processDueLock 处理单条待处理锁行:周期归属不可判定时解锁并转人工,已跨期时认领解锁并条件复机。 +func (s *Service) processDueLock(ctx context.Context, item *DueLock, stats *CycleStats) error { + if item.Kind == DueLockUnresolvable { + return s.processUnresolvableLock(ctx, item, stats) + } + lock := &item.Lock card, err := s.loadCard(ctx, lock.CardID) if err != nil { return err @@ -111,6 +119,42 @@ func (s *Service) processDueLock(ctx context.Context, lock *model.CarrierTraffic return nil } +// processUnresolvableLock 处理周期归属不可判定的锁行:同事务认领解锁并标记异常转人工。 +// +// 只解锁不写复机事件、不调运营商:配置缺失时无法判断是否已跨期,解除通道锁交由既有复机链路 +// 与人工决定;anomaly_flag 与失败原因使运维可见并转人工核对。解锁与异常标记都是条件更新, +// 重复执行不会产生第二次副作用。 +func (s *Service) processUnresolvableLock(ctx context.Context, item *DueLock, stats *CycleStats) error { + lock := &item.Lock + s.logger.Warn("通道阈值锁行周期归属不可判定,解锁并转人工核对", + zap.Uint("lock_id", lock.ID), zap.Uint("card_id", lock.CardID), zap.Uint("carrier_id", lock.CarrierID), + zap.String("reason", item.Reason)) + unlocked := false + err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + claimed, unlockErr := s.unlockInTx(ctx, tx, lock.ID, item.Reason) + if unlockErr != nil { + return unlockErr + } + if !claimed { + return nil + } + unlocked = true + _, anomalyErr := s.markAnomaly(ctx, lock.ID, item.Reason) + return anomalyErr + }) + if err != nil { + return err + } + if !unlocked { + stats.Skipped++ + s.logger.Info("通道阈值不可判定锁已被并发处理,跳过", zap.Uint("lock_id", lock.ID)) + return nil + } + stats.Unlocked++ + stats.Anomaly++ + return nil +} + // recoveryBatchSize 是单次恢复扫描的锁行上限。 const recoveryBatchSize = 200 @@ -158,8 +202,11 @@ func (s *Service) RecoverSubmitted(ctx context.Context, now time.Time) (Recovery } // recoverSubmittedLock 处理单条待确认锁:只查询状态回填,不发起任何停复机调用。 +// +// 入口只处理未决子任务(submitted/unknown/failed):调用失败或结果未知同样可能已在运营商侧生效, +// 必须继续收敛;confirmed 是终态,pending 表示从未对外调用,都不在本扫描范围。 func (s *Service) recoverSubmittedLock(ctx context.Context, lock *model.CarrierTrafficThresholdLock, now time.Time, stats *RecoveryStats) error { - if lock.StopStatus != domain.TaskStatusSubmitted && lock.ResumeStatus != domain.TaskStatusSubmitted { + if !domain.IsUnresolvedTaskStatus(lock.StopStatus) && !domain.IsUnresolvedTaskStatus(lock.ResumeStatus) { return nil } card, err := s.loadCard(ctx, lock.CardID) @@ -178,7 +225,7 @@ func (s *Service) recoverSubmittedLock(ctx context.Context, lock *model.CarrierT } confirmed := 0 unconfirmed := 0 - if lock.StopStatus == domain.TaskStatusSubmitted { + if domain.IsUnresolvedTaskStatus(lock.StopStatus) { switch { case known && status == constants.NetworkStatusOffline: confirmed++ @@ -189,7 +236,7 @@ func (s *Service) recoverSubmittedLock(ctx context.Context, lock *model.CarrierT unconfirmed++ } } - if lock.ResumeStatus == domain.TaskStatusSubmitted { + if domain.IsUnresolvedTaskStatus(lock.ResumeStatus) { switch { case known && status == constants.NetworkStatusOnline: confirmed++ diff --git a/internal/application/carrierthreshold/execute.go b/internal/application/carrierthreshold/execute.go index c2568af..4fbb66c 100644 --- a/internal/application/carrierthreshold/execute.go +++ b/internal/application/carrierthreshold/execute.go @@ -77,7 +77,7 @@ func (s *Service) ExecuteStop(ctx context.Context, lockID uint) error { zap.Uint("lock_id", lockID), zap.Uint("card_id", lock.CardID), zap.String("result", outcome.Result), zap.Error(callErr)) } - _, err = s.markTaskOutcome(ctx, lockID, stopTask, domain.TaskStatusSubmitted, + _, err = s.markTaskOutcome(ctx, lockID, stopTask, []string{domain.TaskStatusSubmitted}, taskResultOf(outcome), outcome.IntegrationID, stopFailureReason(outcome)) if err != nil { return err @@ -142,7 +142,7 @@ func (s *Service) ExecuteResume(ctx context.Context, lockID uint) error { zap.Uint("lock_id", lockID), zap.Uint("card_id", lock.CardID), zap.String("result", outcome.Result), zap.Error(callErr)) } - _, err = s.markTaskOutcome(ctx, lockID, resumeTask, domain.TaskStatusSubmitted, + _, err = s.markTaskOutcome(ctx, lockID, resumeTask, []string{domain.TaskStatusSubmitted}, taskResultOf(outcome), outcome.IntegrationID, resumeFailureReason(outcome)) if err != nil { return err diff --git a/internal/application/carrierthreshold/lock_store.go b/internal/application/carrierthreshold/lock_store.go index 7ae818d..c68609c 100644 --- a/internal/application/carrierthreshold/lock_store.go +++ b/internal/application/carrierthreshold/lock_store.go @@ -122,12 +122,36 @@ func (s *Service) ClaimResumeSubmission(ctx context.Context, lockID uint, now ti return claimed.RowsAffected == 1, nil } -// ScanExpiredLocks 扫描已跨期但仍持锁的锁行,供周期处理解锁与条件复机。 +// DueLockKind 描述一条仍在持锁的锁行的跨期判定结果。 +type DueLockKind string + +const ( + // DueLockExpired 表示已跨期:锁行 period_start 早于该锁行自身运营商按 data_reset_day 算出的当前周期起点。 + DueLockExpired DueLockKind = "expired" + // DueLockUnresolvable 表示周期归属不可判定:锁行引用的运营商已不存在(含软删)或重置日非法。 + DueLockUnresolvable DueLockKind = "unresolvable" +) + +// unresolvableCarrierReason 是周期归属不可判定时写入锁行的可安全原因。 +const unresolvableCarrierReason = "锁行引用的运营商已不存在或上游流量重置日非法,已按跨期解除通道锁并转人工核对" + +// DueLock 是周期处理扫描到的一条待处理锁行。 +type DueLock struct { + // Lock 是持锁锁行本身。 + Lock model.CarrierTrafficThresholdLock + // Kind 是跨期判定结果:已跨期或周期归属不可判定。 + Kind DueLockKind + // Reason 是周期归属不可判定时的可安全原因(Kind 为 DueLockExpired 时为空)。 + Reason string +} + +// ScanDueLocks 扫描需要周期处理的持锁锁行,供周期处理解锁与条件复机。 // // 过期判断按锁行自身运营商的 data_reset_day 计算其当前周期起点,与锁行 period_start 不一致即已跨期, -// 因此换运营商后的旧锁仍按其旧 carrier 的归属被正确识别。锁行引用的运营商已不存在时无法计算周期, -// 本次不返回该锁(不猜测),交由周期处理标记异常转人工。 -func (s *Service) ScanExpiredLocks(ctx context.Context, now time.Time, limit int) ([]model.CarrierTrafficThresholdLock, error) { +// 因此换运营商后的旧锁仍按其旧 carrier 的归属被正确识别。锁行引用的运营商已不存在或重置日非法时 +// 无法计算周期归属,这类行以 DueLockUnresolvable 返回:周期处理必须给出出路(按跨期语义解锁并转人工), +// 绝不能让持锁卡因配置缺失而永久禁止复机。 +func (s *Service) ScanDueLocks(ctx context.Context, now time.Time, limit int) ([]DueLock, error) { if s == nil || s.db == nil { return nil, nil } @@ -144,34 +168,40 @@ func (s *Service) ScanExpiredLocks(ctx context.Context, now time.Time, limit int if err != nil { return nil, err } - expired := make([]model.CarrierTrafficThresholdLock, 0, len(locks)) + due := make([]DueLock, 0, len(locks)) for index := range locks { - resetDay, ok := resetDays[locks[index].CarrierID] + lock := locks[index] + resetDay, ok := resetDays[lock.CarrierID] if !ok { + due = append(due, DueLock{Lock: lock, Kind: DueLockUnresolvable, Reason: unresolvableCarrierReason}) continue } - periodStart, err := domain.PeriodStart(now, resetDay) - if err != nil { + periodStart, periodErr := domain.PeriodStart(now, resetDay) + if periodErr != nil { + due = append(due, DueLock{Lock: lock, Kind: DueLockUnresolvable, Reason: unresolvableCarrierReason}) continue } - if locks[index].PeriodStart.Before(periodStart) { - expired = append(expired, locks[index]) + if lock.PeriodStart.Before(periodStart) { + due = append(due, DueLock{Lock: lock, Kind: DueLockExpired}) } } - return expired, nil + return due, nil } -// ScanSubmittedLocks 扫描存在已提交子任务的锁行,供恢复扫描查询运营商状态回填。 +// ScanSubmittedLocks 扫描存在未决子任务的锁行,供恢复扫描查询运营商状态回填。 // +// 未决集合为 {submitted, unknown, failed}(domain.UnresolvedTaskStatuses):调用失败或结果未知的行 +// 仍可能已在运营商侧生效,因此必须继续用只读状态查询收敛,不得退出链路。 // 已标记异常(转人工)的锁必须退出扫描,否则每次扫描都会重复查询同一笔无法收敛的结果。 func (s *Service) ScanSubmittedLocks(ctx context.Context, limit int) ([]model.CarrierTrafficThresholdLock, error) { if s == nil || s.db == nil { return nil, nil } + unresolved := domain.UnresolvedTaskStatuses() var locks []model.CarrierTrafficThresholdLock if err := s.db.WithContext(ctx). - Where("anomaly_flag = ? AND (stop_status = ? OR resume_status = ?)", - domain.AnomalyFlagNone, domain.TaskStatusSubmitted, domain.TaskStatusSubmitted). + Where("anomaly_flag = ? AND (stop_status IN ? OR resume_status IN ?)", + domain.AnomalyFlagNone, unresolved, unresolved). Order("id ASC").Limit(limit).Find(&locks).Error; err != nil { return nil, errors.Wrap(errors.CodeDatabaseError, err, "扫描待确认通道流量阈值任务失败") } @@ -238,10 +268,10 @@ func (s *Service) loadCard(ctx context.Context, cardID uint) (*model.IotCard, er return &card, nil } -// markTaskOutcome 按 expected 状态条件更新把子任务推进到终态(ENG-CONC-001)。 -// 返回 false 表示记录已被并发推进,调用方必须按幂等处理,不再重复执行外部动作。 -func (s *Service) markTaskOutcome(ctx context.Context, lockID uint, task lockTask, expected, result, integrationID, failureReason string) (bool, error) { - if s == nil || s.db == nil || lockID == 0 { +// markTaskOutcome 按 expected 状态集合条件更新把子任务推进到终态(ENG-CONC-001)。 +// 返回 false 表示记录已被并发推进或已处于终态,调用方必须按幂等处理,不再重复执行外部动作。 +func (s *Service) markTaskOutcome(ctx context.Context, lockID uint, task lockTask, expected []string, result, integrationID, failureReason string) (bool, error) { + if s == nil || s.db == nil || lockID == 0 || len(expected) == 0 { return false, nil } updates := map[string]any{task.statusColumn: result} @@ -252,7 +282,7 @@ func (s *Service) markTaskOutcome(ctx context.Context, lockID uint, task lockTas updates["failure_reason"] = safeFailureReason(failureReason) } result_ := s.db.WithContext(ctx).Model(&model.CarrierTrafficThresholdLock{}). - Where("id = ? AND "+task.statusColumn+" = ?", lockID, expected). + Where("id = ? AND "+task.statusColumn+" IN ?", lockID, expected). Updates(updates) if result_.Error != nil { return false, errors.Wrap(errors.CodeDatabaseError, result_.Error, "回填通道阈值子任务状态失败") @@ -261,8 +291,10 @@ func (s *Service) markTaskOutcome(ctx context.Context, lockID uint, task lockTas } // markTaskConfirmed 由恢复扫描在运营商状态确认成功后把子任务回填为已确认。 +// expected 取未决集合 {submitted, unknown, failed}:失败与结果未知同样可能已在运营商侧生效, +// 必须允许收敛为已确认;confirmed 不在集合内,因此确认只会写入一次。 func (s *Service) markTaskConfirmed(ctx context.Context, lockID uint, task lockTask, integrationID string) (bool, error) { - return s.markTaskOutcome(ctx, lockID, task, domain.TaskStatusSubmitted, domain.TaskStatusConfirmed, integrationID, "") + return s.markTaskOutcome(ctx, lockID, task, domain.UnresolvedTaskStatuses(), domain.TaskStatusConfirmed, integrationID, "") } // markAnomaly 把锁标记为需人工核对并退出自动扫描。 diff --git a/internal/domain/carrierthreshold/lock.go b/internal/domain/carrierthreshold/lock.go index 3ad2434..f016063 100644 --- a/internal/domain/carrierthreshold/lock.go +++ b/internal/domain/carrierthreshold/lock.go @@ -29,3 +29,22 @@ const ( // AnomalyFlagManual 表示查询窗口超期或失败无法自动确认,已转人工核对并退出自动扫描。 AnomalyFlagManual = 1 ) + +// UnresolvedTaskStatuses 返回仍需由恢复扫描查询运营商状态收敛的子任务状态集合。 +// +// 语义:submitted 表示已提交待确认,unknown 表示结果未知,failed 表示调用明确失败但运营商侧 +// 状态仍可能已生效(例如请求已到达而响应超时/异常)——三者都必须继续用只读状态查询确认, +// 因此恢复扫描的查询谓词与超期判定共用本集合,避免「失败或结果未知」的锁行退出收敛链路。 +// pending 表示尚未对运营商发起过调用,confirmed 是终态,都不属于本集合。 +func UnresolvedTaskStatuses() []string { + return []string{TaskStatusSubmitted, TaskStatusUnknown, TaskStatusFailed} +} + +// IsUnresolvedTaskStatus 判断子任务状态是否仍需恢复扫描收敛。 +func IsUnresolvedTaskStatus(status string) bool { + switch status { + case TaskStatusSubmitted, TaskStatusUnknown, TaskStatusFailed: + return true + } + return false +} diff --git a/internal/service/iot_card/stop_resume_service.go b/internal/service/iot_card/stop_resume_service.go index 34106ba..3c1ff77 100644 --- a/internal/service/iot_card/stop_resume_service.go +++ b/internal/service/iot_card/stop_resume_service.go @@ -638,6 +638,10 @@ func (s *StopResumeService) stopCardWithRetry(ctx context.Context, card *model.I if lastErr == nil { lastErr = attemptObserver.recordingErr } + // Integration Log 无法终结时仍必须留下本次停机的审计事实(与重试耗尽路径同一写入)。 + s.recordCardCommandAudit(ctx, card, actionCode, summary+"未完成", attemptObserver.auditResult(lastErr), lastIntegrationID, + map[string]any{"network_status": card.NetworkStatus, "stop_reason": card.StopReason}, + map[string]any{"requested_network_status": constants.NetworkStatusOffline, "stop_reason": stopReason}, lastErr) return carrierthresholddomain.CommandOutcome{ IntegrationID: lastIntegrationID, Result: attemptObserver.auditResult(lastErr), }, lastErr @@ -739,6 +743,10 @@ func (s *StopResumeService) resumeCardWithRetry(ctx context.Context, card *model if lastErr == nil { lastErr = attemptObserver.recordingErr } + // Integration Log 无法终结时仍必须留下本次复机的审计事实(与重试耗尽路径同一写入)。 + s.recordCardCommandAudit(ctx, card, actionCode, summary+"未完成", attemptObserver.auditResult(lastErr), lastIntegrationID, + map[string]any{"network_status": card.NetworkStatus, "stop_reason": card.StopReason, "gateway_extend": card.GatewayExtend}, + map[string]any{"requested_network_status": constants.NetworkStatusOnline}, lastErr) return nil, lastIntegrationID, lastErr } if callErr == nil { diff --git a/migrations/000227_add_carrier_traffic_threshold.down.sql b/migrations/000227_add_carrier_traffic_threshold.down.sql index 3fb955f..915b70a 100644 --- a/migrations/000227_add_carrier_traffic_threshold.down.sql +++ b/migrations/000227_add_carrier_traffic_threshold.down.sql @@ -24,7 +24,7 @@ BEGIN END IF; END $$; -DROP INDEX IF EXISTS idx_carrier_traffic_threshold_lock_submitted; +DROP INDEX IF EXISTS idx_carrier_traffic_threshold_lock_unresolved; DROP INDEX IF EXISTS idx_carrier_traffic_threshold_lock_card; DROP INDEX IF EXISTS idx_carrier_traffic_threshold_lock_status_period; diff --git a/migrations/000227_add_carrier_traffic_threshold.up.sql b/migrations/000227_add_carrier_traffic_threshold.up.sql index 5c59149..c76e3e1 100644 --- a/migrations/000227_add_carrier_traffic_threshold.up.sql +++ b/migrations/000227_add_carrier_traffic_threshold.up.sql @@ -102,11 +102,13 @@ COMMENT ON COLUMN tb_carrier_traffic_threshold_lock.deleted_at IS '软删除时 CREATE INDEX idx_carrier_traffic_threshold_lock_status_period ON tb_carrier_traffic_threshold_lock (status, period_start); -- 四入口拒绝与周期处理按卡查当前周期锁。 CREATE INDEX idx_carrier_traffic_threshold_lock_card ON tb_carrier_traffic_threshold_lock (card_id, period_start); --- 恢复扫描只看 submitted 子任务。 -CREATE INDEX idx_carrier_traffic_threshold_lock_submitted ON tb_carrier_traffic_threshold_lock (stop_status, resume_status) - WHERE deleted_at IS NULL AND (stop_status = 'submitted' OR resume_status = 'submitted'); +-- 恢复扫描只看未决子任务(submitted/unknown/failed 都必须继续用只读状态查询收敛)。 +CREATE INDEX idx_carrier_traffic_threshold_lock_unresolved ON tb_carrier_traffic_threshold_lock (stop_status, resume_status) + WHERE deleted_at IS NULL AND ( + stop_status IN ('submitted', 'unknown', 'failed') OR resume_status IN ('submitted', 'unknown', 'failed') + ); COMMENT ON INDEX idx_carrier_traffic_threshold_lock_status_period IS '通道阈值周期处理扫描的进行中锁索引'; COMMENT ON INDEX idx_carrier_traffic_threshold_lock_card IS '通道阈值按卡查当前周期锁索引'; -COMMENT ON INDEX idx_carrier_traffic_threshold_lock_submitted IS '通道阈值恢复扫描的已提交子任务部分索引'; +COMMENT ON INDEX idx_carrier_traffic_threshold_lock_unresolved IS '通道阈值恢复扫描的未决子任务部分索引(含结果未知与失败,谓词与 domain.UnresolvedTaskStatuses 一致)'; COMMENT ON INDEX uq_carrier_traffic_threshold_lock_key IS '通道阈值周期锁唯一键(软删感知):一个卡在一个运营商计费周期内至多一条,冲突即该周期已处理'; diff --git a/openspec/changes/add-carrier-channel-traffic-thresholds/design.md b/openspec/changes/add-carrier-channel-traffic-thresholds/design.md index 810608a..f24a5dc 100644 --- a/openspec/changes/add-carrier-channel-traffic-thresholds/design.md +++ b/openspec/changes/add-carrier-channel-traffic-thresholds/design.md @@ -43,30 +43,32 @@ period_start = M.AddDate(0, -1, 0) if now < M // 上月重置日 0 点 - 锁唯一键 `(carrier_id, card_id, period_start)`;**postgres 23505 单独捕获为「该周期已处理」并跳过**,不得与既有流量基线 CAS 冲突(`CodeConflict`)混流——两类冲突语义不同,前者幂等跳过,后者重放重试。 - **改阈值不生效于当前周期已持锁卡**:唯一键已占位 = 该周期已处理,当前周期维持拒绝复机;新周期按新阈值判断。 - **停机消费者**(worker,消费停机 Outbox 事件): - 1. 提交认领:锁行 `stop_submitted_at IS NULL` 条件更新(refundchannel `claimChannelSubmission` 范式),至多一次执行;认领失败 = 已提交过 → 只走恢复查询,绝不重复调用。 + 1. 提交认领:锁行 `stop_submitted_at IS NULL AND status='locked'` 条件更新(refundchannel `claimChannelSubmission` 范式),至多一次执行;认领失败 = 已提交过或锁已被周期处理跨期解锁 → 只走恢复查询,绝不重复调用。 2. 卡已 `offline`(其他停因先行)→ 不调 Gateway,直接确认 success。 3. 否则复用 `stopCardWithRetry`(内部 3 次重试、Integration Log、统一审计);成功后写 `network_status=offline`、`stop_reason=channel_threshold`、`stopped_at`。 - 4. 失败/未知:保留锁与任务状态、记安全失败原因;unknown 的 Integration Log 必带 `RecoveryStrategy`。 + 4. 失败/未知:子任务落 `failed`/`unknown` 并记安全失败原因,锁行与历史结果一律保留;**两类结果都不退出收敛链路**——恢复扫描的未决集合为 `{submitted, unknown, failed}`,继续用只读状态查询确认(调用失败或响应超时并不代表运营商侧未生效),超期转人工;unknown 的 Integration Log 必带 `RecoveryStrategy`。 ### 4. 停因与复机拒绝 - 新增 `StopReasonChannelThreshold = "channel_threshold"`(`pkg/constants/iot.go` 停因组)。**不纳入 `isPollingStopReason`**(否则 `EvaluateAndAct` 离线分支会绕过周期逻辑直接自动复机),**不纳入 `isDeviceScopeReason`**(不扩散设备)。复机只由周期处理发起。 -- **持锁拒绝覆盖四个入口**,入口前置检查(锁存在 → 拒绝 + 审计 `denied` + 锁保留): - | 入口 | 覆盖路径 | - |---|---| - | `resumeSingleCard` | `EvaluateAndAct` 自动复机、`ResumeCardIfStopped`、`resumeDeviceCards` 遍历 | - | `ManualStartCard` | 手动复机 | - | `ForceStartCard` | 保护期强制复机(不加此拒绝,保护期一致性检查会强行复机持锁卡) | - | `StartMachineSeparatedCard` | 机卡分离复机 | +- **持锁拒绝覆盖四个入口**,入口前置检查(锁存在 → 拒绝 + 审计 `denied` + 锁保留;顺序先于既有各类前置校验与上游调用): + | 入口 | 覆盖路径 | 拒绝时的返回 | + |---|---|---| + | `resumeSingleCard` | `EvaluateAndAct` 自动复机、`ResumeCardIfStopped`、`resumeDeviceCards` 遍历 | 返回 nil(自动链路不把拒绝当失败重试),但仍写 `denied` 审计 | + | `ManualStartCard` | 手动复机 | 返回 `CodeForbidden` + 中文拒绝文案 | + | `ForceStartCard` | 保护期强制复机(**先持锁拒绝,再走既有保护期一致性检查**,否则会强行复机锁定期内的卡) | 返回 `CodeForbidden` + 中文拒绝文案 | + | `StartMachineSeparatedCard` | 机卡分离复机 | 返回 `CodeForbidden` + 中文拒绝文案 | - **「其他停机锁」定义**:`stop_reason` 非空且非 `channel_threshold`(arrears/manual/carrier_stopped/protect_period 等),或 `gateway_extend` 为风险停机/销户(`isRiskGatewayExtend`)。 ### 5. 周期处理与恢复(两个独立 cron) 均 `@every 1m` + `asynq.Unique(10m)` + 无 payload,照 `TaskTypeRefundChannelRecovery` 形态注册(`cmd/worker/main.go`);独立 cron 而非挂流量轮询同路径,是因为不依赖「停机卡是否继续被轮询」,新周期后最迟 1 分钟处理。 -1. **周期处理**:扫描 `period_start` 已过期的持锁锁行(按**锁行自身 carrier** 的 `data_reset_day` 判断过期,支持换运营商后旧锁归属);条件更新认领解锁(`status=locked → unlocked`,防并发重复);逐卡评估:有效主套餐(`hasValidPackage`)+ 流量未耗尽(`isTrafficExhausted`)+ 实名 OK(`isRealnameOK`)+ 非风险 extend + 无其他停因;全满足写复机 Outbox 事件,任一不满足只解锁。 -2. **恢复扫描**:扫描 `stop_status`/`resume_status` 为 `submitted` 的锁;**只查询网关状态回填,绝不重复发起停复机**;确认成功 → 回填任务状态并补写卡状态(覆盖「Gateway 成功但 DB 更新失败」场景);自提交起超过 **30 分钟**(常量定义)仍不可查 → `anomaly_flag=1` + 安全失败原因,退出扫描转人工,不自动删除锁。 -- **复机消费者**:结构同停机消费者,认领字段 `resume_submitted_at`,复用 `resumeCardWithRetry`,成功后写 `network_status=online`、`stop_reason=""`、`resumed_at`。 +1. **周期处理**:扫描待处理的持锁锁行(`status=locked`),按**锁行自身 carrier** 的 `data_reset_day` 分两种结果处理: + - 已跨期(`period_start` 早于该 carrier 的当前周期起点,支持换运营商后旧锁归属):条件更新认领解锁(`status=locked → unlocked`,防并发重复);逐卡评估有效主套餐(`hasValidPackage`)+ 流量未耗尽(`isTrafficExhausted`)+ 实名 OK(`isRealnameOK`)+ 非风险 extend + 无其他停因;全满足写复机 Outbox 事件,任一不满足只解锁并把原因写入锁行。 + - 周期归属不可判定(锁行引用的 carrier 已不存在/软删,或 `data_reset_day` 非法):**同样认领解锁**(按 spec「新周期对仍持锁卡解除通道锁」的语义)**并同事务 `anomaly_flag=1` + 安全原因**、Warn 日志转人工;不写复机事件、不调运营商。该分支是必须的出路:`ActiveLock` 对配置缺失按仍未生效处理(fail-closed,不放开复机),若周期处理也跳过这些行,持锁卡将永久禁止一切复机且无自动出路。 +2. **恢复扫描**:扫描 `anomaly_flag=0` 且 `stop_status`/`resume_status` 处于**未决集合 `{submitted, unknown, failed}`** 的锁;**只查询网关状态回填,绝不重复发起停复机**;确认成功 → 回填任务状态(`confirmed` 只写一次,条件更新按未决集合为谓词)并补写卡状态(覆盖「Gateway 成功但 DB 更新失败」场景);自提交起超过 **30 分钟**(常量定义)仍不可查 → `anomaly_flag=1` + 安全失败原因,退出扫描转人工,不自动删除锁。 +- **复机消费者**:结构同停机消费者,认领字段 `resume_submitted_at`,复用 `resumeCardWithRetry`,成功后写 `network_status=online`、`resumed_at`,并**只在该卡停因正是 `channel_threshold` 时清除停因**(不覆盖 arrears/manual 等其他停因)。 - 周期处理与恢复 Handler 审计上下文固定 `ActorKind=AuditActorScheduledJob`、`Source=AuditSourceScheduler`。 ### 6. 权限与审计 @@ -75,12 +77,29 @@ period_start = M.AddDate(0, -1, 0) if now < M // 上月重置日 0 点 - 审计并入既有 carrier 更新审计(`writeAudit` 整体快照前后值),不新建独立审计动作;启停即 Update 的一部分,前后值自然覆盖。 - 通道停用:`enabled=false` 只停止新锁创建;已持锁卡当前周期继续拒绝复机,新周期解锁/复机照常;锁与历史保留。 -### 7. 锁表结构 +### 7. 审计决定(ENG-AUDIT-001) + +按 ENG-AUDIT-001「用例 MUST 明确 Audit Event、Domain Ledger、Integration Log、Outbox 的使用决定或 N/A 理由」,本能力的四类事实归属如下: + +- **Audit Event(有,复用既有动作,不新建动作)**:卡状态侧的用户可见结果与关键拒绝——达量停机成功/失败(`iot_card.auto_stop`)、新周期复机成功(`iot_card.auto_start`)、四入口持锁拒绝(同上动作 + `result=denied`)——全部走停复机单一事实源的统一审计写入(与卡状态更新同事务)。 +- **Domain Ledger(N/A)**:本能力不涉及资金与额度,无独立账本事实。 +- **Integration Log(有)**:每次运营商停复机调用与每次恢复扫描的状态查询都写 Integration Log(`stop_card`/`start_card`/`query_card_status`);结果未知必带 `RecoveryStrategy`(既有基座强制)。 +- **Outbox(有)**:达量停机事件与周期复机事件与锁行事实同事务写入(ENG-OUTBOX-001)。 +- **锁行状态迁移(`locked→unlocked`、`anomaly_flag=1`)的 N/A 决定**:这两类迁移是**本能力内部的任务编排状态**,不是用户可见的业务状态变更,因此**不新建独立审计动作**;其承载方式为:锁行自身的 `status`/`failure_reason`/`updated_at` 留痕 + 周期处理计划任务的 Warn 日志(含 lock_id/card_id/carrier_id 与原因)+ 该锁驱动的卡状态变更由上述 Audit Event 覆盖。该决定在此显式登记,四类事实不得互相替代。 +- 读卡/写审计类基础设施故障使结果无法判定时,子任务保持 `submitted`,由恢复扫描按窗口收敛或转人工——不伪造终态。 + +### 8. 锁表结构 `carrier_id`、`card_id`、`period_start`(timestamptz,唯一键三列组合,软删感知)、`status`(`locked`/`unlocked`)、任务状态组(`stop_status`/`resume_status`:`pending`/`submitted`/`confirmed`/`failed`/`unknown`)、提交认领字段(`stop_submitted_at`/`resume_submitted_at`)、`anomaly_flag`、`failure_reason`、时间戳。成对迁移(up/down)。 +### 9. 既有能力的行为变更(Modified Capability) + +持通道阈值锁的卡在锁定期内不再调用上游复机流程,这改变了既有 `package-lifecycle`「套餐状态流转」中「支付后已生效主套餐异步尝试自动复机」与「套餐激活/重置复机」的可观察行为(新增一条前置条件:持通道阈值锁时拒绝复机)。该变更以 `openspec/changes/add-carrier-channel-traffic-thresholds/specs/package-lifecycle/spec.md` 的 MODIFIED delta 显式建模(复述原 Requirement 全文并加入持锁条件),主 Spec 的同步在归档时进行。 + ## Risks / Trade-offs - 达量判定嵌入 `ApplyTrafficObservation` 事务:该事务变重(多一次 carrier 查询 + 锁插入)。收益是并发安全与数据最新;风险是事务失败回滚会连同流量事实一起回滚——既有 CAS 冲突已按此语义处理,行为一致。 - 停机成功写 `stop_reason=channel_threshold` 会覆盖卡上既有停因字段;仅当停机消费者确认 Gateway 成功后写入,且持锁期本就该拒绝其他路径,覆盖可接受。 - 恢复扫描「只查询回填」依赖网关状态查询接口可用;持续不可用 → 30 分钟超期转人工,锁保留,无数据丢失。 +- 停机/复机调用明确失败(`failed`)与结果未知(`unknown`)都继续留在收敛链路上:每分钟只做一次**只读**状态查询,绝不重试外呼;这在极端情况下会让同一笔结果在 30 分钟内被查询数十次,代价换来「调用已生效但响应丢失」场景能被自动纠正。 +- 锁行引用的 carrier 被删除/软删,或 `data_reset_day` 非法时,周期处理按跨期语义解锁并标记异常转人工:解锁意味着该卡不再被通道锁拒绝复机(配置缺失时无法判断是否仍在周期内),异常标记与 Warn 日志保证运维可见。本能力不新增 carrier 删除前置校验(既有 Change 边界),该场景以计划任务 + 锁行留痕转人工。 diff --git a/openspec/changes/add-carrier-channel-traffic-thresholds/proposal.md b/openspec/changes/add-carrier-channel-traffic-thresholds/proposal.md index 36659e8..118a8cb 100644 --- a/openspec/changes/add-carrier-channel-traffic-thresholds/proposal.md +++ b/openspec/changes/add-carrier-channel-traffic-thresholds/proposal.md @@ -17,7 +17,7 @@ - `carrier-channel-traffic-threshold`: 通道流量阈值控制。 ### Modified Capabilities -- 无。 +- `package-lifecycle`: 「套餐状态流转」新增一条复机前置条件——载体在当前计费周期内持有运营商通道流量阈值停机锁时,支付后异步自动复机与套餐激活/重置复机 MUST NOT 调用上游复机流程,并保留该锁(以 `specs/package-lifecycle/spec.md` 的 MODIFIED delta 建模;主 Spec 同步在归档时进行)。 ## Impact @@ -27,5 +27,6 @@ - 新增周期阈值停机锁表:`(carrier_id, card_id, period_start)` 唯一键、任务状态组、提交认领字段、`anomaly_flag`;成对迁移。 - 新增 2 个 Outbox 事件类型:停机、复机各一,worker 消费。 - 新增 2 个 cron 与 worker 装配点:周期处理(解锁/条件复机)、恢复扫描(只查询回填),均 `@every 1m` + `asynq.Unique(10m)`,照 `TaskTypeRefundChannelRecovery` 形态注册。 -- 新增停因常量 `channel_threshold`;4 个复机入口新增持锁拒绝。 +- 新增停因常量 `channel_threshold`;4 个复机入口(`resumeSingleCard`、`ManualStartCard`、`ForceStartCard`、`StartMachineSeparatedCard`)新增持锁拒绝,因而改变了既有套餐生命周期自动复机路径的可观察行为(见 Modified Capabilities)。 +- 停机/复机调用明确失败或结果未知的子任务不退出收敛链路:恢复扫描的未决集合为 `{submitted, unknown, failed}`(只读状态查询,绝不重试外呼),超 30 分钟仍不可确认转人工;锁行引用的运营商缺失或重置日非法时,周期处理按跨期语义解锁并标记异常转人工。 - carrier 管理接口字段级守卫与响应过滤(阈值字段仅平台账号可见可写)。 diff --git a/openspec/changes/add-carrier-channel-traffic-thresholds/specs/package-lifecycle/spec.md b/openspec/changes/add-carrier-channel-traffic-thresholds/specs/package-lifecycle/spec.md new file mode 100644 index 0000000..69a2b26 --- /dev/null +++ b/openspec/changes/add-carrier-channel-traffic-thresholds/specs/package-lifecycle/spec.md @@ -0,0 +1,55 @@ +## MODIFIED Requirements + +### Requirement: 套餐状态流转 + +系统 SHALL 按当前套餐和套餐使用状态控制上架、订购、激活、失效与到期处理。主套餐到期时,系统 MUST 先确定同一载体是否存在待生效的后续主套餐:存在时,后续套餐激活与停复机重新评估 MUST 由同一条顺序流程完成;系统 MUST NOT 依据后续套餐激活前的无套餐快照发起停机。后续套餐成功生效后,系统 MUST 依据最新套餐、流量和实名事实重新判断卡网络状态,且不得遗留 `no_package` 停机。不存在后续套餐或后续套餐经业务校验不能生效时,系统 SHALL 按现有停机规则评估卡状态。后续套餐激活结果未知或任务投递失败不得被当作无后续套餐处理并据此停机,系统 SHALL 保留既有激活恢复与轮询兜底路径。 + +套餐订单支付成功时,系统 MUST 在支付和套餐权益事务提交后,检查本订单是否存在已生效且未挂靠其他主套餐的主套餐权益。存在时,系统 MUST 基于已提交的套餐、流量、实名和停机原因事实异步尝试自动复机;支付回调不得等待上游复机结果。主套餐权益处于待生效、待实名生效或其他非生效状态时,系统 MUST NOT 因本次支付发起自动复机。载体在当前计费周期内持有运营商通道流量阈值停机锁时,系统 MUST NOT 调用上游复机流程,并 MUST 保留该通道阈值锁与既有网络状态。自动复机失败 SHALL 保留既有失败记录与轮询兜底机制。 + +#### Scenario: 支付后已生效主套餐触发自动复机 + +- **GIVEN** 套餐订单支付成功后,本订单存在已生效且未挂靠其他主套餐的主套餐权益,载体因可自动恢复的原因处于停机状态 +- **WHEN** 支付和套餐权益事务提交成功 +- **THEN** 系统异步按最新套餐、流量、实名和停机原因事实检查载体,并在全部复机条件满足时调用既有上游复机流程;支付回调不等待该调用结束 + +#### Scenario: 支付后主套餐未生效不触发自动复机 + +- **GIVEN** 套餐订单支付成功后,本订单主套餐权益仍处于待生效、待实名生效或其他非生效状态 +- **WHEN** 支付和套餐权益事务提交成功 +- **THEN** 系统不因本次支付发起自动复机,并保留后续套餐激活和轮询处理 + +#### Scenario: 支付后复机条件不满足 + +- **GIVEN** 套餐订单支付成功后,本订单存在已生效主套餐权益 +- **WHEN** 载体为手动停机、无有效套餐、流量已耗尽、不满足实名策略,或处于运营商通道流量阈值锁定期 +- **THEN** 系统不调用上游复机流程,保留当前网络状态、既有通道阈值锁与既有停复机处理路径 + +#### Scenario: 套餐状态流转 + +- **GIVEN** 套餐或使用记录处于允许的前置状态 +- **WHEN** 执行状态操作 +- **THEN** 仅发生一次允许的状态变化;不满足前置状态时返回业务错误 + +#### Scenario: 到期主套餐接续后续套餐 + +- **GIVEN** 某载体的当前主套餐到期,且存在满足激活条件的待生效后续主套餐 +- **WHEN** 系统处理该主套餐到期 +- **THEN** 系统先完成后续套餐激活并按最新权益事实重新评估停复机,且不得因到期前的无套餐快照对该载体发起 `no_package` 停机 + +#### Scenario: 到期主套餐无后续可生效套餐 + +- **GIVEN** 某载体的当前主套餐到期,且不存在后续主套餐或队首后续套餐不满足激活条件 +- **WHEN** 系统完成该套餐到期处理 +- **THEN** 系统按当前套餐、流量和实名事实执行既有停机评估 + +#### Scenario: 后续套餐激活结果未知 + +- **GIVEN** 某载体的当前主套餐到期,存在待生效后续主套餐,但激活任务投递或执行结果暂时未知 +- **WHEN** 系统处理该套餐到期 +- **THEN** 系统不得将该未知结果视为不存在后续套餐而依据旧快照发起停机,并保留既有激活恢复与套餐轮询兜底 + +#### Scenario: 卡状态轮询发现缺失的套餐任务 + +- **GIVEN** 启用轮询的卡匹配套餐检查配置,且其 `polling:package` 分片队列项因异常缺失 +- **WHEN** 卡状态轮询成功完成且未命中风险停机 +- **THEN** 系统基于最新卡状态仅补入缺失的套餐任务,不改写已存在套餐任务的执行时间;后续套餐任务仍按既有停复机条件评估该卡 diff --git a/openspec/changes/add-carrier-channel-traffic-thresholds/tasks.md b/openspec/changes/add-carrier-channel-traffic-thresholds/tasks.md index c1cb02d..08673b0 100644 --- a/openspec/changes/add-carrier-channel-traffic-thresholds/tasks.md +++ b/openspec/changes/add-carrier-channel-traffic-thresholds/tasks.md @@ -28,3 +28,12 @@ ## 7. 验证 - [x] 7.1 隔离库验证:平台配置阈值(含 GB 换算)成功 + 审计前后值,代理/企业写被拒且读不到字段;达量触发停机、同周期重复判定唯一冲突不重复锁/不重复停机;持锁期四入口复机拒绝且锁保留;停机 unknown 保留 + RecoveryStrategy + 恢复扫描确认补写卡状态、超期 anomaly 转人工;新周期解锁、条件满足才复机、其他停因(stop_reason 非空非本停因或风险 extend)只解锁;停用后新卡不再锁、已持锁卡新周期照常解锁复机;换运营商旧锁保留、新周期按新 carrier 重置日;第 12 项套餐预警与既有 3 种停因停复机回归不变;迁移 up/down/up。 - [x] 7.2 运行 `gofmt -w`、`go build ./cmd/api ./cmd/worker`、`go run cmd/gendocs/main.go`、`openspec validate add-carrier-channel-traffic-thresholds --strict` 与 `openspec doctor --json`;自动化测试按项目决策为 N/A。 + +## 8. 审查修复 +- [x] 8.1(B1)停机/复机调用明确失败或结果未知的子任务不退出恢复扫描:扫描谓词与超期判定改用未决集合 `{submitted, unknown, failed}`、`markTaskConfirmed` 的 expected 改为未决集合(confirmed 仍只写一次)、部分索引谓词同步为未决集合(000227,已确认生产库未应用)。 +- [x] 8.2(B2)锁行引用的运营商缺失或 `data_reset_day` 非法时给出出路:周期处理按跨期语义认领解锁并同事务标记 `anomaly_flag=1` + 安全原因 + Warn 日志转人工,不写复机事件、不调运营商、不静默跳过;`ActiveLock` 维持配置缺失时不放开复机的 fail-closed 语义。 +- [x] 8.3(R1)恢复 `stopCardWithRetry`/`resumeCardWithRetry` 在 Integration Log 写入失败分支的「未完成」审计,保持既有审计覆盖与返回语义不变。 +- [x] 8.4(R2)按 ENG-AUDIT-001 显式登记锁行状态迁移(`unlocked`、`anomaly_flag=1`)的审计决定(不新建独立审计动作 + 承载方式),并在 design.md 补「审计决定」小节。 +- [x] 8.5(R4)以 `specs/package-lifecycle/spec.md` MODIFIED delta 建模「持通道阈值锁时拒绝复机」对既有套餐生命周期自动复机路径的行为变更,并更新 proposal 的 Modified Capabilities。 +- [x] 8.6 修正 design.md 与实现不符的表述(失败/未知与不可判定 carrier 的收敛路径、四入口拒绝顺序与返回、认领谓词)。 +- [x] 8.7 回归验证:B1(unknown/failed 扫描与收敛、超期 anomaly)、B2(carrier 缺失/重置日非法的解锁与异常)、R1(审计写入)在隔离库实测;gofmt/go build/go vet/openspec validate --strict/context-health 全通过。