Merge pull request 'Report the backend host as the pdbmux_source value' (#31) from benvin/pdbmux-source-fact-fqdn into main
ci/woodpecker/tag/docker Pipeline was successful

Reviewed-on: #31
This commit was merged in pull request #31.
This commit is contained in:
2026-10-10 01:08:14 +11:00
14 changed files with 189 additions and 86 deletions
+8 -6
View File
@@ -24,7 +24,7 @@ not PQL) is forwarded verbatim.
| Path | Behaviour |
|---|---|
| `GET /pdb/query/v4/nodes` | Fan out to all backends, dedupe by `certname`, keep the record with the newer `report_timestamp`, stamped with the winning backend's name (see provenance). An `extract` query with a `function` column is **combined** instead. |
| `GET /pdb/query/v4/nodes` | Fan out to all backends, dedupe by `certname`, keep the record with the newer `report_timestamp`, stamped with the winning backend's host (see provenance). An `extract` query with a `function` column is **combined** instead. |
| `GET /pdb/query/v4/facts` | Fan out to all, and per `certname` keep **all** facts from the backend that owns that node (see merge semantics), plus a synthetic `pdbmux_source` fact naming it. An `extract` query with a `function` column is **combined** instead. |
| `GET /pdb/query/v4/facts/<name>[/<value>]` | Same fan-out and merge as `/facts`, and an `extract` query with a `function` column is **combined** the same way. The path segment is a `name` constraint, so no synthetic `pdbmux_source` record is added — except on the fact's own path, which is **synthesised** from the `/facts` merge (see provenance). |
| `GET /pdb/query/v4/fact-names` | Fan out to all and serve the **union** of the flat name arrays, deduped and re-sorted, re-paged across backends, plus the `pdbmux_source` name while injection is on. `order_by` is only valid on `name`. |
@@ -205,8 +205,10 @@ response itself, so nothing has to query each backend to find out:
- **`/facts`** gains one extra fact record per `certname`, alongside the node's
real facts, in the shape of a real fact record — `certname`, `name`, `value`,
`environment` — with `value` set to the **backend name** from `backends` /
`PDBMUX_BACKENDS`. `environment` is copied from that node's own facts (all
`environment` — with `value` set to the **hostname** of the backend's URL from
`backends` / `PDBMUX_BACKENDS` (no scheme or port), e.g.
`puppetdb.puppet.svc.cluster.local` for `new=http://puppetdb.puppet.svc.cluster.local:8080`.
Two backends may therefore not share a hostname; config validation rejects it. `environment` is copied from that node's own facts (all
four keys are always present, since clients index them directly).
- **`/nodes`** gains a `pdbmux_source` **key** on each merged node record. A
node record carries no facts, so this is a synthetic field, not a fact — the
@@ -268,10 +270,10 @@ response reports, and lets the request's own `query` narrow the result upstream.
That costs one `/facts` fan-out per request — the widest fan-out `pdbmux`
makes — on a rare, user-initiated path.
`/facts/pdbmux_source/<value>` pins the backend name, so it answers with the
`/facts/pdbmux_source/<value>` pins the backend host, so it answers with the
nodes that backend owns. The `<value>` segment never reaches that fan-out: the
synthetic record's value is always a backend name, so a value naming none is
answered `[]` from the configured names alone, with no fan-out at all, and a
synthetic record's value is always a backend host, so any other value — a
backend's name included — is answered `[]` from the configured hosts alone, with no fan-out at all, and a
value naming one filters a record set fetched under a key the value is not part
of. The record set is a property of the estate rather than of the filter, so
every value of it — and the unfiltered path — share one entry and one fetch.
+8 -8
View File
@@ -632,7 +632,7 @@ func TestHandler_LeaderDisconnectDoesNotFailFollowers(t *testing.T) {
body := `[` + node("h1", "2026-01-01T00:00:00.000Z") + `]`
a := newCountingBackend(t, map[string]string{nodesPath: body})
b := newCountingBackend(t, map[string]string{nodesPath: `[]`})
srv, _ := newCachedServer(t, cacheTestConfig(a.srv.URL, b.srv.URL))
srv, _ := newCachedServer(t, cacheTestConfig(withHost(a.srv.URL, hostA), withHost(b.srv.URL, hostB)))
h := srv.Handler()
release := make(chan struct{})
@@ -673,7 +673,7 @@ func TestHandler_LeaderDisconnectDoesNotFailFollowers(t *testing.T) {
}
// The leader built this body, so its provenance names the backend that
// answered the leader's fan-out.
want := `[` + stamped(t, node("h1", "2026-01-01T00:00:00.000Z"), defaultSourceFact, "a") + `]`
want := `[` + stamped(t, node("h1", "2026-01-01T00:00:00.000Z"), defaultSourceFact, hostA) + `]`
if got := strings.TrimSpace(follower.Body.String()); !sameJSON(t, got, want) {
t.Errorf("follower body = %s, want %s", got, want)
}
@@ -1488,16 +1488,16 @@ func stamped(t *testing.T, raw, field, value string) string {
func TestHandler_CachedNodesKeepSourceStamp(t *testing.T) {
a := newCountingBackend(t, map[string]string{nodesPath: `[` + node("h1", "2026-01-02T00:00:00.000Z") + `]`})
b := newCountingBackend(t, map[string]string{nodesPath: `[` + node("h1", "2026-01-01T00:00:00.000Z") + `]`})
srv, _ := newCachedServer(t, cacheTestConfig(a.srv.URL, b.srv.URL))
srv, _ := newCachedServer(t, cacheTestConfig(withHost(a.srv.URL, hostA), withHost(b.srv.URL, hostB)))
h := srv.Handler()
first := doGet(t, h, nodesPath, "")
if got := nodeSources(t, first.Body.Bytes(), defaultSourceFact); got["h1"] != "a" {
if got := nodeSources(t, first.Body.Bytes(), defaultSourceFact); got["h1"] != hostA {
t.Fatalf("first request stamp = %v, want h1 -> a", got)
}
second := doGet(t, h, nodesPath, "")
want := `[` + stamped(t, node("h1", "2026-01-02T00:00:00.000Z"), defaultSourceFact, "a") + `]`
want := `[` + stamped(t, node("h1", "2026-01-02T00:00:00.000Z"), defaultSourceFact, hostA) + `]`
if got := strings.TrimSpace(second.Body.String()); !sameJSON(t, got, want) {
t.Errorf("cached body = %s, want %s", got, want)
}
@@ -1512,7 +1512,7 @@ func TestHandler_CachedNodesKeepSourceStamp(t *testing.T) {
func TestHandler_CachedSourceAgesWithItsData(t *testing.T) {
a := newCountingBackend(t, map[string]string{nodesPath: `[` + node("h1", "2026-01-01T00:00:00.000Z") + `]`})
b := newCountingBackend(t, map[string]string{nodesPath: `[]`})
srv, clk := newCachedServer(t, cacheTestConfig(a.srv.URL, b.srv.URL))
srv, clk := newCachedServer(t, cacheTestConfig(withHost(a.srv.URL, hostA), withHost(b.srv.URL, hostB)))
h := srv.Handler()
doGet(t, h, nodesPath, "")
@@ -1522,7 +1522,7 @@ func TestHandler_CachedSourceAgesWithItsData(t *testing.T) {
b.setBody(nodesPath, `[`+node("h1", "2026-01-01T00:00:00.000Z")+`]`)
cached := doGet(t, h, nodesPath, "")
if got := nodeSources(t, cached.Body.Bytes(), defaultSourceFact); got["h1"] != "a" {
if got := nodeSources(t, cached.Body.Bytes(), defaultSourceFact); got["h1"] != hostA {
t.Errorf("cached stamp = %v, want h1 -> a: the body and its attribution come from the same fetch", got)
}
if got := a.hitCount(nodesPath); got != 1 {
@@ -1531,7 +1531,7 @@ func TestHandler_CachedSourceAgesWithItsData(t *testing.T) {
clk.advance(31 * time.Second)
rebuilt := doGet(t, h, nodesPath, "")
if got := nodeSources(t, rebuilt.Body.Bytes(), defaultSourceFact); got["h1"] != "b" {
if got := nodeSources(t, rebuilt.Body.Bytes(), defaultSourceFact); got["h1"] != hostB {
t.Errorf("rebuilt stamp = %v, want h1 -> b once the entry expired", got)
}
}
+22
View File
@@ -2,6 +2,7 @@ package main
import (
"fmt"
"net/url"
"os"
"path/filepath"
"strconv"
@@ -289,12 +290,22 @@ func (c Config) configHint() string {
return ConfigPath()
}
// backendHost is the upstream's hostname, the value of the source fact.
func backendHost(raw string) string {
u, err := url.Parse(raw)
if err != nil {
return ""
}
return strings.ToLower(u.Hostname())
}
func (c Config) Validate() error {
if len(c.Backends) == 0 {
return fmt.Errorf("no backends configured: set %sBACKENDS to \"name=url,name=url\" or add a backends list to %s",
envPrefix, c.configHint())
}
seen := map[string]bool{}
hosts := map[string]string{}
for _, b := range c.Backends {
if b.Name == "" || b.URL == "" {
return fmt.Errorf("backend requires both name and url: %+v", b)
@@ -303,6 +314,17 @@ func (c Config) Validate() error {
return fmt.Errorf("duplicate backend name %q", b.Name)
}
seen[b.Name] = true
if !c.SourceFactEnabled {
continue
}
host := backendHost(b.URL)
if host == "" {
return fmt.Errorf("backend %q url %q has no host", b.Name, b.URL)
}
if other, dup := hosts[host]; dup {
return fmt.Errorf("backends %q and %q share host %q; %s values would collide", other, b.Name, host, c.SourceFact)
}
hosts[host] = b.Name
}
switch c.Merge {
case mergeFreshness, mergeStatic:
+27 -3
View File
@@ -15,7 +15,7 @@ func testConfigValid() Config {
cfg := DefaultConfig()
cfg.Backends = []Backend{
{Name: "a", URL: "http://localhost:18080"},
{Name: "b", URL: "http://localhost:18081"},
{Name: "b", URL: "http://127.0.0.1:18081"},
}
return cfg
}
@@ -57,7 +57,7 @@ func TestLoad_NoBackendsLoadsButFailsValidation(t *testing.T) {
func TestLoad_BackendsKeepConfiguredOrder(t *testing.T) {
t.Setenv("XDG_CONFIG_HOME", t.TempDir())
clearEnv(t)
t.Setenv(envPrefix+"BACKENDS", "a=http://localhost:18080,b=http://localhost:18081")
t.Setenv(envPrefix+"BACKENDS", "a=http://localhost:18080,b=http://127.0.0.1:18081")
cfg, err := Load("")
if err != nil {
@@ -229,8 +229,20 @@ func TestValidate(t *testing.T) {
}{
{"ok", func(*Config) {}, false},
{"no backends", func(c *Config) { c.Backends = nil }, true},
{"dup name", func(c *Config) { c.Backends = append(c.Backends, Backend{Name: "a", URL: "x"}) }, true},
{"dup name", func(c *Config) { c.Backends = append(c.Backends, Backend{Name: "a", URL: "http://c"}) }, true},
{"missing url", func(c *Config) { c.Backends[0].URL = "" }, true},
{"url without host", func(c *Config) { c.Backends[0].URL = "localhost:18080" }, true},
{"shared host on another port", func(c *Config) { c.Backends[1].URL = "http://localhost:18081" }, true},
{"same host differing by scheme", func(c *Config) { c.Backends[1].URL = "https://localhost" }, true},
{"same host differing by case", func(c *Config) { c.Backends[1].URL = "http://LocalHost:18081" }, true},
{"shared host while source fact disabled", func(c *Config) {
c.SourceFactEnabled = false
c.Backends[1].URL = "http://localhost:18081"
}, false},
{"url without host while source fact disabled", func(c *Config) {
c.SourceFactEnabled = false
c.Backends[0].URL = "localhost:18080"
}, false},
{"bad merge", func(c *Config) { c.Merge = "wrong" }, true},
{"zero timeout", func(c *Config) { c.Timeout = 0 }, true},
{"empty source fact while enabled", func(c *Config) { c.SourceFact = "" }, true},
@@ -512,3 +524,15 @@ func clearEnv(t *testing.T) {
t.Setenv(envPrefix+k, "")
}
}
func TestBackendHost(t *testing.T) {
for raw, want := range map[string]string{
"https://PuppetDB.Example:8081/x": "puppetdb.example",
"http://[::1]:8080": "::1",
"localhost:8080": "",
} {
if got := backendHost(raw); got != want {
t.Errorf("backendHost(%q) = %q, want %q", raw, got, want)
}
}
}
+1 -1
View File
@@ -154,7 +154,7 @@ func startBackend(ctx context.Context, t fatalf, name, netName string) *backend
fail("starting openvoxdb for %s: %v", name, err)
}
return &backend{name: name, url: fmt.Sprintf("http://127.0.0.1:%d", port), pg: pg, pdb: pdb}
return &backend{name: name, url: fmt.Sprintf("http://%s.localhost:%d", name, port), pg: pg, pdb: pdb}
}
// waitForPuppetDBRunning gates on the trapperkeeper status service reporting
+3
View File
@@ -42,6 +42,9 @@ func fixtureTime(offset time.Duration) string {
const (
backendAName = "pdb-a"
backendBName = "pdb-b"
// The source fact's value is the backend host, so each backend gets its own.
backendAHost = backendAName + ".localhost"
backendBHost = backendBName + ".localhost"
)
type nodeFixture struct {
+2 -2
View File
@@ -77,8 +77,8 @@ func TestBackendDeathAndRecovery(t *testing.T) {
if got, _ := factValue(facts, nodeShared, "owner"); got != backendAName {
t.Errorf("owner of %s with %s down = %v, want %q", nodeShared, h.b.name, got, backendAName)
}
if got, _ := factValue(facts, nodeShared, defaultSourceFact); got != backendAName {
t.Errorf("%s for %s with %s down = %v, want %q", defaultSourceFact, nodeShared, h.b.name, got, backendAName)
if got, _ := factValue(facts, nodeShared, defaultSourceFact); got != backendAHost {
t.Errorf("%s for %s with %s down = %v, want %q", defaultSourceFact, nodeShared, h.b.name, got, backendAHost)
}
})
+2 -2
View File
@@ -84,8 +84,8 @@ func TestNodeLookupShowsProvenance(t *testing.T) {
if !strings.Contains(line, defaultSourceFact) {
continue
}
if !strings.Contains(line, backendBName) {
t.Fatalf("node-lookup reports %q for %s, want the owning backend %q", strings.TrimSpace(line), nodeShared, backendBName)
if !strings.Contains(line, backendBHost) {
t.Fatalf("node-lookup reports %q for %s, want the owning backend %q", strings.TrimSpace(line), nodeShared, backendBHost)
}
return
}
+12 -12
View File
@@ -286,8 +286,8 @@ func TestSharedNodeResolvesToTheFresherBackend(t *testing.T) {
if got := node["report_timestamp"]; got != tsSharedOnB {
t.Errorf("merged /nodes report_timestamp for %s = %v, want the fresher %q", nodeShared, got, tsSharedOnB)
}
if got := node[defaultSourceFact]; got != backendBName {
t.Errorf("merged /nodes %s for %s = %v, want %q", defaultSourceFact, nodeShared, got, backendBName)
if got := node[defaultSourceFact]; got != backendBHost {
t.Errorf("merged /nodes %s for %s = %v, want %q", defaultSourceFact, nodeShared, got, backendBHost)
}
rows := get(t, factsPath, query(`["=","certname","`+nodeShared+`"]`)).rows(t)
@@ -769,10 +769,10 @@ func TestSourceFactInjectionAndGating(t *testing.T) {
t.Run("facts carry one source record per node", func(t *testing.T) {
rows := get(t, factsPath, nil).rows(t)
want := map[string]string{
nodeAlpha: backendAName,
nodeBeta: backendBName,
nodeGamma: backendBName,
nodeShared: backendBName, // won on freshness, not on configured order
nodeAlpha: backendAHost,
nodeBeta: backendBHost,
nodeGamma: backendBHost,
nodeShared: backendBHost, // won on freshness, not on configured order
}
got := map[string]int{}
for _, row := range rows {
@@ -839,10 +839,10 @@ func TestSourceFactInjectionAndGating(t *testing.T) {
// records /facts carries rather than the empty set the backends hold.
func TestSourceFactDrilldown(t *testing.T) {
want := map[string]string{
nodeAlpha: backendAName,
nodeBeta: backendBName,
nodeGamma: backendBName,
nodeShared: backendBName, // won on freshness, not on configured order
nodeAlpha: backendAHost,
nodeBeta: backendBHost,
nodeGamma: backendBHost,
nodeShared: backendBHost, // won on freshness, not on configured order
}
path := factsPath + "/" + defaultSourceFact
@@ -882,8 +882,8 @@ func TestSourceFactDrilldown(t *testing.T) {
value string
want []string
}{
{backendAName, []string{nodeAlpha}},
{backendBName, []string{nodeBeta, nodeGamma, nodeShared}},
{backendAHost, []string{nodeAlpha}},
{backendBHost, []string{nodeBeta, nodeGamma, nodeShared}},
{"nosuchbackend", nil},
} {
rows := get(t, path+"/"+tc.value, nil).rows(t)
+31 -29
View File
@@ -259,15 +259,15 @@ func TestHandler_FactsBySourceNamePathSynthesisesOneRecordPerNode(t *testing.T)
for _, merge := range []string{mergeStatic, mergeFreshness} {
t.Run(merge, func(t *testing.T) {
a, b := sourceFactBackends(t)
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, merge))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, merge))
// h3 lives in both: static keeps the first configured backend, freshness
// the one holding its newer report.
wantH3 := "a"
wantH3 := hostA
if merge == mergeFreshness {
wantH3 = "b"
wantH3 = hostB
}
want := map[string]string{"h1": "a", "h2": "b", "h3": wantH3}
want := map[string]string{"h1": hostA, "h2": hostB, "h3": wantH3}
rec := doGet(t, srv.Handler(), sourceFactURL, "")
if rec.Code != http.StatusOK {
@@ -315,20 +315,21 @@ func TestHandler_FactsBySourceNamePathCarriesTheOwnersEnvironment(t *testing.T)
}
}
// The pinned value names a backend, so the sub-route is the drilldown filtered
// by owner; a value naming no backend matches nothing.
// The pinned value is a backend host, so the sub-route is the drilldown filtered
// by owner; any other value, a backend's name included, matches nothing.
func TestHandler_FactsBySourceNameAndValueFiltersByOwner(t *testing.T) {
for _, tc := range []struct {
value string
want map[string]string
}{
{"a", map[string]string{"h1": "a"}},
{"b", map[string]string{"h2": "b", "h3": "b"}},
{hostA, map[string]string{"h1": hostA}},
{hostB, map[string]string{"h2": hostB, "h3": hostB}},
{"a", map[string]string{}},
{"nosuchbackend", map[string]string{}},
} {
t.Run(tc.value, func(t *testing.T) {
a, b := sourceFactBackends(t)
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeFreshness))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeFreshness))
rec := doGet(t, srv.Handler(), sourceFactURL+"/"+tc.value, "")
got, n := sourceValues(t, rec.Body.Bytes(), defaultSourceFact)
@@ -377,15 +378,15 @@ func TestHandler_FactsBySourceNameUnknownValueSkipsFanOut(t *testing.T) {
func TestHandler_FactsBySourceNameValuesFilterOneFetch(t *testing.T) {
a := newCountingBackend(t, map[string]string{factsPath: `[` + fact("h1", "osfamily", "RedHat", "") + `]`})
b := newCountingBackend(t, map[string]string{factsPath: `[` + fact("h2", "osfamily", "Debian", "") + `]`})
srv, _ := newCachedServer(t, cacheTestConfig(a.srv.URL, b.srv.URL))
srv, _ := newCachedServer(t, cacheTestConfig(withHost(a.srv.URL, hostA), withHost(b.srv.URL, hostB)))
for _, tc := range []struct {
path string
want map[string]string
}{
{sourceFactURL + "/a", map[string]string{"h1": "a"}},
{sourceFactURL + "/b", map[string]string{"h2": "b"}},
{sourceFactURL, map[string]string{"h1": "a", "h2": "b"}},
{sourceFactURL + "/" + hostA, map[string]string{"h1": hostA}},
{sourceFactURL + "/" + hostB, map[string]string{"h2": hostB}},
{sourceFactURL, map[string]string{"h1": hostA, "h2": hostB}},
} {
rec := doGet(t, srv.Handler(), tc.path, "")
got, n := sourceValues(t, rec.Body.Bytes(), defaultSourceFact)
@@ -434,12 +435,12 @@ func TestHandler_FactsBySourceNamePathRespectsQuery(t *testing.T) {
a, b := sourceFactBackends(t)
a.bodies[factsPath] = `[` + factEnv("h1", "osfamily", "RedHat", "prod") + `]`
b.bodies[factsPath] = `[]`
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeFreshness))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeFreshness))
const q = `["=","certname","h1"]`
rec := doGet(t, srv.Handler(), sourceFactURL, q)
got, n := sourceValues(t, rec.Body.Bytes(), defaultSourceFact)
if n != 1 || got["h1"] != "a" {
if n != 1 || got["h1"] != hostA {
t.Fatalf("%s?query=%s = %v (%d records), want only h1: %s", sourceFactURL, q, got, n, rec.Body.String())
}
// The query has to reach the backends for them to narrow anything.
@@ -457,11 +458,11 @@ func TestHandler_FactsBySourceNamePathSuppressesUpstream(t *testing.T) {
a.bodies[factsPath] = `[` + factEnv("h1", "osfamily", "RedHat", "prod") + `,` +
fact("h1", defaultSourceFact, "stale", "") + `]`
b.bodies[factsPath] = `[]`
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeFreshness))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeFreshness))
rec := doGet(t, srv.Handler(), sourceFactURL, "")
got, n := sourceValues(t, rec.Body.Bytes(), defaultSourceFact)
if n != 1 || got["h1"] != "a" {
if n != 1 || got["h1"] != hostA {
t.Fatalf("%s = %v (%d records), want exactly the synthetic one: %s", sourceFactURL, got, n, rec.Body.String())
}
}
@@ -714,20 +715,21 @@ func TestHandler_FactsQuerySelectsSourceFact(t *testing.T) {
query string
want map[string]string
}{
{`["=","name","` + src + `"]`, map[string]string{"h1": "a", "h2": "b", "h3": "b"}},
{`["and",["=","name","` + src + `"]]`, map[string]string{"h1": "a", "h2": "b", "h3": "b"}},
{`["and",["=","name","` + src + `"],["=","value","b"]]`, map[string]string{"h2": "b", "h3": "b"}},
{`["and",["=","name","` + src + `"],["not",["~","value","^b$"]]]`, map[string]string{"h1": "a"}},
{`["and",["=","certname","h3"],["=","name","` + src + `"]]`, map[string]string{"h3": "b"}},
{`["and",["or",["=","name","osfamily"],["=","name","` + src + `"]]]`, map[string]string{"h1": "a", "h2": "b", "h3": "b"}},
{`["in","name",["array",["osfamily","` + src + `"]]]`, map[string]string{"h1": "a", "h2": "b", "h3": "b"}},
{`["and",["or",["=","name","` + src + `"]],["=","environment","prod"]]`, map[string]string{"h1": "a"}},
{`["=","name","` + src + `"]`, map[string]string{"h1": hostA, "h2": hostB, "h3": hostB}},
{`["and",["=","name","` + src + `"]]`, map[string]string{"h1": hostA, "h2": hostB, "h3": hostB}},
{`["and",["=","name","` + src + `"],["=","value","` + hostB + `"]]`, map[string]string{"h2": hostB, "h3": hostB}},
{`["and",["=","name","` + src + `"],["=","value","b"]]`, map[string]string{}},
{`["and",["=","name","` + src + `"],["not",["~","value","^b\\.localhost$"]]]`, map[string]string{"h1": hostA}},
{`["and",["=","certname","h3"],["=","name","` + src + `"]]`, map[string]string{"h3": hostB}},
{`["and",["or",["=","name","osfamily"],["=","name","` + src + `"]]]`, map[string]string{"h1": hostA, "h2": hostB, "h3": hostB}},
{`["in","name",["array",["osfamily","` + src + `"]]]`, map[string]string{"h1": hostA, "h2": hostB, "h3": hostB}},
{`["and",["or",["=","name","` + src + `"]],["=","environment","prod"]]`, map[string]string{"h1": hostA}},
}
for _, tc := range tests {
t.Run(tc.query, func(t *testing.T) {
a, b := sourceFactBackends(t)
a.bodies[factsPath] = `[` + factEnv("h1", "osfamily", "RedHat", "prod") + `,` + fact("h1", src, "stale", "") + `]`
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeFreshness))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeFreshness))
rec := doGet(t, srv.Handler(), factsPath, tc.query)
if rec.Code != http.StatusOK {
@@ -796,12 +798,12 @@ func TestHandler_FactsQuerySourceFactNarrowsUpstream(t *testing.T) {
return false
}
}
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeFreshness))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeFreshness))
query := `["and",["=","certname","h3"],["=","environment","dev"],["=","name","` + defaultSourceFact + `"]]`
rec := doGet(t, srv.Handler(), factsPath, query)
if vals, n := sourceValues(t, rec.Body.Bytes(), defaultSourceFact); n != 1 || vals["h3"] != "b" {
t.Fatalf("got %v (%d records), want h3=b: %s", vals, n, rec.Body.String())
if vals, n := sourceValues(t, rec.Body.Bytes(), defaultSourceFact); n != 1 || vals["h3"] != hostB {
t.Fatalf("got %v (%d records), want h3=%s: %s", vals, n, hostB, rec.Body.String())
}
for _, q := range got {
if q == "" {
+12 -6
View File
@@ -53,7 +53,9 @@ var errAllBackendsFailed = errors.New("all backends failed")
const jsonContentType = "application/json;charset=utf-8"
type Server struct {
cfg Config
cfg Config
// hosts maps backend name to the hostname the source fact reports.
hosts map[string]string
client *http.Client
log *log.Logger
@@ -83,6 +85,10 @@ func NewServer(cfg Config, logger *log.Logger) *Server {
client: &http.Client{Timeout: cfg.Timeout},
log: logger,
now: time.Now,
hosts: map[string]string{},
}
for _, b := range cfg.Backends {
s.hosts[b.Name] = backendHost(b.URL)
}
if cfg.cacheEnabled() {
s.nodeCache = newMemoryCache(cfg.FactsTTL, cfg.CacheBytes)
@@ -476,10 +482,10 @@ func (s *Server) serveFactsByName(w http.ResponseWriter, r *http.Request) {
// value would be a fresh cache key, a fresh flight and a fresh whole-estate
// fan-out.
func (s *Server) serveSourceFact(w http.ResponseWriter, r *http.Request, inject *sourceInjector, value string, valued bool) {
// The synthetic record's value is always a backend name, so any other value
// The synthetic record's value is always a backend host, so any other value
// matches zero records. The response is complete rather than degraded, so it
// reports every configured backend.
if valued && !s.hasBackend(value) {
if valued && !s.hasBackendHost(value) {
n := len(s.cfg.Backends)
writeCached(w, cachedResponse{Records: -1, Backends: n, Configured: n})
return
@@ -505,9 +511,9 @@ func (s *Server) serveSourceFact(w http.ResponseWriter, r *http.Request, inject
})
}
func (s *Server) hasBackend(name string) bool {
for _, b := range s.cfg.Backends {
if b.Name == name {
func (s *Server) hasBackendHost(host string) bool {
for _, h := range s.hosts {
if h == host {
return true
}
}
+22
View File
@@ -4,6 +4,7 @@ import (
"encoding/json"
"io"
"log"
"net"
"net/http"
"net/http/httptest"
"net/url"
@@ -163,6 +164,27 @@ func truncate(t *testing.T, body, limit string) string {
return string(out)
}
// Test backends are reached as hostA/hostB, since the source fact's value is the
// backend host and both httptest servers listen on 127.0.0.1.
const (
hostA = "a.localhost"
hostB = "b.localhost"
)
func withHost(rawURL, host string) string {
u, err := url.Parse(rawURL)
if err != nil || u.Port() == "" {
return rawURL
}
u.Host = net.JoinHostPort(host, u.Port())
return u.String()
}
// hostedConfig is testConfig with the backends reached as hostA and hostB.
func hostedConfig(aURL, bURL, merge string) Config {
return testConfig(withHost(aURL, hostA), withHost(bURL, hostB), merge)
}
func testConfig(aURL, bURL, merge string) Config {
return Config{
Listen: ":0",
+5 -3
View File
@@ -13,6 +13,8 @@ import (
// callers need no branch.
type sourceInjector struct {
name string
// hosts maps backend name to the value emitted.
hosts map[string]string
// inject is false when the query shape rules synthesis out. Suppression of an
// upstream fact of the same name does not depend on it.
inject bool
@@ -25,7 +27,7 @@ func (s *Server) newSourceInjector(query string, factEntity bool) *sourceInjecto
if !s.cfg.SourceFactEnabled || s.cfg.SourceFact == "" {
return nil
}
return &sourceInjector{name: s.cfg.SourceFact, inject: injectable(query, factEntity)}
return &sourceInjector{name: s.cfg.SourceFact, hosts: s.hosts, inject: injectable(query, factEntity)}
}
// claims reports whether an upstream record carries the name pdbmux owns, and is
@@ -70,7 +72,7 @@ func (si *sourceInjector) factRecord(certname, backend, environment string) json
Environment string `json:"environment"`
Name string `json:"name"`
Value string `json:"value"`
}{Certname: certname, Environment: environment, Name: si.name, Value: backend})
}{Certname: certname, Environment: environment, Name: si.name, Value: si.hosts[backend]})
if err != nil {
return nil
}
@@ -87,7 +89,7 @@ func (si *sourceInjector) stamp(raw json.RawMessage, backend string) json.RawMes
if json.Unmarshal(raw, &obj) != nil || obj == nil {
return raw
}
value, err := json.Marshal(backend)
value, err := json.Marshal(si.hosts[backend])
if err != nil {
return raw
}
+34 -14
View File
@@ -77,7 +77,7 @@ func TestHandler_FactsSourceFollowsMergeOwner(t *testing.T) {
b := newFakeBackend(t,
`[`+node("h1", "2026-07-01T00:00:00Z")+`,`+node("h2", "2026-07-20T00:00:00Z")+`]`,
`[`+factEnv("h1", "role", "web-b", "staging")+`,`+factEnv("h2", "role", "db-b", "staging")+`]`)
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeFreshness))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeFreshness))
rec := doGet(t, srv.Handler(), factsPath, "")
if rec.Code != http.StatusOK {
@@ -87,7 +87,7 @@ func TestHandler_FactsSourceFollowsMergeOwner(t *testing.T) {
if n != 2 {
t.Fatalf("expected one %s record per certname, got %d: %s", defaultSourceFact, n, rec.Body.String())
}
if got["h1"] != "a" || got["h2"] != "b" {
if got["h1"] != hostA || got["h2"] != hostB {
t.Errorf("provenance must name the backend that won the merge, got %v", got)
}
}
@@ -127,14 +127,14 @@ func TestHandler_NodesSourceStamped(t *testing.T) {
a := newFakeBackend(t,
`[`+node("h1", "2026-07-01T00:00:00Z")+`,`+node("h2", "2026-07-20T00:00:00Z")+`]`, `[]`)
b := newFakeBackend(t, `[`+node("h1", "2026-07-20T00:00:00Z")+`]`, `[]`)
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeStatic))
rec := doGet(t, srv.Handler(), nodesPath, "")
if rec.Code != http.StatusOK {
t.Fatalf("status %d: %s", rec.Code, rec.Body.String())
}
got := nodeSources(t, rec.Body.Bytes(), defaultSourceFact)
if got["h1"] != "b" || got["h2"] != "a" {
if got["h1"] != hostB || got["h2"] != hostA {
t.Errorf("node provenance = %v, want h1=b h2=a", got)
}
}
@@ -186,13 +186,13 @@ func TestHandler_SourceFactNameOverride(t *testing.T) {
a := newFakeBackend(t, `[`+node("h1", "2026-07-20T00:00:00Z")+`]`,
`[`+fact("h1", "role", "web", "")+`]`)
b := newFakeBackend(t, `[]`, `[]`)
cfg := testConfig(a.srv.URL, b.srv.URL, mergeStatic)
cfg := hostedConfig(a.srv.URL, b.srv.URL, mergeStatic)
cfg.SourceFact = "origin_pdb"
srv := newTestServer(cfg)
facts := doGet(t, srv.Handler(), factsPath, "")
got, n := sourceValues(t, facts.Body.Bytes(), "origin_pdb")
if n != 1 || got["h1"] != "a" {
if n != 1 || got["h1"] != hostA {
t.Errorf("override name not honoured: %s", facts.Body.String())
}
if _, n := sourceValues(t, facts.Body.Bytes(), defaultSourceFact); n != 0 {
@@ -200,7 +200,7 @@ func TestHandler_SourceFactNameOverride(t *testing.T) {
}
nodes := doGet(t, srv.Handler(), nodesPath, "")
if got := nodeSources(t, nodes.Body.Bytes(), "origin_pdb"); got["h1"] != "a" {
if got := nodeSources(t, nodes.Body.Bytes(), "origin_pdb"); got["h1"] != hostA {
t.Errorf("override name not honoured on /nodes: %v", got)
}
}
@@ -211,14 +211,14 @@ func TestHandler_UpstreamSourceFactOverridden(t *testing.T) {
a := newFakeBackend(t, `[`+node("h1", "2026-07-20T00:00:00Z")+`]`,
`[`+fact("h1", "role", "web", "")+`,`+fact("h1", defaultSourceFact, "stale-value", "")+`]`)
b := newFakeBackend(t, `[]`, `[]`)
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeStatic))
rec := doGet(t, srv.Handler(), factsPath, "")
got, n := sourceValues(t, rec.Body.Bytes(), defaultSourceFact)
if n != 1 {
t.Fatalf("expected exactly 1 %s record, got %d: %s", defaultSourceFact, n, rec.Body.String())
}
if got["h1"] != "a" {
if got["h1"] != hostA {
t.Errorf("upstream value survived: %v", got)
}
}
@@ -232,10 +232,10 @@ func TestHandler_UpstreamSourceFactSuppressedOnEveryGateState(t *testing.T) {
query string
want string // synthetic value, or "" when the gate blocks injection
}{
{"injection on", "", "a"},
{"injection on", "", hostA},
{"gated by extract", `["extract",["certname","name","value"],["=","certname","h1"]]`, ""},
{"gated by name filter", `["or",["=","name","role"],["=","name","os"]]`, ""},
{"selected by name filter", `["or",["=","name","role"],["=","name","` + defaultSourceFact + `"]]`, "a"},
{"selected by name filter", `["or",["=","name","role"],["=","name","` + defaultSourceFact + `"]]`, hostA},
{"gated by nested extract", `["and",["=","certname","h1"],["extract",["certname"]]]`, ""},
}
for _, tc := range tests {
@@ -243,7 +243,7 @@ func TestHandler_UpstreamSourceFactSuppressedOnEveryGateState(t *testing.T) {
a := newFakeBackend(t, `[`+node("h1", "2026-07-20T00:00:00Z")+`]`,
`[`+fact("h1", "role", "web", "")+`,`+fact("h1", defaultSourceFact, upstream, "")+`]`)
b := newFakeBackend(t, `[]`, `[]`)
srv := newTestServer(testConfig(a.srv.URL, b.srv.URL, mergeStatic))
srv := newTestServer(hostedConfig(a.srv.URL, b.srv.URL, mergeStatic))
rec := doGet(t, srv.Handler(), factsPath, tc.query)
if rec.Code != http.StatusOK {
@@ -560,10 +560,11 @@ func TestMergeFacts_GatedSuppressesUpstream(t *testing.T) {
func TestMergeFacts_SourceOrderedAfterOwnersFacts(t *testing.T) {
a := recs(t, "a", fact("h1", "role", "web-a", ""), fact("h1", "kernel", "Linux", ""))
b := recs(t, "b", fact("h1", "role", "web-b", ""))
merged := mergeFacts([]backendResult{a, b}, nil, &sourceInjector{name: defaultSourceFact, inject: true})
si := &sourceInjector{name: defaultSourceFact, hosts: map[string]string{"a": "a.example", "b": "b.example"}, inject: true}
merged := mergeFacts([]backendResult{a, b}, nil, si)
got := factValues(t, merged)
want := []string{"h1:role=web-a", "h1:kernel=Linux", "h1:" + defaultSourceFact + "=a"}
want := []string{"h1:role=web-a", "h1:kernel=Linux", "h1:" + defaultSourceFact + "=a.example"}
if !slices.Equal(got, want) {
t.Errorf("merged = %v, want %v", got, want)
}
@@ -638,6 +639,25 @@ func TestSourceInjector_SourceQuery(t *testing.T) {
}
}
// The value is the backend URL's hostname alone; the backend name is not used.
func TestSourceInjector_ValueIsBackendHostname(t *testing.T) {
cfg := testConfig("http://puppetdb.puppet.svc.cluster.local:8080", "https://puppetdbapi.service.consul/pdb", mergeStatic)
srv := newTestServer(cfg)
si := srv.newSourceInjector("", true)
for backend, want := range map[string]string{"a": "puppetdb.puppet.svc.cluster.local", "b": "puppetdbapi.service.consul"} {
var rec factFields
if err := json.Unmarshal(si.factRecord("h1", backend, ""), &rec); err != nil {
t.Fatal(err)
}
if rec.Value != want {
t.Errorf("fact value for %s = %q, want %q", backend, rec.Value, want)
}
if got := nodeSources(t, []byte(`[`+string(si.stamp([]byte(`{"certname":"h1"}`), backend))+`]`), defaultSourceFact); got["h1"] != want {
t.Errorf("node stamp for %s = %v, want %q", backend, got, want)
}
}
}
func TestCertnameScope(t *testing.T) {
tests := map[string]string{
`["=","name","x"]`: ``,