Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .changelog/23833.txt
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
```release-note:bug
catalog: Fix missing legacy scalar `Port` field in RPC and streaming subscription responses for multiport services. Older agents that cannot decode `NodeService.Ports` now correctly receive the default port value. The backfill logic is applied at RPC and subscription adapter boundaries via the shared multiport adapter so it is used consistently across `Catalog`, `Health`, `PreparedQuery`, and streaming subscription responses.
```
12 changes: 2 additions & 10 deletions agent/catalog_endpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -413,19 +413,11 @@ func (s *HTTPHandlers) catalogServiceNodes(resp http.ResponseWriter, req *http.R
out.ServiceNodes = make(structs.ServiceNodes, 0)
}
for i, s := range out.ServiceNodes {
var clone = *s

if s.ServiceTags == nil {
clone := *s
clone.ServiceTags = make([]string, 0)
out.ServiceNodes[i] = &clone
}

if clone.ServicePort == 0 && len(clone.ServicePorts) > 0 {
// Populate `port` with default port for backward compatibility
clone.ServicePort = clone.ToNodeService().DefaultPort()
}

out.ServiceNodes[i] = &clone

}
metrics.IncrCounterWithLabels([]string{"client", "api", "success", "catalog_service_nodes"}, 1,
s.nodeMetricsLabels())
Expand Down
32 changes: 32 additions & 0 deletions agent/consul/adapter/multiport_adapter.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
// Copyright IBM Corp. 2024, 2026
// SPDX-License-Identifier: BUSL-1.1

package adapter

import "github.com/hashicorp/consul/agent/structs"

// PopulateLegacyNodeServicePort sets the legacy scalar port from the default
// named port. Older agents do not decode NodeService.Ports from RPC responses,
// so the compatibility value must be populated before the response is sent.
func PopulateLegacyNodeServicePort(service *structs.NodeService) {
if service != nil && service.Port == 0 && len(service.Ports) > 0 {
service.Port = service.DefaultPort()
}
}

// PopulateLegacyServiceNodePorts sets the legacy scalar port on service nodes.
func PopulateLegacyServiceNodePorts(services structs.ServiceNodes) {
for _, service := range services {
if service != nil && service.ServicePort == 0 && len(service.ServicePorts) > 0 {
service.ServicePort = service.ToNodeService().DefaultPort()
}
}
}

// PopulateLegacyCheckServiceNodePorts sets the legacy scalar port on the
// services contained in check-service nodes.
func PopulateLegacyCheckServiceNodePorts(nodes structs.CheckServiceNodes) {
for i := range nodes {
PopulateLegacyNodeServicePort(nodes[i].Service)
}
}
24 changes: 24 additions & 0 deletions agent/consul/adapter/multiport_adapter_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
// Copyright IBM Corp. 2024, 2026
// SPDX-License-Identifier: BUSL-1.1

package adapter

import (
"testing"

"github.com/stretchr/testify/require"

"github.com/hashicorp/consul/agent/structs"
)

