Skip to content

feat(target-allocator): per-node allocation and annotation-based SM/PM routing - #415

Open
Aakash-Dantre wants to merge 13 commits into
aws:mainfrom
Aakash-Dantre:smpm/ta-routing-per-node
Open

Aakash-Dantre wants to merge 13 commits into
aws:mainfrom
Aakash-Dantre:smpm/ta-routing-per-node

Conversation

@Aakash-Dantre

@Aakash-Dantre Aakash-Dantre commented Oct 1, 2026 •

Copy link
Copy Markdown

Summary

Two Target Allocator (TA) capabilities for ServiceMonitor/PodMonitor scraping:

  1. Per-node allocation - a per-node allocation strategy that assigns each target to the agent on the same node as the scraped pod, so scrape traffic stays node-local instead of crossing nodes and AZs.
  2. Annotation-based routing - a monitor annotated cloudwatch.aws.amazon.com/scraper: cluster-scraper is scraped only by the cluster-scraper agent's TA; every other monitor only by the per-node agent's TA.

Supersedes #398 and #399, rebased onto current main. #399 already contained all of #398. Built on the same base as #414, which includes all of #386.

What

Per-node allocation:

  • New per-node strategy in allocation/: assigns each target to the collector on its node. Node-less targets fall back to a configurable strategy, or stay unassigned with a one-time warning if no fallback is set.
  • Collector watch tracks pod scheduling and termination, so targets are reassigned when collectors appear or disappear.

Annotation routing:

  • New spec.targetAllocator.prometheusCR.scraperRole field (TA config key scraper_role).
  • Client-side filter in the TA watcher (annotationRoleMatches, selectsMonitor), applied during monitor discovery in LoadConfig. The cluster-scraper role keeps only annotated monitors and the default (empty) role keeps only unannotated ones, so every monitor is owned by exactly one agent.
  • Routing is by annotation rather than label because both TAs already list every monitor and filter client-side, so a server-side label selector would gain nothing.

Webhook:

  • allocationStrategy: per-node is rejected unless mode: daemonset (as upstream does), since it needs one collector per node.
  • The "we do not recommend enabling Target Allocator when not running as a StatefulSet" warning is no longer shown for the two supported setups: per-node on a DaemonSet, and an agent with a scraperRole. Other non-StatefulSet setups still warn.

