mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
atenet-router: apply review-walkthrough feedback
- Drop the fast-path benchmarks: the before/after numbers live in the PR description; nothing in CI executes benchmarks, so the file only cost maintenance. - Reformat the tests this change adds: nested proto literals one field per line, and the anonymous mock signatures wrapped one parameter per line. No behavior change.
This commit is contained in:
@@ -399,9 +399,20 @@ func TestHandleRequestHeaders_FullLotServesRunningActor(t *testing.T) {
|
||||
|
||||
var resumeCalled bool
|
||||
clientMock := &mockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeCalled = true
|
||||
return &ateapipb.ResumeActorResponse{Actor: &ateapipb.Actor{Status: &ateapipb.ActorStatus{State: ateapipb.ActorState_ACTOR_STATE_RUNNING, WorkerAssignment: &ateapipb.WorkerAssignment{WorkerPodIp: "10.0.0.1"}}}}, nil
|
||||
return &ateapipb.ResumeActorResponse{
|
||||
Actor: &ateapipb.Actor{
|
||||
Status: &ateapipb.ActorStatus{
|
||||
State: ateapipb.ActorState_ACTOR_STATE_RUNNING,
|
||||
WorkerAssignment: &ateapipb.WorkerAssignment{WorkerPodIp: "10.0.0.1"},
|
||||
},
|
||||
},
|
||||
}, nil
|
||||
},
|
||||
}
|
||||
|
||||
@@ -445,7 +456,11 @@ func TestHandleRequestHeaders_FullLotShedsParkedRequest(t *testing.T) {
|
||||
|
||||
var resumeCalls atomic.Int32
|
||||
clientMock := &mockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeCalls.Add(1)
|
||||
return nil, status.Error(codes.ResourceExhausted, "no free workers available")
|
||||
},
|
||||
|
||||
@@ -1,103 +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 ingress
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"log/slog"
|
||||
"testing"
|
||||
|
||||
corev3 "github.com/envoyproxy/go-control-plane/envoy/config/core/v3"
|
||||
"google.golang.org/grpc"
|
||||
|
||||
"github.com/agent-substrate/substrate/cmd/atenet/internal/router/extproc"
|
||||
"github.com/agent-substrate/substrate/internal/atenet"
|
||||
"github.com/agent-substrate/substrate/internal/resources"
|
||||
"github.com/agent-substrate/substrate/pkg/proto/ateapipb"
|
||||
)
|
||||
|
||||
// runningActorResponse is the control-plane reply for an actor that needs no
|
||||
// resume — the fast path every benchmark here exercises.
|
||||
func runningActorResponse() *ateapipb.ResumeActorResponse {
|
||||
return &ateapipb.ResumeActorResponse{
|
||||
Actor: &ateapipb.Actor{
|
||||
Status: &ateapipb.ActorStatus{
|
||||
State: ateapipb.ActorState_ACTOR_STATE_RUNNING,
|
||||
WorkerAssignment: &ateapipb.WorkerAssignment{WorkerPodIp: "10.0.0.52"},
|
||||
},
|
||||
},
|
||||
Resumed: false,
|
||||
}
|
||||
}
|
||||
|
||||
// BenchmarkResumeActorAlreadyRunning measures the resumer's fast path — a
|
||||
// request to an actor the control plane reports RUNNING on the first attempt.
|
||||
// This is on every request's critical path, so the flight bookkeeping around
|
||||
// the single RPC must stay cheap.
|
||||
func BenchmarkResumeActorAlreadyRunning(b *testing.B) {
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
return runningActorResponse(), nil
|
||||
},
|
||||
}
|
||||
resumer := NewActorResumer(mock, withParking(DefaultParkedRequestConfig()))
|
||||
ref := resources.ActorRef{Atespace: "team-a", Name: "bench-actor"}
|
||||
ctx := context.Background()
|
||||
|
||||
b.ReportAllocs()
|
||||
for b.Loop() {
|
||||
if _, _, err := resumer.ResumeActor(ctx, ref); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// BenchmarkHandleRequestHeadersRunningActor measures the whole ingress
|
||||
// handler's fast path, parking admission included — the production cost of one
|
||||
// routed request to a running actor.
|
||||
func BenchmarkHandleRequestHeadersRunningActor(b *testing.B) {
|
||||
prev := slog.Default()
|
||||
slog.SetDefault(slog.New(slog.NewTextHandler(io.Discard, nil)))
|
||||
b.Cleanup(func() { slog.SetDefault(prev) })
|
||||
|
||||
mock := &mockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
return runningActorResponse(), nil
|
||||
},
|
||||
}
|
||||
h := New(mock, DefaultParkedRequestConfig(), nil)
|
||||
|
||||
md := benchRequestMetadata("123e4567-e89b-12d3-a456-426614174000", "team-a")
|
||||
ctx := context.Background()
|
||||
|
||||
b.ReportAllocs()
|
||||
for b.Loop() {
|
||||
if _, err := h.HandleRequestHeaders(ctx, md); err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// benchRequestMetadata mirrors the requestMetadata test helper.
|
||||
func benchRequestMetadata(actorName, atespace string) *extproc.RequestMetadata {
|
||||
return extproc.NewRequestMetadata(
|
||||
[]*corev3.HeaderValue{
|
||||
{Key: atenet.TargetActorHeader, Value: atespace + "/" + actorName},
|
||||
{Key: ":method", Value: "POST"},
|
||||
},
|
||||
nil,
|
||||
)
|
||||
}
|
||||
@@ -586,7 +586,13 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
testActorRef := resources.ActorRef{Atespace: testAtespace, Name: testActorName}
|
||||
runningResp := func() *ateapipb.ResumeActorResponse {
|
||||
return &ateapipb.ResumeActorResponse{
|
||||
Actor: &ateapipb.Actor{Metadata: &ateapipb.ResourceMetadata{Name: testActorName}, Status: &ateapipb.ActorStatus{State: ateapipb.ActorState_ACTOR_STATE_RUNNING, WorkerAssignment: &ateapipb.WorkerAssignment{WorkerPodIp: expectedIP}}},
|
||||
Actor: &ateapipb.Actor{
|
||||
Metadata: &ateapipb.ResourceMetadata{Name: testActorName},
|
||||
Status: &ateapipb.ActorStatus{
|
||||
State: ateapipb.ActorState_ACTOR_STATE_RUNNING,
|
||||
WorkerAssignment: &ateapipb.WorkerAssignment{WorkerPodIp: expectedIP},
|
||||
},
|
||||
},
|
||||
Resumed: true,
|
||||
}
|
||||
}
|
||||
@@ -594,7 +600,11 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
t.Run("FastFlightNeverEntersLot", func(t *testing.T) {
|
||||
synctest.Test(t, func(t *testing.T) {
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
return runningResp(), nil
|
||||
},
|
||||
}
|
||||
@@ -629,7 +639,11 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
var mu sync.Mutex
|
||||
var calls, activeDuringRetry int
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
mu.Lock()
|
||||
calls++
|
||||
n := calls
|
||||
@@ -679,7 +693,11 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
var mu sync.Mutex
|
||||
var calls int
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
mu.Lock()
|
||||
calls++
|
||||
mu.Unlock()
|
||||
@@ -717,7 +735,11 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
var calls int
|
||||
proceed := make(chan struct{})
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
mu.Lock()
|
||||
calls++
|
||||
n := calls
|
||||
@@ -778,7 +800,11 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
cfg := ParkedRequestConfig{Max: 1, Budget: 1 * time.Second}
|
||||
lot := newParkingLot(cfg, nil)
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
return nil, status.Error(codes.ResourceExhausted, "no free workers available")
|
||||
},
|
||||
}
|
||||
@@ -806,7 +832,11 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
var mu sync.Mutex
|
||||
var calls int
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
mu.Lock()
|
||||
calls++
|
||||
n := calls
|
||||
@@ -851,7 +881,11 @@ func TestActorResumer_LotAdmission(t *testing.T) {
|
||||
var mu sync.Mutex
|
||||
var calls int
|
||||
mock := &resumerMockClient{
|
||||
resumeFn: func(ctx context.Context, in *ateapipb.ResumeActorRequest, opts ...grpc.CallOption) (*ateapipb.ResumeActorResponse, error) {
|
||||
resumeFn: func(
|
||||
ctx context.Context,
|
||||
in *ateapipb.ResumeActorRequest,
|
||||
opts ...grpc.CallOption,
|
||||
) (*ateapipb.ResumeActorResponse, error) {
|
||||
mu.Lock()
|
||||
calls++
|
||||
mu.Unlock()
|
||||
|
||||
Reference in New Issue
Block a user