diff --git a/.github/k8s/sam-control-plane-template.yaml b/.github/k8s/sam-control-plane-template.yaml index 7793f819..3b50b44b 100644 --- a/.github/k8s/sam-control-plane-template.yaml +++ b/.github/k8s/sam-control-plane-template.yaml @@ -216,6 +216,9 @@ spec: - path: type: Exact value: /refresh + - path: + type: Exact + value: /nodes/catalog backendRefs: - name: sam-control-plane-${ENV_NAME} port: 8080 diff --git a/.github/workflows/scorecard.yml b/.github/workflows/scorecard.yml index 88fcb0ca..517ee6ac 100644 --- a/.github/workflows/scorecard.yml +++ b/.github/workflows/scorecard.yml @@ -73,6 +73,6 @@ jobs: # Upload the results to GitHub's code scanning dashboard (optional). # Commenting out will disable upload of results to your repo's Code Scanning dashboard - name: "Upload to code-scanning" - uses: github/codeql-action/upload-sarif@v3 + uses: github/codeql-action/upload-sarif@b96794f015dfd88f77b49b1c93e0fa7110f94c63 # v4.38.0 with: sarif_file: results.sarif diff --git a/api/sam.pb.go b/api/sam.pb.go index 65101f11..0845d1d7 100644 --- a/api/sam.pb.go +++ b/api/sam.pb.go @@ -1835,6 +1835,54 @@ func (x *TokenRefreshResponse) GetErrorMessage() string { return "" } +// NodeCatalogReport is the body of POST /nodes/catalog: a node's +// self-reported list of locally registered services. The reporting peer is +// taken from the presented biscuit, never from the body, so a node can only +// ever describe itself. Display-only; carries no authorization weight. +type NodeCatalogReport struct { + state protoimpl.MessageState `protogen:"open.v1"` + Services []*ServiceInfo `protobuf:"bytes,1,rep,name=services,proto3" json:"services,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *NodeCatalogReport) Reset() { + *x = NodeCatalogReport{} + mi := &file_api_sam_proto_msgTypes[24] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *NodeCatalogReport) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*NodeCatalogReport) ProtoMessage() {} + +func (x *NodeCatalogReport) ProtoReflect() protoreflect.Message { + mi := &file_api_sam_proto_msgTypes[24] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use NodeCatalogReport.ProtoReflect.Descriptor instead. +func (*NodeCatalogReport) Descriptor() ([]byte, []int) { + return file_api_sam_proto_rawDescGZIP(), []int{24} +} + +func (x *NodeCatalogReport) GetServices() []*ServiceInfo { + if x != nil { + return x.Services + } + return nil +} + type TokenRevokeRequest struct { state protoimpl.MessageState `protogen:"open.v1"` PeerId string `protobuf:"bytes,1,opt,name=peer_id,json=peerId,proto3" json:"peer_id,omitempty"` @@ -1844,7 +1892,7 @@ type TokenRevokeRequest struct { func (x *TokenRevokeRequest) Reset() { *x = TokenRevokeRequest{} - mi := &file_api_sam_proto_msgTypes[24] + mi := &file_api_sam_proto_msgTypes[25] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1856,7 +1904,7 @@ func (x *TokenRevokeRequest) String() string { func (*TokenRevokeRequest) ProtoMessage() {} func (x *TokenRevokeRequest) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[24] + mi := &file_api_sam_proto_msgTypes[25] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1869,7 +1917,7 @@ func (x *TokenRevokeRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use TokenRevokeRequest.ProtoReflect.Descriptor instead. func (*TokenRevokeRequest) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{24} + return file_api_sam_proto_rawDescGZIP(), []int{25} } func (x *TokenRevokeRequest) GetPeerId() string { @@ -1889,7 +1937,7 @@ type TokenRevokeResponse struct { func (x *TokenRevokeResponse) Reset() { *x = TokenRevokeResponse{} - mi := &file_api_sam_proto_msgTypes[25] + mi := &file_api_sam_proto_msgTypes[26] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1901,7 +1949,7 @@ func (x *TokenRevokeResponse) String() string { func (*TokenRevokeResponse) ProtoMessage() {} func (x *TokenRevokeResponse) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[25] + mi := &file_api_sam_proto_msgTypes[26] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1914,7 +1962,7 @@ func (x *TokenRevokeResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use TokenRevokeResponse.ProtoReflect.Descriptor instead. func (*TokenRevokeResponse) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{25} + return file_api_sam_proto_rawDescGZIP(), []int{26} } func (x *TokenRevokeResponse) GetSuccess() bool { @@ -1945,7 +1993,7 @@ type AgentSecret struct { func (x *AgentSecret) Reset() { *x = AgentSecret{} - mi := &file_api_sam_proto_msgTypes[26] + mi := &file_api_sam_proto_msgTypes[27] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -1957,7 +2005,7 @@ func (x *AgentSecret) String() string { func (*AgentSecret) ProtoMessage() {} func (x *AgentSecret) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[26] + mi := &file_api_sam_proto_msgTypes[27] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -1970,7 +2018,7 @@ func (x *AgentSecret) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentSecret.ProtoReflect.Descriptor instead. func (*AgentSecret) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{26} + return file_api_sam_proto_rawDescGZIP(), []int{27} } func (x *AgentSecret) GetHost() string { @@ -2013,7 +2061,7 @@ type AgentEgress struct { func (x *AgentEgress) Reset() { *x = AgentEgress{} - mi := &file_api_sam_proto_msgTypes[27] + mi := &file_api_sam_proto_msgTypes[28] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2025,7 +2073,7 @@ func (x *AgentEgress) String() string { func (*AgentEgress) ProtoMessage() {} func (x *AgentEgress) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[27] + mi := &file_api_sam_proto_msgTypes[28] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2038,7 +2086,7 @@ func (x *AgentEgress) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentEgress.ProtoReflect.Descriptor instead. func (*AgentEgress) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{27} + return file_api_sam_proto_rawDescGZIP(), []int{28} } func (x *AgentEgress) GetAllow() []string { @@ -2070,7 +2118,7 @@ type AgentIngress struct { func (x *AgentIngress) Reset() { *x = AgentIngress{} - mi := &file_api_sam_proto_msgTypes[28] + mi := &file_api_sam_proto_msgTypes[29] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2082,7 +2130,7 @@ func (x *AgentIngress) String() string { func (*AgentIngress) ProtoMessage() {} func (x *AgentIngress) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[28] + mi := &file_api_sam_proto_msgTypes[29] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2095,7 +2143,7 @@ func (x *AgentIngress) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentIngress.ProtoReflect.Descriptor instead. func (*AgentIngress) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{28} + return file_api_sam_proto_rawDescGZIP(), []int{29} } func (x *AgentIngress) GetType() ServiceType { @@ -2153,7 +2201,7 @@ type AgentBundle struct { func (x *AgentBundle) Reset() { *x = AgentBundle{} - mi := &file_api_sam_proto_msgTypes[29] + mi := &file_api_sam_proto_msgTypes[30] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2165,7 +2213,7 @@ func (x *AgentBundle) String() string { func (*AgentBundle) ProtoMessage() {} func (x *AgentBundle) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[29] + mi := &file_api_sam_proto_msgTypes[30] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2178,7 +2226,7 @@ func (x *AgentBundle) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentBundle.ProtoReflect.Descriptor instead. func (*AgentBundle) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{29} + return file_api_sam_proto_rawDescGZIP(), []int{30} } func (x *AgentBundle) GetVersion() string { @@ -2234,7 +2282,7 @@ type AgentAttachRequest struct { func (x *AgentAttachRequest) Reset() { *x = AgentAttachRequest{} - mi := &file_api_sam_proto_msgTypes[30] + mi := &file_api_sam_proto_msgTypes[31] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2246,7 +2294,7 @@ func (x *AgentAttachRequest) String() string { func (*AgentAttachRequest) ProtoMessage() {} func (x *AgentAttachRequest) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[30] + mi := &file_api_sam_proto_msgTypes[31] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2259,7 +2307,7 @@ func (x *AgentAttachRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentAttachRequest.ProtoReflect.Descriptor instead. func (*AgentAttachRequest) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{30} + return file_api_sam_proto_rawDescGZIP(), []int{31} } func (x *AgentAttachRequest) GetBundle() *AgentBundle { @@ -2283,7 +2331,7 @@ type AgentAttachResponse struct { func (x *AgentAttachResponse) Reset() { *x = AgentAttachResponse{} - mi := &file_api_sam_proto_msgTypes[31] + mi := &file_api_sam_proto_msgTypes[32] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2295,7 +2343,7 @@ func (x *AgentAttachResponse) String() string { func (*AgentAttachResponse) ProtoMessage() {} func (x *AgentAttachResponse) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[31] + mi := &file_api_sam_proto_msgTypes[32] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2308,7 +2356,7 @@ func (x *AgentAttachResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentAttachResponse.ProtoReflect.Descriptor instead. func (*AgentAttachResponse) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{31} + return file_api_sam_proto_rawDescGZIP(), []int{32} } func (x *AgentAttachResponse) GetEgressSocket() string { @@ -2343,7 +2391,7 @@ type AgentDetachRequest struct { func (x *AgentDetachRequest) Reset() { *x = AgentDetachRequest{} - mi := &file_api_sam_proto_msgTypes[32] + mi := &file_api_sam_proto_msgTypes[33] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2355,7 +2403,7 @@ func (x *AgentDetachRequest) String() string { func (*AgentDetachRequest) ProtoMessage() {} func (x *AgentDetachRequest) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[32] + mi := &file_api_sam_proto_msgTypes[33] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2368,7 +2416,7 @@ func (x *AgentDetachRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentDetachRequest.ProtoReflect.Descriptor instead. func (*AgentDetachRequest) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{32} + return file_api_sam_proto_rawDescGZIP(), []int{33} } func (x *AgentDetachRequest) GetAgentId() string { @@ -2388,7 +2436,7 @@ type AgentDetachResponse struct { func (x *AgentDetachResponse) Reset() { *x = AgentDetachResponse{} - mi := &file_api_sam_proto_msgTypes[33] + mi := &file_api_sam_proto_msgTypes[34] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2400,7 +2448,7 @@ func (x *AgentDetachResponse) String() string { func (*AgentDetachResponse) ProtoMessage() {} func (x *AgentDetachResponse) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[33] + mi := &file_api_sam_proto_msgTypes[34] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2413,7 +2461,7 @@ func (x *AgentDetachResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentDetachResponse.ProtoReflect.Descriptor instead. func (*AgentDetachResponse) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{33} + return file_api_sam_proto_rawDescGZIP(), []int{34} } func (x *AgentDetachResponse) GetSuccess() bool { @@ -2443,7 +2491,7 @@ type AgentRefreshRequest struct { func (x *AgentRefreshRequest) Reset() { *x = AgentRefreshRequest{} - mi := &file_api_sam_proto_msgTypes[34] + mi := &file_api_sam_proto_msgTypes[35] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2455,7 +2503,7 @@ func (x *AgentRefreshRequest) String() string { func (*AgentRefreshRequest) ProtoMessage() {} func (x *AgentRefreshRequest) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[34] + mi := &file_api_sam_proto_msgTypes[35] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2468,7 +2516,7 @@ func (x *AgentRefreshRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentRefreshRequest.ProtoReflect.Descriptor instead. func (*AgentRefreshRequest) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{34} + return file_api_sam_proto_rawDescGZIP(), []int{35} } func (x *AgentRefreshRequest) GetAgentId() string { @@ -2496,7 +2544,7 @@ type AgentRefreshResponse struct { func (x *AgentRefreshResponse) Reset() { *x = AgentRefreshResponse{} - mi := &file_api_sam_proto_msgTypes[35] + mi := &file_api_sam_proto_msgTypes[36] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2508,7 +2556,7 @@ func (x *AgentRefreshResponse) String() string { func (*AgentRefreshResponse) ProtoMessage() {} func (x *AgentRefreshResponse) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[35] + mi := &file_api_sam_proto_msgTypes[36] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2521,7 +2569,7 @@ func (x *AgentRefreshResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentRefreshResponse.ProtoReflect.Descriptor instead. func (*AgentRefreshResponse) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{35} + return file_api_sam_proto_rawDescGZIP(), []int{36} } func (x *AgentRefreshResponse) GetSuccess() bool { @@ -2556,7 +2604,7 @@ type AgentStatusRequest struct { func (x *AgentStatusRequest) Reset() { *x = AgentStatusRequest{} - mi := &file_api_sam_proto_msgTypes[36] + mi := &file_api_sam_proto_msgTypes[37] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2568,7 +2616,7 @@ func (x *AgentStatusRequest) String() string { func (*AgentStatusRequest) ProtoMessage() {} func (x *AgentStatusRequest) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[36] + mi := &file_api_sam_proto_msgTypes[37] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2581,7 +2629,7 @@ func (x *AgentStatusRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentStatusRequest.ProtoReflect.Descriptor instead. func (*AgentStatusRequest) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{36} + return file_api_sam_proto_rawDescGZIP(), []int{37} } func (x *AgentStatusRequest) GetAgentId() string { @@ -2603,7 +2651,7 @@ type AgentStatus struct { func (x *AgentStatus) Reset() { *x = AgentStatus{} - mi := &file_api_sam_proto_msgTypes[37] + mi := &file_api_sam_proto_msgTypes[38] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2615,7 +2663,7 @@ func (x *AgentStatus) String() string { func (*AgentStatus) ProtoMessage() {} func (x *AgentStatus) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[37] + mi := &file_api_sam_proto_msgTypes[38] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2628,7 +2676,7 @@ func (x *AgentStatus) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentStatus.ProtoReflect.Descriptor instead. func (*AgentStatus) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{37} + return file_api_sam_proto_rawDescGZIP(), []int{38} } func (x *AgentStatus) GetAgentId() string { @@ -2669,7 +2717,7 @@ type AgentStatusResponse struct { func (x *AgentStatusResponse) Reset() { *x = AgentStatusResponse{} - mi := &file_api_sam_proto_msgTypes[38] + mi := &file_api_sam_proto_msgTypes[39] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2681,7 +2729,7 @@ func (x *AgentStatusResponse) String() string { func (*AgentStatusResponse) ProtoMessage() {} func (x *AgentStatusResponse) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[38] + mi := &file_api_sam_proto_msgTypes[39] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2694,7 +2742,7 @@ func (x *AgentStatusResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use AgentStatusResponse.ProtoReflect.Descriptor instead. func (*AgentStatusResponse) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{38} + return file_api_sam_proto_rawDescGZIP(), []int{39} } func (x *AgentStatusResponse) GetAgents() []*AgentStatus { @@ -2725,7 +2773,7 @@ type IdentityEvidenceResponse struct { func (x *IdentityEvidenceResponse) Reset() { *x = IdentityEvidenceResponse{} - mi := &file_api_sam_proto_msgTypes[39] + mi := &file_api_sam_proto_msgTypes[40] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2737,7 +2785,7 @@ func (x *IdentityEvidenceResponse) String() string { func (*IdentityEvidenceResponse) ProtoMessage() {} func (x *IdentityEvidenceResponse) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[39] + mi := &file_api_sam_proto_msgTypes[40] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2750,7 +2798,7 @@ func (x *IdentityEvidenceResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use IdentityEvidenceResponse.ProtoReflect.Descriptor instead. func (*IdentityEvidenceResponse) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{39} + return file_api_sam_proto_rawDescGZIP(), []int{40} } func (x *IdentityEvidenceResponse) GetPeerId() string { @@ -2811,7 +2859,7 @@ type PeerEvidenceResponse struct { func (x *PeerEvidenceResponse) Reset() { *x = PeerEvidenceResponse{} - mi := &file_api_sam_proto_msgTypes[40] + mi := &file_api_sam_proto_msgTypes[41] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2823,7 +2871,7 @@ func (x *PeerEvidenceResponse) String() string { func (*PeerEvidenceResponse) ProtoMessage() {} func (x *PeerEvidenceResponse) ProtoReflect() protoreflect.Message { - mi := &file_api_sam_proto_msgTypes[40] + mi := &file_api_sam_proto_msgTypes[41] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2836,7 +2884,7 @@ func (x *PeerEvidenceResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use PeerEvidenceResponse.ProtoReflect.Descriptor instead. func (*PeerEvidenceResponse) Descriptor() ([]byte, []int) { - return file_api_sam_proto_rawDescGZIP(), []int{40} + return file_api_sam_proto_rawDescGZIP(), []int{41} } func (x *PeerEvidenceResponse) GetPeerId() string { @@ -3041,7 +3089,9 @@ const file_api_sam_proto_rawDesc = "" + "\rbiscuit_token\x18\x01 \x01(\fR\fbiscuitToken\x12\x1d\n" + "\n" + "expires_at\x18\x02 \x01(\x03R\texpiresAt\x12#\n" + - "\rerror_message\x18\x03 \x01(\tR\ferrorMessage\"-\n" + + "\rerror_message\x18\x03 \x01(\tR\ferrorMessage\"D\n" + + "\x11NodeCatalogReport\x12/\n" + + "\bservices\x18\x01 \x03(\v2\x13.sam.v1.ServiceInfoR\bservices\"-\n" + "\x12TokenRevokeRequest\x12\x17\n" + "\apeer_id\x18\x01 \x01(\tR\x06peerId\"E\n" + "\x13TokenRevokeResponse\x12\x18\n" + @@ -3146,7 +3196,7 @@ func file_api_sam_proto_rawDescGZIP() []byte { } var file_api_sam_proto_enumTypes = make([]protoimpl.EnumInfo, 3) -var file_api_sam_proto_msgTypes = make([]protoimpl.MessageInfo, 46) +var file_api_sam_proto_msgTypes = make([]protoimpl.MessageInfo, 47) var file_api_sam_proto_goTypes = []any{ (EnrollmentStatus)(0), // 0: sam.v1.EnrollmentStatus (ServiceType)(0), // 1: sam.v1.ServiceType @@ -3175,57 +3225,59 @@ var file_api_sam_proto_goTypes = []any{ (*KeysResponse)(nil), // 24: sam.v1.KeysResponse (*TokenRefreshRequest)(nil), // 25: sam.v1.TokenRefreshRequest (*TokenRefreshResponse)(nil), // 26: sam.v1.TokenRefreshResponse - (*TokenRevokeRequest)(nil), // 27: sam.v1.TokenRevokeRequest - (*TokenRevokeResponse)(nil), // 28: sam.v1.TokenRevokeResponse - (*AgentSecret)(nil), // 29: sam.v1.AgentSecret - (*AgentEgress)(nil), // 30: sam.v1.AgentEgress - (*AgentIngress)(nil), // 31: sam.v1.AgentIngress - (*AgentBundle)(nil), // 32: sam.v1.AgentBundle - (*AgentAttachRequest)(nil), // 33: sam.v1.AgentAttachRequest - (*AgentAttachResponse)(nil), // 34: sam.v1.AgentAttachResponse - (*AgentDetachRequest)(nil), // 35: sam.v1.AgentDetachRequest - (*AgentDetachResponse)(nil), // 36: sam.v1.AgentDetachResponse - (*AgentRefreshRequest)(nil), // 37: sam.v1.AgentRefreshRequest - (*AgentRefreshResponse)(nil), // 38: sam.v1.AgentRefreshResponse - (*AgentStatusRequest)(nil), // 39: sam.v1.AgentStatusRequest - (*AgentStatus)(nil), // 40: sam.v1.AgentStatus - (*AgentStatusResponse)(nil), // 41: sam.v1.AgentStatusResponse - (*IdentityEvidenceResponse)(nil), // 42: sam.v1.IdentityEvidenceResponse - (*PeerEvidenceResponse)(nil), // 43: sam.v1.PeerEvidenceResponse - nil, // 44: sam.v1.EnrollRequest.LabelsEntry - nil, // 45: sam.v1.BootstrapEnrollRequest.LabelsEntry - nil, // 46: sam.v1.CommandBackend.EnvEntry - nil, // 47: sam.v1.ServiceAnnounce.LabelsEntry - nil, // 48: sam.v1.PeerEvidenceResponse.LabelsEntry + (*NodeCatalogReport)(nil), // 27: sam.v1.NodeCatalogReport + (*TokenRevokeRequest)(nil), // 28: sam.v1.TokenRevokeRequest + (*TokenRevokeResponse)(nil), // 29: sam.v1.TokenRevokeResponse + (*AgentSecret)(nil), // 30: sam.v1.AgentSecret + (*AgentEgress)(nil), // 31: sam.v1.AgentEgress + (*AgentIngress)(nil), // 32: sam.v1.AgentIngress + (*AgentBundle)(nil), // 33: sam.v1.AgentBundle + (*AgentAttachRequest)(nil), // 34: sam.v1.AgentAttachRequest + (*AgentAttachResponse)(nil), // 35: sam.v1.AgentAttachResponse + (*AgentDetachRequest)(nil), // 36: sam.v1.AgentDetachRequest + (*AgentDetachResponse)(nil), // 37: sam.v1.AgentDetachResponse + (*AgentRefreshRequest)(nil), // 38: sam.v1.AgentRefreshRequest + (*AgentRefreshResponse)(nil), // 39: sam.v1.AgentRefreshResponse + (*AgentStatusRequest)(nil), // 40: sam.v1.AgentStatusRequest + (*AgentStatus)(nil), // 41: sam.v1.AgentStatus + (*AgentStatusResponse)(nil), // 42: sam.v1.AgentStatusResponse + (*IdentityEvidenceResponse)(nil), // 43: sam.v1.IdentityEvidenceResponse + (*PeerEvidenceResponse)(nil), // 44: sam.v1.PeerEvidenceResponse + nil, // 45: sam.v1.EnrollRequest.LabelsEntry + nil, // 46: sam.v1.BootstrapEnrollRequest.LabelsEntry + nil, // 47: sam.v1.CommandBackend.EnvEntry + nil, // 48: sam.v1.ServiceAnnounce.LabelsEntry + nil, // 49: sam.v1.PeerEvidenceResponse.LabelsEntry } var file_api_sam_proto_depIdxs = []int32{ 2, // 0: sam.v1.MeshEvent.type:type_name -> sam.v1.MeshEvent.Type - 44, // 1: sam.v1.EnrollRequest.labels:type_name -> sam.v1.EnrollRequest.LabelsEntry - 45, // 2: sam.v1.BootstrapEnrollRequest.labels:type_name -> sam.v1.BootstrapEnrollRequest.LabelsEntry + 45, // 1: sam.v1.EnrollRequest.labels:type_name -> sam.v1.EnrollRequest.LabelsEntry + 46, // 2: sam.v1.BootstrapEnrollRequest.labels:type_name -> sam.v1.BootstrapEnrollRequest.LabelsEntry 0, // 3: sam.v1.BootstrapEnrollResponse.status:type_name -> sam.v1.EnrollmentStatus 1, // 4: sam.v1.ServiceInfo.type:type_name -> sam.v1.ServiceType - 46, // 5: sam.v1.CommandBackend.env:type_name -> sam.v1.CommandBackend.EnvEntry + 47, // 5: sam.v1.CommandBackend.env:type_name -> sam.v1.CommandBackend.EnvEntry 10, // 6: sam.v1.RegisterServiceRequest.service:type_name -> sam.v1.ServiceInfo 11, // 7: sam.v1.RegisterServiceRequest.command:type_name -> sam.v1.CommandBackend 1, // 8: sam.v1.ServiceAnnounce.type:type_name -> sam.v1.ServiceType - 47, // 9: sam.v1.ServiceAnnounce.labels:type_name -> sam.v1.ServiceAnnounce.LabelsEntry + 48, // 9: sam.v1.ServiceAnnounce.labels:type_name -> sam.v1.ServiceAnnounce.LabelsEntry 18, // 10: sam.v1.PolicyConfigGetResponse.roles:type_name -> sam.v1.PolicyRole 19, // 11: sam.v1.PolicyConfigGetResponse.bindings:type_name -> sam.v1.PolicyBinding 18, // 12: sam.v1.PolicyConfigUpdateRequest.roles:type_name -> sam.v1.PolicyRole 19, // 13: sam.v1.PolicyConfigUpdateRequest.bindings:type_name -> sam.v1.PolicyBinding - 29, // 14: sam.v1.AgentEgress.secrets:type_name -> sam.v1.AgentSecret - 1, // 15: sam.v1.AgentIngress.type:type_name -> sam.v1.ServiceType - 30, // 16: sam.v1.AgentBundle.egress:type_name -> sam.v1.AgentEgress - 31, // 17: sam.v1.AgentBundle.ingress:type_name -> sam.v1.AgentIngress - 32, // 18: sam.v1.AgentAttachRequest.bundle:type_name -> sam.v1.AgentBundle - 31, // 19: sam.v1.AgentStatus.ingress:type_name -> sam.v1.AgentIngress - 40, // 20: sam.v1.AgentStatusResponse.agents:type_name -> sam.v1.AgentStatus - 48, // 21: sam.v1.PeerEvidenceResponse.labels:type_name -> sam.v1.PeerEvidenceResponse.LabelsEntry - 22, // [22:22] is the sub-list for method output_type - 22, // [22:22] is the sub-list for method input_type - 22, // [22:22] is the sub-list for extension type_name - 22, // [22:22] is the sub-list for extension extendee - 0, // [0:22] is the sub-list for field type_name + 10, // 14: sam.v1.NodeCatalogReport.services:type_name -> sam.v1.ServiceInfo + 30, // 15: sam.v1.AgentEgress.secrets:type_name -> sam.v1.AgentSecret + 1, // 16: sam.v1.AgentIngress.type:type_name -> sam.v1.ServiceType + 31, // 17: sam.v1.AgentBundle.egress:type_name -> sam.v1.AgentEgress + 32, // 18: sam.v1.AgentBundle.ingress:type_name -> sam.v1.AgentIngress + 33, // 19: sam.v1.AgentAttachRequest.bundle:type_name -> sam.v1.AgentBundle + 32, // 20: sam.v1.AgentStatus.ingress:type_name -> sam.v1.AgentIngress + 41, // 21: sam.v1.AgentStatusResponse.agents:type_name -> sam.v1.AgentStatus + 49, // 22: sam.v1.PeerEvidenceResponse.labels:type_name -> sam.v1.PeerEvidenceResponse.LabelsEntry + 23, // [23:23] is the sub-list for method output_type + 23, // [23:23] is the sub-list for method input_type + 23, // [23:23] is the sub-list for extension type_name + 23, // [23:23] is the sub-list for extension extendee + 0, // [0:23] is the sub-list for field type_name } func init() { file_api_sam_proto_init() } @@ -3243,7 +3295,7 @@ func file_api_sam_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_api_sam_proto_rawDesc), len(file_api_sam_proto_rawDesc)), NumEnums: 3, - NumMessages: 46, + NumMessages: 47, NumExtensions: 0, NumServices: 0, }, diff --git a/api/sam.proto b/api/sam.proto index b56826ff..13d761d0 100644 --- a/api/sam.proto +++ b/api/sam.proto @@ -252,6 +252,14 @@ message TokenRefreshResponse { string error_message = 3; } +// NodeCatalogReport is the body of POST /nodes/catalog: a node's +// self-reported list of locally registered services. The reporting peer is +// taken from the presented biscuit, never from the body, so a node can only +// ever describe itself. Display-only; carries no authorization weight. +message NodeCatalogReport { + repeated ServiceInfo services = 1; +} + message TokenRevokeRequest { string peer_id = 1; } diff --git a/charts/sam-mesh/README.md b/charts/sam-mesh/README.md index 600e7330..f8c67b84 100644 --- a/charts/sam-mesh/README.md +++ b/charts/sam-mesh/README.md @@ -77,7 +77,7 @@ with no default, because the right GatewayClass is provider-specific The route exposes only the control plane's enrollment surface (`/register`, `/info`, `/keys`, `/routers/lease`, `/policies`, `/enroll`, `/enroll/status`, -`/refresh`) and the console under `gateway.consolePath`; everything else, +`/refresh`, `/nodes/catalog`) and the console under `gateway.consolePath`; everything else, including `/admin` and `/user`, is unrouted. `gateway.adminRoute: true` additionally routes `/admin` — a dev convenience, leave it off in production. diff --git a/charts/sam-mesh/templates/gateway.yaml b/charts/sam-mesh/templates/gateway.yaml index 8d6105e3..83f4b2e0 100644 --- a/charts/sam-mesh/templates/gateway.yaml +++ b/charts/sam-mesh/templates/gateway.yaml @@ -59,6 +59,9 @@ spec: - path: type: Exact value: /refresh + - path: + type: Exact + value: /nodes/catalog backendRefs: - name: {{ $fullName }}-control-plane port: {{ .Values.controlPlane.service.port }} diff --git a/charts/sam-mesh/tests/gateway_test.yaml b/charts/sam-mesh/tests/gateway_test.yaml index 23b69ced..d494932e 100644 --- a/charts/sam-mesh/tests/gateway_test.yaml +++ b/charts/sam-mesh/tests/gateway_test.yaml @@ -28,6 +28,20 @@ tests: - lengthEqual: path: spec.rules count: 3 + # Every path a node calls on its own must be on the allow-list, or the + # node's periodic push silently 404s at the gateway. + - contains: + path: spec.rules[0].matches + content: + path: + type: Exact + value: /refresh + - contains: + path: spec.rules[0].matches + content: + path: + type: Exact + value: /nodes/catalog - it: empty consolePath leaves the console unrouted set: diff --git a/internal/console/public/app.js b/internal/console/public/app.js index 91864650..65295722 100644 --- a/internal/console/public/app.js +++ b/internal/console/public/app.js @@ -244,6 +244,7 @@ async function loadData() { setTableMessage('table-enrollments', 4, 'Restricted to administrators.'); } renderNodesTable(data.enrolled_nodes || []); + renderServicesTable(data.node_catalog || {}, buildLabelsByPeer(data.enrolled_nodes || [])); renderRoutersTable(data.active_routers || []); renderRouterTopography(data.active_routers || []); renderBootstrapTokensTable(data.bootstrap_tokens || []); @@ -332,16 +333,49 @@ function renderUsersTable(users) { `).join(''); } +// Renders an operator-declared labels map (e.g. {component: "stvv", role: +// "producer"}, from sam-node.yaml's labels: key) as a compact key=value +// list - the closest thing to a node mnemonic that exists today, since SAM +// has no dedicated name/alias field. Returns '' if there are none. +function formatLabels(labels) { + const entries = Object.entries(labels || {}); + if (entries.length === 0) { + return ''; + } + return entries.map(([k, v]) => `${k}=${v}`).join(', '); +} + +// A peer ID cell: the raw ID (still the real, authoritative identifier) +// with its operator-declared labels shown underneath when present. +function peerCell(peerID, labels) { + const labelText = formatLabels(labels); + const sub = labelText + ? `
${escapeHTML(labelText)}
` + : ''; + return `${escapeHTML(peerID)}${sub}`; +} + +// Builds a peer ID -> labels lookup from the enrolled_nodes list, so other +// tables (e.g. Services) can show the same labels next to a bare peer ID +// without a second fetch. +function buildLabelsByPeer(nodes) { + const byPeer = {}; + for (const node of nodes || []) { + byPeer[node.PeerID] = node.Labels || {}; + } + return byPeer; +} + function renderNodesTable(nodes) { const tbody = document.getElementById('table-nodes'); if (nodes.length === 0) { tbody.innerHTML = `No enrolled nodes found`; return; } - + tbody.innerHTML = nodes.map(node => ` - ${escapeHTML(node.PeerID)} + ${peerCell(node.PeerID, node.Labels)} ${escapeHTML(node.Role)} ${escapeHTML(node.OwnerID)} @@ -353,6 +387,39 @@ function renderNodesTable(nodes) { `).join(''); } +// node_catalog is {peerID: {services: [{name, type, description}], reported_at}}, +// already restricted server-side to nodes that are still admitted; type is +// the short name ("mcp", "inference", "a2a") rendered by the control plane. +function renderServicesTable(nodeCatalog, labelsByPeer) { + const tbody = document.getElementById('table-services'); + const peerIDs = Object.keys(nodeCatalog || {}); + const rows = []; + for (const peerID of peerIDs) { + const entry = nodeCatalog[peerID] || {}; + const services = entry.services || []; + for (const svc of services) { + if (svc) { + rows.push({ peerID, reportedAt: entry.reported_at, svc }); + } + } + } + + if (rows.length === 0) { + tbody.innerHTML = `No nodes have reported any services yet`; + return; + } + + tbody.innerHTML = rows.map(({ peerID, reportedAt, svc }) => ` + + ${escapeHTML(svc.name || '')} + ${escapeHTML(svc.type || 'unknown')} + ${escapeHTML(svc.description || '')} + ${peerCell(peerID, (labelsByPeer || {})[peerID])} + ${reportedAt ? escapeHTML(new Date(reportedAt).toLocaleString()) : '-'} + + `).join(''); +} + function getStatusBadge(status) { if (status === 0 || status === 'ENROLLMENT_STATUS_PENDING') { return `Pending`; diff --git a/internal/console/public/index.html b/internal/console/public/index.html index 5a6f6bfd..ec55b739 100644 --- a/internal/console/public/index.html +++ b/internal/console/public/index.html @@ -71,6 +71,10 @@

SAM Console

Nodes + + + Services + Enrollments @@ -199,6 +203,32 @@

Enrolled Nodes

+ +
+
+

Mesh Services

+
+
+
+ + + + + + + + + + + + + +
ServiceTypeDescriptionNode IDReported At
Loading...
+
+
+
+
diff --git a/internal/console/public/style.css b/internal/console/public/style.css index 782e4264..e3b4105d 100644 --- a/internal/console/public/style.css +++ b/internal/console/public/style.css @@ -646,6 +646,13 @@ body { border-bottom: none; } +/* An operator-declared label line under a bare peer ID (see peerCell in + app.js) - a class, not an inline style, so CSP style-src can stay strict. */ +.cell-subtext { + color: var(--text-secondary); + font-size: 0.85em; +} + .text-center { text-align: center; } diff --git a/internal/controlplane/catalog.go b/internal/controlplane/catalog.go new file mode 100644 index 00000000..7ff798cd --- /dev/null +++ b/internal/controlplane/catalog.go @@ -0,0 +1,215 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package controlplane + +import ( + "crypto/ed25519" + "encoding/base64" + "fmt" + "io" + "net/http" + "strings" + "time" + + "github.com/google/sam/api" + "github.com/google/sam/internal/identity" + "github.com/google/sam/internal/storage" + "github.com/libp2p/go-libp2p/core/peer" + "google.golang.org/protobuf/proto" +) + +// maxCatalogServices bounds one report so a single admitted node cannot grow +// the in-memory cache without limit. +const maxCatalogServices = 512 + +// nodeCatalogEntry is what HandleNodeCatalog caches per reporting peer. +type nodeCatalogEntry struct { + Services []*api.ServiceInfo + ReportedAt time.Time +} + +// catalogService is the console-facing shape of one reported service: a plain +// struct so the JSON the console reads does not depend on protoc-gen-go's +// struct layout or tags. +type catalogService struct { + Name string `json:"name"` + Type string `json:"type"` + Description string `json:"description"` +} + +// catalogView is HandleAdminStatus's node_catalog value for one peer. +type catalogView struct { + Services []catalogService `json:"services"` + ReportedAt time.Time `json:"reported_at"` +} + +// catalogSnapshot returns a stable copy of the current node service catalog +// cache, safe to range over without holding catalogMu. +func (s *Server) catalogSnapshot() map[string]nodeCatalogEntry { + s.catalogMu.RLock() + defer s.catalogMu.RUnlock() + snap := make(map[string]nodeCatalogEntry, len(s.catalog)) + for k, v := range s.catalog { + snap[k] = v + } + return snap +} + +// catalogViewFor renders the cache for the console, restricted to nodes that +// are still admitted so a banned or expired node's last report disappears +// with its enrollment instead of lingering until the next restart. +func (s *Server) catalogViewFor(nodes []storage.EnrolledNode, now time.Time) map[string]catalogView { + snap := s.catalogSnapshot() + view := make(map[string]catalogView, len(snap)) + for i := range nodes { + node := &nodes[i] + // The cache is keyed by the canonical base58 form from the verified + // biscuit; the stored record may carry another valid encoding of the + // same peer (e.g. CIDv1), so decode before looking up. + pID, err := peer.Decode(node.PeerID) + if err != nil { + logger.Warnw("Skipping enrolled node with undecodable peer ID in catalog view", "peer_id", node.PeerID, "error", err) + continue + } + entry, ok := snap[pID.String()] + if !ok || node.CheckAdmission(now) != nil { + continue + } + services := make([]catalogService, 0, len(entry.Services)) + for _, svc := range entry.Services { + if svc == nil { + continue + } + typeName, err := api.ServiceTypeToString(svc.GetType()) + if err != nil { + typeName = "unknown" + } + services = append(services, catalogService{ + Name: svc.GetName(), + Type: typeName, + Description: svc.GetDescription(), + }) + } + view[node.PeerID] = catalogView{Services: services, ReportedAt: entry.ReportedAt} + } + return view +} + +// dropCatalogEntry forgets a peer's report; called when its enrollment ends. +// Accepts any valid encoding of the peer ID. +func (s *Server) dropCatalogEntry(peerID string) { + pID, err := peer.Decode(peerID) + if err != nil { + // Cache keys always come from a verified biscuit, so an undecodable + // ID cannot have an entry - nothing to evict, but say so. + logger.Warnw("Not evicting catalog entry for undecodable peer ID", "peer_id", peerID, "error", err) + return + } + s.catalogMu.Lock() + delete(s.catalog, pID.String()) + s.catalogMu.Unlock() +} + +// HandleNodeCatalog HTTP POST /nodes/catalog - a node self-reports the +// services it currently has registered locally (the same data +// list_local_services already answers on the node itself), so the control +// plane can show mesh-wide service topology without needing to be a DHT +// participant or open a P2P connection to every enrolled node itself. +// +// The body is an api.NodeCatalogReport. The reporting peer is the one bound +// in the presented Biscuit, so a node can only ever describe itself. +// +// This is a live-status cache, not authoritative state: a node that goes +// offline without ever reporting an empty catalog just leaves its last +// report in place until ReportedAt visibly goes stale or its enrollment +// ends. It is admin-facing display data only and never feeds authorization, +// which is also why a bare bearer Biscuit (no signed challenge, unlike +// /refresh) is accepted here: a replayed token can only repaint a table. +func (s *Server) HandleNodeCatalog(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + http.Error(w, "Method not allowed", http.StatusMethodNotAllowed) + return + } + + authHeader := r.Header.Get("Authorization") + if !strings.HasPrefix(authHeader, "Bearer ") { + http.Error(w, "Missing node Biscuit token in Authorization header", http.StatusUnauthorized) + return + } + biscuitBytes, err := base64.StdEncoding.DecodeString(strings.TrimPrefix(authHeader, "Bearer ")) + if err != nil { + http.Error(w, "Malformed base64 token", http.StatusBadRequest) + return + } + + ctx := r.Context() + validKeys, err := s.store.GetAllValidKeys(ctx) + if err != nil { + logger.Errorf("Failed to retrieve valid signing keys: %v", err) + http.Error(w, "Internal server error", http.StatusInternalServerError) + return + } + var trustedKeys []ed25519.PublicKey + for _, k := range validKeys { + trustedKeys = append(trustedKeys, k.Public) + } + + peerID, err := identity.VerifyAndExtractPeerID(trustedKeys, biscuitBytes, s.config.BiscuitTimeout) + if err != nil { + logger.Warnw("Invalid biscuit presented to /nodes/catalog", "error", err) + http.Error(w, "Invalid biscuit: "+err.Error(), http.StatusUnauthorized) + return + } + + nodeRecord, err := s.store.GetNode(ctx, peerID.String()) + if err == storage.ErrNotFound || (err == nil && nodeRecord == nil) { + http.Error(w, "Node not enrolled", http.StatusUnauthorized) + return + } else if err != nil { + logger.Errorf("Failed to retrieve node record: %v", err) + http.Error(w, "Internal server error", http.StatusInternalServerError) + return + } + if err := nodeRecord.CheckAdmission(time.Now()); err != nil { + http.Error(w, "Node not admitted: "+err.Error(), http.StatusUnauthorized) + return + } + + r.Body = http.MaxBytesReader(w, r.Body, maxRequestBodyBytes) + body, err := io.ReadAll(r.Body) + if err != nil { + http.Error(w, "Failed to read body", http.StatusBadRequest) + return + } + + var req api.NodeCatalogReport + if err := proto.Unmarshal(body, &req); err != nil { + http.Error(w, "Invalid request format", http.StatusBadRequest) + return + } + if len(req.Services) > maxCatalogServices { + http.Error(w, fmt.Sprintf("Too many services in report (max %d)", maxCatalogServices), http.StatusBadRequest) + return + } + + s.catalogMu.Lock() + s.catalog[peerID.String()] = nodeCatalogEntry{ + Services: req.Services, + ReportedAt: time.Now(), + } + s.catalogMu.Unlock() + + w.WriteHeader(http.StatusNoContent) +} diff --git a/internal/controlplane/catalog_test.go b/internal/controlplane/catalog_test.go new file mode 100644 index 00000000..aa28b252 --- /dev/null +++ b/internal/controlplane/catalog_test.go @@ -0,0 +1,354 @@ +// Copyright 2026 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package controlplane + +import ( + "bytes" + "context" + "crypto/ed25519" + "encoding/base64" + "encoding/json" + "net/http" + "testing" + "time" + + "github.com/google/sam/api" + "github.com/google/sam/internal/identity" + "github.com/libp2p/go-libp2p/core/crypto" + "github.com/libp2p/go-libp2p/core/peer" + "google.golang.org/protobuf/proto" +) + +// postCatalog POSTs a raw body to /nodes/catalog under the given +// Authorization header value and returns the response status. +func postCatalog(t *testing.T, cpURL, authHeader string, body []byte) int { + t.Helper() + + req, err := http.NewRequest(http.MethodPost, cpURL+"/nodes/catalog", bytes.NewReader(body)) + if err != nil { + t.Fatalf("NewRequest: %v", err) + } + if authHeader != "" { + req.Header.Set("Authorization", authHeader) + } + req.Header.Set("Content-Type", "application/x-protobuf") + + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatalf("Do: %v", err) + } + defer func() { _ = resp.Body.Close() }() + return resp.StatusCode +} + +func bearer(biscuit []byte) string { + return "Bearer " + base64.StdEncoding.EncodeToString(biscuit) +} + +func catalogBody(t *testing.T, services ...*api.ServiceInfo) []byte { + t.Helper() + body, err := proto.Marshal(&api.NodeCatalogReport{Services: services}) + if err != nil { + t.Fatalf("Marshal: %v", err) + } + return body +} + +// adminNodeCatalog fetches /admin/status and returns its node_catalog value. +func adminNodeCatalog(t *testing.T, cpURL, adminToken string) map[string]catalogView { + t.Helper() + + req, err := http.NewRequest(http.MethodGet, cpURL+"/admin/status", nil) + if err != nil { + t.Fatalf("NewRequest: %v", err) + } + req.Header.Set("Authorization", "Bearer "+adminToken) + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatalf("GET /admin/status: %v", err) + } + defer func() { _ = resp.Body.Close() }() + if resp.StatusCode != http.StatusOK { + t.Fatalf("GET /admin/status: status %d", resp.StatusCode) + } + var status struct { + NodeCatalog map[string]catalogView `json:"node_catalog"` + } + if err := json.NewDecoder(resp.Body).Decode(&status); err != nil { + t.Fatalf("decode /admin/status: %v", err) + } + return status.NodeCatalog +} + +func TestHandleNodeCatalog(t *testing.T) { + t.Parallel() + + srv, store, cpURL := setupTestServer(t, "") + defer func() { + _ = srv.Close() + _ = store.Close() + }() + srv.config.AdminToken = "super-secret-admin-token" + + ctx := context.Background() + priv, biscuitBytes := enrollRefreshTestNode(t, ctx, store) + nodePeer, err := peer.IDFromPrivateKey(priv) + if err != nil { + t.Fatalf("IDFromPrivateKey: %v", err) + } + + body := catalogBody(t, + &api.ServiceInfo{Type: api.ServiceType_SERVICE_TYPE_MCP, Name: "stvv-compliance-docs", Description: "doc lookup"}, + &api.ServiceInfo{Type: api.ServiceType_SERVICE_TYPE_INFERENCE, Name: "llama", Description: "local model"}, + ) + if got := postCatalog(t, cpURL, bearer(biscuitBytes), body); got != http.StatusNoContent { + t.Fatalf("HandleNodeCatalog: got status %d, want %d", got, http.StatusNoContent) + } + + snap := srv.catalogSnapshot() + entry, ok := snap[nodePeer.String()] + if !ok || len(snap) != 1 { + t.Fatalf("expected exactly the reporting peer in the catalog, got %v", snap) + } + if len(entry.Services) != 2 || entry.ReportedAt.IsZero() { + t.Fatalf("unexpected cached entry: %+v", entry) + } + + // A second report replaces the first rather than accumulating. + body = catalogBody(t, &api.ServiceInfo{Type: api.ServiceType_SERVICE_TYPE_A2A, Name: "planner"}) + if got := postCatalog(t, cpURL, bearer(biscuitBytes), body); got != http.StatusNoContent { + t.Fatalf("second report: got status %d, want %d", got, http.StatusNoContent) + } + entry = srv.catalogSnapshot()[nodePeer.String()] + if len(entry.Services) != 1 || entry.Services[0].Name != "planner" { + t.Fatalf("second report must replace the first, got %+v", entry.Services) + } + + // The console sees the plain view with the type rendered as a name. + view := adminNodeCatalog(t, cpURL, srv.config.AdminToken) + got, ok := view[nodePeer.String()] + if !ok || len(view) != 1 { + t.Fatalf("expected the reporting peer in node_catalog, got %v", view) + } + want := []catalogService{{Name: "planner", Type: "a2a", Description: ""}} + if len(got.Services) != 1 || got.Services[0] != want[0] || got.ReportedAt.IsZero() { + t.Fatalf("node_catalog view = %+v, want services %+v", got, want) + } + + // An empty report is valid and clears the node's services. + if got := postCatalog(t, cpURL, bearer(biscuitBytes), catalogBody(t)); got != http.StatusNoContent { + t.Fatalf("empty report: got status %d, want %d", got, http.StatusNoContent) + } + if view := adminNodeCatalog(t, cpURL, srv.config.AdminToken); len(view[nodePeer.String()].Services) != 0 { + t.Fatalf("empty report must clear services, got %+v", view) + } +} + +func TestHandleNodeCatalog_Rejections(t *testing.T) { + t.Parallel() + + srv, store, cpURL := setupTestServer(t, "") + defer func() { + _ = srv.Close() + _ = store.Close() + }() + + ctx := context.Background() + _, biscuitBytes := enrollRefreshTestNode(t, ctx, store) + cpPriv, _, err := store.GetCurrentKey(ctx) + if err != nil { + t.Fatalf("GetCurrentKey: %v", err) + } + + // A biscuit minted by a key the control plane never trusted. + _, rogueKey, err := ed25519.GenerateKey(nil) + if err != nil { + t.Fatalf("GenerateKey: %v", err) + } + strangerPriv, _, err := crypto.GenerateKeyPair(crypto.Ed25519, -1) + if err != nil { + t.Fatalf("GenerateKeyPair: %v", err) + } + strangerPeer, err := peer.IDFromPrivateKey(strangerPriv) + if err != nil { + t.Fatalf("IDFromPrivateKey: %v", err) + } + forged, err := identity.MintBootstrapBiscuitToken(rogueKey, strangerPeer, api.RoleNode, time.Now().Add(api.BiscuitTokenTTL), nil, nil) + if err != nil { + t.Fatalf("MintBootstrapBiscuitToken(rogue): %v", err) + } + // A valid biscuit for a peer that was never enrolled. + unenrolled, err := identity.MintBootstrapBiscuitToken(cpPriv, strangerPeer, api.RoleNode, time.Now().Add(api.BiscuitTokenTTL), nil, nil) + if err != nil { + t.Fatalf("MintBootstrapBiscuitToken(unenrolled): %v", err) + } + + tooMany := make([]*api.ServiceInfo, maxCatalogServices+1) + for i := range tooMany { + tooMany[i] = &api.ServiceInfo{Type: api.ServiceType_SERVICE_TYPE_MCP, Name: "svc"} + } + + ok := catalogBody(t) + tests := []struct { + name string + auth string + body []byte + want int + }{ + {name: "missing authorization", auth: "", body: ok, want: http.StatusUnauthorized}, + {name: "not a bearer token", auth: "Basic abc", body: ok, want: http.StatusUnauthorized}, + {name: "malformed base64", auth: "Bearer %%%not-base64", body: ok, want: http.StatusBadRequest}, + {name: "not a biscuit", auth: bearer([]byte("garbage")), body: ok, want: http.StatusUnauthorized}, + {name: "forged signature", auth: bearer(forged), body: ok, want: http.StatusUnauthorized}, + {name: "unenrolled peer", auth: bearer(unenrolled), body: ok, want: http.StatusUnauthorized}, + {name: "invalid body", auth: bearer(biscuitBytes), body: []byte(`{"services":[]}`), want: http.StatusBadRequest}, + {name: "too many services", auth: bearer(biscuitBytes), body: catalogBody(t, tooMany...), want: http.StatusBadRequest}, + } + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + if got := postCatalog(t, cpURL, tc.auth, tc.body); got != tc.want { + t.Fatalf("got status %d, want %d", got, tc.want) + } + }) + } + if len(srv.catalogSnapshot()) != 0 { + t.Fatalf("rejected reports must not be cached, got %v", srv.catalogSnapshot()) + } + + resp, err := http.Get(cpURL + "/nodes/catalog") + if err != nil { + t.Fatalf("Get: %v", err) + } + _ = resp.Body.Close() + if resp.StatusCode != http.StatusMethodNotAllowed { + t.Fatalf("GET /nodes/catalog: got status %d, want %d", resp.StatusCode, http.StatusMethodNotAllowed) + } +} + +// A node's catalog is only ever as trustworthy as its enrollment: once the +// record is banned or its session lapses, new reports are refused and the +// cached one drops out of the console view. +func TestHandleNodeCatalog_Admission(t *testing.T) { + t.Parallel() + + srv, store, cpURL := setupTestServer(t, "") + defer func() { + _ = srv.Close() + _ = store.Close() + }() + srv.config.AdminToken = "super-secret-admin-token" + + ctx := context.Background() + priv, biscuitBytes := enrollRefreshTestNode(t, ctx, store) + nodePeer, err := peer.IDFromPrivateKey(priv) + if err != nil { + t.Fatalf("IDFromPrivateKey: %v", err) + } + body := catalogBody(t, &api.ServiceInfo{Type: api.ServiceType_SERVICE_TYPE_MCP, Name: "calc"}) + if got := postCatalog(t, cpURL, bearer(biscuitBytes), body); got != http.StatusNoContent { + t.Fatalf("admitted node: got status %d, want %d", got, http.StatusNoContent) + } + if _, ok := adminNodeCatalog(t, cpURL, srv.config.AdminToken)[nodePeer.String()]; !ok { + t.Fatal("admitted node's report missing from node_catalog") + } + + // Session lapsed: the cache still holds the report but the view hides it. + record, err := store.GetNode(ctx, nodePeer.String()) + if err != nil { + t.Fatalf("GetNode: %v", err) + } + record.ExpiresAt = time.Now().Add(-time.Hour) + if err := store.EnrollNode(ctx, record); err != nil { + t.Fatalf("EnrollNode: %v", err) + } + if got := postCatalog(t, cpURL, bearer(biscuitBytes), body); got != http.StatusUnauthorized { + t.Fatalf("expired session: got status %d, want %d", got, http.StatusUnauthorized) + } + if view := adminNodeCatalog(t, cpURL, srv.config.AdminToken); len(view) != 0 { + t.Fatalf("expired node must not appear in node_catalog, got %v", view) + } + + // Banned through the server path: the entry is evicted outright. + record.ExpiresAt = time.Now().Add(time.Hour) + if err := store.EnrollNode(ctx, record); err != nil { + t.Fatalf("EnrollNode: %v", err) + } + if err := srv.banNode(ctx, record); err != nil { + t.Fatalf("banNode: %v", err) + } + if got := postCatalog(t, cpURL, bearer(biscuitBytes), body); got != http.StatusUnauthorized { + t.Fatalf("banned node: got status %d, want %d", got, http.StatusUnauthorized) + } + if snap := srv.catalogSnapshot(); len(snap) != 0 { + t.Fatalf("banning must evict the cached report, got %v", snap) + } +} + +// The cache is keyed by the canonical base58 form from the verified biscuit, +// but the enrollment record's PeerID comes off the wire and may be any valid +// encoding of the same peer (e.g. CIDv1 base32). The view lookup and the +// ban eviction must still hit the entry. +func TestCatalogPeerIDCanonicalization(t *testing.T) { + t.Parallel() + + srv, store, cpURL := setupTestServer(t, "") + defer func() { + _ = srv.Close() + _ = store.Close() + }() + srv.config.AdminToken = "super-secret-admin-token" + + ctx := context.Background() + priv, biscuitBytes := enrollRefreshTestNode(t, ctx, store) + nodePeer, err := peer.IDFromPrivateKey(priv) + if err != nil { + t.Fatalf("IDFromPrivateKey: %v", err) + } + + // Re-enroll the same peer under its CIDv1 base32 encoding, as a raw wire + // string would arrive before the in-flight canonicalization PR lands. + record, err := store.GetNode(ctx, nodePeer.String()) + if err != nil { + t.Fatalf("GetNode: %v", err) + } + cidForm := peer.ToCid(nodePeer).String() + if cidForm == nodePeer.String() { + t.Fatal("test needs a non-canonical encoding, got the canonical one") + } + record.PeerID = cidForm + if err := store.EnrollNode(ctx, record); err != nil { + t.Fatalf("EnrollNode: %v", err) + } + + body := catalogBody(t, &api.ServiceInfo{Type: api.ServiceType_SERVICE_TYPE_MCP, Name: "calc"}) + if got := postCatalog(t, cpURL, bearer(biscuitBytes), body); got != http.StatusNoContent { + t.Fatalf("report: got status %d, want %d", got, http.StatusNoContent) + } + + // The view must join the CIDv1 record with the canonically-keyed entry, + // displayed under the record's own spelling. + view := adminNodeCatalog(t, cpURL, srv.config.AdminToken) + if _, ok := view[cidForm]; !ok || len(view[cidForm].Services) != 1 { + t.Fatalf("CIDv1-enrolled node's report missing from node_catalog, got %v", view) + } + + // Ban via the stored record: eviction must hit the canonical cache key. + if err := srv.banNode(ctx, record); err != nil { + t.Fatalf("banNode: %v", err) + } + if snap := srv.catalogSnapshot(); len(snap) != 0 { + t.Fatalf("banning a CIDv1-enrolled node must evict its cached report, got %v", snap) + } +} diff --git a/internal/controlplane/server.go b/internal/controlplane/server.go index 5f398ba9..565a1c47 100644 --- a/internal/controlplane/server.go +++ b/internal/controlplane/server.go @@ -77,6 +77,13 @@ type Server struct { providersMu sync.RWMutex providers map[string]*oidc.Provider + // catalogMu/catalog cache each node's self-reported local service list + // (see HandleNodeCatalog), keyed by peer ID. In-memory only: this is a + // live-status view, not authoritative state, so it's fine to lose on + // restart - every node re-reports on its own next periodic push. + catalogMu sync.RWMutex + catalog map[string]nodeCatalogEntry + ctx context.Context cancel context.CancelFunc wg sync.WaitGroup @@ -98,6 +105,7 @@ func NewServer(config Options, store storage.Store) (*Server, error) { mesh: NewNopMeshAdapter(), limiter: rate.NewLimiter(rate.Limit(EnrollRateLimit), EnrollBurst), providers: make(map[string]*oidc.Provider), + catalog: make(map[string]nodeCatalogEntry), ctx: ctx, cancel: cancel, }, nil @@ -204,6 +212,7 @@ func (s *Server) RegisterRoutes(mux *http.ServeMux) { mux.HandleFunc("/enroll", s.HandleEnroll) mux.HandleFunc("/enroll/status", s.HandleEnrollStatus) mux.HandleFunc("/refresh", s.HandleRefresh) + mux.HandleFunc("/nodes/catalog", s.HandleNodeCatalog) mux.HandleFunc("/admin/bootstrap-tokens", s.HandleAdminBootstrapTokens) mux.HandleFunc("/admin/enrollments", s.HandleAdminEnrollments) mux.HandleFunc("/admin/enrollments/", s.HandleAdminEnrollmentAction) @@ -2177,6 +2186,7 @@ func (s *Server) banNode(ctx context.Context, node *storage.EnrolledNode) error if err := s.store.SetNodeBanned(ctx, node.PeerID, true); err != nil { return err } + s.dropCatalogEntry(node.PeerID) if node.ClaimsJSON == "" { return nil } diff --git a/internal/controlplane/ui.go b/internal/controlplane/ui.go index cb688fcb..064f5a04 100644 --- a/internal/controlplane/ui.go +++ b/internal/controlplane/ui.go @@ -17,6 +17,7 @@ package controlplane import ( "encoding/json" "net/http" + "time" ) // HandleAdminStatus returns a consolidated JSON state of the control plane. @@ -86,6 +87,7 @@ func (s *Server) HandleAdminStatus(w http.ResponseWriter, r *http.Request) { "enrollment_requests": reqs, "bootstrap_tokens": tokens, "policy_json": policyJSON, + "node_catalog": s.catalogViewFor(nodes, time.Now()), } w.Header().Set("Content-Type", "application/json") diff --git a/internal/node/controlplane.go b/internal/node/controlplane.go index de40ed4a..b05392ee 100644 --- a/internal/node/controlplane.go +++ b/internal/node/controlplane.go @@ -251,3 +251,40 @@ func FetchMeshPolicy(ctx context.Context, controlPlaneURL string, biscuitToken [ return &policyResp, nil } + +// ReportNodeCatalog self-reports this node's locally registered services to +// the control plane's /nodes/catalog endpoint, so an admin can see mesh-wide +// service topology (see catalog.go's HandleNodeCatalog for why this exists +// instead of the control plane discovering it via DHT/P2P itself). +func ReportNodeCatalog(ctx context.Context, controlPlaneURL string, biscuitToken []byte, services []*api.ServiceInfo) error { + if !strings.HasPrefix(controlPlaneURL, "http://") && !strings.HasPrefix(controlPlaneURL, "https://") { + controlPlaneURL = "https://" + controlPlaneURL + } + controlPlaneURL = strings.TrimSuffix(controlPlaneURL, "/") + + payload, err := proto.Marshal(&api.NodeCatalogReport{Services: services}) + if err != nil { + return fmt.Errorf("failed to encode catalog report: %w", err) + } + + urlStr := controlPlaneURL + "/nodes/catalog" + req, err := http.NewRequestWithContext(ctx, "POST", urlStr, bytes.NewReader(payload)) + if err != nil { + return fmt.Errorf("failed to create HTTP request: %w", err) + } + req.Header.Set("Content-Type", "application/x-protobuf") + req.Header.Set("Authorization", "Bearer "+base64.StdEncoding.EncodeToString(biscuitToken)) + + client := &http.Client{Timeout: 10 * time.Second} + resp, err := client.Do(req) + if err != nil { + return fmt.Errorf("HTTP request failed: %w", err) + } + defer resp.Body.Close() //nolint:errcheck + + if resp.StatusCode != http.StatusNoContent { + body, _ := io.ReadAll(io.LimitReader(resp.Body, 4096)) + return fmt.Errorf("control plane returned status %s: %s", resp.Status, string(body)) + } + return nil +} diff --git a/internal/node/controlplane_test.go b/internal/node/controlplane_test.go index e47508e3..38469539 100644 --- a/internal/node/controlplane_test.go +++ b/internal/node/controlplane_test.go @@ -17,6 +17,8 @@ package node import ( "context" "crypto/ed25519" + "encoding/base64" + "io" "net/http" "net/http/httptest" "reflect" @@ -220,3 +222,161 @@ func TestSyncMeshConfig(t *testing.T) { t.Errorf("Expected saved addrs %v, got %v", expectedInfo.RouterAddresses, savedAddrsStr) } } + +func TestReportNodeCatalog(t *testing.T) { + services := []*api.ServiceInfo{ + {Type: api.ServiceType_SERVICE_TYPE_MCP, Name: "stvv-compliance-docs", Description: "doc lookup"}, + } + + var gotAuth, gotContentType string + var gotReq api.NodeCatalogReport + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + t.Errorf("Expected POST, got %s", r.Method) + } + if r.URL.Path != "/nodes/catalog" { + t.Errorf("Expected path /nodes/catalog, got %s", r.URL.Path) + } + gotAuth = r.Header.Get("Authorization") + gotContentType = r.Header.Get("Content-Type") + body, err := io.ReadAll(r.Body) + if err != nil { + t.Errorf("failed to read request body: %v", err) + } + if err := proto.Unmarshal(body, &gotReq); err != nil { + t.Errorf("failed to decode request body: %v", err) + } + w.WriteHeader(http.StatusNoContent) + })) + defer server.Close() + + biscuitToken := []byte("fake-biscuit-bytes") + if err := ReportNodeCatalog(context.Background(), server.URL, biscuitToken, services); err != nil { + t.Fatalf("ReportNodeCatalog failed: %v", err) + } + + wantAuth := "Bearer " + base64.StdEncoding.EncodeToString(biscuitToken) + if gotAuth != wantAuth { + t.Errorf("Expected Authorization header %q, got %q", wantAuth, gotAuth) + } + if gotContentType != "application/x-protobuf" { + t.Errorf("Expected Content-Type application/x-protobuf, got %q", gotContentType) + } + if len(gotReq.Services) != 1 || gotReq.Services[0].Name != "stvv-compliance-docs" || gotReq.Services[0].Type != api.ServiceType_SERVICE_TYPE_MCP { + t.Errorf("Expected relayed services %v, got %v", services, gotReq.Services) + } +} + +func TestReportNodeCatalog_HTTPError(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "node not enrolled or not admitted", http.StatusUnauthorized) + })) + defer server.Close() + + err := ReportNodeCatalog(context.Background(), server.URL, []byte("fake-biscuit-bytes"), nil) + if err == nil { + t.Fatal("Expected error, got nil") + } + if !strings.Contains(err.Error(), "control plane returned status 401") { + t.Errorf("Expected error to mention status 401, got %v", err) + } +} + +// newCatalogTestNode is a SamNode with just enough state for the catalog +// loop: a store holding the control-plane URL, a cached identity, and a +// registry with one service. +func newCatalogTestNode(t *testing.T, controlPlaneURL string) *SamNode { + t.Helper() + store, err := NewStore(t.TempDir()) + if err != nil { + t.Fatalf("NewStore: %v", err) + } + t.Cleanup(func() { _ = store.Close() }) + if controlPlaneURL != "" { + if err := store.SaveControlPlaneURL(controlPlaneURL); err != nil { + t.Fatalf("SaveControlPlaneURL: %v", err) + } + } + node := &SamNode{Store: store, services: newServiceRegistryForTest(&fakeDHT{})} + node.services.insertService(newFakeSvc("calc", api.ServiceType_SERVICE_TYPE_MCP)) + return node +} + +func TestReportNodeCatalog_Preconditions(t *testing.T) { + noURL := newCatalogTestNode(t, "") + noURL.SetIdentityCache([]byte("biscuit")) + if err := noURL.reportNodeCatalog(context.Background()); err == nil || !strings.Contains(err.Error(), "control plane URL") { + t.Errorf("without a control-plane URL: got %v, want a URL error", err) + } + + noIdentity := newCatalogTestNode(t, "http://127.0.0.1:1") + if err := noIdentity.reportNodeCatalog(context.Background()); err == nil || !strings.Contains(err.Error(), "identity") { + t.Errorf("without an identity: got %v, want an identity error", err) + } +} + +// The loop must report once after the initial delay and then keep reporting +// every interval, carrying the live service list each time. +func TestStartCatalogReportLoop(t *testing.T) { + reports := make(chan *api.NodeCatalogReport, 16) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/nodes/catalog" || r.Method != http.MethodPost { + t.Errorf("unexpected request %s %s", r.Method, r.URL.Path) + } + body, _ := io.ReadAll(r.Body) + var report api.NodeCatalogReport + if err := proto.Unmarshal(body, &report); err != nil { + t.Errorf("decode report: %v", err) + } + reports <- &report + w.WriteHeader(http.StatusNoContent) + })) + defer server.Close() + + node := newCatalogTestNode(t, server.URL) + node.SetIdentityCache([]byte("biscuit")) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + node.startCatalogReportLoop(ctx, 10*time.Millisecond, 20*time.Millisecond) + + for i := 0; i < 3; i++ { + select { + case report := <-reports: + if len(report.Services) != 1 || report.Services[0].Name != "calc" { + t.Fatalf("report %d: got services %v, want [calc]", i, report.Services) + } + case <-time.After(5 * time.Second): + t.Fatalf("timed out waiting for report %d", i) + } + } + + // Cancelling stops the loop: no report after the in-flight one settles. + cancel() + time.Sleep(100 * time.Millisecond) + for len(reports) > 0 { + <-reports + } + select { + case <-reports: + t.Fatal("loop kept reporting after context cancellation") + case <-time.After(150 * time.Millisecond): + } +} + +func TestOptionsDefault_CatalogReport(t *testing.T) { + var o Options + o.Default() + if o.CatalogReportInterval != time.Minute { + t.Errorf("CatalogReportInterval = %v, want 1m", o.CatalogReportInterval) + } + if o.CatalogReportInitialDelay != 5*time.Second { + t.Errorf("CatalogReportInitialDelay = %v, want 5s", o.CatalogReportInitialDelay) + } + + custom := Options{CatalogReportInterval: 3 * time.Second, CatalogReportInitialDelay: time.Second} + custom.Default() + if custom.CatalogReportInterval != 3*time.Second || custom.CatalogReportInitialDelay != time.Second { + t.Errorf("Default overwrote explicit values: %+v", custom) + } +} diff --git a/internal/node/node.go b/internal/node/node.go index cf934824..b784588d 100644 --- a/internal/node/node.go +++ b/internal/node/node.go @@ -631,6 +631,10 @@ func (n *SamNode) Start(ctx context.Context) error { // Periodically sync mesh policy n.startPolicySyncLoop(ctx, n.config.PolicySyncInterval) + // Periodically self-report local services to the control plane, so an + // admin can see mesh-wide service topology. + n.startCatalogReportLoop(ctx, n.config.CatalogReportInitialDelay, n.config.CatalogReportInterval) + return nil } @@ -2208,6 +2212,72 @@ func (n *SamNode) syncMeshPolicy(ctx context.Context) error { return nil } +// reportNodeCatalog self-reports this node's local service list to the +// control plane (see internal/controlplane/catalog.go's HandleNodeCatalog), +// so an admin console can show mesh-wide service topology. +func (n *SamNode) reportNodeCatalog(ctx context.Context) error { + controlPlaneURL, err := n.Store.LoadControlPlaneURL() + if err != nil || controlPlaneURL == "" { + return fmt.Errorf("control plane URL not found in store") + } + + token := n.GetIdentity() + if len(token) == 0 { + return fmt.Errorf("node has no identity token to report its catalog") + } + + services := n.ListLocalServices(api.ServiceType_SERVICE_TYPE_UNSPECIFIED) + if err := ReportNodeCatalog(ctx, controlPlaneURL, token, services); err != nil { + return fmt.Errorf("failed to report node catalog: %w", err) + } + return nil +} + +// startCatalogReportLoop reports the local catalog once after initialDelay +// and then every interval, each wait stretched by up to a tenth of interval +// so a fleet started together does not hit the control plane in lockstep. +// A failure is logged at Warn once and at Debug while it persists: the +// usual causes (control plane unreachable, path not routed) do not change +// from one tick to the next. +func (n *SamNode) startCatalogReportLoop(ctx context.Context, initialDelay, interval time.Duration) { + if interval <= 0 { + interval = 1 * time.Minute + } + if initialDelay <= 0 { + initialDelay = 5 * time.Second + } + + go func() { + wait := initialDelay + failures := 0 + for { + timer := time.NewTimer(wait + time.Duration(rand.Int63n(int64(interval/10)+1))) + select { + case <-ctx.Done(): + timer.Stop() + return + case <-timer.C: + } + wait = interval + + err := n.reportNodeCatalog(ctx) + switch { + case err != nil && failures == 0: + logger.Warnf("Node catalog report failed: %v", err) + case err != nil: + logger.Debugf("Node catalog report still failing (%d consecutive): %v", failures+1, err) + case failures > 0: + logger.Infof("Node catalog report recovered after %d failures", failures) + } + if err != nil { + failures++ + } else { + failures = 0 + } + } + }() +} + func (n *SamNode) startPolicySyncLoop(ctx context.Context, interval time.Duration) { if interval <= 0 { interval = 1 * time.Hour diff --git a/internal/node/options.go b/internal/node/options.go index 730ecd1e..dbc5e30e 100644 --- a/internal/node/options.go +++ b/internal/node/options.go @@ -86,6 +86,16 @@ type Options struct { // is tight enough that even simple interpreted-language MCP servers can // miss it on first spawn. BackendProbeTimeout time.Duration + // CatalogReportInterval specifies how often the node self-reports its + // locally registered services to the control plane (POST + // /nodes/catalog), so an admin can see mesh-wide service topology + // without the control plane needing DHT/P2P access to every node + // itself. Zero uses the default. + CatalogReportInterval time.Duration + // CatalogReportInitialDelay is how long after Start the first catalog + // report is sent, so services configured at startup have registered by + // then. Zero uses the default. + CatalogReportInitialDelay time.Duration } // Default applies default values to Options if they are not specified. @@ -136,6 +146,12 @@ func (o *Options) Default() { if o.PolicySyncInterval == 0 { o.PolicySyncInterval = 1 * time.Hour } + if o.CatalogReportInterval <= 0 { + o.CatalogReportInterval = 1 * time.Minute + } + if o.CatalogReportInitialDelay <= 0 { + o.CatalogReportInitialDelay = 5 * time.Second + } if o.PolicySyncJitter <= 0 { o.PolicySyncJitter = 10 * time.Second } diff --git a/site/content/docs/user/control-plane-configuration.md b/site/content/docs/user/control-plane-configuration.md index 3792bd4c..ffabae5f 100644 --- a/site/content/docs/user/control-plane-configuration.md +++ b/site/content/docs/user/control-plane-configuration.md @@ -203,7 +203,19 @@ Administrators can immediately revoke any active session to disable a node's abi --- -## 6. Headless Node Enrollment (Bootstrap Token Flow) +## 6. Node Service Catalog + +Each enrolled node periodically self-reports its locally registered services so the admin console's **Services** view can show what is running where, without the Control Plane joining the P2P mesh. + +* **Endpoint**: `POST /nodes/catalog` +* **Authentication**: The node's own Biscuit as a `Bearer` token. The reporting peer is taken from the verified token, never from the body, so a node can only ever describe itself; reports from banned or expired enrollments are rejected. +* **Payload**: An `api.NodeCatalogReport` protobuf (`application/x-protobuf`) with up to 512 services per report. +* **Cadence**: Every minute by default (`node.Options.CatalogReportInterval`), with jitter. +* **Semantics**: A live-status cache, not authoritative state. It is display-only and never feeds authorization; the cache is in-memory and rebuilt by the nodes' next reports after a Control Plane restart. If you front the Control Plane with a Gateway allow-list, `/nodes/catalog` must be routed like the other node-facing endpoints. + +--- + +## 7. Headless Node Enrollment (Bootstrap Token Flow) To enroll a headless server, router, or background daemon that cannot complete interactive OIDC authentication, SAM supports a **Bootstrap Token** flow. diff --git a/tests/ui/console.spec.js b/tests/ui/console.spec.js index ea2950fe..0d44af25 100644 --- a/tests/ui/console.spec.js +++ b/tests/ui/console.spec.js @@ -312,3 +312,85 @@ test('the served markup carries no inline style attribute', async ({ request }) const html = await (await request.get('/index.html')).text(); expect(html).not.toMatch(/<[^>]+\sstyle=/); }); + +// The Services view has no admin flow to seed it from the browser alone (a +// real node has to self-report), so this only pins what is reachable without +// one: the nav item routes there and the empty state renders without a JS +// error. TestHandleNodeCatalog in internal/controlplane covers the reporting +// endpoint itself, and TestReportNodeCatalog in internal/node covers the +// node-side push. +test('the Services view is reachable and renders its empty state', async ({ page }) => { + await login(page); + + await page.click('.nav-item[data-target="services"]'); + await expect(page).toHaveURL(/#services$/); + await expect(page.locator('#view-services')).toBeVisible(); + await expect(page.locator('#table-services')).toContainText('No nodes have reported any services yet'); +}); + +// A reload must survive on this view too, same as the other deep-linkable +// views already covered above for #routers. +test('the Services view is deep-linkable via the URL hash', async ({ page }) => { + await login(page); + + await page.goto('/#services'); + await expect(page.locator('#view-services')).toBeVisible(); + await expect(page.locator('.nav-item[data-target="services"]')).toHaveAttribute('aria-current', 'page'); +}); + +// Seeding the Services view for real needs an enrolled node self-reporting, so +// render it from an intercepted /admin/status instead: the real response with +// node_catalog and a labelled node spliced in, shaped exactly like +// catalogViewFor in internal/controlplane/catalog.go emits it. +test('reported services render with type, labels and report time', async ({ page }) => { + const PEER = '12D3KooWServicesViewFixturePeer'; + await page.route('**/api/admin/status', async (route) => { + const response = await route.fetch(); + const status = await response.json(); + status.enrolled_nodes = [...(status.enrolled_nodes || []), { + PeerID: PEER, + Role: 'sam:role:node', + OwnerID: 'root-admin', + Labels: { component: 'stvv', region: 'eu-west' }, + }]; + status.node_catalog = { + [PEER]: { + services: [ + { name: 'compliance-docs', type: 'mcp', description: 'doc lookup' }, + // Reported names and descriptions are node-controlled input to an + // admin page; markup in them must render inert. + { name: 'llama', type: 'inference', description: '' }, + ], + reported_at: '2026-09-15T08:00:00Z', + }, + }; + await route.fulfill({ response, json: status }); + }); + + await login(page); + await page.click('.nav-item[data-target="services"]'); + + const rows = page.locator('#table-services tr'); + await expect(rows).toHaveCount(2); + + const first = rows.nth(0); + await expect(first).toContainText('compliance-docs'); + await expect(first).toContainText('mcp'); + await expect(first).toContainText('doc lookup'); + await expect(first).toContainText(PEER); + // The node's labels ride along as the mnemonic under the peer ID. + await expect(first.locator('.cell-subtext')).toHaveText('component=stvv, region=eu-west'); + // reported_at renders as a local time, not the raw RFC 3339 string. + await expect(first.locator('td').nth(4)).not.toHaveText(/2026-09-15T08:00:00Z|-/); + + await expect(rows.nth(1)).toContainText('inference'); + // The description shows as text; nothing was injected into the DOM. + await expect(rows.nth(1)).toContainText('onerror'); + expect(await page.locator('#table-services img').count()).toBe(0); + expect(await page.evaluate(() => window.svcXSS)).toBeUndefined(); + + // The same labels also annotate the peer ID in the Nodes table. + await page.click('.nav-item[data-target="nodes"]'); + const nodeRow = page.locator('#table-nodes tr', { hasText: PEER }); + await expect(nodeRow.locator('.cell-subtext')).toHaveText('component=stvv, region=eu-west'); +});