Included base fixes (from #386)

The first ten commits are all of #386 by @musa-asad, at its current head (2d21b9b6), replayed onto current main with authorship kept. #386 is no longer being driven, so it is folded in here instead of landing separately. They cover:

  • registering the --enable-prometheus-cr-watcher flag (the TA otherwise exits with "unknown flag" when the operator passes it), ORed with prometheus_cr.enabled, plus the Prometheus namespace and evaluation-interval settings the watcher needs
  • defaulting scrape_protocols on config load (test)
  • rolling collector pods when the Prometheus config changes
  • trimming and validating the collector-namespace lookup, including OTELCOL_NAMESPACE, with tests
  • documenting the flag's enable-only behaviour, with a test

While replaying onto main:

Overlap with #414

Both PRs change cmd/amazon-cloudwatch-agent-target-allocator/watcher/promOperator.go. In #414 the SM/PM informers can be absent (if smInformer != nil), while this PR adds the selectsMonitor filter in LoadConfig. Whichever PR merges second needs a rebase that keeps both behaviours: skip an informer that is absent, and apply the routing filter to the ones that are present. That merge needs review and should not be resolved mechanically.

Testing

Unit tests, added in this PR:

  • Webhook: TestTargetAllocatorModeValidation (per-node accepted on daemonset, rejected on deployment and statefulset; no warning for per-node daemonset or a scraper role; warning kept for consistent-hashing outside a StatefulSet)
  • Per-node: TestPerNodeAssignsToMatchingNode, TestPerNodeFallbackAssignsNodelessTarget, TestPerNodeLeavesNodelessTargetUnassigned, TestPerNodeReallocatesWhenCollectorAppears, TestPerNodeCollectorRemovalReassigns, TestPerNodeTwoCollectorsSameNodeTieBreak, TestPerNodeFallbackDuringCollectorChange, plus 11 more TestPerNode* cases; Test_runWatch_UnscheduledThenScheduled, Test_runWatch_TerminatingPodReleased, TestGetNodeName
  • Routing: TestAnnotationRoleMatches (8 cases plus the partition invariant), TestSelectsMonitor, TestLoadConfigScraperRouting
  • From the base fixes: TestScrapeProtocolsDefaultedOnLoad, TestPrometheusConfigChangeBumpsHash, TestEmptyPrometheusHashUnchanged

PR Checklist

  • Commits are squashed into a logical, reviewable set (one commit for a single change) - 13 commits: the ten from Fix Target Allocator startup crashes, PrometheusCR watcher, and Prometheus config pod restart #386 (kept as separate commits, with authorship), per-node allocation, annotation routing, and the webhook change, one commit per change so each can be reviewed on its own.
  • Commits and PR description comply with the contribution guidelines - no internal references in commit messages or this description.
  • make passes locally - go build ./... and make vet pass; go test for cmd/amazon-cloudwatch-agent-target-allocator/... and internal/manifests/collector/... passes. The full make test needs envtest assets that are not available locally; CI runs it.
  • All GitHub Actions checks on the PR are passing - not yet run on this PR; fork PR workflows need a maintainer to approve them. The same code passes on the fork: PR Build (lint + unit tests) and Operator Integration Test (all 5 jobs). The fork run was on the previous revision of this branch, which differs only in comment and test-name wording. Before the rebase onto current main, the integration test's NamespaceAnnotationsTest failed; it passes now.
  • Integration test evidence - Operator Integration Test passed on the previous revision (same code): Daemonset, Deployment, Statefulset and Namespace annotations, and Instrumentation.
  • New or updated integration test coverage - test(otel/pernode): per-node, CRD-bundling and scraper-routing E2E suites amazon-cloudwatch-agent-test#773: TestPerNodeAllocation, TestPerNodeCoverageAcrossNodes, TestScraperRoleWiring, TestClusterScraperTargetAllocatorHealthy, TestAnnotationRoutingPartition.
  • New functionality has unit tests; behaviour changes have a reproducing test - see Testing.
  • Config translation changes include updated golden files - N/A: no translator changes.
  • Breaking or customer-visible changes are called out in the PR description - per-node is opt-in (the TA default stays consistent-hashing), and scraperRole is a new optional field. Behaviour changes: a TA on the default role skips monitors annotated cloudwatch.aws.amazon.com/scraper: cluster-scraper (unannotated monitors are handled as before); the webhook now rejects per-node outside mode: daemonset; and the StatefulSet warning is no longer shown for per-node DaemonSets or agents with a scraperRole. New: the per-node strategy, the scraperRole field, and the cloudwatch.aws.amazon.com/scraper annotation.

// cloudwatch.aws/scraper: cluster-scraper is scraped only by the cluster-scraper agent's Target
// Allocator; all others are scraped only by the per-node agent's Target Allocator.
const (
ScraperAnnotationKey = "cloudwatch.aws/scraper"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The operator's other CloudWatch-specific annotations (auto-annotate-, inject-jmx-, restartedAt) and the CRD API group all use cloudwatch.aws.amazon.com/, and I can't find any existing key under cloudwatch.aws/ in the operator or the chart. Customers will set this one on their own ServiceMonitors and PodMonitors, so it's hard to rename after release. Was cloudwatch.aws/ intentional, or should it be cloudwatch.aws.amazon.com/scraper?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not intentional. I'll switch to cloudwatch.aws.amazon.com/scraper across #415, aws-observability/helm-charts#376 and aws/amazon-cloudwatch-agent-test#773 unless anyone objects.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done: renamed to cloudwatch.aws.amazon.com/scraper in #415 (89d4f7f), aws-observability/helm-charts#376 and aws/amazon-cloudwatch-agent-test#773.


// AmazonCloudWatchAgentTargetAllocatorAllocationStrategyPerNode targets will be allocated to the collector running on the same node as the target.
// Targets without a resolvable node fall back to the configured fallback strategy (consistent-hashing).
AmazonCloudWatchAgentTargetAllocatorAllocationStrategyPerNode AmazonCloudWatchAgentTargetAllocatorAllocationStrategy = "per-node"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should the webhook reject per-node when mode isn't daemonset? Upstream does (internal/webhook/collector_webhook.go#L338 (https://github.com/open-telemetry/opentelemetry-operator/blob/main/internal/webhook/collector_webhook.go)). Here it's accepted, and on a Deployment or StatefulSet each
replica takes every target on its own node while the rest go through the fallback, so load ends up uneven.

Related: with helm-charts#376 the DaemonSet (per-node) and the Deployment cluster-scraper both enable the TA, so the existing warning in collector_webhook.go ("we do not recommend enabling Target Allocator when not running as a StatefulSet") will fire on every apply. What are our plans for that?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed, I'll reject per-node unless mode: daemonset, matching upstream. For the warning, I'd skip it for per-node on a DaemonSet or when scraperRole is set, and keep it otherwise. OK?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in e6fb8f4: per-node is rejected unless mode: daemonset, and the warning is skipped for per-node DaemonSets and agents with a scraperRole. Checked with server-side dry runs on a live cluster.

}
sort.Strings(mapping)
noNode := len(pn.collectors) - len(pn.collectorByNode)
pn.log.Info("per-node: collector node index rebuilt",

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This logs every node-to-collector pair at Info on every collector change. On a large cluster with frequent autoscaling, that's a very long log line on every DaemonSet pod add or remove. Could the full mapping go to V(1) and Info keep just the counts?

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done in f8f0ac6: Info keeps the counts, the full mapping is at V(1).

@Aakash-Dantre
Aakash-Dantre force-pushed the smpm/ta-routing-per-node branch 2 times, most recently from 6574dcf to f8f0ac6 Compare October 2, 2026 13:40
@Aakash-Dantre
Aakash-Dantre force-pushed the smpm/ta-routing-per-node branch from f8f0ac6 to e6fb8f4 Compare October 2, 2026 14:46
musa-asad and others added 13 commits October 7, 2026 14:41
…er startup

The target-allocator declared the enable-prometheus-cr-watcher flag name as a
constant but never registered it on the flag set, while the operator passes
--enable-prometheus-cr-watcher whenever PrometheusCR.enabled is true. Because
args are parsed with pflag.ExitOnError, the unregistered flag caused the binary
to print 'unknown flag' and exit(2), putting the target-allocator pod into
CrashLoopBackOff.

This change registers the flag and ORs it with the YAML prometheus_cr.enabled
setting, and makes three related changes the watcher needs once it starts:

- promOperator: set a non-empty Namespace on the synthetic Prometheus object so
  the prometheus-operator config generator no longer panics with
  'namespace can't be empty' in store.ForNamespace.
- promOperator: set EvaluationInterval so the generated config does not render an
  empty global.evaluation_interval, which the prometheus config parser rejects
  with 'empty duration string'.
- main: create and register service-discovery metrics and pass them to
  discovery.NewManager; passing a nil sdMetrics map makes every SD provider fail
  to register, yielding zero discovered targets.

(cherry picked from commit 1376451)
Add a test asserting that loading a Target Allocator config whose
static scrape job omits scrape_protocols still yields a non-empty
ScrapeProtocols on every loaded scrape config. This is defaulted by the pinned
Prometheus library during yaml.UnmarshalStrict into the prometheus Config type,
so the distributed /scrape_configs payload is never empty and the agent's
prometheus-receiver validation passes. The test fails fast if a future
dependency or load-path change drops this defaulting.

(cherry picked from commit 0405f4d)
The pod-template restart-trigger sha256 was computed from Spec.Config only,
so a change to Spec.Prometheus (rendered into a separate ConfigMap) left the
pod template byte-identical and the workload controller did not roll the pods.

Fold the serialized Spec.Prometheus (PrometheusConfig.Yaml()) into the hash
input when it is non-empty, so a Prometheus-only change bumps the pod-template
annotation and triggers a rolling restart, matching agent-config behavior.
When no Prometheus config is set the hash input is byte-identical to the agent
config alone, leaving non-Prometheus agents unaffected.

(cherry picked from commit d61d693)
configHashInput's error branch is unreachable from the CRD path:
Spec.Prometheus.Config is an *AnyConfig whose Object is a
map[string]interface{} populated by JSON decoding, and every type that
produces is encodable by gopkg.in/yaml.v3. Logging a warning there
implied a runtime condition operators should watch for and act on, which
is misleading.

Replace the slog.Warn with a comment that records three things a reader
cannot recover from the code: that the branch is unreachable from the
CRD path, why the sentinel is deliberately a constant rather than
derived from the error (two different unserializable specs hash equal,
so a Prometheus-only change would not roll pods while the error
persists), and that the hash's stability across operator restarts rests
on gopkg.in/yaml.v3 v3.0.1 sorting map keys -- a yaml library
preserving insertion order here would roll every agent pod on restart.

Removing the call also removes the log/slog import, which was the
package's only use of it.

The non-error path, the IsEmpty() gate and the \x00 separator are
untouched, so the pinned config-hash values asserted in
annotations_test.go are unchanged.

(cherry picked from commit ae7a53c)
The namespace fed to the Prometheus config generator was taken verbatim
from OTELCOL_NAMESPACE, or from the service account namespace file when
that env var was empty. Neither source was trimmed, so a trailing newline
or surrounding whitespace yielded a string that is not a valid Kubernetes
namespace. The file branch also gated on len(ns) > 0, which treats a
whitespace-only file as a real value and so skips the
defaultCollectorNamespace fallback.

Both sources are now whitespace-trimmed, and a value that is blank after
trimming falls through to the next source: env var, then the service
account file, then defaultCollectorNamespace. The resolution moves out of
NewPrometheusCRWatcher into resolveCollectorNamespace, with the service
account path held in a package variable so a test can redirect the read.
The log line still fires exactly once when the env var is unset, and
still reports the resolved value.

(cherry picked from commit c5dee74)
…fault

Adds coverage for the two behaviors NewPrometheusCRWatcher depends on and
that nothing else in the watcher package exercises.

TestResolveCollectorNamespace pins the resolution chain across seven cases:
OTELCOL_NAMESPACE wins when set, both the env var and the service account
namespace file are whitespace-trimmed, a whitespace-only value in either
source falls through to the next source rather than becoming the namespace,
and an absent or blank file falls back to the default. This guards against a
future edit dropping a TrimSpace, reverting to a length-only emptiness check,
or reordering the fallback chain -- any of which yields an empty or blank
namespace, and the config generator panics on a non-empty namespace
requirement.

TestNewPrometheusCRWatcherDefaultsEvaluationInterval pins EvaluationInterval
defaulting to ScrapeInterval. Without that assignment the generated config
carries an empty evaluation_interval, which fails to parse.

The tests live in a new file because promOperator_test.go imports neither
k8s.io/client-go/rest, the allocator config package, nor logr. The test
redirects serviceAccountNamespacePath and restores it via t.Cleanup, so the
package's later test files do not inherit a path into a removed temp dir.
The constructor is pointed at an unroutable host to confirm it does no dialing.

(cherry picked from commit 479e251)
The collector client resolves its watch namespace once at process init
from OTELCOL_NAMESPACE and passes the result to Pods(ns).List and
Pods(ns).Watch. A value carrying surrounding whitespace is not a valid
namespace name, so those calls matched nothing, no collectors were
discovered, and allocation stayed empty with no error surfaced.

Trim the value so padded input resolves to the intended namespace,
matching how the namespace is now resolved on the watcher path.

(cherry picked from commit 02a4bb5)
LoadFromCLI ORs the CLI flag with prometheus_cr.enabled from the config
file, so the flag can only turn the PrometheusCR watcher on. Passing
=false has no effect when the config file enables it, which the
previous help text did not convey. State the OR behavior in the flag's
usage string so the asymmetry is discoverable from --help.

No behavior change: only the usage string passed to flagSet.Bool changes.

(cherry picked from commit ebd1629)
LoadFromCLI OR-es the --enable-prometheus-cr-watcher flag with
prometheus_cr.enabled from the config file, so the flag can only enable
the watcher: passing =false cannot disable one the config file turns on.
Nothing covered that, and replacing the || with a plain assignment would
break only the config-file-true rows while leaving every existing test in
the package green.

Add a table-driven test over the four (config file, CLI) combinations,
plus a committed testdata kubeconfig so the test is hermetic. LoadFromCLI
builds a client config before it reaches the OR, so with no
--kubeconfig-path it falls back to rest.InClusterConfig() and returns
"unable to load in-cluster configuration" early. That passes on a
developer host that happens to have a kubeconfig and fails in CI, so the
flag pointing at testdata/kubeconfig_test.yaml is required rather than
decorative. The fixture's server is 127.0.0.1 and is never dialed.

pflag.ContinueOnError is a test-level choice so a parse error fails the
test instead of exiting the process; the production getFlagSet call keeps
pflag.ExitOnError.

(cherry picked from commit 95adef2)
Extract the hardcoded service account namespace path literal into
defaultServiceAccountNamespacePath per review feedback.

The package-level serviceAccountNamespacePath var remains as the test seam
that redirects the read in unit tests.

There is no behavior change.

(cherry picked from commit 2d21b9b)
Adds a per-node allocation strategy so each CloudWatch Agent scrapes only the
ServiceMonitor/PodMonitor targets on its own node, eliminating cross-node and
cross-AZ scrape traffic. Targets with no resolvable node fall back to
consistent-hashing so every target is still scraped, and unassigned targets are
surfaced via a gauge.

Also changes collector registration for ALL strategies, including the existing
consistent-hashing default: the initial List now skips pods with an empty
spec.NodeName, watch.Modified is handled so a pod scheduled after Add is
picked up, and pods are dropped as soon as they carry a DeletionTimestamp
rather than on Deleted. This releases a terminating agent's targets
immediately instead of at the end of its grace period, at the cost of more
target churn during rollouts.

Verified: go build ./... clean; make impi and make checklicense pass; unit
tests pass across the target-allocator and manifests packages.
…per annotation

Partitions ServiceMonitor/PodMonitor discovery across CloudWatch Agents by the
cloudwatch.aws.amazon.com/scraper annotation on the monitor CR. A monitor annotated
cluster-scraper is scraped only by the cluster-scraper agent's Target
Allocator; all others only by the per-node agent's. The two roles are
complementary, so every monitor is owned by exactly one agent with no
double-scrape and no gap.

Adds spec.targetAllocator.prometheusCR.scraperRole (Target Allocator config
key scraper_role) and a client-side annotation filter applied during monitor
discovery in LoadConfig.

Note: scraperRole is constrained by an enum with no empty member, so the
default per-node role is expressed by omitting the field; setting
scraperRole: "" explicitly is rejected by the API server.

Companion: aws-observability/helm-charts#339.

Verified: go build ./... clean; make impi and make checklicense pass; unit
tests pass across the target-allocator and manifests packages.
per-node assigns each target to the collector on the target's node, so it
needs one collector per node. Reject allocationStrategy per-node unless
mode is daemonset, as upstream does.

Also stop warning about running the Target Allocator outside a StatefulSet
for the two supported setups: per-node on a DaemonSet, and an agent with a
prometheusCR scraperRole. Other non-StatefulSet setups still warn.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants