Fix PBS node identity retry classification

This commit is contained in:
rcourtman 2026-08-08 09:39:02 +01:00
parent 8e2858dac0
commit 8ff64e0dee
7 changed files with 466 additions and 56 deletions

View file

@ -5926,6 +5926,17 @@
"path": "internal/unifiedresources/registry.go",
"kind": "file"
},
{
"repo": "pulse",
"path": "pkg/pbs/client.go",
"kind": "file"
},
{
"repo": "pulse",
"path": "pkg/pbs/client_http_test.go",
"kind": "file",
"evidence_tier": "test-proof"
},
{
"repo": "pulse",
"path": "scripts/release_control/subsystem_lookup_test.py",

View file

@ -405,6 +405,7 @@ only graft them onto fixture data.
21b. `pkg/proxmox/client.go`
21c. `pkg/proxmox/io_counters.go`
22. `pkg/proxmox/zfs.go`
22a. `pkg/pbs/client.go`
23. `internal/monitoring/guest_memory_sources.go`
24. `internal/monitoring/guest_memory_stability.go`
25. `internal/monitoring/monitor_polling_vm.go`
@ -852,6 +853,14 @@ only graft them onto fixture data.
into the alert manager, but threshold selection, override identity, active
alert state, and notification delivery remain alerts-owned. The monitoring
sync bridge must not introduce per-platform evaluator branches.
20. Add or change PBS API transport, optional node identity collection, or PBS
HTTP retry classification through `pkg/pbs/client.go`. HTTP status decisions
must use the concrete client error status rather than rendered error or body
text. Node identity permission suppression is limited to 401 and 403;
429, 5xx, decoding, cancellation, timeout, and network failures remain
transient and retry on the next polling call. Concurrent callers must share
one in-flight `/nodes` request so immediate retry does not create a request
storm.
## Forbidden Paths
@ -972,6 +981,10 @@ only graft them onto fixture data.
`TestCrashedTrueNASAppContainerStillRaisesIncident` and
`TestStoppedTrueNASAppContainerStateStaysExited` in
`internal/truenas/provider_oneshot_containers_test.go`.
15. Keep PBS client HTTP error status structural and the node-name retry policy
explicit. Package proof must cover 401/403 deferral, 429/5xx and network
retry, response bodies containing permission-like text, recovery, caching,
and concurrent single-flight behavior under the race detector.
## Current State
@ -1027,7 +1040,15 @@ probe into a connection failure. Once connectivity is proven, the poll also
captures the hostname the PBS node reports about itself (`GET /nodes`) on
`models.PBSInstance.NodeName` as machine-identity evidence for connected-system
grouping; node-name fetch failure is partial data like the other optional
collections, never a poll failure.
collections, never a poll failure. `pkg/pbs/client.go` preserves HTTP status in
its concrete API error. Only 401 and 403 defer another node-name request for 30
minutes; 429, 5xx, malformed responses, cancellation, timeout, and network
errors retry on the next call regardless of error-body wording. A successful
node name remains cached for the client lifetime, and concurrent callers join
one in-flight request so a transient response produces one bounded request per
polling wave rather than one request per caller. The `GetNodeName` tests in
`pkg/pbs/client_http_test.go` are the focused retry, recovery, cache, and race
proof.
### Host snapshots carry integration provenance; doctor copy is user-facing

View file

@ -5270,7 +5270,8 @@
"internal/monitoring/",
"internal/storagehealth/",
"internal/truenas/",
"pkg/diskinventory/"
"pkg/diskinventory/",
"pkg/pbs/"
],
"owned_files": [
"docker-entrypoint.sh",
@ -5672,6 +5673,20 @@
"tests/integration/tests/43-platform-mock-runtime.spec.ts"
]
},
{
"id": "pbs-client-runtime",
"label": "PBS API transport and node identity proof",
"match_prefixes": [
"pkg/pbs/"
],
"match_files": [],
"allow_same_subsystem_tests": false,
"test_prefixes": [],
"exact_files": [
"pkg/pbs/client_http_test.go",
"pkg/pbs/client_test.go"
]
},
{
"id": "pbs-protection-evidence-runtime",
"label": "Proxmox protection evidence collection proof",

View file

@ -4,6 +4,7 @@ import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
@ -44,6 +45,14 @@ type Client struct {
nodeNameMu sync.Mutex
cachedNodeName string
nodeNameRetryAfter time.Time
nodeNameAttempt *nodeNameAttempt
}
type nodeNameAttempt struct {
done chan struct{}
joined int
name string
err error
}
// nodeNamePermissionRetryInterval throttles /nodes retries after a permission
@ -304,6 +313,45 @@ func shouldFallbackToForm(err error) bool {
return false
}
// apiHTTPError preserves the response status separately from the server body
// so callers can make retry and permission decisions without parsing text.
type apiHTTPError struct {
status int
body string
}
func (e *apiHTTPError) Error() string {
message := fmt.Sprintf("API error %d: %s", e.status, e.body)
if e.status == http.StatusUnauthorized || e.status == http.StatusForbidden {
return "authentication error: " + message
}
return message
}
func pbsHTTPStatus(err error) (int, bool) {
var apiErr *apiHTTPError
if errors.As(err, &apiErr) {
return apiErr.status, true
}
var authErr *authHTTPError
if errors.As(err, &authErr) {
return authErr.status, true
}
return 0, false
}
func isPBSPermissionError(err error) bool {
status, ok := pbsHTTPStatus(err)
return ok && (status == http.StatusUnauthorized || status == http.StatusForbidden)
}
func isPBSNotFoundError(err error) bool {
status, ok := pbsHTTPStatus(err)
return ok && status == http.StatusNotFound
}
// request performs an API request
func (c *Client) request(ctx context.Context, method, path string, data url.Values) (*http.Response, error) {
// Re-authenticate if needed
@ -369,18 +417,10 @@ func (c *Client) request(ctx context.Context, method, path string, data url.Valu
defer resp.Body.Close()
body, err := readResponseBodyLimited(resp.Body)
if err != nil {
return nil, err
return nil, &apiHTTPError{status: resp.StatusCode, body: fmt.Sprintf("read response body: %v", err)}
}
// Create base error
apiErr := fmt.Errorf("API error %d: %s", resp.StatusCode, string(body))
// Wrap with appropriate error type
if resp.StatusCode == 401 || resp.StatusCode == 403 {
return nil, fmt.Errorf("authentication error: %w", apiErr)
}
return nil, apiErr
return nil, &apiHTTPError{status: resp.StatusCode, body: string(body)}
}
return resp, nil
@ -508,7 +548,7 @@ func (c *Client) DeleteUserToken(ctx context.Context, userID, tokenName string)
resp, err := c.delete(ctx, path)
if err != nil {
// Deleting a missing token should be treated as already converged.
if strings.Contains(err.Error(), "API error 404") {
if isPBSNotFoundError(err) {
return nil
}
return fmt.Errorf("delete token: %w", err)
@ -635,17 +675,43 @@ func (c *Client) GetNodeName(ctx context.Context) (string, error) {
c.nodeNameMu.Unlock()
return "", fmt.Errorf("node name unavailable: /nodes permission denied, retry deferred")
}
if attempt := c.nodeNameAttempt; attempt != nil {
attempt.joined++
c.nodeNameMu.Unlock()
select {
case <-ctx.Done():
return "", ctx.Err()
case <-attempt.done:
return attempt.name, attempt.err
}
}
attempt := &nodeNameAttempt{done: make(chan struct{})}
c.nodeNameAttempt = attempt
c.nodeNameMu.Unlock()
name, err := c.fetchNodeName(ctx)
c.nodeNameMu.Lock()
attempt.name = name
attempt.err = err
if err == nil {
c.cachedNodeName = name
c.nodeNameRetryAfter = time.Time{}
} else if isPBSPermissionError(err) {
c.nodeNameRetryAfter = time.Now().Add(nodeNamePermissionRetryInterval)
}
c.nodeNameAttempt = nil
close(attempt.done)
c.nodeNameMu.Unlock()
return name, err
}
func (c *Client) fetchNodeName(ctx context.Context) (string, error) {
log.Debug().Msg("PBS GetNodeName: fetching node name")
resp, err := c.get(ctx, "/nodes")
if err != nil {
if strings.Contains(err.Error(), "403") || strings.Contains(err.Error(), "permission") {
c.nodeNameMu.Lock()
c.nodeNameRetryAfter = time.Now().Add(nodeNamePermissionRetryInterval)
c.nodeNameMu.Unlock()
}
return "", fmt.Errorf("failed to get nodes: %w", err)
}
defer resp.Body.Close()
@ -666,10 +732,6 @@ func (c *Client) GetNodeName(ctx context.Context) (string, error) {
// Return the first (usually only) node name
nodeName := result.Data[0].Node
c.nodeNameMu.Lock()
c.cachedNodeName = nodeName
c.nodeNameRetryAfter = time.Time{}
c.nodeNameMu.Unlock()
log.Debug().Str("nodeName", nodeName).Msg("PBS GetNodeName: found node name")
return nodeName, nil
}
@ -683,8 +745,7 @@ func (c *Client) GetNodeStatus(ctx context.Context) (*NodeStatus, error) {
// We'll gracefully handle the permission error and return nil
statusResp, err := c.get(ctx, "/nodes/localhost/status")
if err != nil {
// Check if this is a permission error (403)
if strings.Contains(err.Error(), "403") || strings.Contains(err.Error(), "permission") {
if isPBSPermissionError(err) {
log.Debug().Msg("PBS GetNodeStatus: permission denied (expected with API tokens) - returning nil")
return nil, nil // Return nil without error for permission issues
}
@ -1640,21 +1701,6 @@ func boolJobField(raw map[string]interface{}, keys ...string) bool {
return false
}
func isPBSPermissionError(err error) bool {
if err == nil {
return false
}
msg := strings.ToLower(err.Error())
return strings.Contains(msg, "403") || strings.Contains(msg, "401") || strings.Contains(msg, "permission") || strings.Contains(msg, "authentication")
}
func isPBSNotFoundError(err error) bool {
if err == nil {
return false
}
return strings.Contains(strings.ToLower(err.Error()), "404")
}
func firstNonEmptyString(values ...string) string {
for _, value := range values {
if strings.TrimSpace(value) != "" {

View file

@ -3,16 +3,25 @@ package pbs
import (
"context"
"encoding/json"
"errors"
"io"
"net/http"
"net/http/httptest"
"net/url"
"slices"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
)
type roundTripFunc func(*http.Request) (*http.Response, error)
func (f roundTripFunc) RoundTrip(req *http.Request) (*http.Response, error) {
return f(req)
}
func TestNewClient_TokenAuth_SetsAuthorizationHeader(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/api2/json/version" {
@ -465,14 +474,172 @@ func TestClient_GetNodeName_CachesResult(t *testing.T) {
}
func TestClient_GetNodeName_PermissionDenialDefersRetry(t *testing.T) {
for _, status := range []int{http.StatusUnauthorized, http.StatusForbidden} {
t.Run(http.StatusText(status), func(t *testing.T) {
var nodesCalls atomic.Int64
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/api2/json/nodes" {
t.Fatalf("unexpected path: %s", r.URL.Path)
}
nodesCalls.Add(1)
w.WriteHeader(status)
_, _ = w.Write([]byte(`{"message":"access denied"}`))
}))
defer server.Close()
client, err := NewClient(ClientConfig{
Host: server.URL,
TokenName: "root@pam!pulse-token",
TokenValue: "secret",
Timeout: 2 * time.Second,
})
if err != nil {
t.Fatalf("NewClient: %v", err)
}
_, firstErr := client.GetNodeName(context.Background())
if firstErr == nil {
t.Fatal("first GetNodeName: expected error")
}
if got, ok := pbsHTTPStatus(firstErr); !ok || got != status {
t.Fatalf("first GetNodeName status = (%d, %v), want (%d, true): %v", got, ok, status, firstErr)
}
for i := 0; i < 2; i++ {
if _, err := client.GetNodeName(context.Background()); err == nil {
t.Fatalf("deferred GetNodeName call %d: expected error", i)
}
}
if got := nodesCalls.Load(); got != 1 {
t.Fatalf("/nodes hit %d times, want 1 (%d defers retries)", got, status)
}
})
}
}
func TestClient_GetNodeName_TransientHTTPFailuresRetryAndRecover(t *testing.T) {
tests := []struct {
name string
status int
body string
}{
{name: "rate limited", status: http.StatusTooManyRequests, body: `{"message":"slow down"}`},
{name: "internal server error", status: http.StatusInternalServerError, body: `{"message":"temporary failure"}`},
{name: "service unavailable body mentions permission", status: http.StatusServiceUnavailable, body: `{"message":"permission service unavailable"}`},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
var nodesCalls atomic.Int64
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/api2/json/nodes" {
t.Fatalf("unexpected path: %s", r.URL.Path)
}
if nodesCalls.Add(1) == 1 {
w.WriteHeader(tc.status)
_, _ = w.Write([]byte(tc.body))
return
}
_ = json.NewEncoder(w).Encode(map[string]any{
"data": []map[string]any{{"node": "pbs-recovered"}},
})
}))
defer server.Close()
client, err := NewClient(ClientConfig{
Host: server.URL,
TokenName: "root@pam!pulse-token",
TokenValue: "secret",
Timeout: 2 * time.Second,
})
if err != nil {
t.Fatalf("NewClient: %v", err)
}
_, firstErr := client.GetNodeName(context.Background())
if firstErr == nil {
t.Fatal("first GetNodeName: expected error")
}
if got, ok := pbsHTTPStatus(firstErr); !ok || got != tc.status {
t.Fatalf("first GetNodeName status = (%d, %v), want (%d, true): %v", got, ok, tc.status, firstErr)
}
client.nodeNameMu.Lock()
retryAfter := client.nodeNameRetryAfter
client.nodeNameMu.Unlock()
if !retryAfter.IsZero() {
t.Fatalf("transient status %d set permission retry deferral to %s", tc.status, retryAfter)
}
name, err := client.GetNodeName(context.Background())
if err != nil {
t.Fatalf("recovery GetNodeName: %v", err)
}
if name != "pbs-recovered" {
t.Fatalf("recovery GetNodeName = %q, want pbs-recovered", name)
}
if got := nodesCalls.Load(); got != 2 {
t.Fatalf("/nodes hit %d times, want 2 (transient failure then recovery)", got)
}
})
}
}
func TestClient_GetNodeName_NetworkFailureRetriesAndRecovers(t *testing.T) {
client, err := NewClient(ClientConfig{
Host: "http://pbs.example.test",
TokenName: "root@pam!pulse-token",
TokenValue: "secret",
Timeout: 2 * time.Second,
})
if err != nil {
t.Fatalf("NewClient: %v", err)
}
var calls atomic.Int64
client.httpClient.Transport = roundTripFunc(func(req *http.Request) (*http.Response, error) {
if calls.Add(1) == 1 {
return nil, errors.New("network permission proxy failure")
}
return &http.Response{
StatusCode: http.StatusOK,
Header: make(http.Header),
Body: io.NopCloser(strings.NewReader(`{"data":[{"node":"pbs-network-recovered"}]}`)),
Request: req,
}, nil
})
if _, err := client.GetNodeName(context.Background()); err == nil {
t.Fatal("first GetNodeName: expected network error")
}
client.nodeNameMu.Lock()
retryAfter := client.nodeNameRetryAfter
client.nodeNameMu.Unlock()
if !retryAfter.IsZero() {
t.Fatalf("network error set permission retry deferral to %s", retryAfter)
}
name, err := client.GetNodeName(context.Background())
if err != nil {
t.Fatalf("recovery GetNodeName: %v", err)
}
if name != "pbs-network-recovered" {
t.Fatalf("recovery GetNodeName = %q, want pbs-network-recovered", name)
}
if got := calls.Load(); got != 2 {
t.Fatalf("transport called %d times, want 2", got)
}
}
func TestClient_GetNodeName_PermissionDeferralExpiresAndRecovers(t *testing.T) {
var nodesCalls atomic.Int64
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/api2/json/nodes" {
t.Fatalf("unexpected path: %s", r.URL.Path)
if nodesCalls.Add(1) == 1 {
w.WriteHeader(http.StatusForbidden)
_, _ = w.Write([]byte(`{"message":"access denied"}`))
return
}
nodesCalls.Add(1)
w.WriteHeader(http.StatusForbidden)
_, _ = w.Write([]byte(`{"message":"permission check failed"}`))
_ = json.NewEncoder(w).Encode(map[string]any{
"data": []map[string]any{{"node": "pbs-after-acl-repair"}},
})
}))
defer server.Close()
@ -486,24 +653,121 @@ func TestClient_GetNodeName_PermissionDenialDefersRetry(t *testing.T) {
t.Fatalf("NewClient: %v", err)
}
for i := 0; i < 3; i++ {
if _, err := client.GetNodeName(context.Background()); err == nil {
t.Fatalf("GetNodeName call %d: expected error", i)
if _, err := client.GetNodeName(context.Background()); err == nil {
t.Fatal("first GetNodeName: expected permission error")
}
if _, err := client.GetNodeName(context.Background()); err == nil {
t.Fatal("deferred GetNodeName: expected error")
}
if got := nodesCalls.Load(); got != 1 {
t.Fatalf("/nodes hit %d times during deferral, want 1", got)
}
client.nodeNameMu.Lock()
client.nodeNameRetryAfter = time.Now().Add(-time.Second)
client.nodeNameMu.Unlock()
name, err := client.GetNodeName(context.Background())
if err != nil {
t.Fatalf("GetNodeName after deferral expiry: %v", err)
}
if name != "pbs-after-acl-repair" {
t.Fatalf("GetNodeName after deferral expiry = %q, want pbs-after-acl-repair", name)
}
if got := nodesCalls.Load(); got != 2 {
t.Fatalf("/nodes hit %d times after recovery, want 2", got)
}
}
func TestClient_GetNodeName_ConcurrentTransientFailureIsSingleFlight(t *testing.T) {
const callers = 16
var nodesCalls atomic.Int64
requestStarted := make(chan struct{})
releaseFirstRequest := make(chan struct{})
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if nodesCalls.Add(1) == 1 {
close(requestStarted)
<-releaseFirstRequest
w.WriteHeader(http.StatusServiceUnavailable)
_, _ = w.Write([]byte(`{"message":"temporary permission backend outage"}`))
return
}
_ = json.NewEncoder(w).Encode(map[string]any{
"data": []map[string]any{{"node": "pbs-concurrent-recovery"}},
})
}))
defer server.Close()
client, err := NewClient(ClientConfig{
Host: server.URL,
TokenName: "root@pam!pulse-token",
TokenValue: "secret",
Timeout: 2 * time.Second,
})
if err != nil {
t.Fatalf("NewClient: %v", err)
}
start := make(chan struct{})
errs := make(chan error, callers)
var wg sync.WaitGroup
for i := 0; i < callers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
<-start
_, err := client.GetNodeName(context.Background())
errs <- err
}()
}
close(start)
<-requestStarted
deadline := time.Now().Add(2 * time.Second)
for {
client.nodeNameMu.Lock()
joined := 0
if client.nodeNameAttempt != nil {
joined = client.nodeNameAttempt.joined
}
client.nodeNameMu.Unlock()
if joined == callers-1 {
break
}
if time.Now().After(deadline) {
close(releaseFirstRequest)
t.Fatalf("only %d of %d concurrent callers joined the in-flight request", joined, callers-1)
}
time.Sleep(time.Millisecond)
}
close(releaseFirstRequest)
wg.Wait()
close(errs)
for err := range errs {
if err == nil {
t.Fatal("concurrent GetNodeName: expected transient error")
}
if got, ok := pbsHTTPStatus(err); !ok || got != http.StatusServiceUnavailable {
t.Fatalf("concurrent GetNodeName status = (%d, %v), want (503, true): %v", got, ok, err)
}
}
if got := nodesCalls.Load(); got != 1 {
t.Fatalf("/nodes hit %d times, want 1 (403 defers retries)", got)
t.Fatalf("/nodes hit %d times for concurrent transient failure, want 1", got)
}
// A transient (non-permission) failure must not defer: reset the deferral
// and confirm the next call goes back to the server.
client.nodeNameMu.Lock()
client.nodeNameRetryAfter = time.Time{}
client.nodeNameMu.Unlock()
if _, err := client.GetNodeName(context.Background()); err == nil {
t.Fatal("expected error after deferral reset")
name, err := client.GetNodeName(context.Background())
if err != nil {
t.Fatalf("recovery GetNodeName: %v", err)
}
if name != "pbs-concurrent-recovery" {
t.Fatalf("recovery GetNodeName = %q, want pbs-concurrent-recovery", name)
}
if _, err := client.GetNodeName(context.Background()); err != nil {
t.Fatalf("cached GetNodeName: %v", err)
}
if got := nodesCalls.Load(); got != 2 {
t.Fatalf("/nodes hit %d times after deferral reset, want 2", got)
t.Fatalf("/nodes hit %d times after recovery and cache read, want 2", got)
}
}

