Skip to content
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@ read it when planning changes to the areas it covers, and keep it current:

- [`docs/internal/metrics.md`](docs/internal/metrics.md) — how metrics are split
(transport vs producer-state vs warpstream-specific) and the franz-go/`kprom`
drop-in-compatibility contract. Update it when changing metrics.
drop-in-compatibility contract. Update it when that categorization or parity
contract changes. Do not catalog individual custom metrics here; their
Prometheus help strings are the source of truth.
- [`docs/internal/tracing.md`](docs/internal/tracing.md) — how the client drives
franz-go's produce-record hooks on its own produce path to support tracing (e.g.
`kotel`) as a drop-in. Update it when changing hook or tracing behaviour.
Expand Down
22 changes: 22 additions & 0 deletions pkg/wgo/agentpool.go
Original file line number Diff line number Diff line change
Expand Up @@ -165,3 +165,25 @@ func diffRemovedAgents(old []int32, newSet map[int32]struct{}) []int32 {
}
return removed
}

// diffAgentMembership counts NodeIDs that appeared or disappeared between two
// sorted, unique agent lists.
func diffAgentMembership(old, new []int32) (added, removed int) {

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

diffAgentMembership assumes sorted, unique input. refresh sorts newAgents but doesn't dedupe it. If a Metadata response ever lists a NodeID twice, [5,5] followed by [5] counts one removed with no real membership change, which inflates agents_changed_total. diffRemovedAgents is unaffected because it goes through agentSet. Probably rare in practice, but a slices.Compact after the sort would make the documented precondition true. (Or derive the counts from the same set-based diff refresh already does.)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Fixed. refresh now deduplicates NodeIDs after sorting. It was worse than described: the pool itself held the duplicate and membership_changed was also counted. Added a test that duplicates a broker in the Metadata response.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

I'm reverting this. A correct Metadata response contains one entry per NodeID, and a duplicate is treated as malformed, so the sorted list already satisfies diffAgentMembership. Compacting also changes fallback routing for a malformed response (hash % len(agents) and the secondary walk), which this PR otherwise doesn't touch. I've removed the slices.Compact call and its test.

i, j := 0, 0
for i < len(old) && j < len(new) {
switch {
case old[i] == new[j]:
i++
j++
case old[i] < new[j]:
removed++
i++
default:
added++
j++
}
}
removed += len(old) - i
added += len(new) - j
return added, removed
}
35 changes: 35 additions & 0 deletions pkg/wgo/agentpool_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,41 @@ func TestBuildLeadersAndTopicIDs(t *testing.T) {

func stringPtr(s string) *string { return &s }

func TestDiffAgentMembership(t *testing.T) {
Comment thread
Copilot marked this conversation as resolved.
Outdated
tests := map[string]struct {
old []int32
new []int32
wantAdded int
wantRemoved int
}{
"unchanged": {
old: []int32{1, 2, 3}, new: []int32{1, 2, 3},
},
"one added": {
old: []int32{1, 2}, new: []int32{1, 2, 3}, wantAdded: 1,
},
"one removed": {
old: []int32{1, 2, 3}, new: []int32{1, 3}, wantRemoved: 1,
},
"replace one": {
old: []int32{1, 2, 3}, new: []int32{1, 4}, wantAdded: 1, wantRemoved: 2,
},
"empty to some": {
old: nil, new: []int32{1, 2}, wantAdded: 2,
},
"all removed": {
old: []int32{1, 2}, new: nil, wantRemoved: 2,
},
}
for name, tc := range tests {
t.Run(name, func(t *testing.T) {
added, removed := diffAgentMembership(tc.old, tc.new)
assert.Equal(t, tc.wantAdded, added)
assert.Equal(t, tc.wantRemoved, removed)
})
}
}

func TestDiffRemovedAgents(t *testing.T) {
tests := map[string]struct {
old []int32
Expand Down
2 changes: 2 additions & 0 deletions pkg/wgo/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -214,6 +214,7 @@ func (c *WarpstreamClient) Produce(ctx context.Context, record *kgo.Record, prom
})
if err != nil {
c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned).Inc()
c.metrics.produceFinalOutcome[produceFinalOutcomeNoAgentAssigned].Inc()
promise(record, err)
return
}
Expand Down Expand Up @@ -308,6 +309,7 @@ func (c *WarpstreamClient) ProduceSync(ctx context.Context, records []*kgo.Recor
// One record had no known candidate. Fail the whole batch
// uniformly: every ok record gets the same error.
c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned).Add(float64(len(okIndices)))
c.metrics.produceFinalOutcome[produceFinalOutcomeNoAgentAssigned].Inc()
Comment thread
koloss2001 marked this conversation as resolved.
Outdated
for _, i := range okIndices {
results[i] = kgo.ProduceResult{Record: records[i], Err: err}
}
Expand Down
4 changes: 4 additions & 0 deletions pkg/wgo/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -251,6 +251,7 @@ func TestWarpstreamClient_ProduceSync(t *testing.T) {
require.Error(t, results[0].Err)
assert.ErrorContains(t, results[0].Err, "no agent assigned")
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned)))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceFinalOutcome[produceFinalOutcomeNoAgentAssigned]))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceRecordsTotal))
// A rejection is not a failure: produceRecordsFailedTotal stays 0.
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsFailedTotal))
Expand Down Expand Up @@ -305,6 +306,7 @@ func TestWarpstreamClient_ProduceSync(t *testing.T) {

assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedRecordTooLarge)))
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned)))
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceFinalOutcome[produceFinalOutcomeNoAgentAssigned]))
assert.Equal(t, float64(2), testutil.ToFloat64(c.metrics.produceRecordsTotal))
// The oversized record is a rejection, not a failure; the ok record succeeds.
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsFailedTotal))
Expand All @@ -329,6 +331,7 @@ func TestWarpstreamClient_ProduceSync(t *testing.T) {
}

