From a2d749801eb56005c433f964afff4259ec97d56c Mon Sep 17 00:00:00 2001 From: Roddie Kieley Date: Wed, 26 Aug 2026 19:01:32 -0230 Subject: [PATCH 1/3] feat: telemetry MetricsStream reverse-scrape hub (JEP-0013 Phase 3) Add the MetricsStream protocol and Go hub so Prometheus can scrape merged exporter OpenMetrics from telemetry without an exporter client yet. Generated Python stubs are included for proto consistency. Co-authored-by: Cursor --- controller/cmd/telemetry/main.go | 49 ++- .../protocol/jumpstarter/v1/telemetry.pb.go | 402 ++++++++++++++++-- .../jumpstarter/v1/telemetry_grpc.pb.go | 47 +- controller/internal/service/metrics_merge.go | 239 +++++++++++ .../internal/service/metrics_merge_test.go | 172 ++++++++ controller/internal/service/metrics_stream.go | 241 +++++++++++ .../internal/service/metrics_stream_test.go | 402 ++++++++++++++++++ controller/internal/service/telemetry_http.go | 135 ++++++ .../internal/service/telemetry_identity.go | 62 +++ .../internal/service/telemetry_service.go | 83 ++-- protocol/proto/jumpstarter/v1/telemetry.proto | 33 +- .../jumpstarter/v1/telemetry_pb2.py | 34 +- .../jumpstarter/v1/telemetry_pb2.pyi | 102 ++++- .../jumpstarter/v1/telemetry_pb2_grpc.py | 51 ++- .../jumpstarter/v1/telemetry_pb2_grpc.pyi | 49 ++- 15 files changed, 1999 insertions(+), 102 deletions(-) create mode 100644 controller/internal/service/metrics_merge.go create mode 100644 controller/internal/service/metrics_merge_test.go create mode 100644 controller/internal/service/metrics_stream.go create mode 100644 controller/internal/service/metrics_stream_test.go create mode 100644 controller/internal/service/telemetry_http.go create mode 100644 controller/internal/service/telemetry_identity.go diff --git a/controller/cmd/telemetry/main.go b/controller/cmd/telemetry/main.go index 1c757d42e..cf2266048 100644 --- a/controller/cmd/telemetry/main.go +++ b/controller/cmd/telemetry/main.go @@ -14,9 +14,10 @@ See the License for the specific language governing permissions and limitations under the License. */ -// jumpstarter-telemetry receives structured log entries from exporters and clients -// via the PushLogs gRPC RPC and writes them to structured stdout for downstream -// log shippers (Promtail, Grafana Alloy, Vector) to forward to Loki. +// jumpstarter-telemetry reverse-scrapes exporter metrics via MetricsStream and +// receives structured log entries via PushLogs. Logs are written to structured +// stdout for downstream log shippers (Promtail, Grafana Alloy, Vector). +// Loki push is a later Phase 3 PR. // // TLS: always enabled. Set EXTERNAL_CERT_PEM and EXTERNAL_KEY_PEM to file paths of // operator-mounted cert/key (e.g. from a cert-manager Secret); when absent a @@ -30,8 +31,7 @@ limitations under the License. // certificate; the controller uses it to advertise the address to exporters via // GetServiceEndpoints. A mismatch causes TLS hostname verification failures. // -// Future phases will add direct Loki push and MetricsStream for reverse-scrape -// of exporter prometheus_client registries. +// HTTP: GET /metrics, /healthz, and /readyz bind separately (default :8080). package main import ( @@ -39,7 +39,9 @@ import ( "flag" "os" "os/signal" + "strings" "syscall" + "time" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/log/zap" @@ -55,9 +57,36 @@ var ( buildDate = "unknown" ) +func splitCSV(s string) []string { + if s == "" { + return nil + } + parts := strings.Split(s, ",") + out := make([]string, 0, len(parts)) + for _, p := range parts { + p = strings.TrimSpace(p) + if p != "" { + out = append(out, p) + } + } + return out +} + func main() { var bindAddr string + var metricsAddr string + var scrapeTimeout time.Duration + var driverTypeEnum string + var exemplarKeys string flag.StringVar(&bindAddr, "grpc-bind", ":9093", "TCP address to bind the gRPC server to") + flag.StringVar(&metricsAddr, "metrics-bind-address", ":8080", + "TCP address for HTTP GET /metrics, /healthz, and /readyz. Use 0 to disable.") + flag.DurationVar(&scrapeTimeout, "scrape-timeout", 7*time.Second, + "Max wait for parallel MetricsStream scrape responses") + flag.StringVar(&driverTypeEnum, "driver-type-enum", strings.Join(service.DefaultDriverTypeEnum, ","), + "Comma-separated allowlist of driver_type values; others are remapped to other") + flag.StringVar(&exemplarKeys, "exemplar-keys", strings.Join(service.DefaultExemplarKeys, ","), + "Comma-separated allowlist of Prometheus exemplar keys") opts := zap.Options{} opts.BindFlags(flag.CommandLine) @@ -71,6 +100,8 @@ func main() { "gitCommit", gitCommit, "buildDate", buildDate, "bindAddr", bindAddr, + "metricsBindAddr", metricsAddr, + "scrapeTimeout", scrapeTimeout, ) ctx, cancel := context.WithCancel(context.Background()) @@ -87,8 +118,12 @@ func main() { } svc := &service.TelemetryService{ - BindAddr: bindAddr, - Signer: signer, + BindAddr: bindAddr, + MetricsBindAddr: metricsAddr, + ScrapeTimeout: scrapeTimeout, + DriverTypeEnum: splitCSV(driverTypeEnum), + ExemplarKeys: splitCSV(exemplarKeys), + Signer: signer, } // Register signal handler before starting the service so no signal diff --git a/controller/internal/protocol/jumpstarter/v1/telemetry.pb.go b/controller/internal/protocol/jumpstarter/v1/telemetry.pb.go index 00e583612..ee0ac7352 100644 --- a/controller/internal/protocol/jumpstarter/v1/telemetry.pb.go +++ b/controller/internal/protocol/jumpstarter/v1/telemetry.pb.go @@ -1,4 +1,4 @@ -// Copyright 2024 The Jumpstarter Authors +// Copyright 2026 The Jumpstarter Authors // Code generated by protoc-gen-go. DO NOT EDIT. // versions: @@ -24,6 +24,289 @@ const ( _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) ) +// Exporter → Telemetry +type MetricsStreamRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Types that are valid to be assigned to Msg: + // + // *MetricsStreamRequest_Register + // *MetricsStreamRequest_ScrapeResponse + Msg isMetricsStreamRequest_Msg `protobuf_oneof:"msg"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MetricsStreamRequest) Reset() { + *x = MetricsStreamRequest{} + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MetricsStreamRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MetricsStreamRequest) ProtoMessage() {} + +func (x *MetricsStreamRequest) ProtoReflect() protoreflect.Message { + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[0] + 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 MetricsStreamRequest.ProtoReflect.Descriptor instead. +func (*MetricsStreamRequest) Descriptor() ([]byte, []int) { + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{0} +} + +func (x *MetricsStreamRequest) GetMsg() isMetricsStreamRequest_Msg { + if x != nil { + return x.Msg + } + return nil +} + +func (x *MetricsStreamRequest) GetRegister() *MetricsRegister { + if x != nil { + if x, ok := x.Msg.(*MetricsStreamRequest_Register); ok { + return x.Register + } + } + return nil +} + +func (x *MetricsStreamRequest) GetScrapeResponse() *MetricsScrapeResponse { + if x != nil { + if x, ok := x.Msg.(*MetricsStreamRequest_ScrapeResponse); ok { + return x.ScrapeResponse + } + } + return nil +} + +type isMetricsStreamRequest_Msg interface { + isMetricsStreamRequest_Msg() +} + +type MetricsStreamRequest_Register struct { + Register *MetricsRegister `protobuf:"bytes,1,opt,name=register,proto3,oneof"` // First message: identify this exporter. +} + +type MetricsStreamRequest_ScrapeResponse struct { + ScrapeResponse *MetricsScrapeResponse `protobuf:"bytes,2,opt,name=scrape_response,json=scrapeResponse,proto3,oneof"` // Subsequent: reply to a scrape. +} + +func (*MetricsStreamRequest_Register) isMetricsStreamRequest_Msg() {} + +func (*MetricsStreamRequest_ScrapeResponse) isMetricsStreamRequest_Msg() {} + +type MetricsRegister struct { + state protoimpl.MessageState `protogen:"open.v1"` + Identity string `protobuf:"bytes,1,opt,name=identity,proto3" json:"identity,omitempty"` // Exporter CRD name (verified against the auth token by the server). + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MetricsRegister) Reset() { + *x = MetricsRegister{} + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[1] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MetricsRegister) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MetricsRegister) ProtoMessage() {} + +func (x *MetricsRegister) ProtoReflect() protoreflect.Message { + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[1] + 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 MetricsRegister.ProtoReflect.Descriptor instead. +func (*MetricsRegister) Descriptor() ([]byte, []int) { + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{1} +} + +func (x *MetricsRegister) GetIdentity() string { + if x != nil { + return x.Identity + } + return "" +} + +type MetricsScrapeResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + MetricsText []byte `protobuf:"bytes,1,opt,name=metrics_text,json=metricsText,proto3" json:"metrics_text,omitempty"` // generate_latest() OpenMetrics output. + Timestamp *timestamppb.Timestamp `protobuf:"bytes,2,opt,name=timestamp,proto3" json:"timestamp,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MetricsScrapeResponse) Reset() { + *x = MetricsScrapeResponse{} + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[2] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MetricsScrapeResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MetricsScrapeResponse) ProtoMessage() {} + +func (x *MetricsScrapeResponse) ProtoReflect() protoreflect.Message { + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[2] + 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 MetricsScrapeResponse.ProtoReflect.Descriptor instead. +func (*MetricsScrapeResponse) Descriptor() ([]byte, []int) { + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{2} +} + +func (x *MetricsScrapeResponse) GetMetricsText() []byte { + if x != nil { + return x.MetricsText + } + return nil +} + +func (x *MetricsScrapeResponse) GetTimestamp() *timestamppb.Timestamp { + if x != nil { + return x.Timestamp + } + return nil +} + +// Telemetry → Exporter +type MetricsStreamResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Types that are valid to be assigned to Msg: + // + // *MetricsStreamResponse_ScrapeRequest + Msg isMetricsStreamResponse_Msg `protobuf_oneof:"msg"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MetricsStreamResponse) Reset() { + *x = MetricsStreamResponse{} + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[3] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MetricsStreamResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MetricsStreamResponse) ProtoMessage() {} + +func (x *MetricsStreamResponse) ProtoReflect() protoreflect.Message { + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[3] + 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 MetricsStreamResponse.ProtoReflect.Descriptor instead. +func (*MetricsStreamResponse) Descriptor() ([]byte, []int) { + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{3} +} + +func (x *MetricsStreamResponse) GetMsg() isMetricsStreamResponse_Msg { + if x != nil { + return x.Msg + } + return nil +} + +func (x *MetricsStreamResponse) GetScrapeRequest() *MetricsScrapeRequest { + if x != nil { + if x, ok := x.Msg.(*MetricsStreamResponse_ScrapeRequest); ok { + return x.ScrapeRequest + } + } + return nil +} + +type isMetricsStreamResponse_Msg interface { + isMetricsStreamResponse_Msg() +} + +type MetricsStreamResponse_ScrapeRequest struct { + ScrapeRequest *MetricsScrapeRequest `protobuf:"bytes,1,opt,name=scrape_request,json=scrapeRequest,proto3,oneof"` +} + +func (*MetricsStreamResponse_ScrapeRequest) isMetricsStreamResponse_Msg() {} + +// Empty request: "send your /metrics now". +type MetricsScrapeRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *MetricsScrapeRequest) Reset() { + *x = MetricsScrapeRequest{} + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[4] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *MetricsScrapeRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*MetricsScrapeRequest) ProtoMessage() {} + +func (x *MetricsScrapeRequest) ProtoReflect() protoreflect.Message { + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[4] + 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 MetricsScrapeRequest.ProtoReflect.Descriptor instead. +func (*MetricsScrapeRequest) Descriptor() ([]byte, []int) { + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{4} +} + // Request to push log entries to the telemetry service. type PushLogsRequest struct { state protoimpl.MessageState `protogen:"open.v1"` @@ -34,7 +317,7 @@ type PushLogsRequest struct { func (x *PushLogsRequest) Reset() { *x = PushLogsRequest{} - mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[0] + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[5] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -46,7 +329,7 @@ func (x *PushLogsRequest) String() string { func (*PushLogsRequest) ProtoMessage() {} func (x *PushLogsRequest) ProtoReflect() protoreflect.Message { - mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[0] + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[5] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -59,7 +342,7 @@ func (x *PushLogsRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use PushLogsRequest.ProtoReflect.Descriptor instead. func (*PushLogsRequest) Descriptor() ([]byte, []int) { - return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{0} + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{5} } func (x *PushLogsRequest) GetEntries() []*LogEntry { @@ -80,7 +363,7 @@ type PushLogsResponse struct { func (x *PushLogsResponse) Reset() { *x = PushLogsResponse{} - mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[1] + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[6] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -92,7 +375,7 @@ func (x *PushLogsResponse) String() string { func (*PushLogsResponse) ProtoMessage() {} func (x *PushLogsResponse) ProtoReflect() protoreflect.Message { - mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[1] + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[6] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -105,7 +388,7 @@ func (x *PushLogsResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use PushLogsResponse.ProtoReflect.Descriptor instead. func (*PushLogsResponse) Descriptor() ([]byte, []int) { - return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{1} + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{6} } func (x *PushLogsResponse) GetAccepted() uint32 { @@ -125,26 +408,27 @@ func (x *PushLogsResponse) GetDropped() uint32 { // A structured log entry from an exporter or client. // Maps directly to a Loki log entry with stream labels and body fields. type LogEntry struct { - state protoimpl.MessageState `protogen:"open.v1"` - Timestamp *timestamppb.Timestamp `protobuf:"bytes,1,opt,name=timestamp,proto3" json:"timestamp,omitempty"` // When the log was emitted. - Severity string `protobuf:"bytes,2,opt,name=severity,proto3" json:"severity,omitempty"` // Log severity: debug, info, warning, error, critical. - Message string `protobuf:"bytes,3,opt,name=message,proto3" json:"message,omitempty"` // Human-readable log message. - Component string `protobuf:"bytes,4,opt,name=component,proto3" json:"component,omitempty"` // Loki stream label: cli, exporter, controller, router, telemetry. - Exporter string `protobuf:"bytes,5,opt,name=exporter,proto3" json:"exporter,omitempty"` // Loki stream label: exporter CRD name (bounded by cluster size). - Lease string `protobuf:"bytes,6,opt,name=lease,proto3" json:"lease,omitempty"` // Log body only (high cardinality): active lease name. - Client string `protobuf:"bytes,7,opt,name=client,proto3" json:"client,omitempty"` // Log body only (high cardinality): client CRD name. - Operation string `protobuf:"bytes,8,opt,name=operation,proto3" json:"operation,omitempty"` // Log body: operation name (flash, power, etc.). - Result string `protobuf:"bytes,9,opt,name=result,proto3" json:"result,omitempty"` // Log body: operation outcome (success, failure, etc.). - DriverType string `protobuf:"bytes,10,opt,name=driver_type,json=driverType,proto3" json:"driver_type,omitempty"` // Log body: driver category (storage, power, network, etc.). - ExtraFields map[string]string `protobuf:"bytes,11,rep,name=extra_fields,json=extraFields,proto3" json:"extra_fields,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` // Additional structured key-value fields. - Namespace string `protobuf:"bytes,12,opt,name=namespace,proto3" json:"namespace,omitempty"` // Loki stream label: Kubernetes namespace (bounded by cluster size). + state protoimpl.MessageState `protogen:"open.v1"` + Timestamp *timestamppb.Timestamp `protobuf:"bytes,1,opt,name=timestamp,proto3" json:"timestamp,omitempty"` // When the log was emitted. + Severity string `protobuf:"bytes,2,opt,name=severity,proto3" json:"severity,omitempty"` // Log severity: debug, info, warning, error, critical. + Message string `protobuf:"bytes,3,opt,name=message,proto3" json:"message,omitempty"` // Human-readable log message. + Component string `protobuf:"bytes,4,opt,name=component,proto3" json:"component,omitempty"` // Loki stream label: cli, exporter, controller, router, telemetry. + Exporter string `protobuf:"bytes,5,opt,name=exporter,proto3" json:"exporter,omitempty"` // Loki stream label: exporter CRD name (bounded by cluster size). + Lease string `protobuf:"bytes,6,opt,name=lease,proto3" json:"lease,omitempty"` // Log body only (high cardinality): active lease name. + Client string `protobuf:"bytes,7,opt,name=client,proto3" json:"client,omitempty"` // Log body only (high cardinality): client CRD name. + Operation string `protobuf:"bytes,8,opt,name=operation,proto3" json:"operation,omitempty"` // Log body: operation name (flash, power, etc.). + Result string `protobuf:"bytes,9,opt,name=result,proto3" json:"result,omitempty"` // Log body: operation outcome (success, failure, etc.). + DriverType string `protobuf:"bytes,10,opt,name=driver_type,json=driverType,proto3" json:"driver_type,omitempty"` // Log body: driver category (storage, power, network, etc.). + ExtraFields map[string]string `protobuf:"bytes,11,rep,name=extra_fields,json=extraFields,proto3" json:"extra_fields,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` // Additional structured key-value fields. + // Capped at 16 entries, 64-char keys, 256-char values. + Namespace string `protobuf:"bytes,12,opt,name=namespace,proto3" json:"namespace,omitempty"` // Loki stream label: Kubernetes namespace (bounded by cluster size). unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } func (x *LogEntry) Reset() { *x = LogEntry{} - mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[2] + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[7] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -156,7 +440,7 @@ func (x *LogEntry) String() string { func (*LogEntry) ProtoMessage() {} func (x *LogEntry) ProtoReflect() protoreflect.Message { - mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[2] + mi := &file_jumpstarter_v1_telemetry_proto_msgTypes[7] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -169,7 +453,7 @@ func (x *LogEntry) ProtoReflect() protoreflect.Message { // Deprecated: Use LogEntry.ProtoReflect.Descriptor instead. func (*LogEntry) Descriptor() ([]byte, []int) { - return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{2} + return file_jumpstarter_v1_telemetry_proto_rawDescGZIP(), []int{7} } func (x *LogEntry) GetTimestamp() *timestamppb.Timestamp { @@ -260,7 +544,20 @@ var File_jumpstarter_v1_telemetry_proto protoreflect.FileDescriptor const file_jumpstarter_v1_telemetry_proto_rawDesc = "" + "\n" + - "\x1ejumpstarter/v1/telemetry.proto\x12\x0ejumpstarter.v1\x1a\x1fgoogle/protobuf/timestamp.proto\"E\n" + + "\x1ejumpstarter/v1/telemetry.proto\x12\x0ejumpstarter.v1\x1a\x1fgoogle/protobuf/timestamp.proto\"\xae\x01\n" + + "\x14MetricsStreamRequest\x12=\n" + + "\bregister\x18\x01 \x01(\v2\x1f.jumpstarter.v1.MetricsRegisterH\x00R\bregister\x12P\n" + + "\x0fscrape_response\x18\x02 \x01(\v2%.jumpstarter.v1.MetricsScrapeResponseH\x00R\x0escrapeResponseB\x05\n" + + "\x03msg\"-\n" + + "\x0fMetricsRegister\x12\x1a\n" + + "\bidentity\x18\x01 \x01(\tR\bidentity\"t\n" + + "\x15MetricsScrapeResponse\x12!\n" + + "\fmetrics_text\x18\x01 \x01(\fR\vmetricsText\x128\n" + + "\ttimestamp\x18\x02 \x01(\v2\x1a.google.protobuf.TimestampR\ttimestamp\"m\n" + + "\x15MetricsStreamResponse\x12M\n" + + "\x0escrape_request\x18\x01 \x01(\v2$.jumpstarter.v1.MetricsScrapeRequestH\x00R\rscrapeRequestB\x05\n" + + "\x03msg\"\x16\n" + + "\x14MetricsScrapeRequest\"E\n" + "\x0fPushLogsRequest\x122\n" + "\aentries\x18\x01 \x03(\v2\x18.jumpstarter.v1.LogEntryR\aentries\"H\n" + "\x10PushLogsResponse\x12\x1a\n" + @@ -283,8 +580,9 @@ const file_jumpstarter_v1_telemetry_proto_rawDesc = "" + "\tnamespace\x18\f \x01(\tR\tnamespace\x1a>\n" + "\x10ExtraFieldsEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + - "\x05value\x18\x02 \x01(\tR\x05value:\x028\x012a\n" + - "\x10TelemetryService\x12M\n" + + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x012\xc3\x01\n" + + "\x10TelemetryService\x12`\n" + + "\rMetricsStream\x12$.jumpstarter.v1.MetricsStreamRequest\x1a%.jumpstarter.v1.MetricsStreamResponse(\x010\x01\x12M\n" + "\bPushLogs\x12\x1f.jumpstarter.v1.PushLogsRequest\x1a .jumpstarter.v1.PushLogsResponseB\xdf\x01\n" + "\x12com.jumpstarter.v1B\x0eTelemetryProtoP\x01Z`github.com/jumpstarter-dev/jumpstarter/controller/internal/protocol/jumpstarter/v1;jumpstarterv1\xa2\x02\x03JXX\xaa\x02\x0eJumpstarter.V1\xca\x02\x0eJumpstarter\\V1\xe2\x02\x1aJumpstarter\\V1\\GPBMetadata\xea\x02\x0fJumpstarter::V1b\x06proto3" @@ -300,25 +598,36 @@ func file_jumpstarter_v1_telemetry_proto_rawDescGZIP() []byte { return file_jumpstarter_v1_telemetry_proto_rawDescData } -var file_jumpstarter_v1_telemetry_proto_msgTypes = make([]protoimpl.MessageInfo, 4) +var file_jumpstarter_v1_telemetry_proto_msgTypes = make([]protoimpl.MessageInfo, 9) var file_jumpstarter_v1_telemetry_proto_goTypes = []any{ - (*PushLogsRequest)(nil), // 0: jumpstarter.v1.PushLogsRequest - (*PushLogsResponse)(nil), // 1: jumpstarter.v1.PushLogsResponse - (*LogEntry)(nil), // 2: jumpstarter.v1.LogEntry - nil, // 3: jumpstarter.v1.LogEntry.ExtraFieldsEntry - (*timestamppb.Timestamp)(nil), // 4: google.protobuf.Timestamp + (*MetricsStreamRequest)(nil), // 0: jumpstarter.v1.MetricsStreamRequest + (*MetricsRegister)(nil), // 1: jumpstarter.v1.MetricsRegister + (*MetricsScrapeResponse)(nil), // 2: jumpstarter.v1.MetricsScrapeResponse + (*MetricsStreamResponse)(nil), // 3: jumpstarter.v1.MetricsStreamResponse + (*MetricsScrapeRequest)(nil), // 4: jumpstarter.v1.MetricsScrapeRequest + (*PushLogsRequest)(nil), // 5: jumpstarter.v1.PushLogsRequest + (*PushLogsResponse)(nil), // 6: jumpstarter.v1.PushLogsResponse + (*LogEntry)(nil), // 7: jumpstarter.v1.LogEntry + nil, // 8: jumpstarter.v1.LogEntry.ExtraFieldsEntry + (*timestamppb.Timestamp)(nil), // 9: google.protobuf.Timestamp } var file_jumpstarter_v1_telemetry_proto_depIdxs = []int32{ - 2, // 0: jumpstarter.v1.PushLogsRequest.entries:type_name -> jumpstarter.v1.LogEntry - 4, // 1: jumpstarter.v1.LogEntry.timestamp:type_name -> google.protobuf.Timestamp - 3, // 2: jumpstarter.v1.LogEntry.extra_fields:type_name -> jumpstarter.v1.LogEntry.ExtraFieldsEntry - 0, // 3: jumpstarter.v1.TelemetryService.PushLogs:input_type -> jumpstarter.v1.PushLogsRequest - 1, // 4: jumpstarter.v1.TelemetryService.PushLogs:output_type -> jumpstarter.v1.PushLogsResponse - 4, // [4:5] is the sub-list for method output_type - 3, // [3:4] is the sub-list for method input_type - 3, // [3:3] is the sub-list for extension type_name - 3, // [3:3] is the sub-list for extension extendee - 0, // [0:3] is the sub-list for field type_name + 1, // 0: jumpstarter.v1.MetricsStreamRequest.register:type_name -> jumpstarter.v1.MetricsRegister + 2, // 1: jumpstarter.v1.MetricsStreamRequest.scrape_response:type_name -> jumpstarter.v1.MetricsScrapeResponse + 9, // 2: jumpstarter.v1.MetricsScrapeResponse.timestamp:type_name -> google.protobuf.Timestamp + 4, // 3: jumpstarter.v1.MetricsStreamResponse.scrape_request:type_name -> jumpstarter.v1.MetricsScrapeRequest + 7, // 4: jumpstarter.v1.PushLogsRequest.entries:type_name -> jumpstarter.v1.LogEntry + 9, // 5: jumpstarter.v1.LogEntry.timestamp:type_name -> google.protobuf.Timestamp + 8, // 6: jumpstarter.v1.LogEntry.extra_fields:type_name -> jumpstarter.v1.LogEntry.ExtraFieldsEntry + 0, // 7: jumpstarter.v1.TelemetryService.MetricsStream:input_type -> jumpstarter.v1.MetricsStreamRequest + 5, // 8: jumpstarter.v1.TelemetryService.PushLogs:input_type -> jumpstarter.v1.PushLogsRequest + 3, // 9: jumpstarter.v1.TelemetryService.MetricsStream:output_type -> jumpstarter.v1.MetricsStreamResponse + 6, // 10: jumpstarter.v1.TelemetryService.PushLogs:output_type -> jumpstarter.v1.PushLogsResponse + 9, // [9:11] is the sub-list for method output_type + 7, // [7:9] is the sub-list for method input_type + 7, // [7:7] is the sub-list for extension type_name + 7, // [7:7] is the sub-list for extension extendee + 0, // [0:7] is the sub-list for field type_name } func init() { file_jumpstarter_v1_telemetry_proto_init() } @@ -326,13 +635,20 @@ func file_jumpstarter_v1_telemetry_proto_init() { if File_jumpstarter_v1_telemetry_proto != nil { return } + file_jumpstarter_v1_telemetry_proto_msgTypes[0].OneofWrappers = []any{ + (*MetricsStreamRequest_Register)(nil), + (*MetricsStreamRequest_ScrapeResponse)(nil), + } + file_jumpstarter_v1_telemetry_proto_msgTypes[3].OneofWrappers = []any{ + (*MetricsStreamResponse_ScrapeRequest)(nil), + } type x struct{} out := protoimpl.TypeBuilder{ File: protoimpl.DescBuilder{ GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_jumpstarter_v1_telemetry_proto_rawDesc), len(file_jumpstarter_v1_telemetry_proto_rawDesc)), NumEnums: 0, - NumMessages: 4, + NumMessages: 9, NumExtensions: 0, NumServices: 1, }, diff --git a/controller/internal/protocol/jumpstarter/v1/telemetry_grpc.pb.go b/controller/internal/protocol/jumpstarter/v1/telemetry_grpc.pb.go index 356b656a2..d5db4a9a7 100644 --- a/controller/internal/protocol/jumpstarter/v1/telemetry_grpc.pb.go +++ b/controller/internal/protocol/jumpstarter/v1/telemetry_grpc.pb.go @@ -1,4 +1,4 @@ -// Copyright 2024 The Jumpstarter Authors +// Copyright 2026 The Jumpstarter Authors // Code generated by protoc-gen-go-grpc. DO NOT EDIT. // versions: @@ -21,16 +21,20 @@ import ( const _ = grpc.SupportPackageIsVersion9 const ( - TelemetryService_PushLogs_FullMethodName = "/jumpstarter.v1.TelemetryService/PushLogs" + TelemetryService_MetricsStream_FullMethodName = "/jumpstarter.v1.TelemetryService/MetricsStream" + TelemetryService_PushLogs_FullMethodName = "/jumpstarter.v1.TelemetryService/PushLogs" ) // TelemetryServiceClient is the client API for TelemetryService service. // // For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream. // -// A service that receives structured logs from exporters and clients. +// A service that reverse-scrapes exporter metrics and receives structured logs. // Implemented by jumpstarter-telemetry; not part of the controller. type TelemetryServiceClient interface { + // Persistent bidirectional stream: telemetry sends scrape requests, + // exporter responds with full metric snapshots (OpenMetrics text). + MetricsStream(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[MetricsStreamRequest, MetricsStreamResponse], error) // Push structured log entries to the telemetry service for Loki ingest. PushLogs(ctx context.Context, in *PushLogsRequest, opts ...grpc.CallOption) (*PushLogsResponse, error) } @@ -43,6 +47,19 @@ func NewTelemetryServiceClient(cc grpc.ClientConnInterface) TelemetryServiceClie return &telemetryServiceClient{cc} } +func (c *telemetryServiceClient) MetricsStream(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[MetricsStreamRequest, MetricsStreamResponse], error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + stream, err := c.cc.NewStream(ctx, &TelemetryService_ServiceDesc.Streams[0], TelemetryService_MetricsStream_FullMethodName, cOpts...) + if err != nil { + return nil, err + } + x := &grpc.GenericClientStream[MetricsStreamRequest, MetricsStreamResponse]{ClientStream: stream} + return x, nil +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type TelemetryService_MetricsStreamClient = grpc.BidiStreamingClient[MetricsStreamRequest, MetricsStreamResponse] + func (c *telemetryServiceClient) PushLogs(ctx context.Context, in *PushLogsRequest, opts ...grpc.CallOption) (*PushLogsResponse, error) { cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) out := new(PushLogsResponse) @@ -57,9 +74,12 @@ func (c *telemetryServiceClient) PushLogs(ctx context.Context, in *PushLogsReque // All implementations must embed UnimplementedTelemetryServiceServer // for forward compatibility. // -// A service that receives structured logs from exporters and clients. +// A service that reverse-scrapes exporter metrics and receives structured logs. // Implemented by jumpstarter-telemetry; not part of the controller. type TelemetryServiceServer interface { + // Persistent bidirectional stream: telemetry sends scrape requests, + // exporter responds with full metric snapshots (OpenMetrics text). + MetricsStream(grpc.BidiStreamingServer[MetricsStreamRequest, MetricsStreamResponse]) error // Push structured log entries to the telemetry service for Loki ingest. PushLogs(context.Context, *PushLogsRequest) (*PushLogsResponse, error) mustEmbedUnimplementedTelemetryServiceServer() @@ -72,6 +92,9 @@ type TelemetryServiceServer interface { // pointer dereference when methods are called. type UnimplementedTelemetryServiceServer struct{} +func (UnimplementedTelemetryServiceServer) MetricsStream(grpc.BidiStreamingServer[MetricsStreamRequest, MetricsStreamResponse]) error { + return status.Error(codes.Unimplemented, "method MetricsStream not implemented") +} func (UnimplementedTelemetryServiceServer) PushLogs(context.Context, *PushLogsRequest) (*PushLogsResponse, error) { return nil, status.Error(codes.Unimplemented, "method PushLogs not implemented") } @@ -96,6 +119,13 @@ func RegisterTelemetryServiceServer(s grpc.ServiceRegistrar, srv TelemetryServic s.RegisterService(&TelemetryService_ServiceDesc, srv) } +func _TelemetryService_MetricsStream_Handler(srv interface{}, stream grpc.ServerStream) error { + return srv.(TelemetryServiceServer).MetricsStream(&grpc.GenericServerStream[MetricsStreamRequest, MetricsStreamResponse]{ServerStream: stream}) +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type TelemetryService_MetricsStreamServer = grpc.BidiStreamingServer[MetricsStreamRequest, MetricsStreamResponse] + func _TelemetryService_PushLogs_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(PushLogsRequest) if err := dec(in); err != nil { @@ -126,6 +156,13 @@ var TelemetryService_ServiceDesc = grpc.ServiceDesc{ Handler: _TelemetryService_PushLogs_Handler, }, }, - Streams: []grpc.StreamDesc{}, + Streams: []grpc.StreamDesc{ + { + StreamName: "MetricsStream", + Handler: _TelemetryService_MetricsStream_Handler, + ServerStreams: true, + ClientStreams: true, + }, + }, Metadata: "jumpstarter/v1/telemetry.proto", } diff --git a/controller/internal/service/metrics_merge.go b/controller/internal/service/metrics_merge.go new file mode 100644 index 000000000..378e54af3 --- /dev/null +++ b/controller/internal/service/metrics_merge.go @@ -0,0 +1,239 @@ +/* +Copyright 2026. + +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 service + +import ( + "bytes" + "errors" + "io" + "sort" + "strings" + + dto "github.com/prometheus/client_model/go" + "github.com/prometheus/common/expfmt" +) + +const ( + labelExporter = "exporter" + labelDriverType = "driver_type" + driverTypeOther = "other" + + // scrapeTimeoutsMetric is incremented once per exporter that does not + // answer a MetricsStream scrape within scrapeTimeout (JEP-0013). + scrapeTimeoutsMetric = "jumpstarter_scrape_timeouts_total" +) + +// DefaultDriverTypeEnum is the JEP-0013 default allowlist for driver_type. +var DefaultDriverTypeEnum = []string{ + "power", "storage", "network", "serial", "console", "video", "composite", +} + +// DefaultExemplarKeys is the JEP-0013 default exemplar allowlist. +var DefaultExemplarKeys = []string{"client", "lease_id"} + +type mergeConfig struct { + exporterName string + driverTypes map[string]struct{} + exemplarKeys map[string]struct{} +} + +type exporterSnapshot struct { + name string + text []byte +} + +func setToMap(values []string) map[string]struct{} { + out := make(map[string]struct{}, len(values)) + for _, v := range values { + if v == "" { + continue + } + out[v] = struct{}{} + } + return out +} + +func (s *TelemetryService) mergeConfigFor(name string) mergeConfig { + types := s.DriverTypeEnum + if len(types) == 0 { + types = DefaultDriverTypeEnum + } + keys := s.ExemplarKeys + if len(keys) == 0 { + keys = DefaultExemplarKeys + } + return mergeConfig{ + exporterName: name, + driverTypes: setToMap(types), + exemplarKeys: setToMap(keys), + } +} + +func parseMetricFamilies(text []byte) ([]*dto.MetricFamily, error) { + dec := expfmt.NewDecoder(bytes.NewReader(text), expfmt.NewFormat(expfmt.TypeOpenMetrics)) + var out []*dto.MetricFamily + for { + mf := new(dto.MetricFamily) + err := dec.Decode(mf) + if errors.Is(err, io.EOF) { + return out, nil + } + if err != nil { + return nil, err + } + out = append(out, mf) + } +} + +func encodeMetricFamilies(w io.Writer, families []*dto.MetricFamily) error { + sorted := append([]*dto.MetricFamily(nil), families...) + sort.Slice(sorted, func(i, j int) bool { + return sorted[i].GetName() < sorted[j].GetName() + }) + enc := expfmt.NewEncoder(w, expfmt.NewFormat(expfmt.TypeOpenMetrics)) + for _, f := range sorted { + if f == nil { + continue + } + if err := enc.Encode(f); err != nil { + return err + } + } + if closer, ok := enc.(expfmt.Closer); ok { + return closer.Close() + } + return nil +} + +func mergeSnapshots(snapshots []exporterSnapshot, extra []*dto.MetricFamily, cfgFor func(string) mergeConfig) []*dto.MetricFamily { + byName := map[string]*dto.MetricFamily{} + for _, f := range extra { + if f == nil || f.GetName() == "" { + continue + } + byName[f.GetName()] = f + } + for _, snap := range snapshots { + if len(snap.text) == 0 { + continue + } + families, err := parseMetricFamilies(snap.text) + if err != nil { + continue + } + cfg := cfgFor(snap.name) + for _, f := range families { + applyFamily(f, cfg) + existing, ok := byName[f.GetName()] + if !ok { + byName[f.GetName()] = f + continue + } + existing.Metric = append(existing.Metric, f.Metric...) + } + } + out := make([]*dto.MetricFamily, 0, len(byName)) + for _, f := range byName { + out = append(out, f) + } + return out +} + +func applyFamily(f *dto.MetricFamily, cfg mergeConfig) { + for _, m := range f.GetMetric() { + applyMetric(m, cfg) + } +} + +func applyMetric(m *dto.Metric, cfg mergeConfig) { + setLabel(m, labelExporter, cfg.exporterName) + if v, ok := getLabel(m, labelDriverType); ok { + setLabel(m, labelDriverType, remapDriverType(v, cfg.driverTypes)) + } + if m.Counter != nil { + m.Counter.Exemplar = filterExemplar(m.Counter.Exemplar, cfg.exemplarKeys) + } + if m.Histogram != nil { + for _, b := range m.Histogram.Bucket { + if b != nil { + b.Exemplar = filterExemplar(b.Exemplar, cfg.exemplarKeys) + } + } + if len(m.Histogram.Exemplars) > 0 { + filtered := make([]*dto.Exemplar, 0, len(m.Histogram.Exemplars)) + for _, ex := range m.Histogram.Exemplars { + if fe := filterExemplar(ex, cfg.exemplarKeys); fe != nil { + filtered = append(filtered, fe) + } + } + m.Histogram.Exemplars = filtered + } + } +} + +func remapDriverType(value string, allow map[string]struct{}) string { + if _, ok := allow[value]; ok { + return value + } + return driverTypeOther +} + +func getLabel(m *dto.Metric, name string) (string, bool) { + for _, lp := range m.GetLabel() { + if lp.GetName() == name { + return lp.GetValue(), true + } + } + return "", false +} + +func setLabel(m *dto.Metric, name, value string) { + for _, lp := range m.GetLabel() { + if lp.GetName() == name { + v := value + lp.Value = &v + return + } + } + m.Label = append(m.Label, labelPair(name, value)) +} + +func labelPair(name, value string) *dto.LabelPair { + n, v := name, value + return &dto.LabelPair{Name: &n, Value: &v} +} + +func filterExemplar(ex *dto.Exemplar, allow map[string]struct{}) *dto.Exemplar { + if ex == nil { + return nil + } + kept := make([]*dto.LabelPair, 0, len(ex.GetLabel())) + for _, lp := range ex.GetLabel() { + if _, ok := allow[lp.GetName()]; !ok { + continue + } + if strings.TrimSpace(lp.GetValue()) == "" { + continue + } + kept = append(kept, lp) + } + if len(kept) == 0 { + return nil + } + ex.Label = kept + return ex +} diff --git a/controller/internal/service/metrics_merge_test.go b/controller/internal/service/metrics_merge_test.go new file mode 100644 index 000000000..5f6dfe70a --- /dev/null +++ b/controller/internal/service/metrics_merge_test.go @@ -0,0 +1,172 @@ +/* +Copyright 2026. + +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 service + +import ( + "bytes" + "strings" + "testing" + + dto "github.com/prometheus/client_model/go" + "github.com/prometheus/common/expfmt" +) + +func TestApplyMetric_OverwritesExporterRemapsDriverTypeAndFiltersExemplars(t *testing.T) { + one := 1.0 + m := &dto.Metric{ + Label: []*dto.LabelPair{ + labelPair("exporter", "spoofed"), + labelPair("driver_type", "tuya"), + labelPair("operation", "on"), + }, + Counter: &dto.Counter{ + Value: &one, + Exemplar: &dto.Exemplar{ + Label: []*dto.LabelPair{ + labelPair("client", "ci"), + labelPair("lease_id", "abc"), + labelPair("trace_id", "drop-me"), + }, + }, + }, + } + applyMetric(m, mergeConfig{ + exporterName: "sidekick", + driverTypes: setToMap(DefaultDriverTypeEnum), + exemplarKeys: setToMap(DefaultExemplarKeys), + }) + + if got, _ := getLabel(m, "exporter"); got != "sidekick" { + t.Errorf("exporter label = %q, want sidekick", got) + } + if got, _ := getLabel(m, "driver_type"); got != "other" { + t.Errorf("driver_type = %q, want other", got) + } + ex := m.Counter.GetExemplar() + if ex == nil { + t.Fatal("expected exemplar to be kept") + } + got := map[string]string{} + for _, lp := range ex.GetLabel() { + got[lp.GetName()] = lp.GetValue() + } + if got["client"] != "ci" || got["lease_id"] != "abc" { + t.Errorf("exemplar labels = %v, want client=ci lease_id=abc", got) + } + if _, ok := got["trace_id"]; ok { + t.Errorf("trace_id should have been dropped, got %v", got) + } +} + +func TestApplyMetric_KeepsAllowlistedDriverType(t *testing.T) { + m := &dto.Metric{ + Label: []*dto.LabelPair{ + labelPair("driver_type", "power"), + }, + } + applyMetric(m, mergeConfig{ + exporterName: "exp-a", + driverTypes: setToMap(DefaultDriverTypeEnum), + exemplarKeys: setToMap(DefaultExemplarKeys), + }) + if got, _ := getLabel(m, "driver_type"); got != "power" { + t.Errorf("driver_type = %q, want power", got) + } + if got, _ := getLabel(m, "exporter"); got != "exp-a" { + t.Errorf("exporter = %q, want exp-a", got) + } +} + +func TestMergeSnapshots_CombinesExportersAndSkipsInvalidText(t *testing.T) { + cfg := func(name string) mergeConfig { + return mergeConfig{ + exporterName: name, + driverTypes: setToMap(DefaultDriverTypeEnum), + exemplarKeys: setToMap(DefaultExemplarKeys), + } + } + textA := []byte(`# TYPE jumpstarter_operations_total counter +# HELP jumpstarter_operations_total Total operations performed. +jumpstarter_operations_total{exporter="a",operation="on",result="success",driver_type="power"} 2.0 +# EOF +`) + textB := []byte(`# TYPE jumpstarter_operations_total counter +jumpstarter_operations_total{exporter="b",operation="off",result="success",driver_type="power"} 3.0 +# EOF +`) + families := mergeSnapshots([]exporterSnapshot{ + {name: "exp-a", text: textA}, + {name: "exp-b", text: textB}, + {name: "bad", text: []byte("not metrics")}, + }, nil, cfg) + + var ops *dto.MetricFamily + for _, f := range families { + if f.GetName() == "jumpstarter_operations_total" { + ops = f + break + } + } + if ops == nil { + t.Fatal("missing jumpstarter_operations_total") + } + if len(ops.Metric) != 2 { + t.Fatalf("got %d series, want 2 (invalid snapshot omitted)", len(ops.Metric)) + } + exporters := map[string]bool{} + for _, m := range ops.Metric { + name, _ := getLabel(m, "exporter") + exporters[name] = true + } + if !exporters["exp-a"] || !exporters["exp-b"] { + t.Errorf("exporters = %v, want exp-a and exp-b", exporters) + } +} + +func TestEncodeMetricFamilies_WritesOpenMetricsEOF(t *testing.T) { + name := "jumpstarter_scrape_timeouts_total" + help := "Exporter MetricsStream scrapes that exceeded scrapeTimeout." + metricType := dto.MetricType_COUNTER + zero := 0.0 + families := []*dto.MetricFamily{{ + Name: &name, + Help: &help, + Type: &metricType, + Metric: []*dto.Metric{{ + Counter: &dto.Counter{Value: &zero}, + }}, + }} + var buf bytes.Buffer + if err := encodeMetricFamilies(&buf, families); err != nil { + t.Fatalf("encode: %v", err) + } + body := buf.String() + if !strings.Contains(body, scrapeTimeoutsMetric) { + t.Errorf("encoded body missing %s:\n%s", scrapeTimeoutsMetric, body) + } + if !strings.Contains(body, "# EOF") { + t.Errorf("OpenMetrics body missing # EOF:\n%s", body) + } + dec := expfmt.NewDecoder(strings.NewReader(body), expfmt.NewFormat(expfmt.TypeOpenMetrics)) + mf := new(dto.MetricFamily) + if err := dec.Decode(mf); err != nil { + t.Fatalf("re-decode: %v", err) + } + if mf.GetName() != scrapeTimeoutsMetric { + t.Errorf("decoded name = %q", mf.GetName()) + } +} diff --git a/controller/internal/service/metrics_stream.go b/controller/internal/service/metrics_stream.go new file mode 100644 index 000000000..b16bb78ef --- /dev/null +++ b/controller/internal/service/metrics_stream.go @@ -0,0 +1,241 @@ +/* +Copyright 2026. + +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 service + +import ( + "context" + "errors" + "sync" + "time" + + pb "github.com/jumpstarter-dev/jumpstarter/controller/internal/protocol/jumpstarter/v1" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" + "sigs.k8s.io/controller-runtime/pkg/log" +) + +var errScrapeTimeout = errors.New("metrics scrape timeout") + +const defaultScrapeTimeout = 7 * time.Second + +// metricsConn is one authenticated MetricsStream. Scrapes are serialized +// because MetricsScrapeRequest has no scrape id (JEP-0013). +type metricsConn struct { + id exporterIdentity + scrapeMu sync.Mutex + mu sync.Mutex + send func(*pb.MetricsStreamResponse) error + // pending is non-nil only while a scrape is waiting for a response. + pending chan *pb.MetricsScrapeResponse + done <-chan struct{} + // started is closed when the Recv loop is running. + started chan struct{} +} + +func (s *TelemetryService) scrapeTimeoutDuration() time.Duration { + if s.ScrapeTimeout > 0 { + return s.ScrapeTimeout + } + return defaultScrapeTimeout +} + +func (s *TelemetryService) initMetricsState() { + s.stateMu.Lock() + defer s.stateMu.Unlock() + if s.conns == nil { + s.conns = make(map[string]*metricsConn) + } +} + +func (s *TelemetryService) registerConn(c *metricsConn) { + s.stateMu.Lock() + defer s.stateMu.Unlock() + if s.conns == nil { + s.conns = make(map[string]*metricsConn) + } + s.conns[c.id.key()] = c +} + +func (s *TelemetryService) unregisterConn(c *metricsConn) { + s.stateMu.Lock() + defer s.stateMu.Unlock() + if s.conns == nil { + return + } + if cur, ok := s.conns[c.id.key()]; ok && cur == c { + delete(s.conns, c.id.key()) + } +} + +func (s *TelemetryService) snapshotConns() []*metricsConn { + s.stateMu.Lock() + defer s.stateMu.Unlock() + out := make([]*metricsConn, 0, len(s.conns)) + for _, c := range s.conns { + out = append(out, c) + } + return out +} + +// MetricsStream reverse-scrapes one exporter. The first message must be +// MetricsRegister whose identity matches the bearer token. +func (s *TelemetryService) MetricsStream(stream pb.TelemetryService_MetricsStreamServer) error { + s.initMetricsState() + id, err := s.authenticateExporter(stream.Context()) + if err != nil { + return err + } + + first, err := stream.Recv() + if err != nil { + return err + } + reg := first.GetRegister() + if reg == nil { + return status.Error(codes.InvalidArgument, "first MetricsStream message must be MetricsRegister") + } + if reg.GetIdentity() != id.name { + return status.Errorf(codes.PermissionDenied, "register identity %q does not match authenticated exporter %q", reg.GetIdentity(), id.name) + } + + done := make(chan struct{}) + conn := &metricsConn{ + id: id, + send: stream.Send, + done: done, + started: make(chan struct{}), + } + s.registerConn(conn) + defer s.unregisterConn(conn) + + logger := log.FromContext(stream.Context()).WithName("telemetry").WithValues( + "exporter", id.name, + "namespace", id.namespace, + ) + logger.Info("MetricsStream registered") + + var closeDone sync.Once + finish := func() { closeDone.Do(func() { close(done) }) } + defer finish() + close(conn.started) + + for { + msg, err := stream.Recv() + if err != nil { + return err + } + resp := msg.GetScrapeResponse() + if resp == nil { + return status.Error(codes.InvalidArgument, "expected MetricsScrapeResponse after register") + } + conn.mu.Lock() + pending := conn.pending + conn.pending = nil + conn.mu.Unlock() + if pending == nil { + continue + } + select { + case pending <- resp: + case <-done: + return nil + case <-stream.Context().Done(): + return stream.Context().Err() + } + } +} + +func (c *metricsConn) scrape(ctx context.Context, timeout time.Duration) ([]byte, error) { + c.scrapeMu.Lock() + defer c.scrapeMu.Unlock() + + reply := make(chan *pb.MetricsScrapeResponse, 1) + c.mu.Lock() + c.pending = reply + c.mu.Unlock() + defer func() { + c.mu.Lock() + if c.pending == reply { + c.pending = nil + } + c.mu.Unlock() + }() + + err := c.send(&pb.MetricsStreamResponse{ + Msg: &pb.MetricsStreamResponse_ScrapeRequest{ + ScrapeRequest: &pb.MetricsScrapeRequest{}, + }, + }) + if err != nil { + return nil, err + } + + timer := time.NewTimer(timeout) + defer timer.Stop() + select { + case resp := <-reply: + if resp == nil { + return nil, errScrapeTimeout + } + return resp.GetMetricsText(), nil + case <-timer.C: + return nil, errScrapeTimeout + case <-ctx.Done(): + return nil, ctx.Err() + case <-c.done: + return nil, context.Canceled + } +} + +func (s *TelemetryService) fanoutScrapes(ctx context.Context) []exporterSnapshot { + s.initMetricsState() + s.initScrapeTimeouts() + timeout := s.scrapeTimeoutDuration() + conns := s.snapshotConns() + if len(conns) == 0 { + return nil + } + + ctx, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + + var mu sync.Mutex + var snaps []exporterSnapshot + var wg sync.WaitGroup + for _, c := range conns { + wg.Add(1) + go func(c *metricsConn) { + defer wg.Done() + text, err := c.scrape(ctx, timeout) + if err != nil { + if errors.Is(err, errScrapeTimeout) || errors.Is(err, context.DeadlineExceeded) { + s.scrapeTimeouts.Inc() + } + log.FromContext(ctx).WithName("telemetry").V(1).Info("exporter scrape omitted", + "exporter", c.id.name, + "error", err.Error(), + ) + return + } + mu.Lock() + snaps = append(snaps, exporterSnapshot{name: c.id.name, text: text}) + mu.Unlock() + }(c) + } + wg.Wait() + return snaps +} diff --git a/controller/internal/service/metrics_stream_test.go b/controller/internal/service/metrics_stream_test.go new file mode 100644 index 000000000..6f4a50cd0 --- /dev/null +++ b/controller/internal/service/metrics_stream_test.go @@ -0,0 +1,402 @@ +/* +Copyright 2026. + +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 service + +import ( + "context" + "errors" + "io" + "net" + "net/http" + "strings" + "testing" + "time" + + pb "github.com/jumpstarter-dev/jumpstarter/controller/internal/protocol/jumpstarter/v1" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/metadata" + "google.golang.org/grpc/status" + "google.golang.org/grpc/test/bufconn" +) + +const bufconnSize = 1024 * 1024 + +func startTestHub(t *testing.T, timeout time.Duration) (*TelemetryService, pb.TelemetryServiceClient, string) { + t.Helper() + signer := testSigner(t) + svc := &TelemetryService{ + Signer: signer, + MetricsBindAddr: "127.0.0.1:0", + ScrapeTimeout: timeout, + DriverTypeEnum: DefaultDriverTypeEnum, + ExemplarKeys: DefaultExemplarKeys, + } + svc.grpcReady.Store(true) + + shutdown, err := svc.startMetricsHTTP() + if err != nil { + t.Fatalf("startMetricsHTTP: %v", err) + } + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + _ = shutdown(ctx) + }) + + lis := bufconn.Listen(bufconnSize) + gs := grpc.NewServer() + pb.RegisterTelemetryServiceServer(gs, svc) + go func() { _ = gs.Serve(lis) }() + t.Cleanup(func() { + gs.Stop() + _ = lis.Close() + }) + + conn, err := grpc.NewClient("passthrough:///bufnet", + grpc.WithContextDialer(func(ctx context.Context, _ string) (net.Conn, error) { + return lis.DialContext(ctx) + }), + grpc.WithTransportCredentials(insecure.NewCredentials()), + ) + if err != nil { + t.Fatalf("grpc.NewClient: %v", err) + } + t.Cleanup(func() { _ = conn.Close() }) + + return svc, pb.NewTelemetryServiceClient(conn), svc.metricsAddr +} + +func exporterStreamCtx(t *testing.T, svc *TelemetryService, subject string) context.Context { + t.Helper() + token, err := svc.Signer.Token(subject) + if err != nil { + t.Fatalf("Token: %v", err) + } + return metadata.NewOutgoingContext( + context.Background(), + metadata.Pairs("authorization", "Bearer "+token), + ) +} + +func waitRegistered(t *testing.T, svc *TelemetryService, n int) { + t.Helper() + deadline := time.Now().Add(2 * time.Second) + for time.Now().Before(deadline) { + conns := svc.snapshotConns() + if len(conns) != n { + time.Sleep(10 * time.Millisecond) + continue + } + ready := true + for _, c := range conns { + if c.started == nil { + ready = false + break + } + select { + case <-c.started: + default: + ready = false + } + if !ready { + break + } + } + if ready { + return + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("timed out waiting for %d MetricsStream connections, have %d", n, len(svc.snapshotConns())) +} + +func httpGet(t *testing.T, url string) (int, string) { + t.Helper() + client := &http.Client{Timeout: 3 * time.Second} + var resp *http.Response + var lastErr error + for i := 0; i < 30; i++ { + resp, lastErr = client.Get(url) + if lastErr == nil { + break + } + time.Sleep(20 * time.Millisecond) + } + if lastErr != nil { + t.Fatalf("GET %s: %v", url, lastErr) + } + defer func() { _ = resp.Body.Close() }() + body, err := io.ReadAll(resp.Body) + if err != nil { + t.Fatalf("read body: %v", err) + } + return resp.StatusCode, string(body) +} + +func TestHealthzAndReadyz(t *testing.T) { + _, _, addr := startTestHub(t, 200*time.Millisecond) + code, body := httpGet(t, "http://"+addr+"/healthz") + if code != http.StatusOK || !strings.Contains(body, "ok") { + t.Fatalf("healthz status=%d body=%q", code, body) + } + code, body = httpGet(t, "http://"+addr+"/readyz") + if code != http.StatusOK || !strings.Contains(body, "ok") { + t.Fatalf("readyz status=%d body=%q", code, body) + } +} + +func TestReadyzNotReadyWhenGRPCDown(t *testing.T) { + svc := &TelemetryService{MetricsBindAddr: "127.0.0.1:0"} + shutdown, err := svc.startMetricsHTTP() + if err != nil { + t.Fatalf("startMetricsHTTP: %v", err) + } + t.Cleanup(func() { + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + _ = shutdown(ctx) + }) + code, _ := httpGet(t, "http://"+svc.metricsAddr+"/readyz") + if code != http.StatusServiceUnavailable { + t.Fatalf("readyz status=%d, want 503", code) + } +} + +func TestMetricsStream_RejectsIdentityMismatch(t *testing.T) { + svc, client, _ := startTestHub(t, time.Second) + ctx := exporterStreamCtx(t, svc, "exporter:jumpstarter:alice:uid1") + stream, err := client.MetricsStream(ctx) + if err != nil { + t.Fatalf("MetricsStream: %v", err) + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_Register{ + Register: &pb.MetricsRegister{Identity: "bob"}, + }, + }); err != nil { + t.Fatalf("Send register: %v", err) + } + _, err = stream.Recv() + if err == nil { + t.Fatal("expected PermissionDenied, got nil") + } + if status.Code(err) != codes.PermissionDenied { + t.Fatalf("code = %v, want PermissionDenied (%v)", status.Code(err), err) + } +} + +func TestMetricsStream_RejectsNonRegisterFirstMessage(t *testing.T) { + svc, client, _ := startTestHub(t, time.Second) + ctx := exporterStreamCtx(t, svc, "exporter:jumpstarter:alice:uid1") + stream, err := client.MetricsStream(ctx) + if err != nil { + t.Fatalf("MetricsStream: %v", err) + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_ScrapeResponse{ + ScrapeResponse: &pb.MetricsScrapeResponse{MetricsText: []byte("# EOF\n")}, + }, + }); err != nil { + t.Fatalf("Send: %v", err) + } + _, err = stream.Recv() + if status.Code(err) != codes.InvalidArgument { + t.Fatalf("code = %v, want InvalidArgument (%v)", status.Code(err), err) + } +} + +func TestMetricsStream_RejectsNonExporterToken(t *testing.T) { + svc, client, _ := startTestHub(t, time.Second) + ctx := exporterStreamCtx(t, svc, "client:jumpstarter:ci:uid1") + stream, err := client.MetricsStream(ctx) + if err != nil { + t.Fatalf("MetricsStream: %v", err) + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_Register{ + Register: &pb.MetricsRegister{Identity: "ci"}, + }, + }); err != nil { + t.Fatalf("Send: %v", err) + } + _, err = stream.Recv() + if status.Code(err) != codes.PermissionDenied { + t.Fatalf("code = %v, want PermissionDenied (%v)", status.Code(err), err) + } +} + +func TestMetricsConn_ScrapeTimeout(t *testing.T) { + done := make(chan struct{}) + c := &metricsConn{ + send: func(*pb.MetricsStreamResponse) error { return nil }, + done: done, + } + _, err := c.scrape(context.Background(), 40*time.Millisecond) + if !errors.Is(err, errScrapeTimeout) { + t.Fatalf("err = %v, want errScrapeTimeout", err) + } +} + +func TestFanout_TimeoutIncrementsCounterWithoutGRPC(t *testing.T) { + svc := &TelemetryService{ScrapeTimeout: 40 * time.Millisecond} + svc.initScrapeTimeouts() + done := make(chan struct{}) + c := &metricsConn{ + id: exporterIdentity{namespace: "jumpstarter", name: "slow"}, + send: func(*pb.MetricsStreamResponse) error { return nil }, + done: done, + } + svc.registerConn(c) + snaps := svc.fanoutScrapes(context.Background()) + if len(snaps) != 0 { + t.Fatalf("timed-out scrape must be omitted, got %d snapshots", len(snaps)) + } + mfs, err := svc.metricsRegistry.Gather() + if err != nil { + t.Fatalf("Gather: %v", err) + } + var value float64 + var found bool + for _, mf := range mfs { + if mf.GetName() == scrapeTimeoutsMetric || mf.GetName()+"_total" == scrapeTimeoutsMetric { + found = true + if len(mf.Metric) > 0 { + value = mf.Metric[0].GetCounter().GetValue() + } + } + } + if !found || value < 1 { + t.Fatalf("scrape timeout counter = %v found=%v, want >= 1", value, found) + } +} + +func TestFanout_TimeoutOmitsExporterAndIncrementsCounter(t *testing.T) { + svc, client, addr := startTestHub(t, 150*time.Millisecond) + ctx := exporterStreamCtx(t, svc, "exporter:jumpstarter:slow:uid1") + stream, err := client.MetricsStream(ctx) + if err != nil { + t.Fatalf("MetricsStream: %v", err) + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_Register{ + Register: &pb.MetricsRegister{Identity: "slow"}, + }, + }); err != nil { + t.Fatalf("Send register: %v", err) + } + // Drain scrape requests so Send on the server does not block, but never reply. + go func() { + for { + if _, err := stream.Recv(); err != nil { + return + } + } + }() + waitRegistered(t, svc, 1) + + code, body := httpGet(t, "http://"+addr+"/metrics") + if code != http.StatusOK { + t.Fatalf("GET /metrics status=%d body=%s", code, body) + } + if strings.Contains(body, "jumpstarter_operations_total") { + t.Errorf("timed-out exporter metrics must be omitted, body:\n%s", body) + } + mfs, err := svc.metricsRegistry.Gather() + if err != nil { + t.Fatalf("Gather: %v", err) + } + var value float64 + var found bool + for _, mf := range mfs { + if mf.GetName() == scrapeTimeoutsMetric || mf.GetName()+"_total" == scrapeTimeoutsMetric { + found = true + if len(mf.Metric) > 0 { + value = mf.Metric[0].GetCounter().GetValue() + } + } + } + if !found || value < 1 { + t.Fatalf("scrape timeout counter = %v found=%v body:\n%s", value, found, body) + } +} + +func TestFanout_MergesOpenMetricsFromConnectedExporter(t *testing.T) { + svc, client, addr := startTestHub(t, time.Second) + ctx := exporterStreamCtx(t, svc, "exporter:jumpstarter:sidekick:uid1") + stream, err := client.MetricsStream(ctx) + if err != nil { + t.Fatalf("MetricsStream: %v", err) + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_Register{ + Register: &pb.MetricsRegister{Identity: "sidekick"}, + }, + }); err != nil { + t.Fatalf("Send register: %v", err) + } + + snapshot := []byte(`# TYPE jumpstarter_operations_total counter +# HELP jumpstarter_operations_total Total operations performed. +jumpstarter_operations_total{exporter="spoofed",operation="on",result="success",driver_type="tuya"} 4.0 +# EOF +`) + errCh := make(chan error, 1) + go func() { + for { + msg, err := stream.Recv() + if err != nil { + errCh <- err + return + } + if msg.GetScrapeRequest() == nil { + continue + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_ScrapeResponse{ + ScrapeResponse: &pb.MetricsScrapeResponse{MetricsText: snapshot}, + }, + }); err != nil { + errCh <- err + return + } + } + }() + waitRegistered(t, svc, 1) + + code, body := httpGet(t, "http://"+addr+"/metrics") + if code != http.StatusOK { + t.Fatalf("GET /metrics status=%d body=%s", code, body) + } + if !strings.Contains(body, `exporter="sidekick"`) { + t.Errorf("authenticated exporter label missing:\n%s", body) + } + if strings.Contains(body, `exporter="spoofed"`) { + t.Errorf("spoofed exporter label must be overwritten:\n%s", body) + } + if !strings.Contains(body, `driver_type="other"`) { + t.Errorf("unknown driver_type must remap to other:\n%s", body) + } + select { + case err := <-errCh: + if err != nil && err != io.EOF && status.Code(err) != codes.Canceled && status.Code(err) != codes.Unavailable { + t.Fatalf("stream goroutine: %v", err) + } + default: + } +} diff --git a/controller/internal/service/telemetry_http.go b/controller/internal/service/telemetry_http.go new file mode 100644 index 000000000..6d12daff4 --- /dev/null +++ b/controller/internal/service/telemetry_http.go @@ -0,0 +1,135 @@ +/* +Copyright 2026. + +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 service + +import ( + "bytes" + "context" + "net" + "net/http" + "time" + + "github.com/prometheus/client_golang/prometheus" + dto "github.com/prometheus/client_model/go" + "github.com/prometheus/common/expfmt" + ctrl "sigs.k8s.io/controller-runtime" +) + +func (s *TelemetryService) initScrapeTimeouts() { + s.stateMu.Lock() + defer s.stateMu.Unlock() + if s.metricsRegistry != nil { + return + } + s.metricsRegistry = prometheus.NewRegistry() + s.scrapeTimeouts = prometheus.NewCounter(prometheus.CounterOpts{ + Name: scrapeTimeoutsMetric, + Help: "Exporter MetricsStream scrapes that exceeded scrapeTimeout.", + }) + s.metricsRegistry.MustRegister(s.scrapeTimeouts) +} + +func metricsHTTPEnabled(addr string) bool { + return addr != "" && addr != "0" +} + +func (s *TelemetryService) startMetricsHTTP() (func(context.Context) error, error) { + addr := s.MetricsBindAddr + if !metricsHTTPEnabled(addr) { + return nil, nil + } + s.initMetricsState() + s.initScrapeTimeouts() + + ln, err := net.Listen("tcp", addr) + if err != nil { + return nil, err + } + s.metricsAddr = ln.Addr().String() + + mux := http.NewServeMux() + mux.HandleFunc("/metrics", s.handleMetrics) + mux.HandleFunc("/healthz", s.handleHealthz) + mux.HandleFunc("/readyz", s.handleReadyz) + + writeTimeout := s.scrapeTimeoutDuration() + 15*time.Second + if writeTimeout < 30*time.Second { + writeTimeout = 30 * time.Second + } + + srv := &http.Server{ + Handler: mux, + ReadHeaderTimeout: 10 * time.Second, + ReadTimeout: 30 * time.Second, + WriteTimeout: writeTimeout, + IdleTimeout: 5 * time.Minute, + } + go func() { + if err := srv.Serve(ln); err != nil && err != http.ErrServerClosed { + ctrl.Log.WithName("telemetry").Error(err, "metrics HTTP server stopped unexpectedly") + } + }() + ctrl.Log.WithName("telemetry").Info("Telemetry metrics HTTP listening", + "addr", s.metricsAddr, + ) + return srv.Shutdown, nil +} + +func (s *TelemetryService) handleHealthz(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte("ok\n")) +} + +// handleReadyz is 200 once the gRPC listener is bound. +// JEP DD-7 also gates on Loki reachability and a connected exporter; Loki is +// Phase 3 PR D, and requiring an exporter would keep an empty lab unready +// (Prometheus could not scrape jumpstarter_scrape_timeouts_total). +func (s *TelemetryService) handleReadyz(w http.ResponseWriter, _ *http.Request) { + if !s.grpcReady.Load() { + http.Error(w, "grpc not ready\n", http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte("ok\n")) +} + +func (s *TelemetryService) handleMetrics(w http.ResponseWriter, r *http.Request) { + s.initScrapeTimeouts() + snaps := s.fanoutScrapes(r.Context()) + + var extra []*dto.MetricFamily + if s.metricsRegistry != nil { + gathered, err := s.metricsRegistry.Gather() + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + extra = gathered + } + + families := mergeSnapshots(snaps, extra, s.mergeConfigFor) + + var buf bytes.Buffer + if err := encodeMetricFamilies(&buf, families); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + + w.Header().Set("Content-Type", string(expfmt.NewFormat(expfmt.TypeOpenMetrics))) + w.WriteHeader(http.StatusOK) + _, _ = w.Write(buf.Bytes()) +} diff --git a/controller/internal/service/telemetry_identity.go b/controller/internal/service/telemetry_identity.go new file mode 100644 index 000000000..e1f23de4e --- /dev/null +++ b/controller/internal/service/telemetry_identity.go @@ -0,0 +1,62 @@ +/* +Copyright 2026. + +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 service + +import ( + "context" + "strings" + + "github.com/jumpstarter-dev/jumpstarter/controller/internal/authentication" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// exporterIdentity is the exporter CRD namespace/name claimed by a bearer token. +type exporterIdentity struct { + namespace string + name string +} + +func (id exporterIdentity) key() string { + return id.namespace + "/" + id.name +} + +func parseExporterSubject(subject string) (exporterIdentity, error) { + parts := strings.SplitN(subject, ":", 4) + if len(parts) != 4 || parts[0] != "exporter" { + return exporterIdentity{}, status.Errorf(codes.PermissionDenied, "token is not an exporter token") + } + if parts[1] == "" || parts[2] == "" { + return exporterIdentity{}, status.Errorf(codes.PermissionDenied, "token has incomplete exporter identity") + } + return exporterIdentity{namespace: parts[1], name: parts[2]}, nil +} + +func (s *TelemetryService) authenticateExporter(ctx context.Context) (exporterIdentity, error) { + token, err := authentication.BearerTokenFromContext(ctx) + if err != nil { + return exporterIdentity{}, err + } + if s.Signer == nil { + return exporterIdentity{}, status.Error(codes.Internal, "telemetry signer is not configured") + } + subject, err := s.Signer.ParseSubject(token) + if err != nil { + return exporterIdentity{}, status.Errorf(codes.Unauthenticated, "invalid token: %v", err) + } + return parseExporterSubject(subject) +} diff --git a/controller/internal/service/telemetry_service.go b/controller/internal/service/telemetry_service.go index 80844b807..e3c60bb77 100644 --- a/controller/internal/service/telemetry_service.go +++ b/controller/internal/service/telemetry_service.go @@ -22,17 +22,17 @@ import ( "fmt" "net" "strings" + "sync" + "sync/atomic" "time" "github.com/grpc-ecosystem/go-grpc-middleware/v2/interceptors/recovery" - "github.com/jumpstarter-dev/jumpstarter/controller/internal/authentication" "github.com/jumpstarter-dev/jumpstarter/controller/internal/oidc" pb "github.com/jumpstarter-dev/jumpstarter/controller/internal/protocol/jumpstarter/v1" + "github.com/prometheus/client_golang/prometheus" "google.golang.org/grpc" - "google.golang.org/grpc/codes" "google.golang.org/grpc/credentials" "google.golang.org/grpc/reflection" - "google.golang.org/grpc/status" ctrl "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/log" ) @@ -69,42 +69,47 @@ var reservedExtraFieldKeys = map[string]struct{}{ type TelemetryService struct { pb.UnimplementedTelemetryServiceServer - // BindAddr is the TCP address to listen on (e.g. ":9093"). + // BindAddr is the TCP address to listen on for gRPC (e.g. ":9093"). BindAddr string - // Signer is used to validate bearer tokens on every PushLogs call. - // Tokens are issued by the controller from the same CONTROLLER_KEY seed, - // so the telemetry binary can verify them locally without a k8s client. + // MetricsBindAddr is the TCP address for HTTP GET /metrics, /healthz, and + // /readyz. Empty or "0" disables the HTTP server (tests). Production + // default is :8080 — a dedicated port, not multiplexed onto gRPC :9093. + MetricsBindAddr string + + // ScrapeTimeout is the fan-out wait for MetricsStream responses (JEP default 7s). + ScrapeTimeout time.Duration + + // DriverTypeEnum allowlist; unknown driver_type values are remapped to "other". + DriverTypeEnum []string + + // ExemplarKeys allowlist applied on the telemetry merge path. + ExemplarKeys []string + + // Signer is used to validate bearer tokens on every PushLogs call and + // MetricsStream. Tokens are issued by the controller from the same + // CONTROLLER_KEY seed, so the telemetry binary can verify them locally + // without a k8s client. Signer *oidc.Signer + + stateMu sync.Mutex + conns map[string]*metricsConn + scrapeTimeouts prometheus.Counter + metricsRegistry *prometheus.Registry + metricsAddr string + grpcReady atomic.Bool } // PushLogs receives a batch of structured log entries and writes them via the // controller-runtime logger (structured JSON to stdout). // Future phase: forward to Loki push API. func (s *TelemetryService) PushLogs(ctx context.Context, req *pb.PushLogsRequest) (*pb.PushLogsResponse, error) { - token, err := authentication.BearerTokenFromContext(ctx) + id, err := s.authenticateExporter(ctx) if err != nil { return nil, err } - - // Validate token and extract the subject (format: exporter:namespace:name:uid). - subject, err := s.Signer.ParseSubject(token) - if err != nil { - return nil, status.Errorf(codes.Unauthenticated, "invalid token: %v", err) - } - - // Only exporter tokens are allowed to push logs. Any other validly-signed - // token (e.g. a client token) is rejected immediately so that the identity - // checks below always have a non-empty claimedName/claimedNamespace. - parts := strings.SplitN(subject, ":", 4) - if len(parts) != 4 || parts[0] != "exporter" { - return nil, status.Errorf(codes.PermissionDenied, "token is not an exporter token") - } - claimedNamespace := parts[1] - claimedName := parts[2] - if claimedNamespace == "" || claimedName == "" { - return nil, status.Errorf(codes.PermissionDenied, "token has incomplete exporter identity") - } + claimedNamespace := id.namespace + claimedName := id.name // Use context-based logger so tests can inject their own via logf.IntoContext. logger := log.FromContext(ctx).WithName("telemetry") @@ -260,9 +265,21 @@ func (s *TelemetryService) Start(ctx context.Context) error { return fmt.Errorf("telemetry: listen %s: %w", s.BindAddr, err) } + s.initMetricsState() + s.initScrapeTimeouts() + s.grpcReady.Store(true) + + httpShutdown, err := s.startMetricsHTTP() + if err != nil { + s.grpcReady.Store(false) + _ = lis.Close() + return fmt.Errorf("telemetry: metrics HTTP listen %s: %w", s.MetricsBindAddr, err) + } + srv := grpc.NewServer( grpc.Creds(creds), grpc.ChainUnaryInterceptor(recovery.UnaryServerInterceptor()), + grpc.ChainStreamInterceptor(recovery.StreamServerInterceptor()), ) pb.RegisterTelemetryServiceServer(srv, s) reflection.Register(srv) @@ -276,12 +293,24 @@ func (s *TelemetryService) Start(ctx context.Context) error { select { case <-ctx.Done(): + s.grpcReady.Store(false) + if httpShutdown != nil { + shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + _ = httpShutdown(shutdownCtx) + } srv.GracefulStop() if err := <-errCh; err != nil && !errors.Is(err, grpc.ErrServerStopped) { return err } return nil case err := <-errCh: + s.grpcReady.Store(false) + if httpShutdown != nil { + shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + _ = httpShutdown(shutdownCtx) + } srv.Stop() return err } diff --git a/protocol/proto/jumpstarter/v1/telemetry.proto b/protocol/proto/jumpstarter/v1/telemetry.proto index 200748201..fbbfd00ee 100644 --- a/protocol/proto/jumpstarter/v1/telemetry.proto +++ b/protocol/proto/jumpstarter/v1/telemetry.proto @@ -6,13 +6,44 @@ package jumpstarter.v1; import "google/protobuf/timestamp.proto"; -// A service that receives structured logs from exporters and clients. +// A service that reverse-scrapes exporter metrics and receives structured logs. // Implemented by jumpstarter-telemetry; not part of the controller. service TelemetryService { + // Persistent bidirectional stream: telemetry sends scrape requests, + // exporter responds with full metric snapshots (OpenMetrics text). + rpc MetricsStream(stream MetricsStreamRequest) returns (stream MetricsStreamResponse); + // Push structured log entries to the telemetry service for Loki ingest. rpc PushLogs(PushLogsRequest) returns (PushLogsResponse); } +// Exporter → Telemetry +message MetricsStreamRequest { + oneof msg { + MetricsRegister register = 1; // First message: identify this exporter. + MetricsScrapeResponse scrape_response = 2; // Subsequent: reply to a scrape. + } +} + +message MetricsRegister { + string identity = 1; // Exporter CRD name (verified against the auth token by the server). +} + +message MetricsScrapeResponse { + bytes metrics_text = 1; // generate_latest() OpenMetrics output. + google.protobuf.Timestamp timestamp = 2; +} + +// Telemetry → Exporter +message MetricsStreamResponse { + oneof msg { + MetricsScrapeRequest scrape_request = 1; + } +} + +// Empty request: "send your /metrics now". +message MetricsScrapeRequest {} + // Request to push log entries to the telemetry service. message PushLogsRequest { repeated LogEntry entries = 1; // Log entries to push. diff --git a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.py b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.py index 30983ccc3..80e4bc2bc 100644 --- a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.py +++ b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.py @@ -25,24 +25,34 @@ from google.protobuf import timestamp_pb2 as google_dot_protobuf_dot_timestamp__pb2 -DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x1ejumpstarter/v1/telemetry.proto\x12\x0ejumpstarter.v1\x1a\x1fgoogle/protobuf/timestamp.proto"E\n\x0fPushLogsRequest\x122\n\x07entries\x18\x01 \x03(\x0b2\x18.jumpstarter.v1.LogEntryR\x07entries"H\n\x10PushLogsResponse\x12\x1a\n\x08accepted\x18\x01 \x01(\rR\x08accepted\x12\x18\n\x07dropped\x18\x02 \x01(\rR\x07dropped"\xe5\x03\n\x08LogEntry\x128\n\ttimestamp\x18\x01 \x01(\x0b2\x1a.google.protobuf.TimestampR\ttimestamp\x12\x1a\n\x08severity\x18\x02 \x01(\tR\x08severity\x12\x18\n\x07message\x18\x03 \x01(\tR\x07message\x12\x1c\n\tcomponent\x18\x04 \x01(\tR\tcomponent\x12\x1a\n\x08exporter\x18\x05 \x01(\tR\x08exporter\x12\x14\n\x05lease\x18\x06 \x01(\tR\x05lease\x12\x16\n\x06client\x18\x07 \x01(\tR\x06client\x12\x1c\n\toperation\x18\x08 \x01(\tR\toperation\x12\x16\n\x06result\x18\t \x01(\tR\x06result\x12\x1f\n\x0bdriver_type\x18\n \x01(\tR\ndriverType\x12L\n\x0cextra_fields\x18\x0b \x03(\x0b2).jumpstarter.v1.LogEntry.ExtraFieldsEntryR\x0bextraFields\x12\x1c\n\tnamespace\x18\x0c \x01(\tR\tnamespace\x1a>\n\x10ExtraFieldsEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x028\x012a\n\x10TelemetryService\x12M\n\x08PushLogs\x12\x1f.jumpstarter.v1.PushLogsRequest\x1a .jumpstarter.v1.PushLogsResponseB\xd1\x01\n\x12com.jumpstarter.v1B\x0eTelemetryProtoP\x01ZRgithub.com/jumpstarter-dev/jumpstarter/controller/internal/protocol/jumpstarter/v1\xa2\x02\x03JXX\xaa\x02\x0eJumpstarter.V1\xca\x02\x0eJumpstarter\\V1\xe2\x02\x1aJumpstarter\\V1\\GPBMetadata\xea\x02\x0fJumpstarter::V1b\x06proto3') +DESCRIPTOR = _descriptor_pool.Default().AddSerializedFile(b'\n\x1ejumpstarter/v1/telemetry.proto\x12\x0ejumpstarter.v1\x1a\x1fgoogle/protobuf/timestamp.proto\"\xae\x01\n\x14MetricsStreamRequest\x12=\n\x08register\x18\x01 \x01(\x0b\x32\x1f.jumpstarter.v1.MetricsRegisterH\x00R\x08register\x12P\n\x0fscrape_response\x18\x02 \x01(\x0b\x32%.jumpstarter.v1.MetricsScrapeResponseH\x00R\x0escrapeResponseB\x05\n\x03msg\"-\n\x0fMetricsRegister\x12\x1a\n\x08identity\x18\x01 \x01(\tR\x08identity\"t\n\x15MetricsScrapeResponse\x12!\n\x0cmetrics_text\x18\x01 \x01(\x0cR\x0bmetricsText\x12\x38\n\ttimestamp\x18\x02 \x01(\x0b\x32\x1a.google.protobuf.TimestampR\ttimestamp\"m\n\x15MetricsStreamResponse\x12M\n\x0escrape_request\x18\x01 \x01(\x0b\x32$.jumpstarter.v1.MetricsScrapeRequestH\x00R\rscrapeRequestB\x05\n\x03msg\"\x16\n\x14MetricsScrapeRequest\"E\n\x0fPushLogsRequest\x12\x32\n\x07\x65ntries\x18\x01 \x03(\x0b\x32\x18.jumpstarter.v1.LogEntryR\x07\x65ntries\"H\n\x10PushLogsResponse\x12\x1a\n\x08\x61\x63\x63\x65pted\x18\x01 \x01(\rR\x08\x61\x63\x63\x65pted\x12\x18\n\x07\x64ropped\x18\x02 \x01(\rR\x07\x64ropped\"\xe5\x03\n\x08LogEntry\x12\x38\n\ttimestamp\x18\x01 \x01(\x0b\x32\x1a.google.protobuf.TimestampR\ttimestamp\x12\x1a\n\x08severity\x18\x02 \x01(\tR\x08severity\x12\x18\n\x07message\x18\x03 \x01(\tR\x07message\x12\x1c\n\tcomponent\x18\x04 \x01(\tR\tcomponent\x12\x1a\n\x08\x65xporter\x18\x05 \x01(\tR\x08\x65xporter\x12\x14\n\x05lease\x18\x06 \x01(\tR\x05lease\x12\x16\n\x06\x63lient\x18\x07 \x01(\tR\x06\x63lient\x12\x1c\n\toperation\x18\x08 \x01(\tR\toperation\x12\x16\n\x06result\x18\t \x01(\tR\x06result\x12\x1f\n\x0b\x64river_type\x18\n \x01(\tR\ndriverType\x12L\n\x0c\x65xtra_fields\x18\x0b \x03(\x0b\x32).jumpstarter.v1.LogEntry.ExtraFieldsEntryR\x0b\x65xtraFields\x12\x1c\n\tnamespace\x18\x0c \x01(\tR\tnamespace\x1a>\n\x10\x45xtraFieldsEntry\x12\x10\n\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n\x05value\x18\x02 \x01(\tR\x05value:\x02\x38\x01\x32\xc3\x01\n\x10TelemetryService\x12`\n\rMetricsStream\x12$.jumpstarter.v1.MetricsStreamRequest\x1a%.jumpstarter.v1.MetricsStreamResponse(\x01\x30\x01\x12M\n\x08PushLogs\x12\x1f.jumpstarter.v1.PushLogsRequest\x1a .jumpstarter.v1.PushLogsResponseB}\n\x12\x63om.jumpstarter.v1B\x0eTelemetryProtoP\x01\xa2\x02\x03JXX\xaa\x02\x0eJumpstarter.V1\xca\x02\x0eJumpstarter\\V1\xe2\x02\x1aJumpstarter\\V1\\GPBMetadata\xea\x02\x0fJumpstarter::V1b\x06proto3') _globals = globals() _builder.BuildMessageAndEnumDescriptors(DESCRIPTOR, _globals) _builder.BuildTopDescriptorsAndMessages(DESCRIPTOR, 'jumpstarter.v1.telemetry_pb2', _globals) if not _descriptor._USE_C_DESCRIPTORS: _globals['DESCRIPTOR']._loaded_options = None - _globals['DESCRIPTOR']._serialized_options = b'\n\022com.jumpstarter.v1B\016TelemetryProtoP\001ZRgithub.com/jumpstarter-dev/jumpstarter/controller/internal/protocol/jumpstarter/v1\242\002\003JXX\252\002\016Jumpstarter.V1\312\002\016Jumpstarter\\V1\342\002\032Jumpstarter\\V1\\GPBMetadata\352\002\017Jumpstarter::V1' + _globals['DESCRIPTOR']._serialized_options = b'\n\022com.jumpstarter.v1B\016TelemetryProtoP\001\242\002\003JXX\252\002\016Jumpstarter.V1\312\002\016Jumpstarter\\V1\342\002\032Jumpstarter\\V1\\GPBMetadata\352\002\017Jumpstarter::V1' _globals['_LOGENTRY_EXTRAFIELDSENTRY']._loaded_options = None _globals['_LOGENTRY_EXTRAFIELDSENTRY']._serialized_options = b'8\001' - _globals['_PUSHLOGSREQUEST']._serialized_start=83 - _globals['_PUSHLOGSREQUEST']._serialized_end=152 - _globals['_PUSHLOGSRESPONSE']._serialized_start=154 - _globals['_PUSHLOGSRESPONSE']._serialized_end=226 - _globals['_LOGENTRY']._serialized_start=229 - _globals['_LOGENTRY']._serialized_end=714 - _globals['_LOGENTRY_EXTRAFIELDSENTRY']._serialized_start=652 - _globals['_LOGENTRY_EXTRAFIELDSENTRY']._serialized_end=714 - _globals['_TELEMETRYSERVICE']._serialized_start=716 - _globals['_TELEMETRYSERVICE']._serialized_end=813 + _globals['_METRICSSTREAMREQUEST']._serialized_start=84 + _globals['_METRICSSTREAMREQUEST']._serialized_end=258 + _globals['_METRICSREGISTER']._serialized_start=260 + _globals['_METRICSREGISTER']._serialized_end=305 + _globals['_METRICSSCRAPERESPONSE']._serialized_start=307 + _globals['_METRICSSCRAPERESPONSE']._serialized_end=423 + _globals['_METRICSSTREAMRESPONSE']._serialized_start=425 + _globals['_METRICSSTREAMRESPONSE']._serialized_end=534 + _globals['_METRICSSCRAPEREQUEST']._serialized_start=536 + _globals['_METRICSSCRAPEREQUEST']._serialized_end=558 + _globals['_PUSHLOGSREQUEST']._serialized_start=560 + _globals['_PUSHLOGSREQUEST']._serialized_end=629 + _globals['_PUSHLOGSRESPONSE']._serialized_start=631 + _globals['_PUSHLOGSRESPONSE']._serialized_end=703 + _globals['_LOGENTRY']._serialized_start=706 + _globals['_LOGENTRY']._serialized_end=1191 + _globals['_LOGENTRY_EXTRAFIELDSENTRY']._serialized_start=1129 + _globals['_LOGENTRY_EXTRAFIELDSENTRY']._serialized_end=1191 + _globals['_TELEMETRYSERVICE']._serialized_start=1194 + _globals['_TELEMETRYSERVICE']._serialized_end=1389 # @@protoc_insertion_point(module_scope) diff --git a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.pyi b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.pyi index 2f28d9f24..968f7ea75 100644 --- a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.pyi +++ b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2.pyi @@ -19,6 +19,103 @@ else: DESCRIPTOR: google.protobuf.descriptor.FileDescriptor +@typing.final +class MetricsStreamRequest(google.protobuf.message.Message): + """Exporter → Telemetry""" + + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + REGISTER_FIELD_NUMBER: builtins.int + SCRAPE_RESPONSE_FIELD_NUMBER: builtins.int + @property + def register(self) -> Global___MetricsRegister: + """First message: identify this exporter.""" + + @property + def scrape_response(self) -> Global___MetricsScrapeResponse: + """Subsequent: reply to a scrape.""" + + def __init__( + self, + *, + register: Global___MetricsRegister | None = ..., + scrape_response: Global___MetricsScrapeResponse | None = ..., + ) -> None: ... + def HasField(self, field_name: typing.Literal["msg", b"msg", "register", b"register", "scrape_response", b"scrape_response"]) -> builtins.bool: ... + def ClearField(self, field_name: typing.Literal["msg", b"msg", "register", b"register", "scrape_response", b"scrape_response"]) -> None: ... + def WhichOneof(self, oneof_group: typing.Literal["msg", b"msg"]) -> typing.Literal["register", "scrape_response"] | None: ... + +Global___MetricsStreamRequest: typing_extensions.TypeAlias = MetricsStreamRequest + +@typing.final +class MetricsRegister(google.protobuf.message.Message): + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + IDENTITY_FIELD_NUMBER: builtins.int + identity: builtins.str + """Exporter CRD name (verified against the auth token by the server).""" + def __init__( + self, + *, + identity: builtins.str = ..., + ) -> None: ... + def ClearField(self, field_name: typing.Literal["identity", b"identity"]) -> None: ... + +Global___MetricsRegister: typing_extensions.TypeAlias = MetricsRegister + +@typing.final +class MetricsScrapeResponse(google.protobuf.message.Message): + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + METRICS_TEXT_FIELD_NUMBER: builtins.int + TIMESTAMP_FIELD_NUMBER: builtins.int + metrics_text: builtins.bytes + """generate_latest() OpenMetrics output.""" + @property + def timestamp(self) -> google.protobuf.timestamp_pb2.Timestamp: ... + def __init__( + self, + *, + metrics_text: builtins.bytes = ..., + timestamp: google.protobuf.timestamp_pb2.Timestamp | None = ..., + ) -> None: ... + def HasField(self, field_name: typing.Literal["timestamp", b"timestamp"]) -> builtins.bool: ... + def ClearField(self, field_name: typing.Literal["metrics_text", b"metrics_text", "timestamp", b"timestamp"]) -> None: ... + +Global___MetricsScrapeResponse: typing_extensions.TypeAlias = MetricsScrapeResponse + +@typing.final +class MetricsStreamResponse(google.protobuf.message.Message): + """Telemetry → Exporter""" + + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + SCRAPE_REQUEST_FIELD_NUMBER: builtins.int + @property + def scrape_request(self) -> Global___MetricsScrapeRequest: ... + def __init__( + self, + *, + scrape_request: Global___MetricsScrapeRequest | None = ..., + ) -> None: ... + def HasField(self, field_name: typing.Literal["msg", b"msg", "scrape_request", b"scrape_request"]) -> builtins.bool: ... + def ClearField(self, field_name: typing.Literal["msg", b"msg", "scrape_request", b"scrape_request"]) -> None: ... + def WhichOneof(self, oneof_group: typing.Literal["msg", b"msg"]) -> typing.Literal["scrape_request"] | None: ... + +Global___MetricsStreamResponse: typing_extensions.TypeAlias = MetricsStreamResponse + +@typing.final +class MetricsScrapeRequest(google.protobuf.message.Message): + """Empty request: "send your /metrics now".""" + + DESCRIPTOR: google.protobuf.descriptor.Descriptor + + def __init__( + self, + ) -> None: ... + +Global___MetricsScrapeRequest: typing_extensions.TypeAlias = MetricsScrapeRequest + @typing.final class PushLogsRequest(google.protobuf.message.Message): """Request to push log entries to the telemetry service.""" @@ -96,6 +193,7 @@ class LogEntry(google.protobuf.message.Message): RESULT_FIELD_NUMBER: builtins.int DRIVER_TYPE_FIELD_NUMBER: builtins.int EXTRA_FIELDS_FIELD_NUMBER: builtins.int + NAMESPACE_FIELD_NUMBER: builtins.int severity: builtins.str """Log severity: debug, info, warning, error, critical.""" message: builtins.str @@ -115,7 +213,9 @@ class LogEntry(google.protobuf.message.Message): driver_type: builtins.str """Log body: driver category (storage, power, network, etc.).""" namespace: builtins.str - """Loki stream label: Kubernetes namespace (bounded by cluster size).""" + """Capped at 16 entries, 64-char keys, 256-char values. + Loki stream label: Kubernetes namespace (bounded by cluster size). + """ @property def timestamp(self) -> google.protobuf.timestamp_pb2.Timestamp: """When the log was emitted.""" diff --git a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.py b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.py index 7be5e24a0..e34c217b6 100644 --- a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.py +++ b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.py @@ -6,7 +6,7 @@ class TelemetryServiceStub: - """A service that receives structured logs from exporters and clients. + """A service that reverse-scrapes exporter metrics and receives structured logs. Implemented by jumpstarter-telemetry; not part of the controller. """ @@ -16,6 +16,11 @@ def __init__(self, channel): Args: channel: A grpc.Channel. """ + self.MetricsStream = channel.stream_stream( + '/jumpstarter.v1.TelemetryService/MetricsStream', + request_serializer=jumpstarter_dot_v1_dot_telemetry__pb2.MetricsStreamRequest.SerializeToString, + response_deserializer=jumpstarter_dot_v1_dot_telemetry__pb2.MetricsStreamResponse.FromString, + _registered_method=True) self.PushLogs = channel.unary_unary( '/jumpstarter.v1.TelemetryService/PushLogs', request_serializer=jumpstarter_dot_v1_dot_telemetry__pb2.PushLogsRequest.SerializeToString, @@ -24,10 +29,18 @@ def __init__(self, channel): class TelemetryServiceServicer: - """A service that receives structured logs from exporters and clients. + """A service that reverse-scrapes exporter metrics and receives structured logs. Implemented by jumpstarter-telemetry; not part of the controller. """ + def MetricsStream(self, request_iterator, context): + """Persistent bidirectional stream: telemetry sends scrape requests, + exporter responds with full metric snapshots (OpenMetrics text). + """ + context.set_code(grpc.StatusCode.UNIMPLEMENTED) + context.set_details('Method not implemented!') + raise NotImplementedError('Method not implemented!') + def PushLogs(self, request, context): """Push structured log entries to the telemetry service for Loki ingest. """ @@ -38,6 +51,11 @@ def PushLogs(self, request, context): def add_TelemetryServiceServicer_to_server(servicer, server): rpc_method_handlers = { + 'MetricsStream': grpc.stream_stream_rpc_method_handler( + servicer.MetricsStream, + request_deserializer=jumpstarter_dot_v1_dot_telemetry__pb2.MetricsStreamRequest.FromString, + response_serializer=jumpstarter_dot_v1_dot_telemetry__pb2.MetricsStreamResponse.SerializeToString, + ), 'PushLogs': grpc.unary_unary_rpc_method_handler( servicer.PushLogs, request_deserializer=jumpstarter_dot_v1_dot_telemetry__pb2.PushLogsRequest.FromString, @@ -52,10 +70,37 @@ def add_TelemetryServiceServicer_to_server(servicer, server): # This class is part of an EXPERIMENTAL API. class TelemetryService: - """A service that receives structured logs from exporters and clients. + """A service that reverse-scrapes exporter metrics and receives structured logs. Implemented by jumpstarter-telemetry; not part of the controller. """ + @staticmethod + def MetricsStream(request_iterator, + target, + options=(), + channel_credentials=None, + call_credentials=None, + insecure=False, + compression=None, + wait_for_ready=None, + timeout=None, + metadata=None): + return grpc.experimental.stream_stream( + request_iterator, + target, + '/jumpstarter.v1.TelemetryService/MetricsStream', + jumpstarter_dot_v1_dot_telemetry__pb2.MetricsStreamRequest.SerializeToString, + jumpstarter_dot_v1_dot_telemetry__pb2.MetricsStreamResponse.FromString, + options, + channel_credentials, + insecure, + call_credentials, + compression, + wait_for_ready, + timeout, + metadata, + _registered_method=True) + @staticmethod def PushLogs(request, target, diff --git a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.pyi b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.pyi index dc163b903..d810b4378 100644 --- a/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.pyi +++ b/python/packages/jumpstarter-protocol/jumpstarter_protocol/jumpstarter/v1/telemetry_pb2_grpc.pyi @@ -25,6 +25,22 @@ class _ServicerContext(grpc.ServicerContext, grpc.aio.ServicerContext): # type: GRPC_GENERATED_VERSION: str GRPC_VERSION: str +_TelemetryServiceMetricsStreamType = typing_extensions.TypeVar( + '_TelemetryServiceMetricsStreamType', + grpc.StreamStreamMultiCallable[ + jumpstarter.v1.telemetry_pb2.MetricsStreamRequest, + jumpstarter.v1.telemetry_pb2.MetricsStreamResponse, + ], + grpc.aio.StreamStreamMultiCallable[ + jumpstarter.v1.telemetry_pb2.MetricsStreamRequest, + jumpstarter.v1.telemetry_pb2.MetricsStreamResponse, + ], + default=grpc.StreamStreamMultiCallable[ + jumpstarter.v1.telemetry_pb2.MetricsStreamRequest, + jumpstarter.v1.telemetry_pb2.MetricsStreamResponse, + ], +) + _TelemetryServicePushLogsType = typing_extensions.TypeVar( '_TelemetryServicePushLogsType', grpc.UnaryUnaryMultiCallable[ @@ -41,13 +57,17 @@ _TelemetryServicePushLogsType = typing_extensions.TypeVar( ], ) -class TelemetryServiceStub(typing.Generic[_TelemetryServicePushLogsType]): - """A service that receives structured logs from exporters and clients. +class TelemetryServiceStub(typing.Generic[_TelemetryServiceMetricsStreamType, _TelemetryServicePushLogsType]): + """A service that reverse-scrapes exporter metrics and receives structured logs. Implemented by jumpstarter-telemetry; not part of the controller. """ @typing.overload def __init__(self: TelemetryServiceStub[ + grpc.StreamStreamMultiCallable[ + jumpstarter.v1.telemetry_pb2.MetricsStreamRequest, + jumpstarter.v1.telemetry_pb2.MetricsStreamResponse, + ], grpc.UnaryUnaryMultiCallable[ jumpstarter.v1.telemetry_pb2.PushLogsRequest, jumpstarter.v1.telemetry_pb2.PushLogsResponse, @@ -56,16 +76,29 @@ class TelemetryServiceStub(typing.Generic[_TelemetryServicePushLogsType]): @typing.overload def __init__(self: TelemetryServiceStub[ + grpc.aio.StreamStreamMultiCallable[ + jumpstarter.v1.telemetry_pb2.MetricsStreamRequest, + jumpstarter.v1.telemetry_pb2.MetricsStreamResponse, + ], grpc.aio.UnaryUnaryMultiCallable[ jumpstarter.v1.telemetry_pb2.PushLogsRequest, jumpstarter.v1.telemetry_pb2.PushLogsResponse, ], ], channel: grpc.aio.Channel) -> None: ... + MetricsStream: _TelemetryServiceMetricsStreamType + """Persistent bidirectional stream: telemetry sends scrape requests, + exporter responds with full metric snapshots (OpenMetrics text). + """ + PushLogs: _TelemetryServicePushLogsType """Push structured log entries to the telemetry service for Loki ingest.""" TelemetryServiceAsyncStub: typing_extensions.TypeAlias = TelemetryServiceStub[ + grpc.aio.StreamStreamMultiCallable[ + jumpstarter.v1.telemetry_pb2.MetricsStreamRequest, + jumpstarter.v1.telemetry_pb2.MetricsStreamResponse, + ], grpc.aio.UnaryUnaryMultiCallable[ jumpstarter.v1.telemetry_pb2.PushLogsRequest, jumpstarter.v1.telemetry_pb2.PushLogsResponse, @@ -73,10 +106,20 @@ TelemetryServiceAsyncStub: typing_extensions.TypeAlias = TelemetryServiceStub[ ] class TelemetryServiceServicer(metaclass=abc.ABCMeta): - """A service that receives structured logs from exporters and clients. + """A service that reverse-scrapes exporter metrics and receives structured logs. Implemented by jumpstarter-telemetry; not part of the controller. """ + @abc.abstractmethod + def MetricsStream( + self, + request_iterator: _MaybeAsyncIterator[jumpstarter.v1.telemetry_pb2.MetricsStreamRequest], + context: _ServicerContext, + ) -> typing.Union[collections.abc.Iterator[jumpstarter.v1.telemetry_pb2.MetricsStreamResponse], collections.abc.AsyncIterator[jumpstarter.v1.telemetry_pb2.MetricsStreamResponse]]: + """Persistent bidirectional stream: telemetry sends scrape requests, + exporter responds with full metric snapshots (OpenMetrics text). + """ + @abc.abstractmethod def PushLogs( self, From 2423e947af0d1eb234b164078b6c8fd7dd65d2f6 Mon Sep 17 00:00:00 2001 From: Roddie Kieley Date: Tue, 1 Sep 2026 20:35:14 -0230 Subject: [PATCH 2/3] fix: surface MetricsStream OpenMetrics parse failures (JEP-0013 Phase 3) Stop silently dropping unparseable exporter snapshots. Log the exporter and error, and increment jumpstarter_metrics_parse_errors_total so reverse-scrape omissions are visible on the same /metrics response. Co-authored-by: Cursor --- controller/internal/service/metrics_merge.go | 9 +- .../internal/service/metrics_merge_test.go | 106 ++++++++++++++++-- .../internal/service/metrics_stream_test.go | 67 +++++++++++ controller/internal/service/telemetry_http.go | 22 +++- .../internal/service/telemetry_service.go | 1 + 5 files changed, 190 insertions(+), 15 deletions(-) diff --git a/controller/internal/service/metrics_merge.go b/controller/internal/service/metrics_merge.go index 378e54af3..484b4e27e 100644 --- a/controller/internal/service/metrics_merge.go +++ b/controller/internal/service/metrics_merge.go @@ -35,6 +35,10 @@ const ( // scrapeTimeoutsMetric is incremented once per exporter that does not // answer a MetricsStream scrape within scrapeTimeout (JEP-0013). scrapeTimeoutsMetric = "jumpstarter_scrape_timeouts_total" + + // metricsParseErrorsMetric is incremented once per exporter snapshot that + // MetricsStream delivered but OpenMetrics parse rejected (JEP-0013). + metricsParseErrorsMetric = "jumpstarter_metrics_parse_errors_total" ) // DefaultDriverTypeEnum is the JEP-0013 default allowlist for driver_type. @@ -119,7 +123,7 @@ func encodeMetricFamilies(w io.Writer, families []*dto.MetricFamily) error { return nil } -func mergeSnapshots(snapshots []exporterSnapshot, extra []*dto.MetricFamily, cfgFor func(string) mergeConfig) []*dto.MetricFamily { +func mergeSnapshots(snapshots []exporterSnapshot, extra []*dto.MetricFamily, cfgFor func(string) mergeConfig, onParseError func(exporter string, err error)) []*dto.MetricFamily { byName := map[string]*dto.MetricFamily{} for _, f := range extra { if f == nil || f.GetName() == "" { @@ -133,6 +137,9 @@ func mergeSnapshots(snapshots []exporterSnapshot, extra []*dto.MetricFamily, cfg } families, err := parseMetricFamilies(snap.text) if err != nil { + if onParseError != nil { + onParseError(snap.name, err) + } continue } cfg := cfgFor(snap.name) diff --git a/controller/internal/service/metrics_merge_test.go b/controller/internal/service/metrics_merge_test.go index 5f6dfe70a..72bdbbbc4 100644 --- a/controller/internal/service/metrics_merge_test.go +++ b/controller/internal/service/metrics_merge_test.go @@ -18,6 +18,7 @@ package service import ( "bytes" + "errors" "strings" "testing" @@ -25,6 +26,24 @@ import ( "github.com/prometheus/common/expfmt" ) +func labeledCounterValue(t *testing.T, mfs []*dto.MetricFamily, name, label, want string) float64 { + t.Helper() + for _, mf := range mfs { + if mf.GetName() != name && mf.GetName()+"_total" != name { + continue + } + for _, m := range mf.Metric { + for _, lp := range m.GetLabel() { + if lp.GetName() == label && lp.GetValue() == want { + return m.GetCounter().GetValue() + } + } + } + } + t.Fatalf("%s{%s=%q} missing from gathered families", name, label, want) + return 0 +} + func TestApplyMetric_OverwritesExporterRemapsDriverTypeAndFiltersExemplars(t *testing.T) { one := 1.0 m := &dto.Metric{ @@ -91,14 +110,32 @@ func TestApplyMetric_KeepsAllowlistedDriverType(t *testing.T) { } } -func TestMergeSnapshots_CombinesExportersAndSkipsInvalidText(t *testing.T) { - cfg := func(name string) mergeConfig { - return mergeConfig{ - exporterName: name, - driverTypes: setToMap(DefaultDriverTypeEnum), - exemplarKeys: setToMap(DefaultExemplarKeys), - } +// pythonOpenMetricsWithExemplar is the prometheus_client OpenMetrics form +// that prometheus/common v0.62's text fallback cannot parse (it treats '#' as +// a timestamp). This is the body exporters send after driver operations. +const pythonOpenMetricsWithExemplar = `# TYPE jumpstarter_operation_duration_seconds histogram +jumpstarter_operation_duration_seconds_bucket{exporter="sidekick",le="0.005",operation="on",result="success",driver_type="power"} 1.0 # {lease_id="lease-1"} 0.00262 1788299670.149 +jumpstarter_operation_duration_seconds_sum{exporter="sidekick",operation="on",result="success",driver_type="power"} 0.00262 +jumpstarter_operation_duration_seconds_count{exporter="sidekick",operation="on",result="success",driver_type="power"} 1.0 +# EOF +` + +func testMergeCfg(name string) mergeConfig { + return mergeConfig{ + exporterName: name, + driverTypes: setToMap(DefaultDriverTypeEnum), + exemplarKeys: setToMap(DefaultExemplarKeys), + } +} + +func TestParseMetricFamilies_OpenMetricsExemplarSuffixFails(t *testing.T) { + _, err := parseMetricFamilies([]byte(pythonOpenMetricsWithExemplar)) + if err == nil { + t.Fatal("expected parse error for OpenMetrics exemplar suffix") } +} + +func TestMergeSnapshots_CombinesExportersAndReportsInvalidText(t *testing.T) { textA := []byte(`# TYPE jumpstarter_operations_total counter # HELP jumpstarter_operations_total Total operations performed. jumpstarter_operations_total{exporter="a",operation="on",result="success",driver_type="power"} 2.0 @@ -108,11 +145,18 @@ jumpstarter_operations_total{exporter="a",operation="on",result="success",driver jumpstarter_operations_total{exporter="b",operation="off",result="success",driver_type="power"} 3.0 # EOF `) + type parseErrCall struct { + exporter string + err error + } + var calls []parseErrCall families := mergeSnapshots([]exporterSnapshot{ {name: "exp-a", text: textA}, {name: "exp-b", text: textB}, {name: "bad", text: []byte("not metrics")}, - }, nil, cfg) + }, nil, testMergeCfg, func(exporter string, err error) { + calls = append(calls, parseErrCall{exporter: exporter, err: err}) + }) var ops *dto.MetricFamily for _, f := range families { @@ -135,6 +179,52 @@ jumpstarter_operations_total{exporter="b",operation="off",result="success",drive if !exporters["exp-a"] || !exporters["exp-b"] { t.Errorf("exporters = %v, want exp-a and exp-b", exporters) } + if len(calls) != 1 { + t.Fatalf("parse-error callbacks = %d, want 1", len(calls)) + } + if calls[0].exporter != "bad" { + t.Errorf("parse-error exporter = %q, want bad", calls[0].exporter) + } + if calls[0].err == nil { + t.Fatal("parse-error callback err is nil") + } +} + +func TestMergeSnapshots_ReportsOpenMetricsExemplarParseError(t *testing.T) { + var gotExporter string + var gotErr error + families := mergeSnapshots([]exporterSnapshot{ + {name: "sidekick", text: []byte(pythonOpenMetricsWithExemplar)}, + }, nil, testMergeCfg, func(exporter string, err error) { + gotExporter = exporter + gotErr = err + }) + for _, f := range families { + if strings.HasPrefix(f.GetName(), "jumpstarter_operation_duration_seconds") { + t.Fatalf("unparseable snapshot must be omitted, got family %s", f.GetName()) + } + } + if gotExporter != "sidekick" { + t.Errorf("parse-error exporter = %q, want sidekick", gotExporter) + } + if gotErr == nil { + t.Fatal("parse-error callback err is nil") + } +} + +func TestRecordMetricsParseError_IncrementsLabeledCounter(t *testing.T) { + svc := &TelemetryService{} + svc.initScrapeTimeouts() + svc.recordMetricsParseError("sidekick", errors.New("parse failed")) + + mfs, err := svc.metricsRegistry.Gather() + if err != nil { + t.Fatalf("Gather: %v", err) + } + value := labeledCounterValue(t, mfs, metricsParseErrorsMetric, labelExporter, "sidekick") + if value != 1 { + t.Fatalf("%s{exporter=sidekick} = %v, want 1", metricsParseErrorsMetric, value) + } } func TestEncodeMetricFamilies_WritesOpenMetricsEOF(t *testing.T) { diff --git a/controller/internal/service/metrics_stream_test.go b/controller/internal/service/metrics_stream_test.go index 6f4a50cd0..6aed15c36 100644 --- a/controller/internal/service/metrics_stream_test.go +++ b/controller/internal/service/metrics_stream_test.go @@ -400,3 +400,70 @@ jumpstarter_operations_total{exporter="spoofed",operation="on",result="success", default: } } + +func TestFanout_UnparseableSnapshotIncrementsParseErrorsOnSameResponse(t *testing.T) { + svc, client, addr := startTestHub(t, time.Second) + ctx := exporterStreamCtx(t, svc, "exporter:jumpstarter:sidekick:uid1") + stream, err := client.MetricsStream(ctx) + if err != nil { + t.Fatalf("MetricsStream: %v", err) + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_Register{ + Register: &pb.MetricsRegister{Identity: "sidekick"}, + }, + }); err != nil { + t.Fatalf("Send register: %v", err) + } + + errCh := make(chan error, 1) + go func() { + for { + msg, err := stream.Recv() + if err != nil { + errCh <- err + return + } + if msg.GetScrapeRequest() == nil { + continue + } + if err := stream.Send(&pb.MetricsStreamRequest{ + Msg: &pb.MetricsStreamRequest_ScrapeResponse{ + ScrapeResponse: &pb.MetricsScrapeResponse{MetricsText: []byte(pythonOpenMetricsWithExemplar)}, + }, + }); err != nil { + errCh <- err + return + } + } + }() + waitRegistered(t, svc, 1) + + code, body := httpGet(t, "http://"+addr+"/metrics") + if code != http.StatusOK { + t.Fatalf("GET /metrics status=%d body=%s", code, body) + } + if strings.Contains(body, "jumpstarter_operation_duration_seconds") { + t.Errorf("unparseable exporter snapshot must be omitted, body:\n%s", body) + } + if !strings.Contains(body, metricsParseErrorsMetric) || !strings.Contains(body, `exporter="sidekick"`) { + t.Errorf("parse-error counter missing from same /metrics response:\n%s", body) + } + + mfs, err := svc.metricsRegistry.Gather() + if err != nil { + t.Fatalf("Gather: %v", err) + } + value := labeledCounterValue(t, mfs, metricsParseErrorsMetric, labelExporter, "sidekick") + if value < 1 { + t.Fatalf("%s{exporter=sidekick} = %v, want >= 1 body:\n%s", metricsParseErrorsMetric, value, body) + } + + select { + case err := <-errCh: + if err != nil && err != io.EOF && status.Code(err) != codes.Canceled && status.Code(err) != codes.Unavailable { + t.Fatalf("stream goroutine: %v", err) + } + default: + } +} diff --git a/controller/internal/service/telemetry_http.go b/controller/internal/service/telemetry_http.go index 6d12daff4..ecd403019 100644 --- a/controller/internal/service/telemetry_http.go +++ b/controller/internal/service/telemetry_http.go @@ -24,7 +24,6 @@ import ( "time" "github.com/prometheus/client_golang/prometheus" - dto "github.com/prometheus/client_model/go" "github.com/prometheus/common/expfmt" ctrl "sigs.k8s.io/controller-runtime" ) @@ -40,7 +39,18 @@ func (s *TelemetryService) initScrapeTimeouts() { Name: scrapeTimeoutsMetric, Help: "Exporter MetricsStream scrapes that exceeded scrapeTimeout.", }) - s.metricsRegistry.MustRegister(s.scrapeTimeouts) + s.parseErrors = prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: metricsParseErrorsMetric, + Help: "Exporter MetricsStream snapshots omitted because OpenMetrics parse failed.", + }, []string{labelExporter}) + s.metricsRegistry.MustRegister(s.scrapeTimeouts, s.parseErrors) +} + +func (s *TelemetryService) recordMetricsParseError(exporter string, err error) { + if s.parseErrors != nil { + s.parseErrors.WithLabelValues(exporter).Inc() + } + ctrl.Log.WithName("telemetry").Error(err, "exporter metrics snapshot omitted", logFieldExporter, exporter) } func metricsHTTPEnabled(addr string) bool { @@ -110,19 +120,19 @@ func (s *TelemetryService) handleReadyz(w http.ResponseWriter, _ *http.Request) func (s *TelemetryService) handleMetrics(w http.ResponseWriter, r *http.Request) { s.initScrapeTimeouts() snaps := s.fanoutScrapes(r.Context()) + // Merge exporter snapshots first so parse failures increment parseErrors + // before hub metrics are gathered into the same /metrics response. + families := mergeSnapshots(snaps, nil, s.mergeConfigFor, s.recordMetricsParseError) - var extra []*dto.MetricFamily if s.metricsRegistry != nil { gathered, err := s.metricsRegistry.Gather() if err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return } - extra = gathered + families = mergeSnapshots(nil, append(gathered, families...), s.mergeConfigFor, nil) } - families := mergeSnapshots(snaps, extra, s.mergeConfigFor) - var buf bytes.Buffer if err := encodeMetricFamilies(&buf, families); err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) diff --git a/controller/internal/service/telemetry_service.go b/controller/internal/service/telemetry_service.go index e3c60bb77..b2be84936 100644 --- a/controller/internal/service/telemetry_service.go +++ b/controller/internal/service/telemetry_service.go @@ -95,6 +95,7 @@ type TelemetryService struct { stateMu sync.Mutex conns map[string]*metricsConn scrapeTimeouts prometheus.Counter + parseErrors *prometheus.CounterVec metricsRegistry *prometheus.Registry metricsAddr string grpcReady atomic.Bool From 26b410d9ca910e4c32a3bb4e53c39ef353b2df3c Mon Sep 17 00:00:00 2001 From: Roddie Kieley Date: Thu, 3 Sep 2026 09:59:33 -0230 Subject: [PATCH 3/3] fix: unparseable becomes unparsable as per typo check Signed-off-by: Roddie Kieley --- controller/internal/service/metrics_merge_test.go | 2 +- controller/internal/service/metrics_stream_test.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/controller/internal/service/metrics_merge_test.go b/controller/internal/service/metrics_merge_test.go index 72bdbbbc4..984612040 100644 --- a/controller/internal/service/metrics_merge_test.go +++ b/controller/internal/service/metrics_merge_test.go @@ -201,7 +201,7 @@ func TestMergeSnapshots_ReportsOpenMetricsExemplarParseError(t *testing.T) { }) for _, f := range families { if strings.HasPrefix(f.GetName(), "jumpstarter_operation_duration_seconds") { - t.Fatalf("unparseable snapshot must be omitted, got family %s", f.GetName()) + t.Fatalf("unparsable snapshot must be omitted, got family %s", f.GetName()) } } if gotExporter != "sidekick" { diff --git a/controller/internal/service/metrics_stream_test.go b/controller/internal/service/metrics_stream_test.go index 6aed15c36..ed2f1356f 100644 --- a/controller/internal/service/metrics_stream_test.go +++ b/controller/internal/service/metrics_stream_test.go @@ -444,7 +444,7 @@ func TestFanout_UnparseableSnapshotIncrementsParseErrorsOnSameResponse(t *testin t.Fatalf("GET /metrics status=%d body=%s", code, body) } if strings.Contains(body, "jumpstarter_operation_duration_seconds") { - t.Errorf("unparseable exporter snapshot must be omitted, body:\n%s", body) + t.Errorf("unparsable exporter snapshot must be omitted, body:\n%s", body) } if !strings.Contains(body, metricsParseErrorsMetric) || !strings.Contains(body, `exporter="sidekick"`) { t.Errorf("parse-error counter missing from same /metrics response:\n%s", body)