mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
Delete ate.scheduler.eligible_workers (#1804)
Closes #1802. ## What this does Deletes the `ate.scheduler.eligible_workers` histogram. Nothing replaces it in this PR. ## Why **It is expensive on the resume path.** `Schedule` built a second slice of the matching workers that scheduling never read, walked it again, and repeated `HasRoom` for each entry. `HasRoom` parses resource quantity strings up to three times for each worker. The metric roughly doubled the per-worker arithmetic of a placement. The filter loop is now one pass over the fleet and builds one slice. **A metric cannot answer the question it was built for.** #564 wanted to tell a full fleet apart from an empty intersection of the constraints, such as a cordoned node behind a node requirement. That needs the selectors, the node requirement and the actor. `docs/metrics/substrate.yaml` bars all three from a metric label and sends them to logs and spans, so no shape of this instrument reaches the answer. A log record at the rejection is the replacement the issue names, and it is not in this PR. ## Also removed `ate.scheduling.constraint` goes with the instrument. No other signal used the attribute. ## Scope of the change - `cmd/ateapi/internal/scheduling/metrics.go` deleted, along with the `WithMeter` option and the second filter pass in `Schedule`. - The instrument and the attribute removed from `docs/metrics/registry/metrics.yaml`. The `pool-keys-paired` exception in `docs/metrics/substrate.yaml` dropped, because the empty pool pair was this instrument's alone. - `docs/observability.md`: the table row and the two label notes removed. - The e2e collector check and the label assertions for the instrument removed. - `TestSchedule_EligibleWorkersMetric` removed. The behaviors it covered (draining workers, a sandbox class mismatch, busy workers) are already in the `TestSchedule` table. ## Testing `go test ./cmd/... ./internal/...` passes. `make verify` passes except `hack/verify/metrics.sh`, which needs Weaver or Docker. Neither is available on this machine, so CI is what checks the registry.
This commit is contained in:
@@ -132,7 +132,7 @@ func NewActorWorkflow(
|
||||
return &ActorWorkflow{
|
||||
store: store,
|
||||
workerCache: workerCache,
|
||||
scheduler: scheduling.New(workerCache, scheduling.WithMeter(otel.Meter("ateapi"))),
|
||||
scheduler: scheduling.New(workerCache),
|
||||
dialer: dialer,
|
||||
sandboxConfigLister: sandboxConfigLister,
|
||||
storageClassLister: storageClassLister,
|
||||
|
||||
@@ -1,109 +0,0 @@
|
||||
// Copyright 2026 Google LLC
|
||||
//
|
||||
// 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 scheduling
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
"go.opentelemetry.io/otel/metric"
|
||||
)
|
||||
|
||||
// Metric identifier for tracking eligible worker counts.
|
||||
const eligibleWorkersMetric = "ate.scheduler.eligible_workers"
|
||||
|
||||
// newEligibleWorkers creates the ate.scheduler.eligible_workers histogram instrument against meter.
|
||||
// Returns (nil, nil) if meter is nil. If registration fails, an error is returned.
|
||||
func newEligibleWorkers(meter metric.Meter) (metric.Int64Histogram, error) {
|
||||
if meter == nil {
|
||||
return nil, nil
|
||||
}
|
||||
hist, err := meter.Int64Histogram(
|
||||
eligibleWorkersMetric,
|
||||
metric.WithUnit("{worker}"),
|
||||
metric.WithDescription("Number of eligible workers available during scheduling given the constraint filters."),
|
||||
metric.WithExplicitBucketBoundaries(0, 1, 2, 3, 5, 10, 20, 50, 100, 250),
|
||||
)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create %s histogram: %w", eligibleWorkersMetric, err)
|
||||
}
|
||||
return hist, nil
|
||||
}
|
||||
|
||||
// WithMeter configures the meter used to create telemetry instruments for the scheduler.
|
||||
// If meter is nil, the option is a no-op. If instrument creation fails, an error is explicitly logged.
|
||||
func WithMeter(meter metric.Meter) Option {
|
||||
return func(s *scheduler) {
|
||||
hist, err := newEligibleWorkers(meter)
|
||||
if err != nil {
|
||||
slog.Error("Failed to register ate.scheduler.eligible_workers histogram", "metric", eligibleWorkersMetric, "error", err)
|
||||
return
|
||||
}
|
||||
s.eligibleWorkers = hist
|
||||
}
|
||||
}
|
||||
|
||||
// Records candidate worker counts grouped by WorkerPool namespace and WorkerPool name,
|
||||
// stamping sandbox class and scheduling constraint attributes on all histogram datapoints.
|
||||
func (s *scheduler) recordEligibleWorkers(ctx context.Context, matching []*ateapipb.Worker, constraints Constraints) {
|
||||
if s.eligibleWorkers == nil {
|
||||
return
|
||||
}
|
||||
|
||||
// Sandbox class and constraint are constant across every key: Applies requires
|
||||
// an exact class match, and the classification is per call. They belong on the
|
||||
// Record call, not in the key.
|
||||
type key struct{ namespace, pool string }
|
||||
eligibleByPool := make(map[key]int64)
|
||||
for _, w := range matching {
|
||||
k := key{namespace: w.GetWorkerNamespace(), pool: w.GetWorkerPool()}
|
||||
if _, ok := eligibleByPool[k]; !ok {
|
||||
eligibleByPool[k] = 0
|
||||
}
|
||||
if s.HasRoom(w, constraints) {
|
||||
eligibleByPool[k]++
|
||||
}
|
||||
}
|
||||
|
||||
// No pool matched the constraints at all. Emit a single zero-valued series so
|
||||
// "nothing is schedulable" stays visible; empty namespace/pool marks it. The
|
||||
// label set matches the per-pool series, so dashboards need no special case.
|
||||
if len(eligibleByPool) == 0 {
|
||||
eligibleByPool[key{}] = 0
|
||||
}
|
||||
|
||||
constraintStr := classifyConstraint(constraints)
|
||||
for k, count := range eligibleByPool {
|
||||
s.eligibleWorkers.Record(ctx, count, metric.WithAttributes(
|
||||
ateattr.WorkerPoolNamespaceKey.String(k.namespace),
|
||||
ateattr.WorkerPoolNameKey.String(k.pool),
|
||||
ateattr.SandboxClassAttribute(constraints.SandboxClass),
|
||||
ateattr.SchedulingConstraintKey.String(constraintStr),
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
func classifyConstraint(c Constraints) string {
|
||||
if len(c.RequiredNodes) > 0 {
|
||||
return ateattr.ConstraintRequiredNodes
|
||||
}
|
||||
if (c.TemplateSelector != nil && !c.TemplateSelector.Empty()) || (c.ActorSelector != nil && !c.ActorSelector.Empty()) {
|
||||
return ateattr.ConstraintSelector
|
||||
}
|
||||
return ateattr.ConstraintNone
|
||||
}
|
||||
@@ -24,7 +24,6 @@ import (
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
"go.opentelemetry.io/otel/metric"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
)
|
||||
|
||||
@@ -77,8 +76,6 @@ type scheduler struct {
|
||||
// intn returns a uniformly distributed random value in [0,n).
|
||||
// Defaults to the global math/rand source
|
||||
intn func(n int) int
|
||||
// Records the number of eligible workers available during scheduling.
|
||||
eligibleWorkers metric.Int64Histogram
|
||||
}
|
||||
|
||||
// Option configures the Scheduler returned by New.
|
||||
@@ -106,21 +103,13 @@ func (s *scheduler) Schedule(ctx context.Context, constraints Constraints) (*ate
|
||||
return nil, fmt.Errorf("while listing workers: %w", err)
|
||||
}
|
||||
|
||||
matching := make([]*ateapipb.Worker, 0, len(workers))
|
||||
var candidates []*ateapipb.Worker
|
||||
for _, worker := range workers {
|
||||
if !s.Applies(worker, constraints) {
|
||||
continue
|
||||
}
|
||||
matching = append(matching, worker)
|
||||
if s.HasRoom(worker, constraints) {
|
||||
if s.Applies(worker, constraints) && s.HasRoom(worker, constraints) {
|
||||
candidates = append(candidates, worker)
|
||||
}
|
||||
}
|
||||
|
||||
// Record telemetry on the number of eligible workers per pool/namespace before returning
|
||||
s.recordEligibleWorkers(ctx, matching, constraints)
|
||||
|
||||
if len(candidates) == 0 {
|
||||
return nil, ErrNoCapacity
|
||||
}
|
||||
|
||||
@@ -19,11 +19,8 @@ import (
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/agent-substrate/substrate/internal/ateattr"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
sdkmetric "go.opentelemetry.io/otel/sdk/metric"
|
||||
"go.opentelemetry.io/otel/sdk/metric/metricdata"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
)
|
||||
|
||||
@@ -419,301 +416,3 @@ func TestAppliesIgnoresRoom(t *testing.T) {
|
||||
|
||||
// firstIntn always picks the first candidate, making Schedule deterministic.
|
||||
func firstIntn(int) int { return 0 }
|
||||
|
||||
func workerWithPool(pod, ns, pool, class, node string, lbls map[string]string, opts ...func(*ateapipb.Worker)) *ateapipb.Worker {
|
||||
w := worker(pod, class, node, lbls, opts...)
|
||||
w.WorkerNamespace = ns
|
||||
w.WorkerPool = pool
|
||||
return w
|
||||
}
|
||||
|
||||
func TestSchedule_EligibleWorkersMetric(t *testing.T) {
|
||||
t.Run("records histogram metric with namespaced attributes for candidates", func(t *testing.T) {
|
||||
reader := sdkmetric.NewManualReader()
|
||||
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
|
||||
meter := provider.Meter("test")
|
||||
|
||||
flt := fleet{
|
||||
workerWithPool("w-1", "ns-a", "pool-1", "gvisor", "node-a", nil),
|
||||
workerWithPool("w-2", "ns-a", "pool-1", "gvisor", "node-a", nil, assigned("demo", "a")),
|
||||
workerWithPool("w-3", "ns-b", "pool-2", "gvisor", "node-b", nil),
|
||||
}
|
||||
|
||||
s := New(flt, WithIntn(firstIntn), WithMeter(meter))
|
||||
_, err := s.Schedule(context.Background(), Constraints{SandboxClass: "gvisor"})
|
||||
if err != nil {
|
||||
t.Fatalf("Schedule() error = %v", err)
|
||||
}
|
||||
|
||||
var rm metricdata.ResourceMetrics
|
||||
if err := reader.Collect(context.Background(), &rm); err != nil {
|
||||
t.Fatalf("reader.Collect() error = %v", err)
|
||||
}
|
||||
|
||||
foundMetric := false
|
||||
for _, sm := range rm.ScopeMetrics {
|
||||
for _, m := range sm.Metrics {
|
||||
if m.Name == "ate.scheduler.eligible_workers" {
|
||||
foundMetric = true
|
||||
histogram, ok := m.Data.(metricdata.Histogram[int64])
|
||||
if !ok {
|
||||
t.Fatalf("metric Data is %T, want metricdata.Histogram[int64]", m.Data)
|
||||
}
|
||||
if len(histogram.DataPoints) == 0 {
|
||||
t.Fatalf("got 0 DataPoints for ate.scheduler.eligible_workers")
|
||||
}
|
||||
for _, dp := range histogram.DataPoints {
|
||||
attrs := dp.Attributes
|
||||
ns, _ := attrs.Value(ateattr.WorkerPoolNamespaceKey)
|
||||
pool, _ := attrs.Value(ateattr.WorkerPoolNameKey)
|
||||
class, _ := attrs.Value(ateattr.SandboxClassKey)
|
||||
constraint, _ := attrs.Value(ateattr.SchedulingConstraintKey)
|
||||
if class.AsString() != "gvisor" {
|
||||
t.Errorf("got sandbox class %q, want %q", class.AsString(), "gvisor")
|
||||
}
|
||||
if constraint.AsString() != "none" {
|
||||
t.Errorf("got constraint %q, want %q", constraint.AsString(), "none")
|
||||
}
|
||||
if ns.AsString() == "ns-a" && pool.AsString() == "pool-1" {
|
||||
if dp.Count != 1 || dp.Sum != 1 {
|
||||
t.Errorf("pool-1 datapoint count=%d sum=%d, want count=1 sum=1", dp.Count, dp.Sum)
|
||||
}
|
||||
}
|
||||
if ns.AsString() == "ns-b" && pool.AsString() == "pool-2" {
|
||||
if dp.Count != 1 || dp.Sum != 1 {
|
||||
t.Errorf("pool-2 datapoint count=%d sum=%d, want count=1 sum=1", dp.Count, dp.Sum)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !foundMetric {
|
||||
t.Fatalf("ate.scheduler.eligible_workers metric not found")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("records 0 eligible workers when fleet has no capacity", func(t *testing.T) {
|
||||
reader := sdkmetric.NewManualReader()
|
||||
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
|
||||
meter := provider.Meter("test")
|
||||
|
||||
flt := fleet{
|
||||
workerWithPool("w-busy", "ns-a", "pool-1", "gvisor", "node-a", nil, assigned("demo", "a")),
|
||||
}
|
||||
|
||||
s := New(flt, WithIntn(firstIntn), WithMeter(meter))
|
||||
_, err := s.Schedule(context.Background(), Constraints{SandboxClass: "gvisor"})
|
||||
if !errors.Is(err, ErrNoCapacity) {
|
||||
t.Fatalf("Schedule() error = %v, want ErrNoCapacity", err)
|
||||
}
|
||||
|
||||
var rm metricdata.ResourceMetrics
|
||||
if err := reader.Collect(context.Background(), &rm); err != nil {
|
||||
t.Fatalf("reader.Collect() error = %v", err)
|
||||
}
|
||||
|
||||
foundMetric := false
|
||||
for _, sm := range rm.ScopeMetrics {
|
||||
for _, m := range sm.Metrics {
|
||||
if m.Name == "ate.scheduler.eligible_workers" {
|
||||
foundMetric = true
|
||||
histogram, ok := m.Data.(metricdata.Histogram[int64])
|
||||
if !ok {
|
||||
t.Fatalf("metric Data is %T, want metricdata.Histogram[int64]", m.Data)
|
||||
}
|
||||
if len(histogram.DataPoints) == 0 {
|
||||
t.Fatalf("got 0 DataPoints")
|
||||
}
|
||||
dp := histogram.DataPoints[0]
|
||||
if dp.Sum != 0 {
|
||||
t.Errorf("datapoint sum = %d, want 0", dp.Sum)
|
||||
}
|
||||
attrs := dp.Attributes
|
||||
ns, _ := attrs.Value(ateattr.WorkerPoolNamespaceKey)
|
||||
pool, _ := attrs.Value(ateattr.WorkerPoolNameKey)
|
||||
constraint, _ := attrs.Value(ateattr.SchedulingConstraintKey)
|
||||
if ns.AsString() != "ns-a" || pool.AsString() != "pool-1" {
|
||||
t.Errorf("got namespace=%q pool=%q, want ns-a / pool-1", ns.AsString(), pool.AsString())
|
||||
}
|
||||
if constraint.AsString() != "none" {
|
||||
t.Errorf("got constraint=%q, want none", constraint.AsString())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !foundMetric {
|
||||
t.Fatalf("ate.scheduler.eligible_workers metric not found")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("records constraint classification attributes correctly", func(t *testing.T) {
|
||||
sel, _ := labels.Parse("env=prod")
|
||||
tests := []struct {
|
||||
name string
|
||||
constraints Constraints
|
||||
wantConstraint string
|
||||
}{
|
||||
{
|
||||
name: "none",
|
||||
constraints: Constraints{SandboxClass: "gvisor"},
|
||||
wantConstraint: ateattr.ConstraintNone,
|
||||
},
|
||||
{
|
||||
name: "selector",
|
||||
constraints: Constraints{SandboxClass: "gvisor", ActorSelector: sel},
|
||||
wantConstraint: ateattr.ConstraintSelector,
|
||||
},
|
||||
{
|
||||
name: "required_nodes",
|
||||
constraints: Constraints{SandboxClass: "gvisor", ActorSelector: sel, RequiredNodes: []string{"node-a"}},
|
||||
wantConstraint: ateattr.ConstraintRequiredNodes,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
reader := sdkmetric.NewManualReader()
|
||||
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
|
||||
meter := provider.Meter("test")
|
||||
|
||||
flt := fleet{
|
||||
workerWithPool("w-1", "ns-a", "pool-1", "gvisor", "node-a", map[string]string{"env": "prod"}),
|
||||
}
|
||||
|
||||
s := New(flt, WithIntn(firstIntn), WithMeter(meter))
|
||||
_, err := s.Schedule(context.Background(), tc.constraints)
|
||||
if err != nil {
|
||||
t.Fatalf("Schedule() error = %v", err)
|
||||
}
|
||||
|
||||
var rm metricdata.ResourceMetrics
|
||||
_ = reader.Collect(context.Background(), &rm)
|
||||
for _, sm := range rm.ScopeMetrics {
|
||||
for _, m := range sm.Metrics {
|
||||
if m.Name == "ate.scheduler.eligible_workers" {
|
||||
histogram := m.Data.(metricdata.Histogram[int64])
|
||||
dp := histogram.DataPoints[0]
|
||||
constraint, _ := dp.Attributes.Value(ateattr.SchedulingConstraintKey)
|
||||
if constraint.AsString() != tc.wantConstraint {
|
||||
t.Errorf("got constraint=%q, want %q", constraint.AsString(), tc.wantConstraint)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("records 0 eligible workers when fleet is completely empty", func(t *testing.T) {
|
||||
reader := sdkmetric.NewManualReader()
|
||||
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
|
||||
meter := provider.Meter("test")
|
||||
|
||||
flt := fleet{}
|
||||
|
||||
s := New(flt, WithIntn(firstIntn), WithMeter(meter))
|
||||
_, err := s.Schedule(context.Background(), Constraints{SandboxClass: "gvisor"})
|
||||
if !errors.Is(err, ErrNoCapacity) {
|
||||
t.Fatalf("Schedule() error = %v, want ErrNoCapacity", err)
|
||||
}
|
||||
|
||||
var rm metricdata.ResourceMetrics
|
||||
_ = reader.Collect(context.Background(), &rm)
|
||||
foundMetric := false
|
||||
for _, sm := range rm.ScopeMetrics {
|
||||
for _, m := range sm.Metrics {
|
||||
if m.Name == "ate.scheduler.eligible_workers" {
|
||||
foundMetric = true
|
||||
histogram := m.Data.(metricdata.Histogram[int64])
|
||||
dp := histogram.DataPoints[0]
|
||||
if dp.Sum != 0 {
|
||||
t.Errorf("datapoint sum = %d, want 0", dp.Sum)
|
||||
}
|
||||
class, _ := dp.Attributes.Value(ateattr.SandboxClassKey)
|
||||
if class.AsString() != "gvisor" {
|
||||
t.Errorf("got sandbox class %q, want gvisor", class.AsString())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !foundMetric {
|
||||
t.Fatalf("ate.scheduler.eligible_workers metric not found")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("records 0 eligible workers on sandbox class mismatch", func(t *testing.T) {
|
||||
reader := sdkmetric.NewManualReader()
|
||||
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
|
||||
meter := provider.Meter("test")
|
||||
|
||||
flt := fleet{
|
||||
workerWithPool("w-1", "ns-a", "pool-1", "gvisor", "node-a", nil),
|
||||
}
|
||||
|
||||
s := New(flt, WithIntn(firstIntn), WithMeter(meter))
|
||||
_, err := s.Schedule(context.Background(), Constraints{SandboxClass: "microvm"})
|
||||
if !errors.Is(err, ErrNoCapacity) {
|
||||
t.Fatalf("Schedule() error = %v, want ErrNoCapacity", err)
|
||||
}
|
||||
|
||||
var rm metricdata.ResourceMetrics
|
||||
_ = reader.Collect(context.Background(), &rm)
|
||||
foundMetric := false
|
||||
for _, sm := range rm.ScopeMetrics {
|
||||
for _, m := range sm.Metrics {
|
||||
if m.Name == "ate.scheduler.eligible_workers" {
|
||||
foundMetric = true
|
||||
histogram := m.Data.(metricdata.Histogram[int64])
|
||||
dp := histogram.DataPoints[0]
|
||||
if dp.Sum != 0 {
|
||||
t.Errorf("datapoint sum = %d, want 0", dp.Sum)
|
||||
}
|
||||
class, _ := dp.Attributes.Value(ateattr.SandboxClassKey)
|
||||
if class.AsString() != "microvm" {
|
||||
t.Errorf("got sandbox class %q, want microvm", class.AsString())
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !foundMetric {
|
||||
t.Fatalf("ate.scheduler.eligible_workers metric not found")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("excludes draining and inactive workers from eligible counts", func(t *testing.T) {
|
||||
reader := sdkmetric.NewManualReader()
|
||||
provider := sdkmetric.NewMeterProvider(sdkmetric.WithReader(reader))
|
||||
meter := provider.Meter("test")
|
||||
|
||||
drainingWorker := workerWithPool("w-draining", "ns-a", "pool-1", "gvisor", "node-a", nil)
|
||||
drainingWorker.Status.State = ateapipb.WorkerState_WORKER_STATE_DRAINING
|
||||
|
||||
flt := fleet{drainingWorker}
|
||||
|
||||
s := New(flt, WithIntn(firstIntn), WithMeter(meter))
|
||||
_, err := s.Schedule(context.Background(), Constraints{SandboxClass: "gvisor"})
|
||||
if !errors.Is(err, ErrNoCapacity) {
|
||||
t.Fatalf("Schedule() error = %v, want ErrNoCapacity", err)
|
||||
}
|
||||
|
||||
var rm metricdata.ResourceMetrics
|
||||
_ = reader.Collect(context.Background(), &rm)
|
||||
foundMetric := false
|
||||
for _, sm := range rm.ScopeMetrics {
|
||||
for _, m := range sm.Metrics {
|
||||
if m.Name == "ate.scheduler.eligible_workers" {
|
||||
foundMetric = true
|
||||
histogram := m.Data.(metricdata.Histogram[int64])
|
||||
dp := histogram.DataPoints[0]
|
||||
if dp.Sum != 0 {
|
||||
t.Errorf("datapoint sum = %d, want 0", dp.Sum)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !foundMetric {
|
||||
t.Fatalf("ate.scheduler.eligible_workers metric not found")
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -474,23 +474,6 @@ groups:
|
||||
stability: development
|
||||
value: error
|
||||
brief: The attempt failed. Only this outcome has an error.type key.
|
||||
- id: ate.scheduling.constraint
|
||||
stability: development
|
||||
brief: The type of constraint in the scheduling request.
|
||||
type:
|
||||
members:
|
||||
- id: none
|
||||
stability: development
|
||||
value: none
|
||||
brief: The request has no constraint.
|
||||
- id: selector
|
||||
stability: development
|
||||
value: selector
|
||||
brief: The request has a label selector for an actor or a template.
|
||||
- id: required_nodes
|
||||
stability: development
|
||||
value: required_nodes
|
||||
brief: The request names specific node VMs.
|
||||
|
||||
- id: registry.ate.router
|
||||
type: attribute_group
|
||||
@@ -878,37 +861,6 @@ groups:
|
||||
requirement_level:
|
||||
conditionally_required: The outcome is error.
|
||||
|
||||
- id: metric.ate.scheduler.eligible_workers
|
||||
type: metric
|
||||
metric_name: ate.scheduler.eligible_workers
|
||||
instrument: histogram
|
||||
unit: "{worker}"
|
||||
stability: development
|
||||
brief: The number of free workers that remain after all the constraint filters, at each scheduling decision.
|
||||
note: >
|
||||
This is an early sign of low capacity. It warns you before the first
|
||||
rejection. The rate of 503 errors tells you only after the users have a
|
||||
fault. If no pool agrees with the constraints, ateapi sends one series
|
||||
with the value 0. Both pool keys are empty in that series. Thus a
|
||||
dashboard shows the "no workers are available" state with the series of
|
||||
the pools. Only this instrument has the empty pair of pool keys.
|
||||
annotations:
|
||||
substrate:
|
||||
emitted_by: [ateapi]
|
||||
golden_signals: [saturation]
|
||||
code_anchor: cmd/ateapi/internal/scheduling/metrics.go
|
||||
buckets: [0, 1, 2, 3, 5, 10, 20, 50, 100, 250]
|
||||
cuj: Resumes fail with a 503 error, but ate.workerpool.workers shows idle workers. Why?
|
||||
attributes:
|
||||
- ref: ate.workerpool.namespace
|
||||
requirement_level: required
|
||||
- ref: ate.workerpool.name
|
||||
requirement_level: required
|
||||
- ref: ate.sandbox.class
|
||||
requirement_level: required
|
||||
- ref: ate.scheduling.constraint
|
||||
requirement_level: required
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Metrics. atelet
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -202,8 +202,7 @@ cardinality_rules:
|
||||
brief: >
|
||||
A component sets ate.workerpool.namespace and ate.workerpool.name
|
||||
together, or it sets no pool key. One key alone does not identify a pool.
|
||||
ate.scheduler.eligible_workers is the one exception. It sends both keys
|
||||
empty to show that no worker is available.
|
||||
No instrument is an exception.
|
||||
enforced: false
|
||||
could_be_enforced_by: >
|
||||
A Weaver Rego policy can make the two keys occur together in this
|
||||
|
||||
@@ -246,7 +246,6 @@ Agent Substrate emits foundational OpenTelemetry system and server metrics to mo
|
||||
| `rpc.server.call.duration` | ateapi & atelet (gRPC servers, via `otelgrpc`) | histogram | per-method gRPC latency, request rate, and errors (labels `rpc.method`, `rpc.response.status_code`) |
|
||||
| `ate.actor.crashes` | ateapi | counter | Number of times actors transitioned to `ACTOR_STATE_CRASHED` with failure reasons (labels `ate.actor.operation.name`, `ate.failure.reason`, `ate.failure.domain`, `ate.template.atespace`, `ate.template.name`, `ate.workerpool.namespace`, `ate.workerpool.name`, `ate.sandbox.class`) |
|
||||
| `atenet.router.route.duration` | atenet-router | histogram | Substrate E2E — Envoy receiving a request to Envoy forwarding it to the resolved worker, excluding actor compute and the response (labels `ate.template.atespace`, `ate.template.name`, `ate.router.outcome`, `ate.router.resume`) |
|
||||
| `ate.scheduler.eligible_workers` | ateapi | histogram | number of eligible unassigned workers available during scheduling given the constraint filters (labels `ate.workerpool.namespace`, `ate.workerpool.name`, `ate.sandbox.class`, `ate.scheduling.constraint`) |
|
||||
| `atelet.snapshot.size` | atelet | histogram | uncompressed size in bytes of each gVisor snapshot image written during checkpoint (labels `file.name`, `ate.template.atespace`, `ate.template.name`) |
|
||||
| `ate.workerpool.desired_workers` | atecontroller | up/down counter | number of worker pods requested for a WorkerPool, from `spec.replicas` (labels
|
||||
`ate.workerpool.namespace`, `ate.workerpool.name`) |
|
||||
@@ -269,18 +268,13 @@ For `atenet.router.route.duration`:
|
||||
* `ate.router.outcome` categorizes the route attempt result: `ok`, `cancelled`, `timeout`, `no_capacity`, `failed_precondition`, `lock_conflict`, `not_found`, `unavailable`, `rate_limited`, or `resume_error`.
|
||||
* `ate.router.resume` indicates the singleflight execution state of actor resumption: `none` (the resume found the actor already running), `triggered` (this request completed a cold activation), `joined` (this request waited on another request's resume, which activated the actor), or `unknown` (the resume did not complete, so whether an activation ran is unknown). `ate.template.atespace` and `ate.template.name` hold `unknown` when the router has no template to name.
|
||||
|
||||
For `ate.scheduler.eligible_workers`:
|
||||
* `ate.scheduling.constraint` categorizes the scheduling request constraint type: `none` (unconstrained), `selector` (actor or template label selectors specified), or `required_nodes` (pinned to specific node VMs).
|
||||
|
||||
For `ate.imagecache.requests`:
|
||||
* `ate.imagecache.outcome` is `hit` when the node holds a complete image record — every layer directory the record names is present — and `miss` when the lookup must pull. A failed lookup is neither: `error` is a failed lookup whatever the cause, and `cancelled` or `timeout` is the caller giving up, as on `ate.router.outcome`. So the hit ratio is `hit / (hit + miss)`, with failures and abandoned lookups out of the denominator.
|
||||
* `error.type` is present only on the `error` outcome, and carries the registry's own HTTP status for its rejection, from a fixed set: `401`, `403`, `404`, `429`, `500`, `502`, `503`, `504`. The set is an allow-list because the registry client reports whatever the remote returned. Each other status, and each failure that carries no status, reports `_OTHER`.
|
||||
|
||||
`ate.workerpool.namespace` and `ate.workerpool.name` identify a pool together, on every instrument that names one. A WorkerPool is a namespaced resource, so the name on its own merges same-named pools from different namespaces into one series. The pair means that capacity (`ate.workerpool.workers`, `ate.workerpool.desired_workers`, `ate.workerpool.ready_workers`) joins to demand (`ate.scheduler.assignment.duration`, `ate.actor.lifecycle.operation.duration`, `ate.actor.crashes`) by pool.
|
||||
|
||||
Two states read differently:
|
||||
* **No keys** means the operation has no pool. The actor-centric instruments omit the pair, so a crash before the actor reached a worker, or the `no_free_worker` outcome, names no pool.
|
||||
* **Both keys empty** means no pool matched. Only `ate.scheduler.eligible_workers` reports it, as one zero-valued series that keeps "nothing is schedulable" on the same chart as the per-pool series.
|
||||
**No keys** means the operation has no pool. The actor-centric instruments omit the pair, so a crash before the actor reached a worker, or the `no_free_worker` outcome, names no pool. No instrument sends the pair with both keys empty.
|
||||
|
||||
The three snapshot labels are orthogonal and mean the same thing on every histogram that carries them:
|
||||
* `ate.snapshot.kind`: which snapshot the operation reads or writes. `local` (node-local, written by a pause), `latest` (the actor's own durable snapshot), `golden` (the template's image), or `boot` (from scratch, so it never appears on the atelet histograms).
|
||||
|
||||
+15
-23
@@ -159,22 +159,21 @@ func ActorStateValue(state ateapipb.ActorState) string {
|
||||
// pool is node state every actor shares. For the same reason it is the only
|
||||
// ate.* label on its counter.
|
||||
const (
|
||||
ActorOperationNameKey = attribute.Key("ate.actor.operation.name")
|
||||
WorkerPoolNamespaceKey = attribute.Key("ate.workerpool.namespace")
|
||||
WorkerPoolNameKey = attribute.Key("ate.workerpool.name")
|
||||
WorkerStateKey = attribute.Key("ate.worker.state")
|
||||
SandboxClassKey = attribute.Key("ate.sandbox.class")
|
||||
SnapshotKindKey = attribute.Key("ate.snapshot.kind")
|
||||
SnapshotScopeKey = attribute.Key("ate.snapshot.scope")
|
||||
SnapshotPhaseKey = attribute.Key("ate.snapshot.phase")
|
||||
ImageCacheOutcomeKey = attribute.Key("ate.imagecache.outcome")
|
||||
SchedulerOutcomeKey = attribute.Key("ate.scheduler.outcome")
|
||||
SchedulingConstraintKey = attribute.Key("ate.scheduling.constraint")
|
||||
RouterResumeKey = attribute.Key("ate.router.resume")
|
||||
RouterOutcomeKey = attribute.Key("ate.router.outcome")
|
||||
FailureReasonKey = attribute.Key("ate.failure.reason")
|
||||
FailureDomainKey = attribute.Key("ate.failure.domain")
|
||||
StatsSourceKey = attribute.Key("ate.stats.source")
|
||||
ActorOperationNameKey = attribute.Key("ate.actor.operation.name")
|
||||
WorkerPoolNamespaceKey = attribute.Key("ate.workerpool.namespace")
|
||||
WorkerPoolNameKey = attribute.Key("ate.workerpool.name")
|
||||
WorkerStateKey = attribute.Key("ate.worker.state")
|
||||
SandboxClassKey = attribute.Key("ate.sandbox.class")
|
||||
SnapshotKindKey = attribute.Key("ate.snapshot.kind")
|
||||
SnapshotScopeKey = attribute.Key("ate.snapshot.scope")
|
||||
SnapshotPhaseKey = attribute.Key("ate.snapshot.phase")
|
||||
ImageCacheOutcomeKey = attribute.Key("ate.imagecache.outcome")
|
||||
SchedulerOutcomeKey = attribute.Key("ate.scheduler.outcome")
|
||||
RouterResumeKey = attribute.Key("ate.router.resume")
|
||||
RouterOutcomeKey = attribute.Key("ate.router.outcome")
|
||||
FailureReasonKey = attribute.Key("ate.failure.reason")
|
||||
FailureDomainKey = attribute.Key("ate.failure.domain")
|
||||
StatsSourceKey = attribute.Key("ate.stats.source")
|
||||
)
|
||||
|
||||
// Values for FailureDomainKey. A strict function of the reason, so it costs no
|
||||
@@ -240,13 +239,6 @@ const (
|
||||
StatsSourceGuestAgent = "guest-agent"
|
||||
)
|
||||
|
||||
// Values for SchedulingConstraintKey.
|
||||
const (
|
||||
ConstraintNone = "none"
|
||||
ConstraintRequiredNodes = "required_nodes"
|
||||
ConstraintSelector = "selector"
|
||||
)
|
||||
|
||||
// Control-plane failure reasons for ate.actor.crashes metric.
|
||||
const (
|
||||
ReasonCorruptedAssignment = string(ateerrors.ReasonCorruptedAssignment)
|
||||
|
||||
@@ -52,7 +52,6 @@ var PlatformMetricPrefixes = []string{
|
||||
"ate_actor_restore_duration",
|
||||
"ate_actor_checkpoint_duration",
|
||||
"atenet_router_route_duration",
|
||||
"ate_scheduler_eligible_workers",
|
||||
}
|
||||
|
||||
// ScrapeAgentGatewayRouterMetrics reads the AgentGateway router's native
|
||||
|
||||
@@ -185,62 +185,6 @@ func TestPlatformMetricsEmitted(t *testing.T) {
|
||||
errs = append(errs, "ate_workerpool_ready_workers validation failed: metric line not found in collector scrape text (no time series emitted by atecontroller callback)")
|
||||
}
|
||||
|
||||
// Verify ate_scheduler_eligible_workers metric carries valid attributes:
|
||||
// - Full labels (namespace, pool, class, constraint) for per-pool candidate lines.
|
||||
// - Necessary base labels (class, constraint) for edge cases when no worker pools match.
|
||||
foundEligibleLine := false
|
||||
foundFullPoolLine := false
|
||||
for _, line := range strings.Split(scrape, "\n") {
|
||||
if strings.HasPrefix(line, "ate_scheduler_eligible_workers") {
|
||||
foundEligibleLine = true
|
||||
nsVal := extractLabelValue(line, "ate_workerpool_namespace")
|
||||
poolVal := extractLabelValue(line, "ate_workerpool_name")
|
||||
classVal := extractLabelValue(line, "ate_sandbox_class")
|
||||
constraintVal := extractLabelValue(line, "ate_scheduling_constraint")
|
||||
|
||||
var lineErrs []string
|
||||
if classVal == "" {
|
||||
lineErrs = append(lineErrs, "ate_sandbox_class label is missing or empty")
|
||||
}
|
||||
if constraintVal == "" {
|
||||
lineErrs = append(lineErrs, "ate_scheduling_constraint label is missing or empty")
|
||||
} else if constraintVal != ateattr.ConstraintNone && constraintVal != ateattr.ConstraintRequiredNodes && constraintVal != ateattr.ConstraintSelector {
|
||||
lineErrs = append(lineErrs, fmt.Sprintf("ate_scheduling_constraint %q is invalid (must be one of {%s, %s, %s})",
|
||||
constraintVal, ateattr.ConstraintNone, ateattr.ConstraintRequiredNodes, ateattr.ConstraintSelector))
|
||||
}
|
||||
|
||||
// Determine line type for error reporting.
|
||||
isPerPoolLine := poolVal != "" || nsVal != ""
|
||||
caseType := "[NORMAL CASE: Per-Pool Candidates Expected]"
|
||||
if !isPerPoolLine {
|
||||
caseType = "[EDGE CASE: No Worker Pools Matched Constraints]"
|
||||
}
|
||||
|
||||
// If the line has pool/namespace labels, verify both are non-empty (full per-pool line).
|
||||
if isPerPoolLine {
|
||||
if nsVal == "" {
|
||||
lineErrs = append(lineErrs, "ate_workerpool_namespace label is missing or empty")
|
||||
}
|
||||
if poolVal == "" {
|
||||
lineErrs = append(lineErrs, "ate_workerpool_name label is missing or empty")
|
||||
}
|
||||
if len(lineErrs) == 0 {
|
||||
foundFullPoolLine = true
|
||||
}
|
||||
}
|
||||
|
||||
if len(lineErrs) > 0 {
|
||||
errs = append(errs, fmt.Sprintf("%s line %q failed label validation:\n - %s\n (Extracted labels: ate_workerpool_namespace=%q, ate_workerpool_name=%q, ate_sandbox_class=%q, ate_scheduling_constraint=%q)",
|
||||
caseType, line, strings.Join(lineErrs, "\n - "), nsVal, poolVal, classVal, constraintVal))
|
||||
}
|
||||
}
|
||||
}
|
||||
if !foundEligibleLine {
|
||||
errs = append(errs, "ate_scheduler_eligible_workers metric line not found in collector scrape output")
|
||||
} else if !foundFullPoolLine {
|
||||
errs = append(errs, "ate_scheduler_eligible_workers [NORMAL CASE] per-pool candidates was not found in collector scrape output; only edge-case 0-count histogram was present")
|
||||
}
|
||||
|
||||
// Verify ate_actor_crashes metric carries valid, non-empty low-cardinality labels for all attributes.
|
||||
foundCrashLine := false
|
||||
for _, line := range strings.Split(scrape, "\n") {
|
||||
|
||||
Reference in New Issue
Block a user