assert.Equal(t, float64(len(records)), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned)))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceFinalOutcome[produceFinalOutcomeNoAgentAssigned]))
assert.Equal(t, float64(len(records)), testutil.ToFloat64(c.metrics.produceRecordsTotal))
// All records are rejections, not failures.
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsFailedTotal))
Expand Down Expand Up @@ -654,6 +657,7 @@ func TestWarpstreamClient_Produce(t *testing.T) {
require.Error(t, err)
assert.ErrorContains(t, err, "no agent assigned")
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned)))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceFinalOutcome[produceFinalOutcomeNoAgentAssigned]))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.produceRecordsTotal))
// A rejection is not a failure: produceRecordsFailedTotal stays 0.
assert.Equal(t, float64(0), testutil.ToFloat64(c.metrics.produceRecordsFailedTotal))
Expand Down
58 changes: 44 additions & 14 deletions pkg/wgo/hedger.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package wgo

import (
"context"
"errors"
"fmt"
"slices"
"sync"
Expand All @@ -24,6 +25,21 @@ type HedgerConfig struct {
MaxHedgeAgents int
}

type hedgeTrigger int8

const (
hedgeTriggerLatency hedgeTrigger = iota
hedgeTriggerPrimaryFailure
hedgeTriggerDemotedProbe
hedgeTriggerCount
)

const (
hedgeTriggerLabelLatency = "latency"
hedgeTriggerLabelPrimaryFailure = "primary_failure"
hedgeTriggerLabelDemotedProbe = "demoted_probe"
)

// Hedger orchestrates produce attempts across multiple agents for the
// same batch of partitions. The whole "retry on a different agent,
// possibly fire a hedge to race the primary, give up when
Expand Down Expand Up @@ -131,20 +147,25 @@ func (h *Hedger) ProduceSync(ctx context.Context, primaryID int32, routedPartiti
}
}

// Check the hedging delay to apply to this request.
delay, shouldHedge := h.shouldHedge(time.Now(), primaryID, routedPartitions)