func TestPopulateLegacyNodeServicePort_PreservesExplicitPort(t *testing.T) {
service := &structs.NodeService{
Port: 7000,
Ports: structs.ServicePorts{
{Name: "http", Port: 8080, Default: true},
},
}

PopulateLegacyNodeServicePort(service)
require.Equal(t, 7000, service.Port)
}
12 changes: 10 additions & 2 deletions agent/consul/catalog_endpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,18 +10,19 @@ import (
"strings"
"time"

"github.com/hashicorp/go-metrics"
"github.com/hashicorp/go-metrics/prometheus"
hashstructure_v2 "github.com/mitchellh/hashstructure/v2"

"github.com/hashicorp/go-bexpr"
"github.com/hashicorp/go-hclog"
"github.com/hashicorp/go-memdb"
"github.com/hashicorp/go-metrics"
"github.com/hashicorp/go-metrics/prometheus"
"github.com/hashicorp/go-uuid"

"github.com/hashicorp/consul/acl"
"github.com/hashicorp/consul/acl/resolver"
"github.com/hashicorp/consul/agent/configentry"
"github.com/hashicorp/consul/agent/consul/adapter"
"github.com/hashicorp/consul/agent/consul/state"
"github.com/hashicorp/consul/agent/structs"
"github.com/hashicorp/consul/ipaddr"
Expand Down Expand Up @@ -876,6 +877,7 @@ func (c *Catalog) ServiceNodes(args *structs.ServiceSpecificRequest, reply *stru
return err
}
reply.ServiceNodes = raw.(structs.ServiceNodes)
adapter.PopulateLegacyServiceNodePorts(reply.ServiceNodes)

return c.srv.sortNodesByDistanceFrom(args.Source, reply.ServiceNodes)
})
Expand Down Expand Up @@ -972,6 +974,9 @@ func (c *Catalog) NodeServices(args *structs.NodeSpecificRequest, reply *structs
return err
}
reply.NodeServices.Services = raw.(map[string]*structs.NodeService)
for _, service := range reply.NodeServices.Services {
adapter.PopulateLegacyNodeServicePort(service)
}
}

return nil
Expand Down Expand Up @@ -1083,6 +1088,9 @@ func (c *Catalog) NodeServiceList(args *structs.NodeSpecificRequest, reply *stru
return err
}
reply.NodeServices.Services = raw.([]*structs.NodeService)
for _, service := range reply.NodeServices.Services {
adapter.PopulateLegacyNodeServicePort(service)
}

return nil
})
Expand Down
11 changes: 7 additions & 4 deletions agent/consul/health_endpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,18 @@ import (
"fmt"
"sort"

"github.com/hashicorp/go-metrics"
hashstructure_v2 "github.com/mitchellh/hashstructure/v2"

"github.com/hashicorp/go-bexpr"
"github.com/hashicorp/go-hclog"
"github.com/hashicorp/go-memdb"
"github.com/hashicorp/go-metrics"

"github.com/hashicorp/consul/acl"
"github.com/hashicorp/consul/agent/configentry"
"github.com/hashicorp/consul/agent/consul/adapter"
"github.com/hashicorp/consul/agent/consul/state"
"github.com/hashicorp/consul/agent/structs"
"github.com/hashicorp/go-bexpr"
"github.com/hashicorp/go-hclog"
"github.com/hashicorp/go-memdb"
)

// Health endpoint is used to query the health information
Expand Down Expand Up @@ -347,6 +349,7 @@ func (h *Health) ServiceNodes(args *structs.ServiceSpecificRequest, reply *struc
thisReply.Index = sgIdx
}

adapter.PopulateLegacyCheckServiceNodePorts(thisReply.Nodes)
*reply = thisReply
return nil
})
Expand Down
116 changes: 116 additions & 0 deletions agent/consul/multiport_adapter_rpc_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,116 @@
// Copyright (c) HashiCorp, Inc.
// SPDX-License-Identifier: BUSL-1.1

package consul

import (
"os"
"testing"

"github.com/stretchr/testify/require"

msgpackrpc "github.com/hashicorp/consul-net-rpc/net-rpc-msgpackrpc"

"github.com/hashicorp/consul/agent/structs"
"github.com/hashicorp/consul/testrpc"
)

