mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
atenet/router: make parkOutcome a named type
Replace the bare-string park-outcome constants with a parkOutcome string type and rename the classifier to parkOutcomeFor, so the parking.wait.duration outcome label is type-checked rather than an arbitrary string.
This commit is contained in:
@@ -162,8 +162,7 @@ func (s *ExtProcServer) handleRequestHeaders(
|
||||
|
||||
slog.InfoContext(ctx, "ResumeActor", slog.Any("actor", actorRef))
|
||||
actor, err := s.resumer.ResumeActor(ctx, actorRef)
|
||||
release(parkOutcome(err))
|
||||
|
||||
release(parkOutcomeFor(err))
|
||||
|
||||
if err != nil {
|
||||
return nil, metadata, "", "", "", mapResumeError(actorRef, err)
|
||||
|
||||
@@ -115,11 +115,11 @@ func (m *parkingMetrics) addActive(ctx context.Context, delta int64) {
|
||||
m.active.Add(ctx, delta)
|
||||
}
|
||||
|
||||
func (m *parkingMetrics) recordWait(ctx context.Context, d time.Duration, outcome string) {
|
||||
func (m *parkingMetrics) recordWait(ctx context.Context, d time.Duration, outcome parkOutcome) {
|
||||
if m == nil || m.wait == nil {
|
||||
return
|
||||
}
|
||||
m.wait.Record(ctx, d.Seconds(), metric.WithAttributes(attribute.String("outcome", outcome)))
|
||||
m.wait.Record(ctx, d.Seconds(), metric.WithAttributes(attribute.String("outcome", string(outcome))))
|
||||
}
|
||||
|
||||
func (m *parkingMetrics) recordRejected(ctx context.Context) {
|
||||
|
||||
@@ -30,12 +30,16 @@ const (
|
||||
defaultParkingMaxParked = 2048
|
||||
)
|
||||
|
||||
// Parking-wait outcome labels, recorded on the parking.wait.duration histogram.
|
||||
// parkOutcome is the terminal disposition of a parked request. It is recorded
|
||||
// as the `outcome` label on the parking.wait.duration histogram.
|
||||
type parkOutcome string
|
||||
|
||||
// Park-wait outcomes, recorded on the parking.wait.duration histogram.
|
||||
const (
|
||||
parkOutcomeServed = "served" // resume succeeded and the request was routed
|
||||
parkOutcomeTimeout = "timeout" // the request's deadline elapsed while parked
|
||||
parkOutcomeCanceled = "canceled" // the client disconnected while parked
|
||||
parkOutcomeError = "error" // resume failed (including park-budget exhaustion)
|
||||
parkOutcomeServed parkOutcome = "served" // resume succeeded and the request was routed
|
||||
parkOutcomeTimeout parkOutcome = "timeout" // the request's deadline elapsed while parked
|
||||
parkOutcomeCanceled parkOutcome = "canceled" // the client disconnected while parked
|
||||
parkOutcomeError parkOutcome = "error" // resume failed (including park-budget exhaustion)
|
||||
)
|
||||
|
||||
// parkingConfig controls how the router parks resume-gated requests.
|
||||
@@ -87,9 +91,9 @@ func newParkingLot(cfg parkingConfig, m *parkingMetrics) *parkingLot {
|
||||
// ok=false means the lot is full and the request should be shed without
|
||||
// waiting. When parking is disabled every request is admitted and no slot
|
||||
// accounting or metrics are recorded.
|
||||
func (l *parkingLot) enter(ctx context.Context) (release func(outcome string), ok bool) {
|
||||
func (l *parkingLot) enter(ctx context.Context) (release func(outcome parkOutcome), ok bool) {
|
||||
if l == nil || !l.cfg.enabled {
|
||||
return func(string) {}, true
|
||||
return func(parkOutcome) {}, true
|
||||
}
|
||||
|
||||
for {
|
||||
@@ -108,7 +112,7 @@ func (l *parkingLot) enter(ctx context.Context) (release func(outcome string), o
|
||||
l.metrics.addActive(ctx, 1)
|
||||
|
||||
var once sync.Once
|
||||
return func(outcome string) {
|
||||
return func(outcome parkOutcome) {
|
||||
once.Do(func() {
|
||||
atomic.AddInt64(&l.active, -1)
|
||||
l.metrics.addActive(ctx, -1)
|
||||
@@ -138,10 +142,10 @@ func (l *parkingLot) status() ParkingStatus {
|
||||
}
|
||||
}
|
||||
|
||||
// parkOutcome classifies a completed resume attempt for the wait-duration
|
||||
// parkOutcomeFor classifies a completed resume attempt for the wait-duration
|
||||
// metric. A budget-exhausted park surfaces the underlying capacity error and is
|
||||
// reported as parkOutcomeError.
|
||||
func parkOutcome(err error) string {
|
||||
func parkOutcomeFor(err error) parkOutcome {
|
||||
switch {
|
||||
case err == nil:
|
||||
return parkOutcomeServed
|
||||
|
||||
@@ -111,7 +111,7 @@ func TestParkingLot_ConcurrentEntryRespectsCapacity(t *testing.T) {
|
||||
|
||||
var admitted int64
|
||||
var mu sync.Mutex
|
||||
releases := make([]func(string), 0, capacity)
|
||||
releases := make([]func(parkOutcome), 0, capacity)
|
||||
|
||||
var wg sync.WaitGroup
|
||||
wg.Add(goroutines)
|
||||
@@ -142,11 +142,11 @@ func TestParkingLot_ConcurrentEntryRespectsCapacity(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestParkOutcome(t *testing.T) {
|
||||
func TestParkOutcomeFor(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
err error
|
||||
want string
|
||||
want parkOutcome
|
||||
}{
|
||||
{"nil is served", nil, parkOutcomeServed},
|
||||
{"canceled", context.Canceled, parkOutcomeCanceled},
|
||||
@@ -155,8 +155,8 @@ func TestParkOutcome(t *testing.T) {
|
||||
}
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
if got := parkOutcome(tc.err); got != tc.want {
|
||||
t.Errorf("parkOutcome(%v) = %q, want %q", tc.err, got, tc.want)
|
||||
if got := parkOutcomeFor(tc.err); got != tc.want {
|
||||
t.Errorf("parkOutcomeFor(%v) = %q, want %q", tc.err, got, tc.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user