Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ This client has been designed around the following non-negotiable assumptions:
1. **Warpstream-specific.** Hedging the same batch across agents only works because any agent can serve any partition. Pointed at vanilla Kafka, the secondary leg would fail with `NotLeaderForPartition`.
2. **At-least-once delivery only.** Duplicates are tolerable. Any code that assumes exactly-once or in-partition record ordering must stay on franz-go.
3. **No transactional or idempotent producer support.** `DisableIdempotentWrite()` semantics are baked in — no `producerId`/`producerEpoch`/`baseSequence` handshake.
4. **Produce never blocks on Metadata.** The agent pool is refreshed on a timer and also on demand. A routing miss fails that Produce immediately and asks for a refresh; a later Produce can use the updated pool. A refresh that excluded a leader asks for another fetch even when that Produce already succeeded on a stand-in. Neither waits on the fetch. On-demand refreshes are coalesced and paced by `OnDemandMetadataRefreshInterval` (default 1s).
4. **Produce never blocks on Metadata.** The agent pool is refreshed on a timer and also on demand. A routing miss fails that Produce immediately and asks for a refresh; a later Produce can use the updated pool. A refresh that excluded a leader asks for another fetch on its own, with no Produce in flight. Those follow-ups start at `OnDemandMetadataRefreshInterval` (default 1s) and double up to `MetadataRefreshInterval` (default 10s) while the exclusion continues. The chain stops when a refresh excludes nobody: that fetch does not ask for another, the periodic refresh runs, and the gap returns to the floor. Neither the Produce nor the follow-up waits on the fetch. On-demand refreshes are coalesced.

## How it works

Expand All @@ -48,7 +48,7 @@ For every record the client asks a `PartitionAssignmentStrategy` for an ordered
- **Deterministic.** Given the same Metadata view, every client instance picks the same primary and the same secondary for a given partition. Hedge load is predictable and analysable instead of randomly smeared across agents.
- **State-aware.** A wrapper around the base strategy (the **Demoter**, see below) can mark an agent as demoted so it's elided from the candidate list or surfaced as a probe.

A partition's leader can briefly go missing from that map during a Metadata refresh, even though nothing is actually wrong. When the topic is known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. That call succeeds without waiting. The refresh that excluded the leader asks for another fetch on its own. A topic is known when another partition still has a leader, and also when this refresh listed the topic's partitions and kept none of the leaders it named. A partition Metadata reported with no leader (`Leader` below 0) is left unroutable, including when a sibling partition still has one, so that produce fails and still gets an on-demand refresh. A topic that has never appeared in Metadata, or that came back with an error, is left alone for the same reason.
A partition's leader can briefly go missing from that map during a Metadata refresh, even though nothing is actually wrong. When the topic is known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. That call succeeds without waiting. The refresh that excluded the leader asks for another fetch on its own, even when no Produce is in flight. Further exclusions double the wait from `OnDemandMetadataRefreshInterval` up to `MetadataRefreshInterval`. Repeated routing misses climb the same way. While on-demand refreshes keep being requested, the periodic refresh does not run, so the gap stays at the ceiling until they stop. A refresh that excludes nobody stops the follow-ups, and the next periodic refresh puts the wait back. A topic is known when another partition still has a leader, and also when this refresh listed the topic's partitions and kept none of the leaders it named. A partition Metadata reported with no leader (`Leader` below 0) is left unroutable, including when a sibling partition still has one, so that produce fails and still gets an on-demand refresh. A topic that has never appeared in Metadata, or that came back with an error, is left alone for the same reason.

### Buffering: linger by destination agent, not by partition

Expand Down
79 changes: 63 additions & 16 deletions pkg/wgo/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -377,6 +377,35 @@ func (c *WarpstreamClient) Close() {
})
}

// refreshBackoff paces on-demand Metadata fetches. The delay starts at floor,
Comment thread
koloss2001 marked this conversation as resolved.
// grows after each on-demand fetch up to ceiling, and returns to floor when a
// periodic fetch runs.
Comment thread
koloss2001 marked this conversation as resolved.
type refreshBackoff struct {
floor time.Duration
ceiling time.Duration
delay time.Duration
}

func newRefreshBackoff(floor, ceiling time.Duration) refreshBackoff {
return refreshBackoff{floor: floor, ceiling: ceiling, delay: floor}
}

func (b *refreshBackoff) current() time.Duration { return b.delay }

func (b *refreshBackoff) advance() {
// time.Duration is an int64. Doubling past MaxInt64 wraps negative.
// At or below half the ceiling, delay*2 still fits and stays within it.
if b.delay > b.ceiling/2 {
b.delay = b.ceiling
return
}
b.delay *= 2
}

