diff --git a/app/models/check_aggregator_test.go b/app/models/check_aggregator_test.go index e14e1f2..c87fe48 100644 --- a/app/models/check_aggregator_test.go +++ b/app/models/check_aggregator_test.go @@ -161,6 +161,85 @@ func TestApplyRemoteCheckResult_DirectWhenQuorumOne(t *testing.T) { "the region result row is always inserted even on the legacy path") } +func TestApplyRemoteCheckResultFromWorkerDeduplicatesFailureReplay(t *testing.T) { + models.Drop() + models.Migrate() + + mon, check := seedAggregatorWorld(t, 1, 5) + seedRegion(t, "ru-msk") + worker := models.WorkerNode{ + WorkerID: "replay-worker", + RegionCode: "ru-msk", + AuthToken: "replay-token", + } + require.NoError(t, models.DB().Create(&worker).Error) + report := makeReport(check.ID, mon.ID, "ERR") + report.JobID = "same-failed-execution" + + require.NoError(t, models.ApplyRemoteCheckResultFromWorker(report, worker.RegionCode, &worker)) + require.NoError(t, models.ApplyRemoteCheckResultFromWorker(report, worker.RegionCode, &worker)) + + got := loadCheck(t, check.ID) + assert.Equal(t, 1, got.Fails) + assert.EqualValues(t, 1, countPendingResults(t, check.ID)) + var attempts int64 + require.NoError(t, models.DB().Model(&models.CheckAttempt{}).Where("job_id = ?", report.JobID).Count(&attempts).Error) + assert.EqualValues(t, 1, attempts) +} + +func TestApplyRemoteCheckResultFromWorkerDeduplicatesQuorumReplay(t *testing.T) { + models.Drop() + models.Migrate() + + mon, check := seedAggregatorWorld(t, 2, 5) + seedRegion(t, "ru-msk") + worker := models.WorkerNode{ + WorkerID: "quorum-replay-worker", + RegionCode: "ru-msk", + AuthToken: "quorum-replay-token", + } + require.NoError(t, models.DB().Create(&worker).Error) + report := makeReport(check.ID, mon.ID, "OK") + report.JobID = "same-quorum-execution" + + require.NoError(t, models.ApplyRemoteCheckResultFromWorker(report, worker.RegionCode, &worker)) + require.NoError(t, models.ApplyRemoteCheckResultFromWorker(report, worker.RegionCode, &worker)) + + got := loadCheck(t, check.ID) + assert.Equal(t, "UNK", got.State, "quorum state remains owned by the aggregator") + assert.EqualValues(t, 1, countPendingResults(t, check.ID)) + var attempts int64 + require.NoError(t, models.DB().Model(&models.CheckAttempt{}).Where("job_id = ?", report.JobID).Count(&attempts).Error) + assert.EqualValues(t, 1, attempts) +} + +func TestApplyRemoteCheckResultFromWorkerDeduplicatesFailReplay(t *testing.T) { + models.Drop() + models.Migrate() + + mon, check := seedAggregatorWorld(t, 1, 5) + seedRegion(t, "ru-msk") + worker := models.WorkerNode{ + WorkerID: "fail-replay-worker", + RegionCode: "ru-msk", + AuthToken: "fail-replay-token", + } + require.NoError(t, models.DB().Create(&worker).Error) + report := makeReport(check.ID, mon.ID, "FAIL") + report.JobID = "same-fail-execution" + + require.NoError(t, models.ApplyRemoteCheckResultFromWorker(report, worker.RegionCode, &worker)) + require.NoError(t, models.ApplyRemoteCheckResultFromWorker(report, worker.RegionCode, &worker)) + + got := loadCheck(t, check.ID) + assert.Equal(t, "FAIL", got.State) + assert.Equal(t, 0, got.Fails, "FAIL resets consecutive ERR tracking") + assert.EqualValues(t, 1, countPendingResults(t, check.ID)) + var attempts int64 + require.NoError(t, models.DB().Model(&models.CheckAttempt{}).Where("job_id = ?", report.JobID).Count(&attempts).Error) + assert.EqualValues(t, 1, attempts) +} + // TestApplyRemoteCheckResult_BuffersWhenQuorumN pins the new path: // when RequireQuorum > 1, ApplyRemoteCheckResult does NOT touch // Check.State — it only inserts the CheckRegionResult row. The check diff --git a/app/models/check_jobs.go b/app/models/check_jobs.go index 80d09d9..82f27b9 100644 --- a/app/models/check_jobs.go +++ b/app/models/check_jobs.go @@ -275,91 +275,141 @@ func ApplyRemoteCheckResultFromWorkerTx(tx *gorm.DB, report wire.CheckResultRepo return nil, fmt.Errorf("apply check result: nil transaction") } now := time.Now() - if worker != nil { - handled := false - if err := ApplyDiagnosticResultTx(tx, report, worker, now, &handled); err != nil { - return nil, err - } - if handled { - return nil, nil - } + + handled, err := handleDiagnosticResult(tx, report, worker, now) + if err != nil { + return nil, err } - check := Check{} - if err := tx.Preload("Monitor").First(&check, report.CheckID).Error; err != nil { - log.Println("worker: check not found:", report.CheckID, err) + if handled { + return nil, nil + } + + check, err := loadCheckWithMonitor(tx, report.CheckID) + if err != nil { return nil, err } - // Always persist the per-region result first so the aggregator can - // pick it up regardless of which path we take next. We rely on - // StoreCheckRegionResult to default AggregatedAt=NULL (the column - // type is *time.Time, so a zero value writes SQL NULL). + skipped, err := recordCheckAttempt(tx, check, report, worker, now) + if err != nil { + return nil, err + } + if skipped { + return nil, nil + } + if err := StoreCheckRegionResultTx(tx, report, regionCode, now); err != nil { log.Println("worker: error storing region result:", report.CheckID, err) return nil, err } - // Quorum-enabled checks: write nothing to Check.State here. The - // aggregator will compute the aggregate state once the window has - // elapsed (or enough regions have reported) and stamp AggregatedAt on - // the contributing CheckRegionResult rows. if check.QuorumEnabled() { return nil, nil } + if err := updateCheckState(tx, &check, report, now); err != nil { + return nil, err + } + + if err := handleCheckStateSideEffects(tx, &check, report.State, worker, now); err != nil { + return nil, err + } + return check.Monitor, nil +} + +func handleDiagnosticResult(tx *gorm.DB, report wire.CheckResultReport, worker *WorkerNode, now time.Time) (bool, error) { + if worker == nil { + return false, nil + } + handled := false + if err := ApplyDiagnosticResultTx(tx, report, worker, now, &handled); err != nil { + return false, err + } + return handled, nil +} + +func loadCheckWithMonitor(tx *gorm.DB, checkID int64) (Check, error) { + check := Check{} + if err := tx.Preload("Monitor").First(&check, checkID).Error; err != nil { + log.Println("worker: check not found:", checkID, err) + return check, err + } + return check, nil +} + +func recordCheckAttempt(tx *gorm.DB, check Check, report wire.CheckResultReport, worker *WorkerNode, now time.Time) (bool, error) { + if worker == nil { + return false, nil + } + payload, _ := json.Marshal(report) + attempt := CheckAttempt{JobID: report.JobID, CheckID: check.ID, MonitorID: check.MonitorID, WorkerNodeID: &worker.ID, Kind: AttemptKindRegular, State: AttemptStateFinished, ResultState: report.State, Result: payload, StartedAt: &now, FinishedAt: &now, Deweighted: worker.NetworkProblemActive(now)} + if attempt.JobID == "" { + attempt.JobID = uuid.New().String() + } + created := tx.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "job_id"}}, + DoNothing: true, + }).Create(&attempt) + if created.Error != nil { + return false, created.Error + } + return created.RowsAffected == 0, nil +} + +func updateCheckState(tx *gorm.DB, check *Check, report wire.CheckResultReport, now time.Time) error { + update := buildCheckStateUpdate(report, now) + if err := tx.Model(check).UpdateColumns(update).Error; err != nil { + log.Println("worker: error updating check:", report.CheckID, err) + return err + } + return nil +} + +func buildCheckStateUpdate(report wire.CheckResultReport, now time.Time) map[string]interface{} { update := map[string]interface{}{ colState: report.State, colLastEnd: now, colWarnings: pq.StringArray(report.Warnings), colInfos: pq.StringArray(report.Infos), } - - if report.State == "OK" { + switch report.State { + case stateOK: update["was_up"] = now update["last_ok"] = now update["fails"] = 0 update["error"] = gorm.Expr("NULL") - } else { + case stateERR: update["last_fail"] = now update["fails"] = gorm.Expr("fails + 1") if report.Error != nil { update["error"] = *report.Error } + default: + update["last_fail"] = now + update["fails"] = 0 + if report.Error != nil { + update["error"] = *report.Error + } } - if report.ExpiresAt != nil { t, err := time.Parse(time.RFC3339, *report.ExpiresAt) if err == nil { update["expires"] = t } } + return update +} - if err := tx.Model(&check).UpdateColumns(update).Error; err != nil { - log.Println("worker: error updating check:", report.CheckID, err) - return nil, err +func handleCheckStateSideEffects(tx *gorm.DB, check *Check, state string, worker *WorkerNode, now time.Time) error { + if worker == nil { + return nil } - if worker != nil { - payload, _ := json.Marshal(report) - attempt := CheckAttempt{JobID: report.JobID, CheckID: check.ID, MonitorID: check.MonitorID, WorkerNodeID: &worker.ID, Kind: AttemptKindRegular, State: AttemptStateFinished, ResultState: report.State, Result: payload, StartedAt: &now, FinishedAt: &now, Deweighted: worker.NetworkProblemActive(now)} - if attempt.JobID == "" { - attempt.JobID = uuid.New().String() - } - // A duplicate websocket/HTTP delivery must not create another attempt. - if err := tx.Where("job_id = ?", attempt.JobID).FirstOrCreate(&attempt).Error; err != nil { - return nil, err - } - switch report.State { - case stateERR, stateFail: - if err := StartConfirmationTx(tx, check.ID, worker.ID, now); err != nil { - return nil, err - } - case stateOK: - if err := RecoverDiagnosticTx(tx, check.ID, now); err != nil { - return nil, err - } - } + switch state { + case stateERR, stateFail: + return StartConfirmationTx(tx, check.ID, worker.ID, now) + case stateOK: + return RecoverDiagnosticTx(tx, check.ID, now) } - return check.Monitor, nil + return nil } // StoreRemoteCheckMetrics persists TSDB points reported by a distributed worker.