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..484b4e27e --- /dev/null +++ b/controller/internal/service/metrics_merge.go @@ -0,0 +1,246 @@ +/* +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" + + // 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. +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, onParseError func(exporter string, err error)) []*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 { + if onParseError != nil { + onParseError(snap.name, err) + } + 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..72bdbbbc4 --- /dev/null +++ b/controller/internal/service/metrics_merge_test.go @@ -0,0 +1,262 @@ +/* +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" + "strings" + "testing" + + dto "github.com/prometheus/client_model/go" + "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{ + 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) + } +} + +// 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 +# EOF +`) + textB := []byte(`# TYPE jumpstarter_operations_total counter +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, testMergeCfg, func(exporter string, err error) { + calls = append(calls, parseErrCall{exporter: exporter, err: err}) + }) + + 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) + } + 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) { + 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..6aed15c36 --- /dev/null +++ b/controller/internal/service/metrics_stream_test.go @@ -0,0 +1,469 @@ +/* +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: + } +} + +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 new file mode 100644 index 000000000..ecd403019 --- /dev/null +++ b/controller/internal/service/telemetry_http.go @@ -0,0 +1,145 @@ +/* +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" + "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.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 { + 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()) + // 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) + + if s.metricsRegistry != nil { + gathered, err := s.metricsRegistry.Gather() + if err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + families = mergeSnapshots(nil, append(gathered, families...), s.mergeConfigFor, nil) + } + + 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..b2be84936 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,48 @@ 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 + parseErrors *prometheus.CounterVec + 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 +266,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 +294,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, diff --git a/python/packages/jumpstarter/jumpstarter/exporter/exporter.py b/python/packages/jumpstarter/jumpstarter/exporter/exporter.py index da57a1099..4b8684f42 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.exporter_name, + token=self.token, + ) logger.info("Telemetry log handler attached") async def _retry_rpc( @@ -942,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: @@ -978,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: @@ -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,17 @@ 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: + 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, @@ -1546,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 57ad81c2c..447ac92ee 100644 --- a/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py +++ b/python/packages/jumpstarter/jumpstarter/exporter/exporter_telemetry_test.py @@ -35,6 +35,11 @@ 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.exporter_name = "test-exporter" + exp.token = "test-token" + exp.namespace = "" tls = MagicMock() tls.insecure = False @@ -100,6 +105,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 +119,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 +138,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 +212,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