fix(router): sync both usage keys from Redis even when one increment is zero

A stream counted before its usage is known increments TPM by zero, so the
worker that served it never refreshed its local TPM value from Redis and
the first byte headers reported the token count another worker had already
consumed. Both pipeline operations now always run, matching the pre-change
callback, so the returned values refresh both worker local keys

Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
yassin 2026-09-16 19:45:48 +00:00
parent cb26454034
commit 7ff8945da4
2 changed files with 13 additions and 3 deletions

View file

@ -8281,7 +8281,6 @@ class Router:
pipeline_operations: Final[list[RedisPipelineIncrementOperation]] = [
RedisPipelineIncrementOperation(key=key, increment_value=increment_value, ttl=RoutingArgs.ttl.value)
for key, increment_value in ((tpm_key, total_tokens), (rpm_key, rpm_increment))
if increment_value > 0
]
post_increment_values: Final = await self.cache.async_increment_cache_pipeline(
increment_list=pipeline_operations,

View file

@ -1041,7 +1041,7 @@ async def test_acompletion_stream_counts_request_before_headers_and_tokens_once_
headers = _ratelimit_headers(stream)
assert headers["x-ratelimit-remaining-tokens"] == 1000
assert headers["x-ratelimit-remaining-requests"] == 99
assert await router.get_model_group_usage("gpt-5-mini") == (None, 1)
assert await router.get_model_group_usage("gpt-5-mini") == (0, 1)
chunks = [chunk async for chunk in stream]
total_tokens = chunks[-1].usage.total_tokens
@ -1078,7 +1078,7 @@ async def test_deployment_callback_on_success_adds_only_uncounted_tokens():
)
assert tpm_key is not None
assert await router.get_model_group_usage("gpt-5-mini") == (40, None)
assert await router.get_model_group_usage("gpt-5-mini") == (40, 0)
class _GatedIncrementCache(DualCache):
@ -1228,6 +1228,17 @@ async def test_headers_on_fresh_worker_reflect_shared_redis_usage():
assert headers["x-ratelimit-remaining-requests"] == 96
assert headers["x-ratelimit-remaining-tokens"] == 1000 - tokens_on_a - response.usage.total_tokens
counted_tokens = tokens_on_a + response.usage.total_tokens
for _ in range(2):
response = await worker_a.acompletion(model="gpt-5-mini", messages=messages, mock_response="pong")
counted_tokens += response.usage.total_tokens
stream = await worker_b.acompletion(model="gpt-5-mini", messages=messages, mock_response="pong", stream=True)
stream_headers = _ratelimit_headers(stream)
assert stream_headers["x-ratelimit-remaining-requests"] == 93
assert stream_headers["x-ratelimit-remaining-tokens"] == 1000 - counted_tokens
assert [chunk async for chunk in stream]
@pytest.mark.asyncio
async def test_get_model_group_io_token_usage_sums_across_deployments():