mirror of
https://github.com/rcourtman/Pulse.git
synced 2026-10-03 12:47:49 +00:00
Integrate reviewed listener-owned health and telemetry transition repair
Change-source: pulse-maintainer
This commit is contained in:
commit
858abc579a
12 changed files with 402 additions and 21 deletions
|
|
@ -3686,6 +3686,16 @@ fields. No agent-lifecycle behavior keyed off them — agent update targeting
|
|||
and command admission are unaffected — and the extension-point expectations
|
||||
on the system-settings boundary are otherwise unchanged.
|
||||
|
||||
### Shared telemetry settings preserve sender lifecycle ownership
|
||||
|
||||
Telemetry preference saves invoke the sender callback only after a persisted
|
||||
explicit boolean changes the effective runtime value. Repeated saves and
|
||||
null/omitted preferences do not restart it. This lifecycle belongs to telemetry
|
||||
reporting, not agent registration or command admission. A telemetry `startup`
|
||||
event also follows a genuine sender re-enable or ID reset and therefore cannot
|
||||
serve as a count of agent or server restarts. Settings transition/persistence
|
||||
tests pin this shared boundary without changing agent lifecycle authority.
|
||||
|
||||
### Shared system-settings boundary gained an SSH backoff reset side effect
|
||||
|
||||
The shared `internal/api` system-settings surface this subsystem consumes
|
||||
|
|
|
|||
|
|
@ -4766,6 +4766,17 @@ ignores them without a validation error and never writes them back into
|
|||
`internal/api/system_settings_telemetry_test.go` and the response snapshot in
|
||||
`internal/api/contract_test.go` pin that payload shape.
|
||||
|
||||
### Telemetry preference saves preserve the sender lifecycle
|
||||
|
||||
The telemetry preference callback runs only after durable persistence of an
|
||||
explicit boolean that changes the effective runtime value. Resubmitting that
|
||||
value still saves it, including correction of a stale disk preference, without
|
||||
restarting or stopping the sender. Null and omitted preferences preserve the
|
||||
runtime value and the stored preference; failed saves never invoke the toggle.
|
||||
`TestTelemetryUpdate_OnlyPreferenceTransitionsToggle` pins ordering and both
|
||||
transitions; the stale-disk/null and persistence-failure tests pin those edges.
|
||||
The settings payload and administration/tenant authority are unchanged.
|
||||
|
||||
### System settings save clears the temperature SSH failure backoff
|
||||
|
||||
A successful `POST /api/system-settings` save now also calls
|
||||
|
|
|
|||
|
|
@ -1520,8 +1520,11 @@ persisted receiver row. The previous-release fields are direct adjacent-release
|
|||
observations, not 30-day update counters, and are the only valid basis for a
|
||||
before/after release-health cohort in the adoption report.
|
||||
The service-health self-check must preserve the explicit bound TCP address and
|
||||
port; wildcard IPv6/unspecified-family listeners try IPv4 and IPv6 loopback,
|
||||
while an IPv4-only wildcard remains IPv4-only. No external address discovery,
|
||||
port. It inspects the IPv6-only option of a bound IPv6 wildcard socket: only a
|
||||
proven dual-stack listener tries IPv4 then IPv6 loopback. IPv6-only or
|
||||
uninspectable IPv6 wildcards remain IPv6-only, while IPv4-only and nil-IP
|
||||
wildcards remain IPv4-only. Another socket on the same port is not evidence of
|
||||
this listener's health. No external address discovery,
|
||||
proxy, redirect, or remote frontend asset fetch is permitted. One bounded
|
||||
settling retry may follow a failed observation in the telemetry background
|
||||
runner: at most two five-second attempts, separated by one second. Each attempt
|
||||
|
|
@ -1532,6 +1535,14 @@ the fixed `timeout` bucket; other failure classes remain visible. Tests must
|
|||
cover unavailable IPv6 with working IPv4, IPv6-only serving, a stalled family,
|
||||
startup recovery, final timeouts, explicit addresses, HTTPS, redirect refusal,
|
||||
bounded bodies/assets, closed-category serialization and background execution.
|
||||
`service_health_socket_test.go` also pins real HTTP/HTTPS wildcard socket modes,
|
||||
both healthy/unhealthy same-port isolation directions, and unavailable socket
|
||||
inspection. Socket-mode inspection does not add a telemetry field or export
|
||||
its result, target or error. Telemetry preference saves invoke the live toggle
|
||||
only for persisted explicit boolean transitions; unchanged/null/omitted values
|
||||
do not restart the sender. A `startup` event also follows re-enabling telemetry
|
||||
or resetting its ID, so its count alone does not establish process restarts or
|
||||
a release regression. Existing payload confidentiality remains unchanged.
|
||||
That same outbound usage telemetry floor now also permits only content-free Pulse
|
||||
Patrol control and governed Pulse Intelligence operations adoption flags and
|
||||
counters inside the same rotating 30-day telemetry window:
|
||||
|
|
|
|||
|
|
@ -2805,6 +2805,16 @@ fields. Persisted `system.json` files that still carry the legacy keys load
|
|||
cleanly with the keys ignored, so tenant workspace preservation and recovery
|
||||
flows that copy `system.json` forward are unaffected.
|
||||
|
||||
### Shared telemetry preference updates preserve recovery ownership
|
||||
|
||||
Telemetry preference updates remain durable before the sender callback. An
|
||||
explicit boolean equal to the effective runtime value still repairs a stale
|
||||
disk preference but must not restart the sender. Null/omitted values leave that
|
||||
preference alone, and a failed save neither mutates runtime nor invokes the
|
||||
callback. The document shape and recovery ownership are unchanged.
|
||||
`system_settings_telemetry_test.go` pins repeated saves, transitions, durable
|
||||
callback ordering, stale preferences, null input and failed persistence.
|
||||
|
||||
### Shared system-settings boundary gained an SSH backoff reset side effect
|
||||
|
||||
The shared `internal/api` system-settings surface this subsystem consumes
|
||||
|
|
|
|||
|
|
@ -55,7 +55,7 @@ type SystemSettingsHandler struct {
|
|||
// runtime, so the router can reconfigure the monitor's checker and
|
||||
// collector without a restart.
|
||||
guestDockerInventoryToggleFunc func()
|
||||
mtMonitor interface {
|
||||
mtMonitor interface {
|
||||
GetMonitor(string) (*monitoring.Monitor, error)
|
||||
}
|
||||
defaultMonitor SystemSettingsMonitor
|
||||
|
|
@ -1069,10 +1069,12 @@ func (h *SystemSettingsHandler) HandleUpdateSystemSettings(w http.ResponseWriter
|
|||
Bool("enabled", settings.EnableProxmoxGuestDockerInventory).
|
||||
Msg("Proxmox guest Docker inventory opt-in changed via settings")
|
||||
}
|
||||
if _, ok := rawRequest["telemetryEnabled"]; ok && settings.TelemetryEnabled != nil {
|
||||
h.config.TelemetryEnabled = *settings.TelemetryEnabled
|
||||
// A null/omitted preference is not a transition, even if disk and runtime
|
||||
// differ. Repeated saves must not restart the telemetry sender either.
|
||||
if updates.TelemetryEnabled != nil && h.config.TelemetryEnabled != *updates.TelemetryEnabled {
|
||||
h.config.TelemetryEnabled = *updates.TelemetryEnabled
|
||||
if h.telemetryToggleFunc != nil {
|
||||
h.telemetryToggleFunc(*settings.TelemetryEnabled)
|
||||
h.telemetryToggleFunc(*updates.TelemetryEnabled)
|
||||
}
|
||||
}
|
||||
if _, ok := rawRequest["publicURL"]; ok {
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@ package api
|
|||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
|
|
@ -367,6 +368,8 @@ func TestTelemetryUpdate_NoMutationOnPersistFailure(t *testing.T) {
|
|||
EnvOverrides: make(map[string]bool),
|
||||
}
|
||||
handler, persistence, token := setupTelemetryTest(t, cfg)
|
||||
toggleCalled := false
|
||||
handler.SetTelemetryToggleFunc(func(bool) { toggleCalled = true })
|
||||
|
||||
initial := config.DefaultSystemSettings()
|
||||
if err := persistence.SaveSystemSettings(*initial); err != nil {
|
||||
|
|
@ -401,6 +404,129 @@ func TestTelemetryUpdate_NoMutationOnPersistFailure(t *testing.T) {
|
|||
if !cfg.TelemetryEnabled {
|
||||
t.Error("TelemetryEnabled should still be true after persistence failure")
|
||||
}
|
||||
if toggleCalled {
|
||||
t.Fatal("failed persistence invoked the telemetry toggle")
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelemetryUpdate_OnlyPreferenceTransitionsToggle(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
cfg := &config.Config{
|
||||
DataPath: tempDir, ConfigPath: tempDir, TelemetryEnabled: true,
|
||||
EnvOverrides: make(map[string]bool),
|
||||
}
|
||||
handler, persistence, token := setupTelemetryTest(t, cfg)
|
||||
settings := config.DefaultSystemSettings()
|
||||
enabled := true
|
||||
settings.TelemetryEnabled = &enabled
|
||||
if err := persistence.SaveSystemSettings(*settings); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var toggles []bool
|
||||
handler.SetTelemetryToggleFunc(func(value bool) {
|
||||
// The sender must observe the committed preference, not a proposed
|
||||
// value that could still fail persistence.
|
||||
persisted, err := persistence.LoadSystemSettings()
|
||||
if err != nil || persisted.TelemetryEnabled == nil || *persisted.TelemetryEnabled != value {
|
||||
t.Fatalf("toggle preceded persistence: settings=%#v error=%v", persisted, err)
|
||||
}
|
||||
if cfg.TelemetryEnabled != value {
|
||||
t.Fatalf("toggle preceded runtime mutation: enabled=%v want=%v", cfg.TelemetryEnabled, value)
|
||||
}
|
||||
toggles = append(toggles, value)
|
||||
})
|
||||
for i, value := range []bool{true, true, false, false, true, true} {
|
||||
body, err := json.Marshal(map[string]interface{}{"telemetryEnabled": value})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/system-settings", bytes.NewReader(body))
|
||||
req.Header.Set("X-API-Token", token)
|
||||
rec := httptest.NewRecorder()
|
||||
handler.HandleUpdateSystemSettings(rec, req)
|
||||
if rec.Code != http.StatusOK || cfg.TelemetryEnabled != value {
|
||||
t.Fatalf("save %d returned %d, enabled=%v: %s", i, rec.Code, cfg.TelemetryEnabled, rec.Body.String())
|
||||
}
|
||||
persisted, err := persistence.LoadSystemSettings()
|
||||
if err != nil || persisted.TelemetryEnabled == nil || *persisted.TelemetryEnabled != value {
|
||||
t.Fatalf("save %d did not preserve requested preference: settings=%#v error=%v", i, persisted, err)
|
||||
}
|
||||
wantToggleCount := 0
|
||||
if i >= 2 {
|
||||
wantToggleCount++
|
||||
}
|
||||
if i >= 4 {
|
||||
wantToggleCount++
|
||||
}
|
||||
if len(toggles) != wantToggleCount {
|
||||
t.Fatalf("save %d toggles = %v, want %d transitions", i, toggles, wantToggleCount)
|
||||
}
|
||||
}
|
||||
if len(toggles) != 2 || toggles[0] != false || toggles[1] != true {
|
||||
t.Fatalf("runtime toggles = %v, want only the disable and re-enable transitions", toggles)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelemetryUpdate_UnchangedPreferenceRepairsStaleDiskWithoutToggle(t *testing.T) {
|
||||
for _, enabled := range []bool{true, false} {
|
||||
t.Run(fmt.Sprintf("enabled=%v", enabled), func(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
cfg := &config.Config{
|
||||
DataPath: tempDir, ConfigPath: tempDir, TelemetryEnabled: enabled,
|
||||
EnvOverrides: make(map[string]bool),
|
||||
}
|
||||
handler, persistence, token := setupTelemetryTest(t, cfg)
|
||||
settings := config.DefaultSystemSettings()
|
||||
stale := !enabled
|
||||
settings.TelemetryEnabled = &stale
|
||||
if err := persistence.SaveSystemSettings(*settings); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
handler.SetTelemetryToggleFunc(func(bool) { t.Fatal("unchanged effective preference invoked toggle") })
|
||||
body, err := json.Marshal(map[string]interface{}{"telemetryEnabled": enabled})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/system-settings", bytes.NewReader(body))
|
||||
req.Header.Set("X-API-Token", token)
|
||||
rec := httptest.NewRecorder()
|
||||
handler.HandleUpdateSystemSettings(rec, req)
|
||||
if rec.Code != http.StatusOK || cfg.TelemetryEnabled != enabled {
|
||||
t.Fatalf("unchanged save returned %d, enabled=%v: %s", rec.Code, cfg.TelemetryEnabled, rec.Body.String())
|
||||
}
|
||||
persisted, err := persistence.LoadSystemSettings()
|
||||
if err != nil || persisted.TelemetryEnabled == nil || *persisted.TelemetryEnabled != enabled {
|
||||
t.Fatalf("unchanged save did not repair disk: settings=%#v error=%v", persisted, err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelemetryUpdate_NullPreservesEffectiveValueWithStalePreference(t *testing.T) {
|
||||
tempDir := t.TempDir()
|
||||
cfg := &config.Config{
|
||||
DataPath: tempDir, ConfigPath: tempDir, TelemetryEnabled: true,
|
||||
EnvOverrides: make(map[string]bool),
|
||||
}
|
||||
handler, persistence, token := setupTelemetryTest(t, cfg)
|
||||
settings := config.DefaultSystemSettings()
|
||||
stale := false
|
||||
settings.TelemetryEnabled = &stale
|
||||
if err := persistence.SaveSystemSettings(*settings); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
handler.SetTelemetryToggleFunc(func(bool) { t.Fatal("null preference invoked toggle") })
|
||||
req := httptest.NewRequest(http.MethodPost, "/api/system-settings", strings.NewReader(`{"telemetryEnabled":null}`))
|
||||
req.Header.Set("X-API-Token", token)
|
||||
rec := httptest.NewRecorder()
|
||||
handler.HandleUpdateSystemSettings(rec, req)
|
||||
if rec.Code != http.StatusOK || !cfg.TelemetryEnabled {
|
||||
t.Fatalf("null preference changed runtime: status=%d enabled=%v", rec.Code, cfg.TelemetryEnabled)
|
||||
}
|
||||
persisted, err := persistence.LoadSystemSettings()
|
||||
if err != nil || persisted.TelemetryEnabled == nil || *persisted.TelemetryEnabled {
|
||||
t.Fatalf("null preference changed disk: settings=%#v error=%v", persisted, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestTelemetryUpdate_UnrelatedUpdateDoesNotToggle(t *testing.T) {
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@ import (
|
|||
"net/url"
|
||||
"regexp"
|
||||
"strings"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/telemetry"
|
||||
|
|
@ -168,11 +169,14 @@ func localServiceHealthBaseURLs(listener net.Listener, tlsEnabled bool) []string
|
|||
hosts := []string{tcpAddr.IP.String()}
|
||||
if tcpAddr.IP == nil || tcpAddr.IP.IsUnspecified() {
|
||||
hosts = []string{"127.0.0.1"}
|
||||
// An IPv6 wildcard can be dual-stack or IPv6-only. IPv4-first also
|
||||
// works when IPv6 loopback is disabled. An IPv4-only wildcard must
|
||||
// not inspect a different IPv6 listener that happens to share its port.
|
||||
if tcpAddr.IP == nil || tcpAddr.IP.To4() == nil {
|
||||
hosts = append(hosts, "::1")
|
||||
if tcpAddr.IP != nil && tcpAddr.IP.To4() == nil {
|
||||
// Only a proven dual-stack socket owns both loopback families.
|
||||
// An IPv6-only socket may share its port with an unrelated IPv4
|
||||
// server. Unknown socket modes stay conservatively IPv6-only.
|
||||
hosts = []string{"::1"}
|
||||
if serviceHealthListenerAcceptsIPv4(listener) {
|
||||
hosts = []string{"127.0.0.1", "::1"}
|
||||
}
|
||||
}
|
||||
} else if tcpAddr.Zone != "" && tcpAddr.IP.To4() == nil {
|
||||
hosts[0] += "%" + tcpAddr.Zone
|
||||
|
|
@ -189,6 +193,27 @@ func localServiceHealthBaseURLs(listener net.Listener, tlsEnabled bool) []string
|
|||
return baseURLs
|
||||
}
|
||||
|
||||
// IPv4-first works when a dual-stack listener serves Pulse over IPv4 but IPv6
|
||||
// loopback is disabled. Inspect the bound socket rather than inferring its
|
||||
// mode from the wildcard address or from another server answering on its port.
|
||||
func serviceHealthListenerAcceptsIPv4(listener net.Listener) bool {
|
||||
socket, ok := listener.(syscall.Conn)
|
||||
if !ok {
|
||||
return false
|
||||
}
|
||||
raw, err := socket.SyscallConn()
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
dualStack := false
|
||||
if err := raw.Control(func(fd uintptr) {
|
||||
dualStack = serviceHealthSocketIsDualStack(fd)
|
||||
}); err != nil {
|
||||
return false
|
||||
}
|
||||
return dualStack
|
||||
}
|
||||
|
||||
func frontendAssetPaths(index []byte) []string {
|
||||
matches := frontendAssetReferencePattern.FindAllSubmatch(index, serviceHealthAssetLimit)
|
||||
paths := make([]string, 0, len(matches))
|
||||
|
|
|
|||
|
|
@ -18,8 +18,8 @@ import (
|
|||
"github.com/rcourtman/pulse-go-rewrite/internal/telemetry"
|
||||
)
|
||||
|
||||
// Override only Addr, retaining the real listener for serving. This models an
|
||||
// IPv6 wildcard reported on an install where IPv6 loopback is unavailable.
|
||||
// Address-only listeners have no inspectable socket. IPv6 wildcards must stay
|
||||
// conservative rather than treating an IPv4 response as evidence of their mode.
|
||||
type serviceHealthAddrListener struct {
|
||||
net.Listener
|
||||
addr net.Addr
|
||||
|
|
@ -49,8 +49,8 @@ func TestServiceHealthWildcardFamiliesAndExplicitAddresses(t *testing.T) {
|
|||
tls bool
|
||||
want []string
|
||||
}{
|
||||
{"IPv6 wildcard", &net.TCPAddr{IP: net.IPv6unspecified, Port: 7655}, false, []string{"http://127.0.0.1:7655", "http://[::1]:7655"}},
|
||||
{"unspecified family", &net.TCPAddr{Port: 7655}, false, []string{"http://127.0.0.1:7655", "http://[::1]:7655"}},
|
||||
{"uninspectable IPv6 wildcard", &net.TCPAddr{IP: net.IPv6unspecified, Port: 7655}, false, []string{"http://[::1]:7655"}},
|
||||
{"unspecified family", &net.TCPAddr{Port: 7655}, false, []string{"http://127.0.0.1:7655"}},
|
||||
{"IPv4 wildcard", &net.TCPAddr{IP: net.IPv4zero, Port: 7655}, false, []string{"http://127.0.0.1:7655"}},
|
||||
{"explicit IPv4", &net.TCPAddr{IP: net.ParseIP("192.0.2.1"), Port: 7655}, false, []string{"http://192.0.2.1:7655"}},
|
||||
{"explicit IPv6", &net.TCPAddr{IP: net.IPv6loopback, Port: 7655}, true, []string{"https://[::1]:7655"}},
|
||||
|
|
@ -72,6 +72,15 @@ func TestServiceHealthWildcardFamiliesAndExplicitAddresses(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestServiceHealthIPv4SurvivesUnavailableIPv6(t *testing.T) {
|
||||
listener, err := net.Listen("tcp", "[::]:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = listener.Close() })
|
||||
targets := localServiceHealthBaseURLs(listener, false)
|
||||
if len(targets) != 2 {
|
||||
t.Fatalf("required dual-stack listener targets = %v", targets)
|
||||
}
|
||||
for _, failure := range []error{errors.New("IPv6 disabled"), errors.New("IPv6 unreachable"), context.DeadlineExceeded} {
|
||||
client := &http.Client{Transport: serviceHealthRoundTripper(func(r *http.Request) (*http.Response, error) {
|
||||
if r.URL.Hostname() == "::1" {
|
||||
|
|
@ -87,7 +96,6 @@ func TestServiceHealthIPv4SurvivesUnavailableIPv6(t *testing.T) {
|
|||
if old.Healthy || !old.Observed {
|
||||
t.Fatalf("IPv6-only control = %#v", old)
|
||||
}
|
||||
targets := localServiceHealthBaseURLs(serviceHealthAddrListener{addr: &net.TCPAddr{IP: net.IPv6unspecified, Port: 7655}}, false)
|
||||
got := serviceHealthProbe(targets, client, time.Second, 0)()
|
||||
if !got.Observed || !got.Healthy || got.FailureCategory != "" {
|
||||
t.Fatalf("wildcard with available IPv4 = %#v", got)
|
||||
|
|
@ -98,9 +106,9 @@ func TestServiceHealthIPv4SurvivesUnavailableIPv6(t *testing.T) {
|
|||
func TestServiceHealthWildcardOnRealIPv4AndIPv6OnlyListeners(t *testing.T) {
|
||||
for _, network := range []string{"tcp4", "tcp6"} {
|
||||
t.Run(network, func(t *testing.T) {
|
||||
address := "127.0.0.1:0"
|
||||
address := "0.0.0.0:0"
|
||||
if network == "tcp6" {
|
||||
address = "[::1]:0"
|
||||
address = "[::]:0"
|
||||
}
|
||||
listener, err := net.Listen(network, address)
|
||||
if err != nil {
|
||||
|
|
@ -114,9 +122,7 @@ func TestServiceHealthWildcardOnRealIPv4AndIPv6OnlyListeners(t *testing.T) {
|
|||
done := make(chan struct{})
|
||||
go func() { defer close(done); _ = server.Serve(listener) }()
|
||||
t.Cleanup(func() { _ = server.Close(); <-done })
|
||||
bound := listener.Addr().(*net.TCPAddr)
|
||||
wildcard := serviceHealthAddrListener{Listener: listener, addr: &net.TCPAddr{IP: net.IPv6unspecified, Port: bound.Port}}
|
||||
got := newServiceHealthProbe(wildcard, false)()
|
||||
got := newServiceHealthProbe(listener, false)()
|
||||
if !got.Observed || !got.Healthy || got.FailureCategory != "" {
|
||||
t.Fatalf("real %s server through wildcard = %#v", network, got)
|
||||
}
|
||||
|
|
|
|||
7
pkg/server/service_health_socket_other.go
Normal file
7
pkg/server/service_health_socket_other.go
Normal file
|
|
@ -0,0 +1,7 @@
|
|||
//go:build !unix && !windows
|
||||
|
||||
package server
|
||||
|
||||
func serviceHealthSocketIsDualStack(_ uintptr) bool {
|
||||
return false
|
||||
}
|
||||
153
pkg/server/service_health_socket_test.go
Normal file
153
pkg/server/service_health_socket_test.go
Normal file
|
|
@ -0,0 +1,153 @@
|
|||
package server
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"reflect"
|
||||
"sync/atomic"
|
||||
"syscall"
|
||||
"testing"
|
||||
|
||||
"github.com/rcourtman/pulse-go-rewrite/internal/telemetry"
|
||||
)
|
||||
|
||||
func startServiceHealthListener(t *testing.T, listener net.Listener, tlsEnabled bool, handler http.Handler) {
|
||||
t.Helper()
|
||||
server := httptest.NewUnstartedServer(handler)
|
||||
_ = server.Listener.Close()
|
||||
server.Listener = listener
|
||||
if tlsEnabled {
|
||||
server.StartTLS()
|
||||
} else {
|
||||
server.Start()
|
||||
}
|
||||
t.Cleanup(server.Close)
|
||||
}
|
||||
|
||||
func serveHealthyServiceHealthFixture(w http.ResponseWriter, r *http.Request) {
|
||||
response, _ := healthyServiceHealthResponse(r)
|
||||
defer response.Body.Close()
|
||||
_, _ = io.Copy(w, response.Body)
|
||||
}
|
||||
|
||||
// Adapted from PR2351's real HTTP/HTTPS wildcard regression. Keep main's
|
||||
// bounded alternate-family retry, but only on sockets that own both families.
|
||||
func TestServiceHealthProbeHandlesWildcardListeners(t *testing.T) {
|
||||
for _, test := range []struct {
|
||||
name, network, address, wantHost string
|
||||
wantTargets int
|
||||
}{
|
||||
{"ipv4", "tcp4", "0.0.0.0:0", "127.0.0.1", 1},
|
||||
{"dual-stack", "tcp", "[::]:0", "127.0.0.1", 2},
|
||||
{"ipv6-only", "tcp6", "[::]:0", "::1", 1},
|
||||
} {
|
||||
for _, tlsEnabled := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("%s/tls=%v", test.name, tlsEnabled), func(t *testing.T) {
|
||||
listener, err := net.Listen(test.network, test.address)
|
||||
if err != nil {
|
||||
t.Fatalf("required %s listener: %v", test.name, err)
|
||||
}
|
||||
startServiceHealthListener(t, listener, tlsEnabled, http.HandlerFunc(serveHealthyServiceHealthFixture))
|
||||
baseURLs := localServiceHealthBaseURLs(listener, tlsEnabled)
|
||||
if len(baseURLs) != test.wantTargets {
|
||||
t.Fatalf("probe targets = %v, want %d owned families", baseURLs, test.wantTargets)
|
||||
}
|
||||
parsed, err := url.Parse(baseURLs[0])
|
||||
if err != nil || parsed.Hostname() != test.wantHost {
|
||||
t.Fatalf("probe targets = %v, want bound listener through %s", baseURLs, test.wantHost)
|
||||
}
|
||||
got := newServiceHealthProbe(listener, tlsEnabled)()
|
||||
if !got.Observed || !got.Healthy || got.FailureCategory != "" {
|
||||
t.Fatalf("reachable listener reported unhealthy: %#v", got)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestServiceHealthWildcardNeverObservesOtherSocketOnSamePort(t *testing.T) {
|
||||
for _, network := range []string{"tcp4", "tcp6"} {
|
||||
for _, tlsEnabled := range []bool{false, true} {
|
||||
for _, ownHealthy := range []bool{false, true} {
|
||||
t.Run(fmt.Sprintf("%s/tls=%v/healthy=%v", network, tlsEnabled, ownHealthy), func(t *testing.T) {
|
||||
ownAddress, otherNetwork, otherHost := "0.0.0.0:0", "tcp6", "::1"
|
||||
if network == "tcp6" {
|
||||
ownAddress, otherNetwork, otherHost = "[::]:0", "tcp4", "127.0.0.1"
|
||||
}
|
||||
listener, err := net.Listen(network, ownAddress)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = listener.Close() })
|
||||
port := listener.Addr().(*net.TCPAddr).Port
|
||||
other, err := net.Listen(otherNetwork, net.JoinHostPort(otherHost, fmt.Sprint(port)))
|
||||
if err != nil {
|
||||
t.Fatalf("required separate socket on same port: %v", err)
|
||||
}
|
||||
var otherRequests atomic.Int32
|
||||
startServiceHealthListener(t, other, tlsEnabled, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
otherRequests.Add(1)
|
||||
if ownHealthy {
|
||||
http.Error(w, "unrelated server unavailable", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
serveHealthyServiceHealthFixture(w, r)
|
||||
}))
|
||||
startServiceHealthListener(t, listener, tlsEnabled, http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if !ownHealthy {
|
||||
http.Error(w, "own server unavailable", http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
serveHealthyServiceHealthFixture(w, r)
|
||||
}))
|
||||
got := newServiceHealthProbe(listener, tlsEnabled)()
|
||||
wantCategory := ""
|
||||
if !ownHealthy {
|
||||
wantCategory = telemetry.ServiceHealthFailureAPIStatus
|
||||
}
|
||||
if !got.Observed || got.Healthy != ownHealthy || got.FailureCategory != wantCategory || otherRequests.Load() != 0 {
|
||||
t.Fatalf("observation=%#v want healthy=%v category=%q; unrelated requests=%d", got, ownHealthy, wantCategory, otherRequests.Load())
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
type serviceHealthSocketErrorListener struct {
|
||||
serviceHealthAddrListener
|
||||
raw syscall.RawConn
|
||||
err error
|
||||
}
|
||||
|
||||
func (l serviceHealthSocketErrorListener) SyscallConn() (syscall.RawConn, error) { return l.raw, l.err }
|
||||
|
||||
type serviceHealthControlError struct{}
|
||||
|
||||
func (serviceHealthControlError) Control(func(uintptr)) error {
|
||||
return errors.New("control unavailable")
|
||||
}
|
||||
func (serviceHealthControlError) Read(func(uintptr) bool) error {
|
||||
return errors.New("read unavailable")
|
||||
}
|
||||
func (serviceHealthControlError) Write(func(uintptr) bool) error {
|
||||
return errors.New("write unavailable")
|
||||
}
|
||||
|
||||
func TestServiceHealthUnknownSocketModeStaysOnIPv6(t *testing.T) {
|
||||
addressOnly := serviceHealthAddrListener{addr: &net.TCPAddr{IP: net.IPv6unspecified, Port: 7655}}
|
||||
for _, listener := range []net.Listener{
|
||||
addressOnly,
|
||||
serviceHealthSocketErrorListener{serviceHealthAddrListener: addressOnly, err: errors.New("socket unavailable")},
|
||||
serviceHealthSocketErrorListener{serviceHealthAddrListener: addressOnly, raw: serviceHealthControlError{}},
|
||||
} {
|
||||
if got := localServiceHealthBaseURLs(listener, false); !reflect.DeepEqual(got, []string{"http://[::1]:7655"}) {
|
||||
t.Fatalf("uninspectable socket targets = %v", got)
|
||||
}
|
||||
}
|
||||
}
|
||||
10
pkg/server/service_health_socket_unix.go
Normal file
10
pkg/server/service_health_socket_unix.go
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
//go:build unix
|
||||
|
||||
package server
|
||||
|
||||
import "golang.org/x/sys/unix"
|
||||
|
||||
func serviceHealthSocketIsDualStack(fd uintptr) bool {
|
||||
v6Only, err := unix.GetsockoptInt(int(fd), unix.IPPROTO_IPV6, unix.IPV6_V6ONLY)
|
||||
return err == nil && v6Only == 0
|
||||
}
|
||||
10
pkg/server/service_health_socket_windows.go
Normal file
10
pkg/server/service_health_socket_windows.go
Normal file
|
|
@ -0,0 +1,10 @@
|
|||
//go:build windows
|
||||
|
||||
package server
|
||||
|
||||
import "golang.org/x/sys/windows"
|
||||
|
||||
func serviceHealthSocketIsDualStack(fd uintptr) bool {
|
||||
v6Only, err := windows.GetsockoptInt(windows.Handle(fd), windows.IPPROTO_IPV6, windows.IPV6_V6ONLY)
|
||||
return err == nil && v6Only == 0
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue