fix(worker): deduplicate check result replays
Все проверки выполнены успешно
CI / test (push) Successful in 2m7s
Docker / Build and publish worker image (push) Successful in 22m17s
Все проверки выполнены успешно
CI / test (push) Successful in 2m7s
Docker / Build and publish worker image (push) Successful in 22m17s
Этот коммит содержится в:
@@ -161,6 +161,85 @@ func TestApplyRemoteCheckResult_DirectWhenQuorumOne(t *testing.T) {
|
|||||||
"the region result row is always inserted even on the legacy path")
|
"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:
|
// TestApplyRemoteCheckResult_BuffersWhenQuorumN pins the new path:
|
||||||
// when RequireQuorum > 1, ApplyRemoteCheckResult does NOT touch
|
// when RequireQuorum > 1, ApplyRemoteCheckResult does NOT touch
|
||||||
// Check.State — it only inserts the CheckRegionResult row. The check
|
// Check.State — it only inserts the CheckRegionResult row. The check
|
||||||
|
|||||||
@@ -275,91 +275,141 @@ func ApplyRemoteCheckResultFromWorkerTx(tx *gorm.DB, report wire.CheckResultRepo
|
|||||||
return nil, fmt.Errorf("apply check result: nil transaction")
|
return nil, fmt.Errorf("apply check result: nil transaction")
|
||||||
}
|
}
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
if worker != nil {
|
|
||||||
handled := false
|
handled, err := handleDiagnosticResult(tx, report, worker, now)
|
||||||
if err := ApplyDiagnosticResultTx(tx, report, worker, now, &handled); err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
|
||||||
if handled {
|
|
||||||
return nil, nil
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
check := Check{}
|
if handled {
|
||||||
if err := tx.Preload("Monitor").First(&check, report.CheckID).Error; err != nil {
|
return nil, nil
|
||||||
log.Println("worker: check not found:", report.CheckID, err)
|
}
|
||||||
|
|
||||||
|
check, err := loadCheckWithMonitor(tx, report.CheckID)
|
||||||
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
// Always persist the per-region result first so the aggregator can
|
skipped, err := recordCheckAttempt(tx, check, report, worker, now)
|
||||||
// pick it up regardless of which path we take next. We rely on
|
if err != nil {
|
||||||
// StoreCheckRegionResult to default AggregatedAt=NULL (the column
|
return nil, err
|
||||||
// type is *time.Time, so a zero value writes SQL NULL).
|
}
|
||||||
|
if skipped {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
|
||||||
if err := StoreCheckRegionResultTx(tx, report, regionCode, now); err != nil {
|
if err := StoreCheckRegionResultTx(tx, report, regionCode, now); err != nil {
|
||||||
log.Println("worker: error storing region result:", report.CheckID, err)
|
log.Println("worker: error storing region result:", report.CheckID, err)
|
||||||
return nil, 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() {
|
if check.QuorumEnabled() {
|
||||||
return nil, nil
|
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{}{
|
update := map[string]interface{}{
|
||||||
colState: report.State,
|
colState: report.State,
|
||||||
colLastEnd: now,
|
colLastEnd: now,
|
||||||
colWarnings: pq.StringArray(report.Warnings),
|
colWarnings: pq.StringArray(report.Warnings),
|
||||||
colInfos: pq.StringArray(report.Infos),
|
colInfos: pq.StringArray(report.Infos),
|
||||||
}
|
}
|
||||||
|
switch report.State {
|
||||||
if report.State == "OK" {
|
case stateOK:
|
||||||
update["was_up"] = now
|
update["was_up"] = now
|
||||||
update["last_ok"] = now
|
update["last_ok"] = now
|
||||||
update["fails"] = 0
|
update["fails"] = 0
|
||||||
update["error"] = gorm.Expr("NULL")
|
update["error"] = gorm.Expr("NULL")
|
||||||
} else {
|
case stateERR:
|
||||||
update["last_fail"] = now
|
update["last_fail"] = now
|
||||||
update["fails"] = gorm.Expr("fails + 1")
|
update["fails"] = gorm.Expr("fails + 1")
|
||||||
if report.Error != nil {
|
if report.Error != nil {
|
||||||
update["error"] = *report.Error
|
update["error"] = *report.Error
|
||||||
}
|
}
|
||||||
|
default:
|
||||||
|
update["last_fail"] = now
|
||||||
|
update["fails"] = 0
|
||||||
|
if report.Error != nil {
|
||||||
|
update["error"] = *report.Error
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
if report.ExpiresAt != nil {
|
if report.ExpiresAt != nil {
|
||||||
t, err := time.Parse(time.RFC3339, *report.ExpiresAt)
|
t, err := time.Parse(time.RFC3339, *report.ExpiresAt)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
update["expires"] = t
|
update["expires"] = t
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
return update
|
||||||
|
}
|
||||||
|
|
||||||
if err := tx.Model(&check).UpdateColumns(update).Error; err != nil {
|
func handleCheckStateSideEffects(tx *gorm.DB, check *Check, state string, worker *WorkerNode, now time.Time) error {
|
||||||
log.Println("worker: error updating check:", report.CheckID, err)
|
if worker == nil {
|
||||||
return nil, err
|
return nil
|
||||||
}
|
}
|
||||||
if worker != nil {
|
switch state {
|
||||||
payload, _ := json.Marshal(report)
|
case stateERR, stateFail:
|
||||||
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)}
|
return StartConfirmationTx(tx, check.ID, worker.ID, now)
|
||||||
if attempt.JobID == "" {
|
case stateOK:
|
||||||
attempt.JobID = uuid.New().String()
|
return RecoverDiagnosticTx(tx, check.ID, now)
|
||||||
}
|
|
||||||
// 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
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
return check.Monitor, nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// StoreRemoteCheckMetrics persists TSDB points reported by a distributed worker.
|
// StoreRemoteCheckMetrics persists TSDB points reported by a distributed worker.
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user