// observeAttempts records the attempt depth of a resolved produce call,
// split by outcome: 1 = resolved on the primary, N = resolved after N-1
// hedge waves.
observeAttempts := func(result ProduceResult, attempts int) {
Comment thread
stephclay marked this conversation as resolved.
Outdated
if result.succeeded() {
h.metrics.produceRequestsAttemptsSuccess.Observe(float64(attempts))
} else {
h.metrics.produceRequestsAttemptsFailure.Observe(float64(attempts))
return
}
h.metrics.produceRequestsAttemptsFailure.Observe(float64(attempts))
// Caller cancellation is not a terminal cluster outcome.
if callerErr := ctx.Err(); callerErr != nil && errors.Is(result.error(), callerErr) {
return
}
Comment thread
stephclay marked this conversation as resolved.
h.observeProduceFinalOutcome(shouldHedge)
}

// Check the hedging delay to apply to this request.
delay, shouldHedge := h.shouldHedge(time.Now(), primaryID, routedPartitions)

// The rest of the Hedger works with unrouted partitions, because it will be
// responsible to route partitions to other candidate agents during hedging
// and retries.
Expand Down Expand Up @@ -172,8 +193,8 @@ func (h *Hedger) ProduceSync(ctx context.Context, primaryID int32, routedPartiti
// primary and returns whichever produces a usable outcome first. Shared by
// the delay==0 fast path and the timer.C branch below, so a future edit to
// this behavior can't land in only one of them.
raceWithPrimary := func() ProduceResult {
hedged := h.runHedgingAttemptsAndRaceWithPrimary(workCtx, primaryID, partitions, candidates, primaryCh)
raceWithPrimary := func(trigger hedgeTrigger) ProduceResult {
hedged := h.runHedgingAttemptsAndRaceWithPrimary(workCtx, primaryID, partitions, candidates, primaryCh, trigger)
observeAttempts(hedged.result, hedged.attempts)
return hedged.result
}
Expand All @@ -183,7 +204,7 @@ func (h *Hedger) ProduceSync(ctx context.Context, primaryID int32, routedPartiti
// primaryCh is a coin flip. Skip the timer and always start
// the fallback race; either leg can still win.
if delay == 0 {
Comment thread
stephclay marked this conversation as resolved.
Outdated
return raceWithPrimary()
return raceWithPrimary(hedgeTriggerDemotedProbe)
Comment thread
stephclay marked this conversation as resolved.
Outdated
}

timer := time.NewTimer(delay)
Expand All @@ -196,12 +217,12 @@ func (h *Hedger) ProduceSync(ctx context.Context, primaryID int32, routedPartiti
return primaryResult
}

hedged := h.runHedgingAttempts(workCtx, primaryID, partitions, candidates)
hedged := h.runHedgingAttempts(workCtx, primaryID, partitions, candidates, hedgeTriggerPrimaryFailure)
result := selectProduceResult(primaryResult, hedged.result)
observeAttempts(result, hedged.attempts)
return result
case <-timer.C:
return raceWithPrimary()
return raceWithPrimary(hedgeTriggerLatency)
}
}

Expand All @@ -213,12 +234,20 @@ func (h *Hedger) ProduceSync(ctx context.Context, primaryID int32, routedPartiti
}

// The primary has failed. Try secondaries.
hedged := h.runHedgingAttempts(workCtx, primaryID, partitions, candidates)
hedged := h.runHedgingAttempts(workCtx, primaryID, partitions, candidates, hedgeTriggerPrimaryFailure)
result := selectProduceResult(primaryResult, hedged.result)
observeAttempts(result, hedged.attempts)
return result
}

func (h *Hedger) observeProduceFinalOutcome(shouldHedge bool) {
reason := produceFinalOutcomeAllCandidatesExhausted
Comment thread
koloss2001 marked this conversation as resolved.
Outdated
if !shouldHedge {
reason = produceFinalOutcomeHedgingSuppressed
}
h.metrics.produceFinalOutcome[reason].Inc()
}

