JEP-0013 Phase 3 - operator image - #1061
Conversation
Add the MetricsStream protocol and Go hub so Prometheus can scrape merged exporter OpenMetrics from telemetry without an exporter client yet. Generated Python stubs are included for proto consistency. Co-authored-by: Cursor <cursoragent@cursor.com>
Stop silently dropping unparseable exporter snapshots. Log the exporter and error, and increment jumpstarter_metrics_parse_errors_total so reverse-scrape omissions are visible on the same /metrics response. Co-authored-by: Cursor <cursoragent@cursor.com>
…e 3) Build /telemetry into the controller image and have the operator mount cert-manager TLS, advertise the CA, and expose scrape flags plus the HTTP metrics port so reverse-scrape can run in-cluster.
|
Warning Review limit reachedNext included review available in 54 minutes. View limit detailsLimit details: You’ve used all 2 included reviews currently available. You've used all free OSS reviews for now. Wait for the free limit to reset to keep reviewing this public repository. Review configuration: ⚙️ Run configurationConfiguration used: Organization UI Review profile: CHILL Plan: Team Run ID: ⛔ Files ignored due to path filters (2)
📒 Files selected for processing (23)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
raballew
left a comment
There was a problem hiding this comment.
Just a few comments for now. Waiting for merge of the first PR.
| return []string{ | ||
| fmt.Sprintf("--grpc-bind=:%d", telemetryPort), | ||
| fmt.Sprintf("-metrics-bind-address=:%d", telemetryMetricsPort), | ||
| fmt.Sprintf("-scrape-timeout=%s", timeout), | ||
| fmt.Sprintf("-driver-type-enum=%s", strings.Join(enum, ",")), | ||
| fmt.Sprintf("-exemplar-keys=%s", strings.Join(keys, ",")), | ||
| } |
There was a problem hiding this comment.
Mixed use of dash and double-dash
| // Max wait for parallel exporter MetricsStream responses during a /metrics fan-out. | ||
| // Should be lower than the Prometheus scrape_timeout. | ||
| // +kubebuilder:default="7s" | ||
| ScrapeTimeout *metav1.Duration `json:"scrapeTimeout,omitempty"` |
There was a problem hiding this comment.
Also set an upper limit to avoid keeping connections open for very long times.
| srv := grpc.NewServer( | ||
| grpc.Creds(creds), | ||
| grpc.ChainUnaryInterceptor(recovery.UnaryServerInterceptor()), | ||
| grpc.ChainStreamInterceptor(recovery.StreamServerInterceptor()), | ||
| ) |
There was a problem hiding this comment.
Adding grpc.MaxRecvMsgSize to the server options would provide an explicit limit especially if you are running 1000s of exporters that are scraped concurrently (each with up to 4 MB per message). Otherwise this might cause some OOM killed processes.
| 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") |
There was a problem hiding this comment.
entry.Message, entry.Component, entry.Severity, entry.Lease, entry.Client, entry.Operation, entry.Result, entry.DriverType have no size bound. Should we truncate them too?
| // Allowlist of keys to include in Prometheus exemplars. Unlisted keys are omitted. | ||
| // +kubebuilder:default={"client","lease_id"} | ||
| ExemplarKeys []string `json:"exemplarKeys,omitempty"` | ||
|
|
||
| // Allowed driver_type label values. Unlisted types are remapped to "other". | ||
| // +kubebuilder:default={"power","storage","network","serial","console","video","composite"} | ||
| DriverTypeEnum []string `json:"driverTypeEnum,omitempty"` | ||
|
|
There was a problem hiding this comment.
A user with CR write access can set either list to an arbitrarily large array for ExemplarKeys []string and DriverTypeEnum []string. Adding // +kubebuilder:validation:MaxItems=<insert a good limit here> and // +kubebuilder:validation:MaxLength=<insert a good limit here> per item would enforce a bound.
Summary
JEP-0013 Phase 3 PR B**: make the MetricsStream hub from PR A runnable in-cluster.
Depends on PR A: #1060
Do not merge until A is on
main, then rebase this branch ontomain. Related PRs C–E will be linked here as they are opened./telemetryincontroller/Containerfileandmake build/docker-build-ci(operator alreadyCommand: ["/telemetry"];mainstill only shipsmanager+router)./healthzand/readyzprobes on that port (replacing TCP probes on gRPC:9093).spec.telemetry.metrics.scrapeTimeout,driverTypeEnum,exemplarKeys. Defaults match JEP-0013 (7s, the driver-type enum,client,lease_id). NoServiceMonitor(Phase 5).GRPC_TELEMETRY_ENDPOINTso the telemetry process advertises the in-cluster Service DNS.TLS: #1023 already landed telemetry TLS on
main. This PR keeps that path (cert-manager or manualspec.telemetry.grpc.tls.certSecret, rolling-restart hash, non-fatal missing CA). It does not reimplement TLS.Lab: in-cluster reverse-scrape of
GET /metricson Service port 8080 worked (port-forward). After driver ops, Python OpenMetrics exemplars still fail GoparseMetricFamilies; the hub from A omits that snapshot and incrementsjumpstarter_metrics_parse_errors_total.spec.telemetry.metrics.exemplarKeysonly allowlists keys on the merge path after a successful parse — it does not fix exemplar decode. That remains A / JEP DD-3, not this PR.DEMO
The lab walkthrough for the stacked Phase 3 work (A+B+C+D+E plus lab-only Route) is on the fork demo branch, not this PR:
See MetricsStream during the lease for the DD-3 exemplar limitation observed on
:8080.How this PR fits the series
jep-0013-phase3-metricsstreamjep-0013-phase3-operator-imagejep-0013-phase3-exporter-metricsstreamjep-0013-phase3-loki-pushjep-0013-phase3-client-pushlogsUnique work vs PR A: the operator/image commit on this branch (
435eebee). If GitHub shows A's commits, that is stacking againstmain; review that unique commit only.Out of scope
identity(C) — after#1058, identity must beexporter_name, notMetadata.nameServiceMonitor(Phase 5)NOTE
Open #1027 also touches telemetry operator files (log-ingest e2e).