fix(disagg): dedup per-layer enqueues so high-cache-hit can't hang (lever A)
E2E at high cache-hit + concurrency exposed a vicious cycle: a context stays active until finish, but the notifier re-fires every layer on every subsequent forward and note_enqueued counted each, so _enqueued (the finish target) grew by the layer count each forward (target=3042-5538 observed) faster than 4 workers can drain -> finish never reaches it -> 30s timeout -> context stays active -> repeat. TTFT p99 = 116s. Fix: note_enqueued(layer_id) dedups per layer (target caps at the layer count); the first fire for a layer is from the request's own forward so its event is correct. Also guard register() against overwriting an active room (was leaking + re-registering, registered=917 for ~100 reqs). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -102,6 +102,20 @@ class TestPerLayerTransferContext(unittest.TestCase):
|
||||
ctx.submit_layer(0, None)
|
||||
self.assertEqual(ctx.finish(), -9)
|
||||
|
||||
def test_repeated_enqueues_dedup_so_finish_cannot_hang(self):
|
||||
# the notifier re-fires every layer on every subsequent forward while the ctx
|
||||
# is still active; note_enqueued must dedup so target stays capped at the layer
|
||||
# count and finish() completes (the high-cache-hit hang was target>>layers).
|
||||
eng = _FakeEngine()
|
||||
ctx = _ctx(eng, num_layers=3)
|
||||
for _forward in range(5): # 5 forwards all re-fire layers 0,1,2
|
||||
for L in range(3):
|
||||
self.assertEqual(ctx.note_enqueued(L), _forward == 0) # True only 1st
|
||||
if _forward == 0:
|
||||
ctx.submit_layer(L, None)
|
||||
self.assertEqual(len(eng.submits), 3) # 3 unique despite 15 note_enqueued
|
||||
self.assertEqual(ctx.finish(timeout=0.5), 0) # completes, no 30s timeout
|
||||
|
||||
def test_finish_falls_back_for_notifier_missed_layers(self):
|
||||
# only layer 0 fired by the notifier; finish must SYNCHRONOUSLY transfer the
|
||||
# missing layers 1,2 (e.g. an MTP buffer not hooked) -> all 3 moved, success.
|
||||
@@ -155,8 +169,12 @@ class _MockCtx:
|
||||
self.finished = False
|
||||
self.failed = False
|
||||
|
||||
def note_enqueued(self):
|
||||
def note_enqueued(self, layer_id):
|
||||
if layer_id in getattr(self, "_enq_layers", set()):
|
||||
return False
|
||||
self._enq_layers = getattr(self, "_enq_layers", set()) | {layer_id}
|
||||
self.enqueued += 1
|
||||
return True
|
||||
|
||||
def submit_layer(self, layer_id, event):
|
||||
self.submitted.append((layer_id, event))
|
||||
|
||||
Reference in New Issue
Block a user