diff --git a/backend/internal/service/channel_monitor_runner.go b/backend/internal/service/channel_monitor_runner.go index 748ffc3638..dbf03c2395 100644 --- a/backend/internal/service/channel_monitor_runner.go +++ b/backend/internal/service/channel_monitor_runner.go @@ -2,6 +2,7 @@ package service import ( "context" + "errors" "log/slog" "math/rand/v2" "sync" @@ -139,14 +140,14 @@ func (r *ChannelMonitorRunner) Start() { } // Schedule 为指定监控创建(或重置)独立定时任务。 -// - m.Enabled=false → 等同于 Unschedule(m.ID) +// - m.Enabled=false 或 APIKeyDecryptFailed=true → 等同于 Unschedule(m.ID) // - 已存在的任务会先被取消再重建(适用于 IntervalSeconds 变更场景) // - 新任务立即触发首次检测,之后按 IntervalSeconds 周期触发 func (r *ChannelMonitorRunner) Schedule(m *ChannelMonitor) { if r == nil || m == nil { return } - if !m.Enabled { + if !m.Enabled || m.APIKeyDecryptFailed { r.Unschedule(m.ID) return } @@ -293,7 +294,7 @@ func (r *ChannelMonitorRunner) releaseInFlight(id int64) { r.inFlightMu.Unlock() } -// runOne 执行单个监控的检测。所有错误只记日志,不熔断。 +// runOne 执行单个监控的检测。普通错误只记日志;API key 解密失败会撤销任务。 // 任务结束时(含 panic recover)必须释放 in-flight 槽。 func (r *ChannelMonitorRunner) runOne(id int64, name string) { ctx, cancel := context.WithTimeout(context.Background(), monitorRequestTimeout+monitorPingTimeout+monitorRunOneBuffer) @@ -309,6 +310,9 @@ func (r *ChannelMonitorRunner) runOne(id int64, name string) { }() if _, err := r.svc.RunCheck(ctx, id); err != nil { + if errors.Is(err, ErrChannelMonitorAPIKeyDecryptFailed) { + r.Unschedule(id) + } slog.Warn("channel_monitor: run check failed", "monitor_id", id, "name", name, "error", err) } diff --git a/backend/internal/service/channel_monitor_runner_test.go b/backend/internal/service/channel_monitor_runner_test.go index 5eed3c2092..e1a9b3cccb 100644 --- a/backend/internal/service/channel_monitor_runner_test.go +++ b/backend/internal/service/channel_monitor_runner_test.go @@ -159,6 +159,70 @@ func TestSchedule_DisabledRedirectsToUnschedule(t *testing.T) { stoppedWithin(t, r, 3*time.Second) } +func TestSchedule_DecryptFailedRedirectsToUnschedule(t *testing.T) { + svc := &stubMonitorSvc{runCalled: make(chan int64, 4)} + r := newRunnerForTest(svc) + r.Start() + + r.Schedule(&ChannelMonitor{ID: 10, Enabled: true, IntervalSeconds: 60}) + waitFor(t, time.Second, "task registered", func() bool { return runnerTaskCount(r) == 1 }) + + r.Schedule(&ChannelMonitor{ID: 10, Enabled: true, IntervalSeconds: 60, APIKeyDecryptFailed: true}) + if got := runnerTaskCount(r); got != 0 { + t.Fatalf("expected tasks empty after decrypt-failed re-Schedule, got %d", got) + } + + stoppedWithin(t, r, 3*time.Second) +} + +func TestSchedule_RepairedAPIKeyCanBeScheduled(t *testing.T) { + svc := &stubMonitorSvc{runCalled: make(chan int64, 1)} + r := newRunnerForTest(svc) + r.Start() + + r.Schedule(&ChannelMonitor{ID: 11, Enabled: true, IntervalSeconds: 60, APIKeyDecryptFailed: true}) + if got := runnerTaskCount(r); got != 0 { + t.Fatalf("expected no task for decrypt-failed monitor, got %d", got) + } + + r.Schedule(&ChannelMonitor{ID: 11, Enabled: true, IntervalSeconds: 60, APIKey: "replacement-key"}) + if got := runnerTaskCount(r); got != 1 { + t.Fatalf("expected repaired monitor to be scheduled, got %d tasks", got) + } + select { + case id := <-svc.runCalled: + if id != 11 { + t.Fatalf("expected repaired monitor id=11 to fire, got %d", id) + } + case <-time.After(2 * time.Second): + t.Fatal("expected repaired monitor to fire immediately") + } + + stoppedWithin(t, r, 3*time.Second) +} + +func TestRunOne_DecryptFailureUnschedulesTask(t *testing.T) { + svc := &stubMonitorSvc{ + runCalled: make(chan int64, 1), + runErr: ErrChannelMonitorAPIKeyDecryptFailed, + } + r := newRunnerForTest(svc) + r.Start() + + r.Schedule(&ChannelMonitor{ID: 12, Enabled: true, IntervalSeconds: 60}) + select { + case id := <-svc.runCalled: + if id != 12 { + t.Fatalf("expected failing monitor id=12 to fire, got %d", id) + } + case <-time.After(2 * time.Second): + t.Fatal("expected failing monitor to fire immediately") + } + waitFor(t, time.Second, "decrypt-failed task unscheduled", func() bool { return runnerTaskCount(r) == 0 }) + + stoppedWithin(t, r, 3*time.Second) +} + // TestSchedule_InvalidIntervalSkipped 验证 IntervalSeconds<=0 不会注册任务(防御性检查)。 func TestSchedule_InvalidIntervalSkipped(t *testing.T) { svc := &stubMonitorSvc{} @@ -201,6 +265,25 @@ func TestStart_LoadsAllEnabledMonitors(t *testing.T) { stoppedWithin(t, r, 3*time.Second) } +func TestStart_SkipsDecryptFailedMonitor(t *testing.T) { + svc := &stubMonitorSvc{ + enabled: []*ChannelMonitor{ + {ID: 4, Enabled: true, IntervalSeconds: 60, APIKeyDecryptFailed: true}, + }, + } + r := newRunnerForTest(svc) + r.Start() + + if got := runnerTaskCount(r); got != 0 { + t.Fatalf("expected no task for decrypt-failed startup monitor, got %d", got) + } + if got := svc.runCount.Load(); got != 0 { + t.Fatalf("expected decrypt-failed startup monitor not to run, got %d calls", got) + } + + stoppedWithin(t, r, 3*time.Second) +} + // TestStop_DrainsAllGoroutines 验证 Stop 会等待所有调度 goroutine 退出(无游离)。 func TestStop_DrainsAllGoroutines(t *testing.T) { svc := &stubMonitorSvc{}