Replay every unanimous upstream status, not just 4xx
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful

openvoxdb does not reserve 5xx for its own faults: the same malformed
query is a 400 on /nodes and a 500 on /facts, and /metrics answers a flat
403, so a 4xx-only replay rule made pdbmux's behaviour depend on the
route. The meta and metrics handlers held their own copy of the gateway
error and bypassed the replay entirely.

- Replay any status from 400 up that every backend agreed on, with the
  backend's own body and content type.
- Keep 502 for backends disagreeing on the status, or a backend that
  answered nothing at all.
- Route /pdb/meta, /metrics and the pass-through path through the same
  rule as the merged query handlers.
- Count a unanimous 5xx as a failed round and let it fall back to a stale
  cache entry; only a unanimous 4xx stays exempt from both.
- Answer successful queries with openvoxdb's application/json;charset=utf-8.
This commit is contained in:
2026-09-13 13:35:21 +10:00
parent 5207a79ad4
commit 121bfacc2f
11 changed files with 541 additions and 126 deletions
+201 -29
View File
@@ -14,11 +14,16 @@ import (
const badOrderBy = `Unrecognized column 'bogus' specified in :order_by`
// upstreamJSONContentType is spelled out rather than taken from jsonContentType,
// so the assertion is against what openvoxdb answers and not against whatever
// pdbmux happens to be configured with.
const upstreamJSONContentType = "application/json;charset=utf-8"
func upErr(status int, body string) error {
return newUpstreamError(status, "text/plain; charset=utf-8", []byte(body))
}
func TestUnanimousClientError_AgreementRule(t *testing.T) {
func TestUnanimousUpstreamError_AgreementRule(t *testing.T) {
refused := errors.New("dial tcp: connection refused")
for _, tc := range []struct {
name string
@@ -29,18 +34,25 @@ func TestUnanimousClientError_AgreementRule(t *testing.T) {
{"all agree on 415", []error{upErr(415, "bad media"), upErr(415, "bad media")}, 415},
{"differing bodies still agree", []error{upErr(400, "one"), upErr(400, "two")}, 400},
{"disagreeing 4xx", []error{upErr(400, "bad"), upErr(404, "gone")}, 0},
// Both client-shaped, so only the same-status rule refuses these.
{"agree on shape, disagree on code", []error{upErr(400, "x"), upErr(415, "y")}, 0},
{"a majority agrees", []error{upErr(400, "x"), upErr(400, "y"), upErr(422, "z")}, 0},
{"the first differs", []error{upErr(422, "z"), upErr(400, "x"), upErr(400, "y")}, 0},
{"4xx with a 5xx", []error{upErr(400, "bad"), upErr(500, "boom")}, 0},
{"4xx with a transport failure", []error{upErr(400, "bad"), refused}, 0},
{"all 403", []error{upErr(403, "denied"), upErr(403, "denied")}, 0},
{"all 404", []error{upErr(404, "gone"), upErr(404, "gone")}, 0},
{"all 429", []error{upErr(429, "slow down"), upErr(429, "slow down")}, 0},
{"all 408", []error{upErr(408, "too slow"), upErr(408, "too slow")}, 0},
{"all 500", []error{upErr(500, "boom"), upErr(500, "boom")}, 0},
{"all 503", []error{upErr(503, "unavailable"), upErr(503, "unavailable")}, 0},
{"5xx with a transport failure", []error{upErr(500, "boom"), refused}, 0},
// Statuses the old 4xx-only rule swallowed. Every one of them is what
// both backends actually said, so every one of them is the answer.
{"all 403", []error{upErr(403, "denied"), upErr(403, "denied")}, 403},
{"all 404", []error{upErr(404, "gone"), upErr(404, "gone")}, 404},
{"all 429", []error{upErr(429, "slow down"), upErr(429, "slow down")}, 429},
{"all 408", []error{upErr(408, "too slow"), upErr(408, "too slow")}, 408},
{"all 500", []error{upErr(500, "boom"), upErr(500, "boom")}, 500},
{"all 503", []error{upErr(503, "unavailable"), upErr(503, "unavailable")}, 503},
{"disagreeing 5xx", []error{upErr(500, "boom"), upErr(503, "later")}, 0},
// Nothing below 400 is an error a backend explained, and a 3xx's meaning
// lives in a Location header the fan-out never kept.
{"all 302", []error{upErr(302, ""), upErr(302, "")}, 0},
{"all 204", []error{upErr(204, ""), upErr(204, "")}, 0},
{"no backends", nil, 0},
{"a backend succeeded", []error{upErr(400, "bad"), nil}, 0},
} {
@@ -49,7 +61,7 @@ func TestUnanimousClientError_AgreementRule(t *testing.T) {
for i, err := range tc.errs {
results[i] = backendResult{name: fmt.Sprintf("b%d", i), err: err}
}
got := unanimousClientError(results)
got := unanimousUpstreamError(backendUpstreamErrors(results))
switch {
case tc.want == 0 && got != nil:
t.Fatalf("want no replay, got HTTP %d", got.status)
@@ -64,17 +76,33 @@ func TestUnanimousClientError_AgreementRule(t *testing.T) {
// Two backends can explain the same rejection differently; the reply is the
// first in configured order so a retried query is answered the same way twice.
func TestUnanimousClientError_PicksFirstInConfiguredOrder(t *testing.T) {
func TestUnanimousUpstreamError_PicksFirstInConfiguredOrder(t *testing.T) {
results := []backendResult{
{name: "a", err: upErr(400, "from a")},
{name: "b", err: upErr(400, "from b")},
}
got := unanimousClientError(results)
got := unanimousUpstreamError(backendUpstreamErrors(results))
if got == nil || !strings.Contains(string(got.body), "from a") {
t.Fatalf("body = %q, want the first backend's", got)
}
}
// Only a 4xx says the estate is well and the query was wrong. A 5xx is replayed
// too, but it must keep counting as a backend fault.
func TestClientRefusal_OnlyFourXX(t *testing.T) {
for status, want := range map[int]bool{400: true, 404: true, 429: true, 499: true, 500: false, 503: false} {
if got := clientRefusal(upErr(status, "x")); got != want {
t.Errorf("clientRefusal(%d) = %v, want %v", status, got, want)
}
}
if clientRefusal(errAllBackendsFailed) {
t.Error("a transport failure is not a client refusal")
}
if clientRefusal(unanimousUpstreamError(nil)) {
t.Error("an absent unanimous error is not a client refusal")
}
}
func TestUpstreamError_BodyIsCapped(t *testing.T) {
ue := newUpstreamError(400, "text/plain", []byte(strings.Repeat("x", upstreamBodyLimit*2)))
if len(ue.body) != upstreamBodyLimit {
@@ -112,31 +140,78 @@ func TestHandler_MergedReplaysNonBadRequestStatus(t *testing.T) {
}
}
// A unanimous 403 is pdbmux's own credentials being refused, not the client's
// query, so it must not be handed back as the client's fault.
func TestHandler_MergedForbiddenStays502(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
a.reject, a.rejectBody = http.StatusForbidden, "certificate not allowed"
b.reject, b.rejectBody = http.StatusForbidden, "certificate not allowed"
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
// A status every backend agreed on is the estate's own answer whatever it is,
// so the statuses the 4xx-only rule used to swallow now reach the client.
func TestHandler_MergedReplaysEveryUnanimousStatus(t *testing.T) {
for _, tc := range []struct {
status int
body string
}{
{http.StatusForbidden, "certificate not allowed"},
{http.StatusNotFound, "Not Found"},
{http.StatusTooManyRequests, "slow down"},
{http.StatusInternalServerError, "Output of convert-query-params does not match schema"},
{http.StatusServiceUnavailable, "starting up"},
} {
t.Run(http.StatusText(tc.status), func(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
a.reject, a.rejectBody = tc.status, tc.body
b.reject, b.rejectBody = tc.status, tc.body
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
if rec := doGet(t, srv.Handler(), nodesPath, ""); rec.Code != http.StatusBadGateway {
t.Fatalf("status = %d, want 502", rec.Code)
rec := doGet(t, srv.Handler(), nodesPath, "")
if rec.Code != tc.status {
t.Fatalf("status = %d, want the upstream %d: %s", rec.Code, tc.status, rec.Body.String())
}
if !strings.Contains(rec.Body.String(), tc.body) {
t.Errorf("body = %q, want the upstream explanation", rec.Body.String())
}
})
}
}
// Every backend answering 404 is ambiguous between a bad path and a record
// nobody holds, so it stays a gateway error on a merged route.
func TestHandler_MergedNotFoundStays502(t *testing.T) {
// A backend that answered nothing at all leaves no unanimity to replay, whatever
// the others said.
func TestHandler_MergedServerErrorWithUnreachableStays502(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
a.reject, a.rejectBody = http.StatusInternalServerError, "boom"
dead := newFakeBackend(t, `[]`, `[]`)
deadURL := dead.srv.URL
dead.srv.Close()
srv := newTestServer(testConfig(a.srv.URL, deadURL, mergeStatic))
if rec := doGet(t, srv.Handler(), nodesPath, ""); rec.Code != http.StatusBadGateway {
t.Fatalf("status = %d, want 502 when a backend never answered", rec.Code)
}
}
func TestHandler_MergedDisagreeingServerErrorsStay502(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
a.reject, a.rejectBody = http.StatusNotFound, "Not Found"
b.reject, b.rejectBody = http.StatusNotFound, "Not Found"
a.reject, a.rejectBody = http.StatusInternalServerError, "boom"
b.reject, b.rejectBody = http.StatusServiceUnavailable, "later"
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
if rec := doGet(t, srv.Handler(), nodesPath, ""); rec.Code != http.StatusBadGateway {
t.Fatalf("status = %d, want 502", rec.Code)
t.Fatalf("status = %d, want 502 when backends disagree", rec.Code)
}
}
// A unanimous 5xx is replayed, but it is the backends admitting a fault, so it
// still has to read as degraded service.
func TestHandler_UnanimousServerErrorIsStillADegradedRound(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
a.reject, a.rejectBody = http.StatusInternalServerError, "boom"
b.reject, b.rejectBody = http.StatusInternalServerError, "boom"
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
if rec := doGet(t, srv.Handler(), nodesPath, ""); rec.Code != http.StatusInternalServerError {
t.Fatalf("status = %d, want the upstream 500", rec.Code)
}
if hr := health(t, srv); hr.Query.PartialRounds == 0 {
t.Errorf("query health = %+v, want the 5xx counted as a failed round", hr.Query)
}
}
@@ -255,7 +330,7 @@ func TestHandler_ReportSubResourceRejectionReplayedButNotItsAbsence(t *testing.T
b := newFakeBackend(t, `[]`, `[]`)
if rec := doGet(t, newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic)).Handler(),
reportsPath+"/nope/events", ""); rec.Code != http.StatusNotFound {
t.Fatalf("status = %d, want pdbmux's own 404 when nobody holds the report", rec.Code)
t.Fatalf("status = %d, want 404 when nobody holds the report", rec.Code)
}
a.reject, a.rejectBody = http.StatusBadRequest, badOrderBy
@@ -316,7 +391,7 @@ func TestHandler_RejectedQueryIsNotCached(t *testing.T) {
func TestHandler_AllBackendsFailedIsNotCached(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
a.fail, b.fail = true, true
a.dead, b.dead = true, true
srv, _ := newCachedServer(t, cacheTestConfig(a.srv.URL, b.srv.URL))
if rec := doGet(t, srv.Handler(), factsPath, ""); rec.Code != http.StatusBadGateway {
@@ -493,3 +568,100 @@ func hostOf(t *testing.T, raw string) string {
}
return u.Host
}
// Every backend answering 500 is the estate's answer on the first-holder route
// as much as on a merged one.
func TestHandler_ReportSubResourceReplaysUnanimousServerError(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
a.reject, a.rejectBody = http.StatusInternalServerError, "boom"
b.reject, b.rejectBody = http.StatusInternalServerError, "boom"
rec := doGet(t, newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic)).Handler(),
reportsPath+"/nope/events", "")
if rec.Code != http.StatusInternalServerError {
t.Fatalf("status = %d, want the upstream 500: %s", rec.Code, rec.Body.String())
}
}
// One backend saying 404 and one saying nothing leaves it unknown whether the
// report exists, which is the gateway error's job to say.
func TestHandler_ReportSubResourceUnreachableBackendIs502(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
b.dead = true
rec := doGet(t, newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic)).Handler(),
reportsPath+"/nope/events", "")
if rec.Code != http.StatusBadGateway {
t.Fatalf("status = %d, want 502 when a backend never answered", rec.Code)
}
}
// The pass-through route used to replay whichever backend answered first, even
// when the others disagreed or never answered. It follows the same rule now.
func TestProxyUnmerged_ReplayRule(t *testing.T) {
for _, tc := range []struct {
name string
aStatus int
bStatus int
bDead bool
want int
}{
{"unanimous 500 is replayed", 500, 500, false, 500},
{"unanimous 404 is replayed", 404, 404, false, 404},
{"unanimous 403 is replayed", 403, 403, false, 403},
{"disagreement is a gateway error", 400, 404, false, http.StatusBadGateway},
{"a silent backend is a gateway error", 404, 0, true, http.StatusBadGateway},
} {
t.Run(tc.name, func(t *testing.T) {
a := newFakeBackend(t, `[]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
a.reject, a.rejectBody = tc.aStatus, "upstream said no"
b.reject, b.rejectBody = tc.bStatus, "upstream said no"
b.dead = tc.bDead
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
rec := doGet(t, srv.Handler(), resourcesPath, `["=","type","File"]`)
if rec.Code != tc.want {
t.Fatalf("status = %d, want %d: %s", rec.Code, tc.want, rec.Body.String())
}
})
}
}
// pdbmux has to be indistinguishable from a PuppetDB, and openvoxdb answers a
// query with a charset on the content type.
func TestHandler_SuccessContentTypeMatchesUpstream(t *testing.T) {
a := newFakeBackend(t, `[`+node("h1", "2026-07-20T00:00:00Z")+`]`, `[]`)
b := newFakeBackend(t, `[]`, `[]`)
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
h := srv.Handler()
for _, path := range []string{nodesPath, factsPath, factNamesPath, reportsPath} {
t.Run(path, func(t *testing.T) {
rec := doGet(t, h, path, "")
if rec.Code != http.StatusOK {
t.Fatalf("status = %d: %s", rec.Code, rec.Body.String())
}
if got := rec.Header().Get("Content-Type"); got != upstreamJSONContentType {
t.Errorf("Content-Type = %q, want %q", got, upstreamJSONContentType)
}
})
}
}
func TestHandler_CachedSuccessContentTypeMatchesUpstream(t *testing.T) {
a := newFakeBackend(t, `[]`, `[`+fact("h1", "role", "web", "")+`]`)
b := newFakeBackend(t, `[]`, `[]`)
srv, _ := newCachedServer(t, cacheTestConfig(a.srv.URL, b.srv.URL))
h := srv.Handler()
// The miss builds the entry and the hit replays it; both are the client's view.
for _, label := range []string{"miss", "hit"} {
rec := doGet(t, h, factsPath, "")
if got := rec.Header().Get("Content-Type"); got != upstreamJSONContentType {
t.Errorf("%s Content-Type = %q, want %q", label, got, upstreamJSONContentType)
}
}
}