Files
worker/internal/distworker/selfcheck_test.go
Gleb Tv 2c884c5612
Некоторые проверки не удались
CI / test (push) Successful in 2m5s
Docker / Build and publish worker image (push) Failing after 31s
refactor: adopt worker module path
2026-07-13 17:56:12 +03:00

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 }