From a2d749801eb56005c433f964afff4259ec97d56c Mon Sep 17 00:00:00 2001 From: Roddie Kieley Date: Wed, 26 Aug 2026 19:01:32 -0230 Subject: [PATCH 1/4] 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/4] 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 b84f952430ac6adfe6bc78c94c0788a0a6cf92ad Mon Sep 17 00:00:00 2001 From: Roddie Kieley Date: Wed, 26 Aug 2026 21:06:58 -0230 Subject: [PATCH 3/4] feat: exporter MetricsStream client for reverse-scrape (JEP-0013 Phase 3) Open a long-lived MetricsStream after GetServiceEndpoints so the telemetry hub can pull generate_latest() snapshots without replacing local /metrics. --- .../jumpstarter/exporter/exporter.py | 25 +- .../exporter/exporter_telemetry_test.py | 33 +++ .../jumpstarter/exporter/exporter_test.py | 1 + .../jumpstarter/exporter/metrics_stream.py | 86 ++++++ .../exporter/metrics_stream_test.py | 278 ++++++++++++++++++ 5 files changed, 419 insertions(+), 4 deletions(-) create mode 100644 python/packages/jumpstarter/jumpstarter/exporter/metrics_stream.py create mode 100644 python/packages/jumpstarter/jumpstarter/exporter/metrics_stream_test.py diff --git a/python/packages/jumpstarter/jumpstarter/exporter/exporter.py b/python/packages/jumpstarter/jumpstarter/exporter/exporter.py index da57a1099..857b4e626 100644 --- a/python/packages/jumpstarter/jumpstarter/exporter/exporter.py +++ b/python/packages/jumpstarter/jumpstarter/exporter/exporter.py @@ -34,6 +34,7 @@ from jumpstarter.config.tls import TLSConfigV1Alpha1 from jumpstarter.exporter.hooks import HookExecutor from jumpstarter.exporter.lease_context import LeaseContext +from jumpstarter.exporter.metrics_stream import MetricsStreamClient from jumpstarter.exporter.session import Session from jumpstarter.exporter.telemetry import TelemetryLogHandler from jumpstarter.logging import clear_log_context, set_log_context @@ -379,6 +380,9 @@ class Exporter(AsyncContextManagerMixin, Metadata): _telemetry_channel: grpc.aio.Channel | None = field(init=False, default=None) """gRPC channel to the telemetry service. Closed on exporter shutdown.""" + _metrics_stream: MetricsStreamClient | None = field(init=False, default=None) + """Optional MetricsStream client. Started next to the log flush loop when telemetry is configured.""" + _status_drain_active: bool = field(init=False, default=False) """True only while serve()'s task group is running and the drain task is active. @@ -558,8 +562,9 @@ async def _setup_telemetry(self) -> None: """Discover and connect to the optional telemetry service. Calls GetServiceEndpoints on the controller. When a telemetry endpoint is - returned, a gRPC channel is created and TelemetryLogHandler is attached - to the root Python logger. Safe to call multiple times; subsequent calls + returned, a gRPC channel is created, TelemetryLogHandler is attached + to the root Python logger, and a MetricsStreamClient is prepared for + reverse-scrape. Safe to call multiple times; subsequent calls are no-ops if a handler is already configured. """ if self._telemetry_handler is not None: @@ -609,6 +614,11 @@ async def _setup_telemetry(self) -> None: handler.setLevel(_severity_to_level(ep.min_severity)) logging.getLogger().addHandler(handler) self._telemetry_handler = handler + self._metrics_stream = MetricsStreamClient( + stub, + identity=self.name, + token=self.token, + ) logger.info("Telemetry log handler attached") async def _retry_rpc( @@ -1280,6 +1290,7 @@ async def serve(self): if self._telemetry_channel is not None: await self._telemetry_channel.close() self._telemetry_channel = None + self._metrics_stream = None async def _run_control_plane( self, @@ -1294,8 +1305,7 @@ async def _run_control_plane( self._pending_status_request = None self._status_drain_active = True tg.start_soon(self._drain_status_reports) - if self._telemetry_handler is not None: - tg.start_soon(self._telemetry_handler.flush_loop) + self._start_telemetry_tasks(tg) tg.start_soon( self._retry_stream, "Status", @@ -1313,6 +1323,13 @@ async def _run_control_plane( if await self._apply_status(message, tg): break + def _start_telemetry_tasks(self, tg: TaskGroup) -> None: + """Start PushLogs flush and MetricsStream next to the control-plane tasks.""" + if self._telemetry_handler is not None: + tg.start_soon(self._telemetry_handler.flush_loop) + if self._metrics_stream is not None: + tg.start_soon(self._metrics_stream.run) + async def _apply_status( self, status: jumpstarter_pb2.StatusResponse, diff --git a/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py b/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py index 57ad81c2c..129009230 100644 --- a/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py +++ b/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py @@ -35,6 +35,10 @@ def make_bare_exporter(): exp = Exporter.__new__(Exporter) exp._telemetry_handler = None exp._telemetry_channel = None + exp._metrics_stream = None + exp.labels = {"jumpstarter.dev/name": "test-exporter"} + exp.token = "test-token" + exp.namespace = "" tls = MagicMock() tls.insecure = False @@ -100,6 +104,7 @@ async def _ctx(): await exp._setup_telemetry() assert exp._telemetry_handler is None + assert exp._metrics_stream is None async def test_setup_telemetry_no_endpoints_returns_early(): @@ -113,6 +118,7 @@ async def test_setup_telemetry_no_endpoints_returns_early(): await exp._setup_telemetry() assert exp._telemetry_handler is None + assert exp._metrics_stream is None async def test_setup_telemetry_insecure_uses_insecure_channel(): @@ -131,6 +137,9 @@ async def test_setup_telemetry_insecure_uses_insecure_channel(): mock_insecure.assert_called_once_with(ep.endpoint) assert exp._telemetry_handler is not None + assert exp._metrics_stream is not None + assert exp._metrics_stream.identity == "test-exporter" + assert exp._metrics_stream.token == "test-token" async def test_setup_telemetry_env_var_insecure(monkeypatch): @@ -202,3 +211,27 @@ async def test_setup_telemetry_applies_min_severity(): assert exp._telemetry_handler is not None assert exp._telemetry_handler.level == logging.WARNING + assert exp._metrics_stream is not None + + +def test_start_telemetry_tasks_starts_flush_loop_and_metrics_stream(): + """MetricsStream runs next to PushLogs flush_loop (JEP-0013 reverse-scrape).""" + exp = make_bare_exporter() + handler = MagicMock() + stream = MagicMock() + exp._telemetry_handler = handler + exp._metrics_stream = stream + tg = MagicMock() + + exp._start_telemetry_tasks(tg) + + started = [c.args[0] for c in tg.start_soon.call_args_list] + assert handler.flush_loop in started + assert stream.run in started + + +def test_start_telemetry_tasks_skips_when_telemetry_disabled(): + exp = make_bare_exporter() + tg = MagicMock() + exp._start_telemetry_tasks(tg) + tg.start_soon.assert_not_called() diff --git a/python/packages/jumpstarter/jumpstarter/exporter/exporter_test.py b/python/packages/jumpstarter/jumpstarter/exporter/exporter_test.py index 924b00427..24e72690a 100644 --- a/python/packages/jumpstarter/jumpstarter/exporter/exporter_test.py +++ b/python/packages/jumpstarter/jumpstarter/exporter/exporter_test.py @@ -76,6 +76,7 @@ def _make_base_exporter(**overrides): "_request_lease_release": AsyncMock(), "_telemetry_handler": None, "_telemetry_channel": None, + "_metrics_stream": None, } defaults.update(overrides) exporter = Exporter.__new__(Exporter) diff --git a/python/packages/jumpstarter/jumpstarter/exporter/metrics_stream.py b/python/packages/jumpstarter/jumpstarter/exporter/metrics_stream.py new file mode 100644 index 000000000..20ec53ee7 --- /dev/null +++ b/python/packages/jumpstarter/jumpstarter/exporter/metrics_stream.py @@ -0,0 +1,86 @@ +"""Exporter MetricsStream client — reverse-scrape snapshots for jumpstarter-telemetry. + +Does not replace the local HTTP GET /metrics loopback bind (CLI :0); reverse-scrape +is an additional path using the same generate_latest() bytes. +""" + +from __future__ import annotations + +import logging +from typing import Any + +from anyio import sleep +from google.protobuf.timestamp_pb2 import Timestamp +from jumpstarter_protocol import telemetry_pb2, telemetry_pb2_grpc + +from jumpstarter.metrics import MetricsRegistry, get_registry + +logger = logging.getLogger(__name__) + +_BACKOFF_BASE = 1.0 +_BACKOFF_CAP = 30.0 + + +class MetricsStreamClient: + """Long-lived MetricsStream: register, answer scrapes, reconnect with backoff.""" + + def __init__( + self, + stub: telemetry_pb2_grpc.TelemetryServiceStub, + *, + identity: str, + token: str = "", + registry: MetricsRegistry | None = None, + backoff_base: float = _BACKOFF_BASE, + backoff_cap: float = _BACKOFF_CAP, + ) -> None: + self.identity = identity + self.token = token + self._stub = stub + self._registry = registry + self._backoff_base = backoff_base + self._backoff_cap = backoff_cap + + def _metrics_text(self) -> bytes: + reg = self._registry if self._registry is not None else get_registry() + return reg.generate_latest() + + def _call_kwargs(self) -> dict[str, Any]: + if self.token: + return {"metadata": [("authorization", f"Bearer {self.token}")]} + return {} + + async def run(self) -> None: + """Reconnect forever until cancelled. Local counters are never reset.""" + delay = self._backoff_base + while True: + try: + await self._session() + delay = self._backoff_base + except Exception as exc: + logger.debug("MetricsStream disconnected: %s; retry in %ss", exc, delay) + await sleep(delay) + delay = min(delay * 2, self._backoff_cap) + continue + await sleep(delay) + + async def _session(self) -> None: + call = self._stub.MetricsStream(**self._call_kwargs()) + await call.write( + telemetry_pb2.MetricsStreamRequest( + register=telemetry_pb2.MetricsRegister(identity=self.identity), + ) + ) + async for resp in call: + if resp.WhichOneof("msg") != "scrape_request": + continue + ts = Timestamp() + ts.GetCurrentTime() + await call.write( + telemetry_pb2.MetricsStreamRequest( + scrape_response=telemetry_pb2.MetricsScrapeResponse( + metrics_text=self._metrics_text(), + timestamp=ts, + ), + ) + ) diff --git a/python/packages/jumpstarter/jumpstarter/exporter/metrics_stream_test.py b/python/packages/jumpstarter/jumpstarter/exporter/metrics_stream_test.py new file mode 100644 index 000000000..1a57a86e8 --- /dev/null +++ b/python/packages/jumpstarter/jumpstarter/exporter/metrics_stream_test.py @@ -0,0 +1,278 @@ +"""Tests for the exporter MetricsStream client (JEP-0013 Phase 3 PR C). + +Copyright 2026. The Jumpstarter Authors. + +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. +""" + +from __future__ import annotations + +import urllib.request +from unittest.mock import MagicMock, patch + +import grpc +import pytest +from anyio import EndOfStream, create_memory_object_stream, create_task_group, fail_after, sleep +from jumpstarter_protocol import telemetry_pb2 + +from jumpstarter.exporter.metrics_stream import MetricsStreamClient +from jumpstarter.metrics.registry import MetricsRegistry +from jumpstarter.metrics.server import start_metrics_server + +pytestmark = pytest.mark.anyio + +IDENTITY = "board-42" +TOKEN = "exporter-jwt" + + +class FakeMetricsCall: + """In-memory stand-in for a grpc.aio stream-stream MetricsStream call.""" + + def __init__(self): + self.writes: list[telemetry_pb2.MetricsStreamRequest] = [] + self._send, self._recv = create_memory_object_stream[telemetry_pb2.MetricsStreamResponse](32) + + async def write(self, msg: telemetry_pb2.MetricsStreamRequest) -> None: + self.writes.append(msg) + + def __aiter__(self): + return self + + async def __anext__(self) -> telemetry_pb2.MetricsStreamResponse: + try: + return await self._recv.receive() + except EndOfStream as exc: + raise StopAsyncIteration from exc + + async def push_scrape(self) -> None: + await self._send.send( + telemetry_pb2.MetricsStreamResponse( + scrape_request=telemetry_pb2.MetricsScrapeRequest(), + ) + ) + + async def close_from_server(self) -> None: + await self._send.aclose() + + +class RecordingStub: + def __init__(self, call: FakeMetricsCall | None = None, *, error: Exception | None = None): + self.call = call + self.error = error + self.calls: list[dict] = [] + + def MetricsStream(self, *args, **kwargs): + assert not args, "MetricsStream client must use write/read style, not a request iterator" + self.calls.append(kwargs) + if self.error is not None: + raise self.error + assert self.call is not None + return self.call + + +async def _wait_until(pred, *, timeout: float = 2.0) -> None: + with fail_after(timeout): + while not pred(): + await sleep(0.01) + + +def _record_op(reg: MetricsRegistry, *, exporter: str = IDENTITY) -> None: + reg.record_operation( + exporter=exporter, + operation="power", + result="success", + driver_type="power", + duration_seconds=0.05, + ) + + +def test_module_documents_local_http_metrics_is_not_replaced(): + """CLI loopback :0 scrape from JEP-0013 Phase 2 stays; reverse-scrape is additive.""" + import jumpstarter.exporter.metrics_stream as mod + + assert "does not replace" in (mod.__doc__ or "").lower() + + +class TestMetricsStreamProtocol: + async def test_first_message_is_register_with_exporter_identity(self): + call = FakeMetricsCall() + stub = RecordingStub(call) + client = MetricsStreamClient(stub, identity=IDENTITY, token=TOKEN) + + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: len(call.writes) >= 1) + tg.cancel_scope.cancel() + + first = call.writes[0] + assert first.WhichOneof("msg") == "register" + assert first.register.identity == IDENTITY + + async def test_sends_bearer_token_metadata(self): + call = FakeMetricsCall() + stub = RecordingStub(call) + client = MetricsStreamClient(stub, identity=IDENTITY, token=TOKEN) + + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: stub.calls) + tg.cancel_scope.cancel() + + assert stub.calls[0]["metadata"] == [("authorization", f"Bearer {TOKEN}")] + + async def test_scrape_replies_with_generate_latest_bytes(self): + reg = MetricsRegistry() + _record_op(reg) + expected = reg.generate_latest() + + call = FakeMetricsCall() + client = MetricsStreamClient(RecordingStub(call), identity=IDENTITY, token=TOKEN, registry=reg) + + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: len(call.writes) >= 1) + await call.push_scrape() + await _wait_until(lambda: len(call.writes) >= 2) + tg.cancel_scope.cancel() + + scrape = call.writes[1] + assert scrape.WhichOneof("msg") == "scrape_response" + assert scrape.scrape_response.metrics_text == expected + assert scrape.scrape_response.HasField("timestamp") + assert scrape.scrape_response.timestamp.seconds > 0 + + async def test_empty_registry_still_sends_scrape_response(self): + reg = MetricsRegistry() + empty = reg.generate_latest() + assert empty # OpenMetrics text is still produced with no samples + + call = FakeMetricsCall() + client = MetricsStreamClient(RecordingStub(call), identity=IDENTITY, token="", registry=reg) + + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: len(call.writes) >= 1) + await call.push_scrape() + await _wait_until(lambda: len(call.writes) >= 2) + tg.cancel_scope.cancel() + + scrape = call.writes[1] + assert scrape.WhichOneof("msg") == "scrape_response" + assert scrape.scrape_response.metrics_text == empty + + async def test_scrape_bytes_match_local_http_metrics(self): + """Reverse-scrape payload must be the same bytes as GET /metrics (lab loopback).""" + reg = MetricsRegistry() + _record_op(reg) + listen, shutdown = start_metrics_server("127.0.0.1:0", registry=reg) + try: + http_body = urllib.request.urlopen(f"http://{listen}/metrics").read() + finally: + assert shutdown is not None + shutdown() + + call = FakeMetricsCall() + client = MetricsStreamClient(RecordingStub(call), identity=IDENTITY, token=TOKEN, registry=reg) + + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: len(call.writes) >= 1) + await call.push_scrape() + await _wait_until(lambda: len(call.writes) >= 2) + tg.cancel_scope.cancel() + + assert call.writes[1].scrape_response.metrics_text == http_body + + async def test_cancel_stops_run(self): + call = FakeMetricsCall() + client = MetricsStreamClient(RecordingStub(call), identity=IDENTITY, token=TOKEN) + + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: len(call.writes) >= 1) + tg.cancel_scope.cancel() + + # Leaving the task group without timeout is the assertion: cancel unwinds. + + async def test_reconnect_keeps_local_counters(self): + """Telemetry unavailable: counters stay local; next scrape is full current state.""" + reg = MetricsRegistry() + _record_op(reg) + expected = reg.generate_latest() + + failing = RecordingStub( + error=grpc.aio.AioRpcError( + code=grpc.StatusCode.UNAVAILABLE, + initial_metadata=None, + trailing_metadata=None, + ) + ) + call = FakeMetricsCall() + succeeding = RecordingStub(call) + stubs = iter([failing, succeeding]) + + def metrics_stream(*args, **kwargs): + return next(stubs).MetricsStream(*args, **kwargs) + + stub = MagicMock() + stub.MetricsStream = metrics_stream + + client = MetricsStreamClient( + stub, identity=IDENTITY, token=TOKEN, registry=reg, backoff_base=0.01, backoff_cap=0.01 + ) + + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: len(call.writes) >= 1) + await call.push_scrape() + await _wait_until(lambda: len(call.writes) >= 2) + tg.cancel_scope.cancel() + + assert call.writes[0].WhichOneof("msg") == "register" + assert call.writes[1].scrape_response.metrics_text == expected + assert b"jumpstarter_operations_total" in expected + + async def test_reconnect_uses_backoff(self): + delays: list[float] = [] + + async def fake_sleep(seconds: float) -> None: + delays.append(seconds) + await sleep(0) + + call = FakeMetricsCall() + attempts = {"n": 0} + + def metrics_stream(*args, **kwargs): + attempts["n"] += 1 + if attempts["n"] == 1: + raise grpc.aio.AioRpcError( + code=grpc.StatusCode.UNAVAILABLE, + initial_metadata=None, + trailing_metadata=None, + ) + return call + + stub = MagicMock() + stub.MetricsStream = metrics_stream + client = MetricsStreamClient( + stub, identity=IDENTITY, token=TOKEN, backoff_base=0.5, backoff_cap=2.0 + ) + + with patch("jumpstarter.exporter.metrics_stream.sleep", fake_sleep): + async with create_task_group() as tg: + tg.start_soon(client.run) + await _wait_until(lambda: len(call.writes) >= 1) + tg.cancel_scope.cancel() + + assert delays[0] == 0.5 + assert attempts["n"] >= 2 From 841250fabf7c4eb5b3f544d39a6cadd5b0bebc23 Mon Sep 17 00:00:00 2001 From: Roddie Kieley Date: Wed, 2 Sep 2026 14:10:02 -0230 Subject: [PATCH 4/4] fix: register MetricsStream with exporter_name after #1058 #1058 stopped sending jumpstarter.dev/name on Register, so identity=self.name was "unknown" and the hub rejected the stream. Co-authored-by: Cursor --- .../jumpstarter/jumpstarter/exporter/exporter.py | 12 ++++++++---- .../jumpstarter/exporter/exporter_telemetry_test.py | 1 + 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/python/packages/jumpstarter/jumpstarter/exporter/exporter.py b/python/packages/jumpstarter/jumpstarter/exporter/exporter.py index 857b4e626..4b8684f42 100644 --- a/python/packages/jumpstarter/jumpstarter/exporter/exporter.py +++ b/python/packages/jumpstarter/jumpstarter/exporter/exporter.py @@ -616,7 +616,7 @@ async def _setup_telemetry(self) -> None: self._telemetry_handler = handler self._metrics_stream = MetricsStreamClient( stub, - identity=self.name, + identity=self.exporter_name, token=self.token, ) logger.info("Telemetry log handler attached") @@ -952,7 +952,7 @@ async def session(self): """ with Session( uuid=self.uuid, - labels=self.labels, + labels=self._session_labels(), root_device=self.device_factory(), motd=self.motd, ) as session: @@ -988,7 +988,7 @@ async def session_for_lease(self): logger.info("Creating new session for lease") with Session( uuid=self.uuid, - labels=self.labels, + labels=self._session_labels(), root_device=self.device_factory(), motd=self.motd, ) as session: @@ -1323,6 +1323,10 @@ async def _run_control_plane( if await self._apply_status(message, tg): break + def _session_labels(self) -> dict[str, str]: + """Labels for local Session/metrics. Not sent on Register (#1058).""" + return {**self.labels, "jumpstarter.dev/name": self.exporter_name} + def _start_telemetry_tasks(self, tg: TaskGroup) -> None: """Start PushLogs flush and MetricsStream next to the control-plane tasks.""" if self._telemetry_handler is not None: @@ -1563,7 +1567,7 @@ async def serve_standalone_tcp( hook_path_str = str(hook_path) with Session( uuid=self.uuid, - labels=self.labels, + labels=self._session_labels(), root_device=self.device_factory(), motd=self.motd, ) as session: diff --git a/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py b/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py index 129009230..447ac92ee 100644 --- a/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py +++ b/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py @@ -37,6 +37,7 @@ def make_bare_exporter(): exp._telemetry_channel = None exp._metrics_stream = None exp.labels = {"jumpstarter.dev/name": "test-exporter"} + exp.exporter_name = "test-exporter" exp.token = "test-token" exp.namespace = ""