fix(stream): disjoint key namespaces, O(N) fence drain

- Seq key moved from `{prefix}:stream:{user}:{msg}:seq` to
  `{prefix}:streamseq:{user}:{msg}`. A message_id containing a literal
  `:seq` suffix would otherwise make its stream key equal to another
  message's seq key, which under a user-controlled message_id would
  collide INCR against an XADD on the same Redis key.

- clearResumeFence swap-and-iterate instead of shift() in a loop.
  Array.shift is O(n) per call, so the old drain was O(n²). Practical
  frame counts during a fence window are small (dozens worst-case) so
  this was never a user-visible problem, but the fix is two lines and
  removes a Big-O footgun. Concurrency behavior unchanged: the batch
  is captured via reference-swap so frames arriving during an await
  continue buffering into the fresh empty array still present in the
  map, and the outer while loop drains those in the next iteration.

Not addressed:
- Concurrent-emitter out-of-seq live frames — already documented as a
  known limitation with a block comment at get_event_emitter; fix
  requires distributed locking.
- Prune resumeSeqByMessageId on per-message terminal events —
  deliberately removed two rounds ago because it caused continuation-
  reuses-message_id to replay duplicate content; the trade-off
  (bounded int-per-message memory vs. correctness) is already
  captured in that commit message.
This commit is contained in:
Claude 2026-04-15 07:24:45 +00:00
parent 1eae0f793a
commit 1964c39ac1
No known key found for this signature in database
2 changed files with 23 additions and 18 deletions

View file

@ -224,7 +224,10 @@ def _stream_key(user_id: str, message_id: str) -> str:
def _stream_seq_key(user_id: str, message_id: str) -> str:
return f'{REDIS_KEY_PREFIX}:stream:{user_id}:{message_id}:seq'
# Distinct top-level namespace (`streamseq`, not `stream`) so a
# message_id containing delimiter-like characters can't collide the
# seq key of one message with the stream key of another.
return f'{REDIS_KEY_PREFIX}:streamseq:{user_id}:{message_id}'
async def _stream_seq_allocate(user_id: str, message_id: str):

View file

@ -648,29 +648,31 @@
};
// Drop fence AND flush buffered events (happy path + timeout).
// Drain via shift(), keeping the queue present in the map so live
// frames arriving during an await inside chatEventHandler keep
// buffering into the same queue instead of racing ahead. Flushed
// events are marked `_replayed` before dispatch so chatEventHandler's
// fence-buffering branch short-circuits and doesn't re-queue them
// (which would be an infinite loop).
// Swap the queue array out for an empty one each batch so live frames
// arriving during an await keep buffering (into the fresh empty
// array that's still in the map) instead of racing ahead. Iterate
// the captured batch linearly, then loop until no new frames came
// in. O(N) drain instead of O(N²) from shift().
const clearResumeFence = async (messageId) => {
const timer = resumeFenceTimerByMessageId.get(messageId);
if (timer) {
clearTimeout(timer);
resumeFenceTimerByMessageId.delete(messageId);
}
const queue = resumeQueueByMessageId.get(messageId);
if (!queue) return;
while (queue.length > 0) {
const event = queue.shift();
if (event && typeof event === 'object') {
event._replayed = true;
}
try {
await chatEventHandler(event);
} catch (e) {
console.error('resume fence flush error', e);
if (!resumeQueueByMessageId.has(messageId)) return;
while (true) {
const batch = resumeQueueByMessageId.get(messageId);
if (!batch || batch.length === 0) break;
resumeQueueByMessageId.set(messageId, []);
for (const event of batch) {
if (event && typeof event === 'object') {
event._replayed = true;
}
try {
await chatEventHandler(event);
} catch (e) {
console.error('resume fence flush error', e);
}
}
}
resumeQueueByMessageId.delete(messageId);