func (h *Hedger) withCoverageCheck(res ProduceResult, nodeID int32, requested []encodedTopicPartitionRecords) ProduceResult {
if res.err != nil {
return res
Expand All @@ -238,10 +267,10 @@ func (h *Hedger) withCoverageCheck(res ProduceResult, nodeID int32, requested []
// cancel as soon as this function returns, which unwinds the losing leg.
// The reported attempts is the depth of the winning leg: 1 when the
// primary wins, the fallback's own depth when the fallback wins.
func (h *Hedger) runHedgingAttemptsAndRaceWithPrimary(workCtx context.Context, primaryID int32, partitions []encodedTopicPartitionRecords, candidates *hedgerCandidates, primaryCh <-chan ProduceResult) hedgerProduceResult {
func (h *Hedger) runHedgingAttemptsAndRaceWithPrimary(workCtx context.Context, primaryID int32, partitions []encodedTopicPartitionRecords, candidates *hedgerCandidates, primaryCh <-chan ProduceResult, trigger hedgeTrigger) hedgerProduceResult {
fallbackCh := make(chan hedgerProduceResult, 1)
go func() {
fallbackCh <- h.runHedgingAttempts(workCtx, primaryID, partitions, candidates)
fallbackCh <- h.runHedgingAttempts(workCtx, primaryID, partitions, candidates, trigger)
}()

// Either side's "wait for the other" branch is bounded: both legs
Expand Down Expand Up @@ -278,8 +307,9 @@ func (h *Hedger) runHedgingAttemptsAndRaceWithPrimary(workCtx context.Context, p
// MaxHedgeAgents, or ctx is canceled. The reported attempts is the total
// attempt depth: 1 (the primary, which this function models as the first
// tried agent) plus one per hedge wave dispatched.
func (h *Hedger) runHedgingAttempts(workCtx context.Context, primaryID int32, partitions []encodedTopicPartitionRecords, candidates *hedgerCandidates) (out hedgerProduceResult) {
func (h *Hedger) runHedgingAttempts(workCtx context.Context, primaryID int32, partitions []encodedTopicPartitionRecords, candidates *hedgerCandidates, trigger hedgeTrigger) (out hedgerProduceResult) {
h.metrics.hedgeAttemptsTotal.Inc()
h.metrics.hedgeTriggers[trigger].Inc()
// Count a win only when the result is successful AND we weren't
// preempted by ctx cancellation (e.g. the racing variant aborting
// the fallback because the primary won).
Expand Down
56 changes: 56 additions & 0 deletions pkg/wgo/hedger_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,8 @@ func TestHedger_ProduceSync(t *testing.T) {
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeWinsTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(m.produceRequestsPrimaryTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceRequestsHedgeTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeAllCandidatesExhausted]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeHedgingSuppressed]))
// Primary won → one success observation of attempt depth 1.
count, sum := histogramCountSum(t, m.produceRequestsAttemptsSuccess.(prometheus.Histogram))
assert.Equal(t, uint64(1), count)
Expand Down Expand Up @@ -241,6 +243,9 @@ func TestHedger_ProduceSync(t *testing.T) {
}
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeAttemptsTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeWinsTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerLatency]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerPrimaryFailure]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerDemotedProbe]))
assert.Equal(t, float64(1), testutil.ToFloat64(m.produceRequestsPrimaryTotal))
assert.GreaterOrEqual(t, testutil.ToFloat64(m.produceRequestsHedgeTotal), float64(1))
})
Expand All @@ -262,6 +267,9 @@ func TestHedger_ProduceSync(t *testing.T) {
assert.Contains(t, callIDs, secondaryID)
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeAttemptsTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeWinsTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerLatency]))
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerPrimaryFailure]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerDemotedProbe]))
assert.Equal(t, float64(1), testutil.ToFloat64(m.produceRequestsPrimaryTotal))
assert.GreaterOrEqual(t, testutil.ToFloat64(m.produceRequestsHedgeTotal), float64(1))
// Primary failed, one hedge wave resolved it → success attempt depth 2.
Expand Down Expand Up @@ -343,6 +351,9 @@ func TestHedger_ProduceSync(t *testing.T) {
require.NoError(t, capture.get(topic, partition).err)
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeAttemptsTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeWinsTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerLatency]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerPrimaryFailure]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerDemotedProbe]))
assert.Equal(t, float64(1), testutil.ToFloat64(m.produceRequestsPrimaryTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceRequestsHedgeTotal))
})
Expand Down Expand Up @@ -414,13 +425,52 @@ func TestHedger_ProduceSync(t *testing.T) {
// fallback candidate.
assert.Equal(t, float64(1), testutil.ToFloat64(m.produceRequestsPrimaryTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceRequestsHedgeTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeAllCandidatesExhausted]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeHedgingSuppressed]))
// No hedge candidate → only the primary attempt was made before
// giving up, recorded under the failure outcome.
count, sum := histogramCountSum(t, m.produceRequestsAttemptsFailure.(prometheus.Histogram))
assert.Equal(t, uint64(1), count)
assert.Equal(t, float64(1), sum)
})