func TestMultiportAdapter_RPCBackwardCompatibility(t *testing.T) {
if testing.Short() {
t.Skip("too slow for testing.Short")
}

t.Parallel()
dir, server := testServer(t)
defer os.RemoveAll(dir)
defer server.Shutdown()
codec := rpcClient(t, server)
defer codec.Close()

testrpc.WaitForLeader(t, server.RPC, "dc1")

ports := structs.ServicePorts{
{Name: "http", Port: 8080, Default: true},
{Name: "admin", Port: 9090},
}
registerReq := structs.RegisterRequest{
Datacenter: "dc1",
Node: "node-1",
Address: "127.0.0.1",
Service: &structs.NodeService{
ID: "web-1",
Service: "web",
Ports: ports,
},
}
var registerResp struct{}
require.NoError(t, msgpackrpc.CallWithCodec(codec, "Catalog.Register", &registerReq, &registerResp))

t.Run("Catalog.ServiceNodes", func(t *testing.T) {
req := structs.ServiceSpecificRequest{Datacenter: "dc1", ServiceName: "web"}
var resp structs.IndexedServiceNodes
require.NoError(t, msgpackrpc.CallWithCodec(codec, "Catalog.ServiceNodes", &req, &resp))
require.Len(t, resp.ServiceNodes, 1)
require.Equal(t, 8080, resp.ServiceNodes[0].ServicePort)
require.Equal(t, ports, resp.ServiceNodes[0].ServicePorts)
})

t.Run("Catalog.NodeServices", func(t *testing.T) {
req := structs.NodeSpecificRequest{Datacenter: "dc1", Node: "node-1"}
var resp structs.IndexedNodeServices
require.NoError(t, msgpackrpc.CallWithCodec(codec, "Catalog.NodeServices", &req, &resp))
require.NotNil(t, resp.NodeServices)
require.Equal(t, 8080, resp.NodeServices.Services["web-1"].Port)
require.Equal(t, ports, resp.NodeServices.Services["web-1"].Ports)
})

t.Run("Catalog.NodeServiceList", func(t *testing.T) {
req := structs.NodeSpecificRequest{Datacenter: "dc1", Node: "node-1"}
var resp structs.IndexedNodeServiceList
require.NoError(t, msgpackrpc.CallWithCodec(codec, "Catalog.NodeServiceList", &req, &resp))
require.Len(t, resp.NodeServices.Services, 1)
require.Equal(t, 8080, resp.NodeServices.Services[0].Port)
require.Equal(t, ports, resp.NodeServices.Services[0].Ports)
})

t.Run("Health.ServiceNodes", func(t *testing.T) {
req := structs.ServiceSpecificRequest{Datacenter: "dc1", ServiceName: "web"}
var resp structs.IndexedCheckServiceNodes
require.NoError(t, msgpackrpc.CallWithCodec(codec, "Health.ServiceNodes", &req, &resp))
require.Len(t, resp.Nodes, 1)
require.Equal(t, 8080, resp.Nodes[0].Service.Port)
require.Equal(t, ports, resp.Nodes[0].Service.Ports)
})

query := structs.PreparedQuery{
Name: "web-query",
Service: structs.ServiceQuery{
Service: "web",
},
}
applyReq := structs.PreparedQueryRequest{
Datacenter: "dc1",
Op: structs.PreparedQueryCreate,
Query: &query,
}
var queryID string
require.NoError(t, msgpackrpc.CallWithCodec(codec, "PreparedQuery.Apply", &applyReq, &queryID))

t.Run("PreparedQuery.Execute", func(t *testing.T) {
req := structs.PreparedQueryExecuteRequest{Datacenter: "dc1", QueryIDOrName: queryID}
var resp structs.PreparedQueryExecuteResponse
require.NoError(t, msgpackrpc.CallWithCodec(codec, "PreparedQuery.Execute", &req, &resp))
require.Len(t, resp.Nodes, 1)
require.Equal(t, 8080, resp.Nodes[0].Service.Port)
require.Equal(t, ports, resp.Nodes[0].Service.Ports)
})

t.Run("PreparedQuery.ExecuteRemote", func(t *testing.T) {
req := structs.PreparedQueryExecuteRemoteRequest{Datacenter: "dc1", Query: query}
var resp structs.PreparedQueryExecuteResponse
require.NoError(t, msgpackrpc.CallWithCodec(codec, "PreparedQuery.ExecuteRemote", &req, &resp))
require.Len(t, resp.Nodes, 1)
require.Equal(t, 8080, resp.Nodes[0].Service.Port)
require.Equal(t, ports, resp.Nodes[0].Service.Ports)
})
}
3 changes: 3 additions & 0 deletions agent/consul/prepared_query_endpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"github.com/hashicorp/go-uuid"

