Repository navigation
fix(chat): lock heartbeat, per-thread drain isolation, debounce drain (#190) - #261
Conversation
…#190) Port upstream eccc6b91 (#656), 076fe5dc (#659), 6cb933eb (#832) and 5b538f6f (#821): - _LockHeartbeat renews held locks every DEFAULT_LOCK_TTL_MS / 3 for drop/queue/debounce/burst, capped by the new ConcurrencyConfig.max_lock_lifetime_ms (default 600000). stop() waits for an in-flight extend before release_lock (_with_held_lock). - Drains dispatch each message under its own thread_id with skipped filtered to that thread (channel-scoped lock isolation, security). - Debounce uses max_queue_size, accumulates superseded messages into context.skipped, and loops after dispatch; the Python-only 20-iteration cap is removed. - _set_mention_flags detects mentions on skipped messages; skipped mentions route to on_mention when mention handlers exist. - StateAdapter.extend_lock documents the token-compare contract. Divergence: lifetime cap measured on a monotonic clock (see docs/UPSTREAM_SYNC.md). Part of #184
|
Warning Review limit reachedYou've used all free OSS reviews for now. Wait for the free limit to reset to keep reviewing this public repository. Next included review available in 17 minutes. View limit detailsLimit details: You’ve used the included review currently available. Review configuration: ⚙️ Run configurationConfiguration used: Repository UI Review profile: CHILL Plan: Advanced Run ID: 📒 Files selected for processing (9)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
…beat crash logging; drain coverage (#190)
…t, token-checked drain ownership (#190)
…e exit, monotonic held_until max, cap-gated ownership confirm (#190)
…nd dequeued messages back (#190)
…ection, cancel stalled extend, capacity-safe hand-back (#190)
# Conflicts: # scripts/fidelity_target.json
Mirror chat@4.41.1 withHeldLock / startLockHeartbeat / debounceLoop / drainQueue. Ownership is checked only at the top of each drain/debounce iteration; a collected batch is dispatched (upstream parity, loss-free). Removed (no Python-specific hazard): monotonic clock, confirm_ownership, may_start_step + _renewing, post-collection recheck, _hand_back, bounded collection, settle_in_flight, bounded/cancelling stop(), max() on held_until, (lock, acquired_at) plumbing, done-callbacks. Kept (Python-specific): client-side held_until seed (Postgres stamps expires_at with the DB clock), shielded extend task, non-yielding idle stop(), max_lock_lifetime_ms validation, crash logging in _run, camelCase threadId rehydrate fallback.
A slow extend_lock pushed every later tick back by its duration, because the loop slept a full TTL/3 after each extend returned. Upstream's setInterval keeps ticking at started_at + k*TTL/3 and skips ticks that fire while an extend is in flight. Match that schedule so a slow backend cannot push the next extend past the lock's expiry. Found by gpt-6-astra review round 1 on PR #261.
setInterval fires on a monotonic timer. Tick deadlines computed on the epoch clock let a wall-clock step back push the next extend past the lock's TTL on the backend. Ticks now use _monotonic_ms(); the lifetime cap and held_until stay on epoch time, as upstream's Date.now(). Found by gpt-6-astra review round 2 on PR #261.
|
Merge gate: CI green (Lint & Type Check, test (3.12), test (3.13), CodeQL, Analyze (python), Analyze (actions) all pass on 3118274); local Codex review (gpt-6-astra, xhigh, --base origin/main) on 3118274: no actionable findings, and the full suite passed in the reviewer's run; PR re-centered on upstream withHeldLock/heartbeat design after earlier rounds accreted Python-only machinery (see ## Merge gate). Merging with --admin (Protect Main requires a code-owner approval). |
Summary
Ports upstream's
Chatconcurrency rework. Before this change, a handler running longer than 30s lost its thread lock, so a second message could run concurrently on the same thread. With a channel-scoped lock, a queued message could also be answered in the lock holder's thread instead of its own._LockHeartbeat+Chat._with_held_lockrenew a held lock everyDEFAULT_LOCK_TTL_MS / 3(10s) while the handler runs, fordrop(including theon_lock_conflictforce path),queue,debounceandburst. Renewal stops after the newConcurrencyConfig.max_lock_lifetime_ms(default600_000); the lock then lapses one TTL later. An extend returningFalsemarks ownership lost. So does an extend that raises once the wall clock has passedheld_until; one that raises earlier only logs a warning.stop()cancels the loop and waits for any in-flight extend beforerelease_lock._drain_queue/_debounce_loopno longer takethread_id. Each drained message is dispatched under its ownthread_id,skippedis filtered to that thread, andtotal_since_last_handler = len(skipped) + 1. Both loops checkis_ownership_lost()and leave the queue to the new holder. The inlineextend_lock/ "Lock lost ... aborting" branches are gone.max_queue_sizecapacity for both the busy path and self-enqueue (was 1). Each tick drains all pending messages, adds superseded ones toskipped, dispatches with aMessageContext, resets, and loops again, so a message that arrives mid-handler is processed without a new webhook. The Python-onlymax_iterations = 20is removed;max_lock_lifetime_msbounds the loop. As upstream, the drop-newest early return still excludes debounce.extend_lockis gone; the heartbeat covers the window._set_mention_flags(adapter, message, context) -> boolfillsis_mentionon the message and eachcontext.skipped(keeping today'sorsemantics; [4.41/C2a] Core mentions & message model: tri-state is_mention, mention regex, Author.email/is_system, Message.reply_to #192 flips it). Routing is nowmessage.is_mention or (has_mention and self._mention_handlers).ConcurrencyConfig.max_lock_lifetime_ms: int = 600_000, resolved in all three config branches (an explicitNonefalls back to the default, like??). TheStateAdapter.extend_lockdocstring states the token-compare contract; memory/Redis/Postgres already comply._rehydrate_messagenow also reads camelCasethreadId. Drains dispatch under the message's own thread id now, so an untagged camelCase dict would otherwise have dispatched to"".Upstream commits mapped
eccc6b91fix(chat): detect mentions in skipped queued messages (#656)_set_mention_flagsovercontext.skipped076fe5dcfix(chat): preserve skipped mention routing (#659)max_queue_size, drain-all + skipped accumulation,has_mention and self._mention_handlersrouting6cb933ebfix(chat): isolate channel-scoped queue dispatch by thread (#832)latest.thread_id, thread-filteredskipped5b538f6ffix(chat): keep thread locks alive during long handlers (#821)_LockHeartbeat,_with_held_lock,max_lock_lifetime_ms,is_ownership_lost()checks, debounce loops after dispatch,extend_lockcontractTests ported (
tests/test_chat_faithful.py, fromchat.test.ts)it.each(["queue","burst","debounce"])"should keep the %s lock alive while a handler is running" →test_should_keep_the_lock_alive_while_a_handler_is_running[queue|burst|debounce]activeConversation()assertions are left out; they belong to [4.41/C3] Conversation context + AI tool scoping (read & write guards, strict_scope) #195.is_mention, [4.41/C2a] Core mentions & message model: tri-state is_mention, mention regex, Author.email/is_system, Message.reply_to #192).tests/_fake_clock.pystands in forvi.useFakeTimers()andinstallTokenLockMock. It swapschat_sdk.chat._sleep,_now_msand_monotonic_msfor virtual time (withstep_wall_clock()to move only the epoch clock), and provides a token-checked, expiry-aware lock on the same clock. There are no real sleeps in the new tests: the 90-second hung-handler scenario runs in about a millisecond.Python-specific (
tests/test_chat_lock_heartbeat.py):setIntervaltranslation): a slow extend (10s -> 26s) is followed by the next at 30s, not 36s; a 60s wall-clock step back does not move ticks.drop: a second message arriving after the TTL raisesLockError, and the heartbeat task is cancelled after release. The force-acquired lock (not the stale holder's) is the one renewed.stop(): waits for an in-flight extend beforerelease_lock; cancelling the caller does not cancel the extend; an idlestop()does not yield (back-to-back queue messages are both handled); a handler error still stops the heartbeat and releases exactly once.held_untilkeeps ownership. Pastheld_until, the drain stops, the queued message stays for the next holder, andrelease_lockstill runs. AFalseextend stops renewal after one call.held_untilseed: a database clock 31s behind the app does not make a fresh lock look lapsed.skippedresets.to_json()-tagged, untagged snake_case, and untagged camelCase.max_lock_lifetime_msresolves in every config branch and is validated; a crashed renewal loop is logged.Each kept Python-specific item was broken on purpose (idle stop that yields, seeding from
expires_at, no shield instop(), no shield in the loop, no crash log, fixed-delay ticks, epoch-clock ticks); each change made a test fail.Fidelity
--strictat the pin stays 733/733.Regenerated after merging
origin/main(bd5e45a). Against main's committed report:Divergences (recorded in
docs/UPSTREAM_SYNC.md, CHANGELOG "Python-specific")held_untilseed (the one behavioral divergence). Seeded from_now_ms() + DEFAULT_LOCK_TTL_MSwhen the heartbeat starts, notLock.expires_at: the Python Postgres backend stampsexpires_atwith the database clock (upstream state-pg uses the client clock), so a database clock 20-30s behind the app would make a fresh lock look lapsed. Each successful extend sets_now_ms() + TTL, as upstream.setInterval+if (inFlight) return). Each extend runs as its own task behindasyncio.shield(asyncio cancellation interrupts in-flight I/O; a JS promise cannot be cancelled).stop()awaits the in-flight extend (shielded) only if one exists; idle, it returns without yielding (awaiting a cancelled task costs an event-loop turn thatclearIntervaldoes not)._runlogs an unexpected exception aterror(the interval callback's.catch).max_lock_lifetime_msvalidation (config, not behavioral). Non-negativeintatChatinit; JS coerces a numeric string where Python would raiseTypeErroron the first tick.docs/UPSTREAM_SYNC.mditem 4 also lists the parity-accepted races (the heartbeat is best-effort liveness, not a fencing token), so a review finding that re-raises one can be closed by citation.Other doc changes: the
UPSTREAM_SYNC.md"must stay 1:1" item no longer claims a 20-iteration debounce cap, and it records themax_iterationsremoval.docs/ARCHITECTURE.mdnow describes the heartbeat and the new drain/debounce flow. The stale Python-onlytest_burst_preserves_skipped_unlike_debounceis removed: debounce now keeps skipped context too, and what the test still checked duplicated the faithful burst-collapse test.Consumer impact (high)
dropstrategy: a second message that arrives 30s or more into a long handler (for example a long Slack/Teams stream) now raisesLockError(logged "Could not acquire lock on thread") instead of running concurrently. Downstream consumers such as chinchill should confirm their strategy (queue/debounce/burst, oron_lock_conflict).contextwhoseskippedholds the superseded messages (it used to beNone). Messages that arrive mid-handler are processed afterwards.lock_scope="channel"): queued and debounced messages are answered in their own thread, andskippedonly contains messages from that thread. As upstream, a pending message from another thread that is not the latest is dropped rather than surfaced.on_mention.message-queued/message-debounce-resetgainqueue_depth; drain logs carry the message's ownthread_idpluslock_key; the "Lock lost ... aborting" warnings are replaced.Merge gate
Re-centered on upstream (ebed4dd). Five earlier gpt-6-astra rounds (R1-R5 on 0541ff2..067b7d0) had grown
chat.pyby +674 lines of Python-only machinery, each piece answering one finding: a monotonic lifetime clock,confirm_ownership(),may_start_step(), a post-collection ownership recheck,_hand_back()re-enqueue, bounded collection,settle_in_flight(), a bounded/cancellingstop(), amax()onheld_until, the(lock, acquired_at)tuple plumbing and task done-callbacks. R5 still reported 3 P2s, all inside_hand_back(ordering, non-atomic capacity check, Redis whole-list PEXPIRE reset). Per the Divergence Policy (match upstream unless a Python-specific hazard makes it wrong here), all of that was removed:heartbeat.is_ownership_lost()only at the top of each iteration (chat@4.41.1chat.tsdebounceLoop / drainQueue) and dispatch a batch they already collected. Dequeue is destructive, so dispatching is the only loss-free choice: dropping loses messages, and handing back needs an atomic push-front-if-room primitive theStateAdapterprotocol lacks. The cost is a brief overlap with a new holder, which upstream accepts. No check inside the drain can close that race without fencing tokens on handler side effects.skipped, as upstream.stop()isstopped = True; cancel; await shield(in_flight)with no bound, as upstream'sawait inFlight.chat.pydiff: +674 -> +444;test_chat_lock_heartbeat.py: 1258 -> 623 lines.gpt-6-astra on the re-centered diff: 3 rounds. Final verdict: clean (no actionable findings).
setIntervalis fixed-rate, so a slow backend could push the next extend past the lock's expiry. Ticks now run on fixed deadlines and skip ticks that came due while an extend was in flight. Test:TestHeartbeatCadence::test_slow_extend_keeps_the_fixed_rate_tick_schedule(failed before the fix).setIntervalis monotonic, so a wall-clock step back could delay renewal past the TTL. Ticks now use_monotonic_ms(); the lifetime cap andheld_untilstay on epoch time, as upstream'sDate.now(). Test:TestHeartbeatCadence::test_wall_clock_step_back_does_not_delay_the_next_tick(failed before the fix).docs/UPSTREAM_SYNC.mditem 4); the R5 hand-back P2s are moot because the hand-back is gone.CI on 3118274: all green (Lint & Type Check, test 3.12, test 3.13, CodeQL). Local validation green: ruff, ruff format, audit (0 hard failures),
--check-docs,--strictat chat@4.31.0 (all TS tests have Python equivalents), pyrefly 0 errors, pytest 6208 passed / 24 skipped.Closes #190
Part of #184