Некоторые проверки не удались
CI / test (push) Successful in 2m5s
Docker / Build and publish worker image (push) Failing after 31s
626 строки
24 KiB
Go
626 строки
24 KiB
Go
package distworker
|
|
|
|
import (
|
|
"encoding/json"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
|
|
"rocketgit.ru/rsmon/worker/internal/wire"
|
|
)
|
|
|
|
func TestRunMainAPIHTTPCheck_OK(t *testing.T) {
|
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
}))
|
|
defer srv.Close()
|
|
|
|
ok := runMainAPIHTTPCheck(srv.URL)
|
|
assert.True(t, ok, "200 OK should be a successful probe")
|
|
}
|
|
|
|
func TestRunMainAPIHTTPCheck_404IsDown(t *testing.T) {
|
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusNotFound)
|
|
}))
|
|
defer srv.Close()
|
|
|
|
ok := runMainAPIHTTPCheck(srv.URL)
|
|
assert.False(t, ok, "normal HTTP check flow treats non-200 as down")
|
|
}
|
|
|
|
func TestRunMainAPIHTTPCheck_5xxIsDown(t *testing.T) {
|
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusInternalServerError)
|
|
}))
|
|
defer srv.Close()
|
|
|
|
ok := runMainAPIHTTPCheck(srv.URL)
|
|
assert.False(t, ok, "5xx should be treated as down")
|
|
}
|
|
|
|
func TestRunMainAPIHTTPCheck_NetworkErrorIsDown(t *testing.T) {
|
|
ok := runMainAPIHTTPCheck("http://127.0.0.1:1")
|
|
assert.False(t, ok, "connection refused should be treated as down")
|
|
}
|
|
|
|
func TestSelfcheckTarget(t *testing.T) {
|
|
assert.Equal(t, "https://rsmon.ru/up", selfcheckTarget("https://rsmon.ru"))
|
|
assert.Equal(t, "https://rsmon.ru/up", selfcheckTarget("https://rsmon.ru/"))
|
|
assert.Equal(t, "https://api.example.com/api/v1/up", selfcheckTarget("https://api.example.com/api/v1"))
|
|
assert.Equal(t, "", selfcheckTarget(""), "empty base URL yields empty target")
|
|
}
|
|
|
|
func TestParseTelegramChatID(t *testing.T) {
|
|
id, err := parseTelegramChatID("200318758")
|
|
require.NoError(t, err)
|
|
assert.Equal(t, int64(200318758), id)
|
|
|
|
_, err = parseTelegramChatID("")
|
|
assert.Error(t, err)
|
|
|
|
_, err = parseTelegramChatID("not-a-number")
|
|
assert.Error(t, err)
|
|
}
|
|
|
|
func TestSMTPCredToModelCarriesPlaintextSecret(t *testing.T) {
|
|
c := wire.SMTPCredential{
|
|
ID: 11,
|
|
Name: "primary",
|
|
Server: "smtp.example.com",
|
|
Port: 587,
|
|
Login: "alerts@example.com",
|
|
Password: "smtp-password-xyz",
|
|
FromName: "RSMon",
|
|
FromAddress: "alerts@example.com",
|
|
}
|
|
cred := smtpCredToModel(&c)
|
|
require.NotNil(t, cred)
|
|
assert.Equal(t, "plain:smtp-password-xyz", cred.SecretEnc)
|
|
assert.Equal(t, int64(11), cred.ID)
|
|
assert.Equal(t, models_credentialKindSMTP(), cred.Kind)
|
|
|
|
got, err := cred.GetSecret()
|
|
require.NoError(t, err)
|
|
assert.Equal(t, "smtp-password-xyz", got,
|
|
"worker-side credential must yield plaintext via models.GetSecret")
|
|
}
|
|
|
|
func TestTelegramCredToModelCarriesPlaintextToken(t *testing.T) {
|
|
c := wire.TelegramCredential{
|
|
ID: 22,
|
|
Name: "main-bot",
|
|
BotName: "rsmon_alerts_bot",
|
|
Token: "bot-token-9876543210:ABCDEFG",
|
|
APIURL: "https://api.telegram.org",
|
|
}
|
|
cred := telegramCredToModel(c)
|
|
require.NotNil(t, cred)
|
|
assert.Equal(t, "plain:bot-token-9876543210:ABCDEFG", cred.SecretEnc)
|
|
assert.Equal(t, int64(22), cred.ID)
|
|
assert.Equal(t, models_credentialKindTelegram(), cred.Kind)
|
|
|
|
got, err := cred.GetSecret()
|
|
require.NoError(t, err)
|
|
assert.Equal(t, "bot-token-9876543210:ABCDEFG", got)
|
|
}
|
|
|
|
// models_credentialKindSMTP/Telegram are local shims to avoid an import-cycle
|
|
// in this isolated test file. The real constants live in app/models and are
|
|
// asserted at runtime through the wire translation functions.
|
|
func models_credentialKindSMTP() string { return "smtp" }
|
|
func models_credentialKindTelegram() string { return "telegram" }
|
|
|
|
func TestApplyInitStoresSystemContacts(t *testing.T) {
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
|
|
assert.Empty(t, r.SystemContacts(), "no system contacts before init")
|
|
|
|
want := []wire.SystemContact{
|
|
{ID: 1, Kind: "email", Value: "ops@example.com", Name: "ops"},
|
|
{ID: 2, Kind: "telegram_group", Value: "200318758", Name: "alerts"},
|
|
}
|
|
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Concurrency: 2,
|
|
SystemContacts: want,
|
|
})
|
|
|
|
got := r.SystemContacts()
|
|
require.Len(t, got, 2)
|
|
assert.Equal(t, want[0].ID, got[0].ID)
|
|
assert.Equal(t, "email", got[0].Kind)
|
|
assert.Equal(t, "ops@example.com", got[0].Value)
|
|
assert.Equal(t, "telegram_group", got[1].Kind)
|
|
assert.Equal(t, "200318758", got[1].Value)
|
|
|
|
r.applyInit(&wire.WorkerInit{WorkerID: "w-1", Concurrency: 2})
|
|
assert.Empty(t, r.SystemContacts(),
|
|
"subsequent init with empty SystemContacts should replace cache")
|
|
}
|
|
|
|
func TestRunSelfcheckOnce_FailureStartsDownStateNoAlert(t *testing.T) {
|
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusInternalServerError)
|
|
}))
|
|
defer srv.Close()
|
|
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.config.URL = srv.URL
|
|
creds := &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp.example.com", Port: 587, Login: "a@a", FromAddress: "a@a", Password: "x"}},
|
|
}
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: creds,
|
|
SystemContacts: []wire.SystemContact{{ID: 1, Kind: "email", Value: "ops@example.com"}},
|
|
})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
require.NotNil(t, state.consensus.downSince,
|
|
"first down probe must stamp downSince")
|
|
assert.False(t, state.consensus.alertActive,
|
|
"no alert expected before selfcheckConsensusWait elapses")
|
|
}
|
|
|
|
func TestRunSelfcheckOnce_AlertFiredAfterConsensusWait(t *testing.T) {
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.config.URL = "http://127.0.0.1:1"
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
downSince := time.Now().Add(-selfcheckConsensusWait - time.Second)
|
|
state.consensus.downSince = &downSince
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
assert.True(t, state.consensus.alertActive,
|
|
"alert must be active after selfcheckConsensusWait elapses")
|
|
}
|
|
|
|
func TestRunSelfcheckOnce_RecoveryClearsAlert(t *testing.T) {
|
|
srvOK := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
}))
|
|
defer srvOK.Close()
|
|
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
state.consensus.alertActive = true
|
|
downSince := time.Now().Add(-20 * time.Minute)
|
|
state.consensus.downSince = &downSince
|
|
|
|
r.config.URL = srvOK.URL
|
|
// First up tick: marks the up consensus start; recovery alert
|
|
// waits for selfcheckConsensusWait. Assert the alert is still
|
|
// active so we can distinguish "up just arrived" from "up held".
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
assert.True(t, state.consensus.alertActive,
|
|
"recovery alert should not fire on the first up tick")
|
|
assert.Equal(t, consensusUp, state.consensus.lastVerdict)
|
|
|
|
// Tick again with the up verdict now held past
|
|
// selfcheckConsensusWait so the recovery alert fires.
|
|
state.consensus.downSince = ptrTime(time.Now().Add(-selfcheckConsensusWait - time.Second))
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
assert.False(t, state.consensus.alertActive,
|
|
"recovery alert must clear alertActive after selfcheckConsensusWait up")
|
|
assert.Nil(t, state.consensus.downSince)
|
|
}
|
|
|
|
func TestRunSelfcheckOnce_RequiresSelfVote(t *testing.T) {
|
|
// When self has not yet produced a verdict (nil), the
|
|
// consensus must refuse to call the cluster down regardless of
|
|
// what peers report.
|
|
up := true
|
|
peers := []peerObservation{{WorkerID: "p1", Up: true, ObservedAt: time.Now()}}
|
|
verdict := decideConsensus(nil, peers, 2)
|
|
assert.Equal(t, consensusNoQuorum, verdict,
|
|
"missing self vote must prevent a down consensus")
|
|
|
|
verdict = decideConsensus(&up, peers, 2)
|
|
assert.Equal(t, consensusUp, verdict,
|
|
"self up + peer up must reach consensus up")
|
|
}
|
|
|
|
func TestSendSystemAlert_NoCredentials(t *testing.T) {
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
SystemContacts: []wire.SystemContact{{ID: 1, Kind: "email", Value: "ops@example.com"}},
|
|
})
|
|
|
|
r.sendSystemAlert(t.Context(), true, "test message")
|
|
}
|
|
|
|
func TestSendSystemAlert_NoContacts(t *testing.T) {
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
})
|
|
|
|
r.sendSystemAlert(t.Context(), true, "test message")
|
|
}
|
|
|
|
func TestWorkerInit_SystemContactsJSONRoundtrip(t *testing.T) {
|
|
init := wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Concurrency: 4,
|
|
SystemContacts: []wire.SystemContact{
|
|
{ID: 1, Kind: "email", Value: "ops@example.com", Name: "ops"},
|
|
{ID: 2, Kind: "telegram_private", Value: "200318758", Name: "alerts"},
|
|
},
|
|
}
|
|
|
|
raw, err := json.Marshal(init)
|
|
require.NoError(t, err)
|
|
assert.Contains(t, string(raw), `"system_contacts":`)
|
|
|
|
var decoded wire.WorkerInit
|
|
require.NoError(t, json.Unmarshal(raw, &decoded))
|
|
require.Len(t, decoded.SystemContacts, 2)
|
|
assert.Equal(t, init.SystemContacts[0], decoded.SystemContacts[0])
|
|
assert.Equal(t, init.SystemContacts[1], decoded.SystemContacts[1])
|
|
}
|
|
|
|
func TestWorkerInit_PeersJSONRoundtrip(t *testing.T) {
|
|
init := wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Concurrency: 4,
|
|
Peers: []wire.PeerInfo{
|
|
{WorkerID: "w-2", URL: "http://127.0.0.1:27402", RegionCode: "local"},
|
|
{WorkerID: "w-3", URL: "http://127.0.0.1:27403", RegionCode: "local", Login: "ops", Password: "x"},
|
|
},
|
|
}
|
|
|
|
raw, err := json.Marshal(init)
|
|
require.NoError(t, err)
|
|
body := string(raw)
|
|
assert.Contains(t, body, `"peers":[`)
|
|
assert.Contains(t, body, `"worker_id":"w-2"`)
|
|
assert.Contains(t, body, `"url":"http://127.0.0.1:27402"`)
|
|
assert.Contains(t, body, `"login":"ops"`)
|
|
assert.Contains(t, body, `"password":"x"`)
|
|
|
|
var decoded wire.WorkerInit
|
|
require.NoError(t, json.Unmarshal(raw, &decoded))
|
|
require.Len(t, decoded.Peers, 2)
|
|
assert.Equal(t, init.Peers[0], decoded.Peers[0])
|
|
assert.Equal(t, init.Peers[1], decoded.Peers[1])
|
|
}
|
|
|
|
func TestWorkerInit_PeersOmittedWhenEmpty(t *testing.T) {
|
|
init := wire.WorkerInit{WorkerID: "w-1", Concurrency: 1}
|
|
raw, err := json.Marshal(init)
|
|
require.NoError(t, err)
|
|
assert.NotContains(t, string(raw), `"peers"`,
|
|
"empty peers slice must be omitted so legacy workers stay wire-compatible")
|
|
}
|
|
|
|
// TestRunSelfcheckOnce_ConfiguredPeersNoFreshPeer_NoAlert verifies that
|
|
// in a multi-worker install a transient local probe failure (or
|
|
// success) does NOT fire an alert or reset the state machine when
|
|
// peers are configured but have not produced a fresh observation yet.
|
|
// The first peer poll lands 2s after startup, and after that on the
|
|
// peer-poll interval, so the early ticks are a real "lone self vote"
|
|
// window in production.
|
|
func TestRunSelfcheckOnce_ConfiguredPeersNoFreshPeer_NoAlert(t *testing.T) {
|
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusInternalServerError)
|
|
}))
|
|
defer srv.Close()
|
|
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.config.URL = srv.URL
|
|
creds := &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
}
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: creds,
|
|
SystemContacts: []wire.SystemContact{
|
|
{ID: 1, Kind: "email", Value: "ops@example.com"},
|
|
},
|
|
// Two peers configured but neither has produced a
|
|
// fresh observation yet — the selfcheck is the lone
|
|
// voter.
|
|
Peers: []wire.PeerInfo{
|
|
{WorkerID: "w-2", URL: "http://127.0.0.1:27402"},
|
|
{WorkerID: "w-3", URL: "http://127.0.0.1:27403"},
|
|
},
|
|
})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
|
|
assert.Nil(t, state.consensus.downSince,
|
|
"with configured peers and no fresh peer observations, the down timer must not start")
|
|
assert.False(t, state.consensus.alertActive,
|
|
"no alert can fire without a peer-backed down verdict")
|
|
}
|
|
|
|
// TestRunSelfcheckOnce_ConfiguredPeersWithFreshDownPeer_FiresAfterWait
|
|
// exercises the happy path: self and one peer both see the master as
|
|
// down, the consensus verdict is down, and after selfcheckConsensusWait
|
|
// the alert fires exactly once.
|
|
func TestRunSelfcheckOnce_ConfiguredPeersWithFreshDownPeer_FiresAfterWait(t *testing.T) {
|
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusInternalServerError)
|
|
}))
|
|
defer srv.Close()
|
|
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.config.URL = srv.URL
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
SystemContacts: []wire.SystemContact{
|
|
{ID: 1, Kind: "email", Value: "ops@example.com"},
|
|
},
|
|
Peers: []wire.PeerInfo{
|
|
{WorkerID: "w-2", URL: "http://127.0.0.1:27402"},
|
|
},
|
|
})
|
|
|
|
// Seed a fresh down observation from the peer as if its
|
|
// selfcheck tick has just completed.
|
|
r.peerCache.put(peerObservation{
|
|
WorkerID: "w-2",
|
|
Up: false,
|
|
ObservedAt: time.Now(),
|
|
})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
|
|
// First tick: stamps downSince, does not yet fire.
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
require.NotNil(t, state.consensus.downSince,
|
|
"peer-backed down verdict must stamp downSince on first tick")
|
|
assert.False(t, state.consensus.alertActive,
|
|
"alert must not fire before selfcheckConsensusWait elapses")
|
|
|
|
// Second tick, simulate the 5-minute wait having elapsed
|
|
// by rewinding downSince past the wait threshold. This
|
|
// avoids waiting wall-clock time in a unit test.
|
|
state.consensus.downSince = ptrTime(time.Now().Add(-selfcheckConsensusWait - time.Second))
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
assert.True(t, state.consensus.alertActive,
|
|
"down consensus held for selfcheckConsensusWait must fire the alert exactly once")
|
|
}
|
|
|
|
func TestRunSelfcheckOnce_NonLeaderSkipsDuplicateAlert(t *testing.T) {
|
|
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusInternalServerError)
|
|
}))
|
|
defer srv.Close()
|
|
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.config.URL = srv.URL
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "worker-local-2",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
SystemContacts: []wire.SystemContact{{ID: 1, Kind: "email", Value: "ops@example.com"}},
|
|
Peers: []wire.PeerInfo{{WorkerID: "worker-local-1", URL: "http://127.0.0.1:27401"}},
|
|
})
|
|
r.peerCache.put(peerObservation{WorkerID: "worker-local-1", Up: false, ObservedAt: time.Now()})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
state.consensus.downSince = ptrTime(time.Now().Add(-selfcheckConsensusWait - time.Second))
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
|
|
assert.True(t, state.consensus.alertActive, "non-leader still marks the incident handled locally")
|
|
assert.Empty(t, r.RecentNotifications(10), "non-leader must not deliver duplicate system-contact notifications")
|
|
}
|
|
|
|
func TestFormatSelfcheckLeaderMessage_IncludesOldAndNewLeader(t *testing.T) {
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{WorkerID: "worker-local-2", RegionCode: "local", Concurrency: 1})
|
|
|
|
out := r.formatSelfcheckLeaderMessage("worker-local-1", "http://localhost:7401/up", consensusDown, 0, 2, 2, 2, 2, 1,
|
|
"worker-local-2", []string{"worker-local-2", "worker-local-3"})
|
|
|
|
assert.Contains(t, out, "Изменился мастер воркер оповещений: worker-local-1 -> worker-local-2")
|
|
assert.Contains(t, out, "Notification leader: worker-local-2")
|
|
assert.Contains(t, out, "Voters: worker-local-2,worker-local-3")
|
|
}
|
|
|
|
// TestRunSelfcheckOnce_RecoveryRequiresHeldUpConsensus verifies that
|
|
// after a peer-backed down alert, recovery only fires after the up
|
|
// verdict has itself been held for selfcheckConsensusWait, and that
|
|
// the recovery flow respects the multi-worker gate (an isolated up
|
|
// tick with no fresh peer must NOT clear the alert).
|
|
func TestRunSelfcheckOnce_RecoveryRequiresHeldUpConsensus(t *testing.T) {
|
|
srvOK := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
}))
|
|
defer srvOK.Close()
|
|
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
SystemContacts: []wire.SystemContact{
|
|
{ID: 1, Kind: "email", Value: "ops@example.com"},
|
|
},
|
|
Peers: []wire.PeerInfo{
|
|
{WorkerID: "w-2", URL: "http://127.0.0.1:27402"},
|
|
},
|
|
})
|
|
|
|
// Pretend the peer also sees the master as up — the only
|
|
// way the up verdict is allowed past the multi-worker gate.
|
|
r.peerCache.put(peerObservation{
|
|
WorkerID: "w-2",
|
|
Up: true,
|
|
ObservedAt: time.Now(),
|
|
})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
state.consensus.alertActive = true
|
|
state.consensus.downSince = ptrTime(time.Now().Add(-20 * time.Minute))
|
|
r.config.URL = srvOK.URL
|
|
|
|
// First up tick: stamps the recovery timer; alert stays
|
|
// active because the up verdict is not yet held for
|
|
// selfcheckConsensusWait.
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
assert.True(t, state.consensus.alertActive,
|
|
"recovery alert must not fire on the first up tick")
|
|
assert.Equal(t, consensusUp, state.consensus.lastVerdict)
|
|
|
|
// Simulate the 5-minute recovery wait having elapsed.
|
|
state.consensus.downSince = ptrTime(time.Now().Add(-selfcheckConsensusWait - time.Second))
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
assert.False(t, state.consensus.alertActive,
|
|
"recovery alert must clear alertActive after selfcheckConsensusWait up consensus")
|
|
assert.Nil(t, state.consensus.downSince)
|
|
}
|
|
|
|
// TestRunSelfcheckOnce_RecoveryIsolatedUpDoesNotClearAlert verifies
|
|
// that an isolated up tick (configured peers, no fresh peer
|
|
// observation) cannot unilaterally clear an active alert. The
|
|
// alert must remain active until a peer-backed up verdict is held.
|
|
func TestRunSelfcheckOnce_RecoveryIsolatedUpDoesNotClearAlert(t *testing.T) {
|
|
srvOK := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
w.WriteHeader(http.StatusOK)
|
|
}))
|
|
defer srvOK.Close()
|
|
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
SystemContacts: []wire.SystemContact{
|
|
{ID: 1, Kind: "email", Value: "ops@example.com"},
|
|
},
|
|
Peers: []wire.PeerInfo{
|
|
{WorkerID: "w-2", URL: "http://127.0.0.1:27402"},
|
|
},
|
|
})
|
|
// Note: no fresh peer observation is seeded — the selfcheck
|
|
// is the lone voter in this scenario.
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
state.consensus.alertActive = true
|
|
state.consensus.downSince = ptrTime(time.Now().Add(-20 * time.Minute))
|
|
r.config.URL = srvOK.URL
|
|
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
|
|
assert.True(t, state.consensus.alertActive,
|
|
"isolated up tick must not clear an active alert when peers are configured")
|
|
}
|
|
|
|
// TestRunSelfcheckOnce_NoQuorumDoesNotResetDownTimer verifies that a
|
|
// NoQuorum verdict neither stamps nor advances the 5-minute down
|
|
// timer. A brief blip in peer reachability must not page or reset
|
|
// progress toward the consensus wait.
|
|
func TestRunSelfcheckOnce_NoQuorumDoesNotResetDownTimer(t *testing.T) {
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{
|
|
WorkerID: "w-1",
|
|
Credentials: &wire.NotificationCredentials{
|
|
SMTP: []wire.SMTPCredential{{ID: 1, Name: "primary", Server: "smtp", Port: 25, Login: "a", FromAddress: "a", Password: "x"}},
|
|
},
|
|
SystemContacts: []wire.SystemContact{
|
|
{ID: 1, Kind: "email", Value: "ops@example.com"},
|
|
},
|
|
Peers: []wire.PeerInfo{
|
|
{WorkerID: "w-2", URL: "http://127.0.0.1:27402"},
|
|
},
|
|
})
|
|
// Peer observation is older than peerStatusMaxAge so the
|
|
// selfcheck is alone in the vote.
|
|
r.peerCache.put(peerObservation{
|
|
WorkerID: "w-2",
|
|
Up: false,
|
|
ObservedAt: time.Now().Add(-10 * time.Minute),
|
|
})
|
|
|
|
state := &selfcheckState{consensus: &consensusState{}}
|
|
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
assert.Nil(t, state.consensus.downSince,
|
|
"first NoQuorum must not stamp the down timer")
|
|
assert.False(t, state.consensus.alertActive)
|
|
|
|
// A subsequent NoQuorum must also not advance the timer;
|
|
// the existing stamp (from a previous verdict) is held.
|
|
previous := ptrTime(time.Now().Add(-3 * time.Minute))
|
|
state.consensus.downSince = previous
|
|
r.runSelfcheckOnce(t.Context(), state)
|
|
require.NotNil(t, state.consensus.downSince)
|
|
assert.Equal(t, *previous, *state.consensus.downSince,
|
|
"NoQuorum must not advance the down timer")
|
|
}
|
|
|
|
func TestFormatSelfcheckMessage_IncludesWorkerAndConsensus(t *testing.T) {
|
|
r := newTestRunner(t, 4, 1, func(interface{}) interface{} { return []wire.CheckResultReport{} })
|
|
r.applyInit(&wire.WorkerInit{WorkerID: "worker-local-1", RegionCode: "local", Concurrency: 1})
|
|
|
|
out := r.formatSelfcheckMessage(selfcheckDownMessage, "http://localhost:7401/up", consensusDown, 0, 3, 3, 2, 2, 2,
|
|
"worker-local-1", []string{"worker-local-1", "worker-local-2", "worker-local-3"})
|
|
|
|
assert.Contains(t, out, "Нет связи с основным api")
|
|
assert.Contains(t, out, "Worker: worker-local-1")
|
|
assert.Contains(t, out, "Region: local")
|
|
assert.Contains(t, out, "Target: http://localhost:7401/up")
|
|
assert.Contains(t, out, "Raft status: lightweight peer quorum")
|
|
assert.Contains(t, out, "Notification leader: worker-local-1")
|
|
assert.Contains(t, out, "Consensus: down")
|
|
assert.Contains(t, out, "Votes: up=0 down=3 total=3 min_votes=2")
|
|
assert.Contains(t, out, "Workers seen: 3 current votes, 2 fresh peers of 2 configured peers")
|
|
assert.Contains(t, out, "Voters: worker-local-1,worker-local-2,worker-local-3")
|
|
assert.Contains(t, out, "Hold time: 5m0s")
|
|
}
|
|
|
|
func TestConsensusNotificationLeader(t *testing.T) {
|
|
voters := consensusVoters("worker-local-2", []peerObservation{{WorkerID: "worker-local-3"}, {WorkerID: "worker-local-1"}})
|
|
assert.Equal(t, []string{"worker-local-1", "worker-local-2", "worker-local-3"}, voters)
|
|
assert.Equal(t, "worker-local-1", consensusNotificationLeader(voters))
|
|
}
|
|
|
|
func TestConsensusState_NotificationLeaderChanged(t *testing.T) {
|
|
state := &consensusState{}
|
|
old, changed := state.notificationLeaderChanged("worker-local-1")
|
|
assert.False(t, changed)
|
|
assert.Empty(t, old)
|
|
|
|
old, changed = state.notificationLeaderChanged("worker-local-1")
|
|
assert.False(t, changed)
|
|
assert.Empty(t, old)
|
|
|
|
old, changed = state.notificationLeaderChanged("worker-local-2")
|
|
assert.True(t, changed)
|
|
assert.Equal(t, "worker-local-1", old)
|
|
}
|
|
|
|
func ptrTime(t time.Time) *time.Time { return &t }
|