diff --git a/README.md b/README.md index f7e5553..d087ffa 100644 --- a/README.md +++ b/README.md @@ -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/[/]` | 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/` pins the backend name, so it answers with the +`/facts/pdbmux_source/` pins the backend host, so it answers with the nodes that backend owns. The `` 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. diff --git a/cache_test.go b/cache_test.go index cae13f1..80b6cd3 100644 --- a/cache_test.go +++ b/cache_test.go @@ -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) } } diff --git a/config.go b/config.go index f7508c2..ca2ef29 100644 --- a/config.go +++ b/config.go @@ -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: diff --git a/config_test.go b/config_test.go index 89c2d9d..aca37a4 100644 --- a/config_test.go +++ b/config_test.go @@ -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) + } + } +} diff --git a/e2e_backend_test.go b/e2e_backend_test.go index 72eefb3..41aff11 100644 --- a/e2e_backend_test.go +++ b/e2e_backend_test.go @@ -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 diff --git a/e2e_fixture_test.go b/e2e_fixture_test.go index 20e9fb3..1328152 100644 --- a/e2e_fixture_test.go +++ b/e2e_fixture_test.go @@ -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 { diff --git a/e2e_health_test.go b/e2e_health_test.go index be42a80..dc2db20 100644 --- a/e2e_health_test.go +++ b/e2e_health_test.go @@ -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) } }) diff --git a/e2e_nodelookup_test.go b/e2e_nodelookup_test.go index 7420557..07e53d5 100644 --- a/e2e_nodelookup_test.go +++ b/e2e_nodelookup_test.go @@ -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 } diff --git a/e2e_query_test.go b/e2e_query_test.go index 6778f08..4a5ffc7 100644 --- a/e2e_query_test.go +++ b/e2e_query_test.go @@ -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) diff --git a/factroutes_test.go b/factroutes_test.go index 67635b5..f7e9f75 100644 --- a/factroutes_test.go +++ b/factroutes_test.go @@ -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 == "" { diff --git a/server.go b/server.go index bc2f1b7..3cff7a0 100644 --- a/server.go +++ b/server.go @@ -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 } } diff --git a/server_test.go b/server_test.go index 83d7933..5b06e5c 100644 --- a/server_test.go +++ b/server_test.go @@ -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", diff --git a/source.go b/source.go index 9c83da2..d4f1d74 100644 --- a/source.go +++ b/source.go @@ -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 } diff --git a/source_test.go b/source_test.go index 70895f9..9d22614 100644 --- a/source_test.go +++ b/source_test.go @@ -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"]`: ``,