t.Run("hedging suppressed and primary fails: cascade is classified as primary failure", func(t *testing.T) {
producer := newMockDirectProducer()
producer.respFn = successResp
producer.errs[primaryID] = kerr.RequestTimedOut

m := newMetrics(prometheus.NewPedanticRegistry())
h := NewHedger(producer, NewAverageAgentStatsTracker(), stratPrimaryAndSecondary, health, cfg, 0, 1<<20, m, nil)

capture := newResultCapture()
runHedger(h, context.Background(), makeReq(capture))

require.NoError(t, capture.get(topic, partition).err)
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeAttemptsTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerLatency]))
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerPrimaryFailure]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerDemotedProbe]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeHedgingSuppressed]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeAllCandidatesExhausted]))
})

t.Run("hedging suppressed and primary fails with no secondary: hedging_suppressed_and_primary_failed", func(t *testing.T) {
emptyStrat := &mockPartitionAssignmentStrategy{
candidates: map[partitionKey][]Agent{{topic, partition}: healthyAgents(primaryID)},
}
producer := newMockDirectProducer()
producer.errs[primaryID] = kerr.RequestTimedOut
m := newMetrics(prometheus.NewPedanticRegistry())
h := NewHedger(producer, NewAverageAgentStatsTracker(), emptyStrat, health, cfg, 0, 1<<20, m, nil)

capture := newResultCapture()
runHedger(h, context.Background(), makeReq(capture))

require.Error(t, capture.get(topic, partition).err)
assert.Equal(t, float64(1), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeHedgingSuppressed]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeAllCandidatesExhausted]))
})

t.Run("partial fallback coverage: any partition exhausting candidates stops the whole attempt", func(t *testing.T) {
// Two partitions; partition 0 has a fallback (secondaryID),
// partition 1 has none. Primary fails. runHedgingAttempt bails
Expand Down Expand Up @@ -517,6 +567,9 @@ func TestHedger_ProduceSync(t *testing.T) {
assert.True(t, sawFallback)
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeAttemptsTotal))
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeWinsTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerLatency]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerPrimaryFailure]))
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeTriggers[hedgeTriggerDemotedProbe]))
})

t.Run("demoted primary is the only candidate: hedge fires but finds no fallback; primary's success wins", func(t *testing.T) {
Expand Down Expand Up @@ -834,6 +887,9 @@ func TestHedger_ProduceSync(t *testing.T) {
assert.ErrorIs(t, err, context.Canceled)
assert.Equal(t, float64(1), testutil.ToFloat64(m.hedgeAttemptsTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.hedgeWinsTotal))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeAllCandidatesExhausted]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeHedgingSuppressed]))
assert.Equal(t, float64(0), testutil.ToFloat64(m.produceFinalOutcome[produceFinalOutcomeNoAgentAssigned]))
})

t.Run("fallback wins: primary leg is canceled via workCtx instead of running until its per-attempt deadline", func(t *testing.T) {
Expand Down
Loading
Loading