"github.com/hashicorp/consul/acl"
"github.com/hashicorp/consul/agent/consul/adapter"
"github.com/hashicorp/consul/agent/consul/state"
"github.com/hashicorp/consul/agent/structs"
"github.com/hashicorp/consul/agent/structs/aclfilter"
Expand Down Expand Up @@ -489,6 +490,7 @@ func (p *PreparedQuery) Execute(args *structs.PreparedQueryExecuteRequest,
}
}

adapter.PopulateLegacyCheckServiceNodePorts(reply.Nodes)
return nil
}

Expand Down Expand Up @@ -539,6 +541,7 @@ func (p *PreparedQuery) ExecuteRemote(args *structs.PreparedQueryExecuteRemoteRe
reply.Nodes = reply.Nodes[:args.Limit]
}

adapter.PopulateLegacyCheckServiceNodePorts(reply.Nodes)
return nil
}

Expand Down
14 changes: 13 additions & 1 deletion agent/consul/state/catalog_events.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
memdb "github.com/hashicorp/go-memdb"

"github.com/hashicorp/consul/acl"
"github.com/hashicorp/consul/agent/consul/adapter"
"github.com/hashicorp/consul/agent/consul/stream"
"github.com/hashicorp/consul/agent/structs"
"github.com/hashicorp/consul/proto/private/pbcommon"
Expand Down Expand Up @@ -63,12 +64,23 @@ func (e EventPayloadCheckServiceNode) Subject() stream.Subject {
}

func (e EventPayloadCheckServiceNode) ToSubscriptionEvent(idx uint64) *pbsubscribe.Event {
value := e.Value
if value != nil {
valueCopy := *value
if value.Service != nil {
serviceCopy := *value.Service
valueCopy.Service = &serviceCopy
}
adapter.PopulateLegacyCheckServiceNodePorts(structs.CheckServiceNodes{valueCopy})
value = &valueCopy
}

return &pbsubscribe.Event{
Index: idx,
Payload: &pbsubscribe.Event_ServiceHealth{
ServiceHealth: &pbsubscribe.ServiceHealthUpdate{
Op: e.Op,
CheckServiceNode: pbservice.NewCheckServiceNodeFromStructs(e.Value),
CheckServiceNode: pbservice.NewCheckServiceNodeFromStructs(value),
},
},
}
Expand Down
18 changes: 18 additions & 0 deletions agent/consul/state/catalog_events_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,24 @@ func TestServiceHealthSnapshot(t *testing.T) {
prototest.AssertDeepEqual(t, expected, buf.events, cmpEvents)
}

func TestEventPayloadCheckServiceNode_ToSubscriptionEvent_PopulatesLegacyPort(t *testing.T) {
service := &structs.NodeService{
Ports: structs.ServicePorts{
{Name: "http", Port: 8080, Default: true},
},
}
payload := EventPayloadCheckServiceNode{
Op: pbsubscribe.CatalogOp_Register,
Value: &structs.CheckServiceNode{
Service: service,
},
}

event := payload.ToSubscriptionEvent(1)
require.Equal(t, int32(8080), event.GetServiceHealth().GetCheckServiceNode().GetService().GetPort())
require.Zero(t, service.Port, "converting the event must not mutate the state-store value")
}

func TestServiceHealthSnapshot_ConnectTopic(t *testing.T) {
netutil.GetAgentBindAddrFunc = netutil.GetMockGetAgentBindAddrFunc("0.0.0.0")
store := NewStateStore(nil)
Expand Down
4 changes: 0 additions & 4 deletions agent/health_endpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -270,10 +270,6 @@ func (s *HTTPHandlers) healthServiceNodes(resp http.ResponseWriter, req *http.Re
clone.Tags = make([]string, 0)
out.Nodes[i].Service = &clone
}

if out.Nodes[i].Service != nil && out.Nodes[i].Service.Port == 0 && len(out.Nodes[i].Service.Ports) > 0 {
out.Nodes[i].Service.Port = out.Nodes[i].Service.DefaultPort()
}
}
return out.Nodes, nil
}
Expand Down
Loading