func (b *refreshBackoff) reset() {
b.delay = b.floor
}

// startBackgroundRefresh owns every post-startup AgentPool.Refresh: the
// periodic ticker and on-demand nudges from triggerRefresh. One owner means
// Refresh is never concurrent. refreshCtx (cancelled by Close) stops the loop,
Expand All @@ -389,20 +418,42 @@ func (c *WarpstreamClient) startBackgroundRefresh() {
defer c.refreshWG.Done()
ticker := time.NewTicker(c.cfg.MetadataRefreshInterval)
defer ticker.Stop()
backoff := newRefreshBackoff(c.cfg.OnDemandMetadataRefreshInterval, c.cfg.MetadataRefreshInterval)
for {
periodic := false
select {
case <-c.refreshCtx.Done():
return
case <-c.refreshNowCh:
ticker.Reset(c.cfg.MetadataRefreshInterval)
startedAt := time.Now()
c.refreshPool(metadataRefreshTriggerOnDemand)
if !c.waitRefreshCooldown(time.Since(startedAt)) {
return
}
case <-ticker.C:
// A nudge queued at the same instant wins. At the ceiling the
// cooldown ends as the tick lands, and taking the tick would
// reset the delay.
select {
case <-c.refreshNowCh:
default:
periodic = true
}
}

if periodic {
c.refreshPool(metadataRefreshTriggerPeriodic)
// A nudge queued during this fetch must not start another after Close.
if c.refreshCtx.Err() != nil {
return
}
backoff.reset()
continue
}
Comment thread
koloss2001 marked this conversation as resolved.
ticker.Reset(c.cfg.MetadataRefreshInterval)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Question, not a blocker: the backoff resets only on a periodic tick, and this line pushes the tick back on every on-demand fetch. A workload that keeps producing routing misses, such as one that produces to an unknown topic, never gets a periodic tick. It stays at the 10s ceiling indefinitely. Before this change, those misses refreshed every 1s. Is that the trade-off that you want? If so, a short note in the README "How it works" section would help, because AGENTS.md asks for it when refresh behavior changes.

— via Claude Code

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.

Yes, this is intentional and follows the #95 discussion about sustained 1s all-topic Metadata refreshes creating excessive load. The first fetch can run immediately; subsequent gaps are 1s → 2s → 4s → 8s → 10s, then remain at 10s while nudges continue. Metadata still refreshes at the normal cadence.

The trade-off also applies to unknown topics: the backoff is shared across the client, so persistent misses for topic A can leave a newly encountered topic B waiting for the next 10s fetch, without its own fast retry sequence.

For now, I think that is acceptable. We can capture a future design for bounded resets when relevant metadata changes, I'll see if #87 has all the metrics to measure this and if not will add it there.

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.

The README is updated

startedAt := time.Now()
c.refreshPool(metadataRefreshTriggerOnDemand)
c.waitRefreshCooldown(backoff.current(), time.Since(startedAt))
// Return before select can take the nudge the fetch just queued.
if c.refreshCtx.Err() != nil {
return
}
backoff.advance()
Comment thread
cursor[bot] marked this conversation as resolved.
}
}()
}
Expand Down Expand Up @@ -431,8 +482,7 @@ func (c *WarpstreamClient) refreshPool(trigger metadataRefreshTrigger) {
}
c.demoter.Refresh(c.pool.Agents())
// Nudge so a stand-in does not wait for the next periodic tick. Produce does
// not block. After a periodic tick the follow-up starts at once; later
// repeats wait out OnDemandMetadataRefreshInterval (default 1s).
// not block. The refresh loop paces the follow-up.
c.noteLeaderDrops(dropped, true)
}

Expand All @@ -452,21 +502,18 @@ func (c *WarpstreamClient) noteLeaderDrops(dropped leaderDrops, nudge bool) {
}
}

// waitRefreshCooldown enforces the configured minimum between on-demand
// refresh start times. Time spent fetching already counts toward the interval.
// Returns false during close so the refresh loop exits without spinning.
func (c *WarpstreamClient) waitRefreshCooldown(elapsed time.Duration) bool {
remaining := c.cfg.OnDemandMetadataRefreshInterval - elapsed
// waitRefreshCooldown sleeps out the rest of delay. Time already spent
// fetching counts. Close cancels refreshCtx and the wait returns.
func (c *WarpstreamClient) waitRefreshCooldown(delay, elapsed time.Duration) {
remaining := delay - elapsed
if remaining <= 0 {
return c.refreshCtx.Err() == nil
return
}
t := time.NewTimer(remaining)
defer t.Stop()
select {
case <-c.refreshCtx.Done():
return false
case <-t.C:
return true
}
}

Expand Down
Loading
Loading