View file

@ -273,6 +273,7 @@ class CanonicalCompletionGuardTest(unittest.TestCase):
"proxmox-backup-identity-monitoring",
"container-entrypoint-runtime",
"mock-runtime-fixtures",
"pbs-client-runtime",
"pbs-protection-evidence-runtime",
"diskinventory-collection-trust",
"agent-fleet-diagnostics-runtime",
@ -380,6 +381,36 @@ class CanonicalCompletionGuardTest(unittest.TestCase):
],
)
def test_pbs_client_runtime_change_requires_monitoring_contract(self):
required = infer_impacted_subsystems(["pkg/pbs/client.go"])
self.assertEqual(set(required), {"monitoring"})
monitoring = required["monitoring"]
self.assertEqual(
monitoring["contract"],
"docs/release-control/v6/internal/subsystems/monitoring.md",
)
self.assertEqual(
monitoring["touched_runtime_files"],
["pkg/pbs/client.go"],
)
self.assertEqual(
monitoring["verification_requirements"],
[
{
"id": "pbs-client-runtime",
"label": "PBS API transport and node identity proof",
"touched_runtime_files": ["pkg/pbs/client.go"],
"allow_same_subsystem_tests": False,
"test_prefixes": [],
"exact_files": [
"pkg/pbs/client_http_test.go",
"pkg/pbs/client_test.go",
],
}
],
)
def test_host_agent_ingest_runtime_change_requires_monitoring_and_agent_lifecycle_contracts(
self,
):

View file

@ -4198,6 +4198,28 @@ class SubsystemLookupTest(unittest.TestCase):
],
)
def test_lookup_paths_assigns_pbs_client_runtime_to_monitoring(self) -> None:
result = lookup_paths(["pkg/pbs/client.go", "pkg/pbs/client_http_test.go"])
self.assertEqual(result["unowned_runtime_files"], [])
for file_entry in result["files"]:
self.assertEqual(len(file_entry["matches"]), 1)
match = file_entry["matches"][0]
self.assertEqual(match["subsystem"], "monitoring")
self.assertEqual(
match["contract"],
"docs/release-control/v6/internal/subsystems/monitoring.md",
)
self.assertEqual(match["lane_context"]["lane_id"], "L13")
self.assertEqual(
match["verification_requirement"]["id"],
"pbs-client-runtime",
)
self.assertEqual(
match["verification_requirement"]["exact_files"],
["pkg/pbs/client_http_test.go", "pkg/pbs/client_test.go"],
)
def test_lookup_paths_assigns_monitoring_metrics_history_runtime(self) -> None:
result = lookup_paths(["internal/monitoring/metrics_history.go"])
self.assertEqual(result["unowned_runtime_files"], [])