Skip to content

fix(chat): lock heartbeat, per-thread drain isolation, debounce drain (#190) - #261

Merged
patrick-chinchill merged 12 commits into
mainfrom
sync/4.41-c1a
Sep 30, 2026
Merged

patrick-chinchill merged 12 commits into
mainfrom
sync/4.41-c1a

Conversation

@patrick-chinchill

@patrick-chinchill patrick-chinchill commented Sep 30, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

Ports upstream's Chat concurrency 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.

  • Lock heartbeat. _LockHeartbeat + Chat._with_held_lock renew a held lock every DEFAULT_LOCK_TTL_MS / 3 (10s) while the handler runs, for drop (including the on_lock_conflict force path), queue, debounce and burst. Renewal stops after the new ConcurrencyConfig.max_lock_lifetime_ms (default 600_000); the lock then lapses one TTL later. An extend returning False marks ownership lost. So does an extend that raises once the wall clock has passed held_until; one that raises earlier only logs a warning. stop() cancels the loop and waits for any in-flight extend before release_lock.
  • Per-thread drain isolation (security, 6cb933eb). _drain_queue / _debounce_loop no longer take thread_id. Each drained message is dispatched under its own thread_id, skipped is filtered to that thread, and total_since_last_handler = len(skipped) + 1. Both loops check is_ownership_lost() and leave the queue to the new holder. The inline extend_lock / "Lock lost ... aborting" branches are gone.
  • Debounce drain. The debounce queue now uses max_queue_size capacity for both the busy path and self-enqueue (was 1). Each tick drains all pending messages, adds superseded ones to skipped, dispatches with a MessageContext, resets, and loops again, so a message that arrives mid-handler is processed without a new webhook. The Python-only max_iterations = 20 is removed; max_lock_lifetime_ms bounds the loop. As upstream, the drop-newest early return still excludes debounce.
  • Burst. The post-sleep one-off extend_lock is gone; the heartbeat covers the window.
  • Skipped mentions. _set_mention_flags(adapter, message, context) -> bool fills is_mention on the message and each context.skipped (keeping today's or semantics; [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 now message.is_mention or (has_mention and self._mention_handlers).
  • Types. ConcurrencyConfig.max_lock_lifetime_ms: int = 600_000, resolved in all three config branches (an explicit None falls back to the default, like ??). The StateAdapter.extend_lock docstring states the token-compare contract; memory/Redis/Postgres already comply.
  • Rehydration. The dict fallback in _rehydrate_message now also reads camelCase threadId. Drains dispatch under the message's own thread id now, so an untagged camelCase dict would otherwise have dispatched to "".

Upstream commits mapped

Commit Tag Ported as
eccc6b91 fix(chat): detect mentions in skipped queued messages (#656) chat@4.32.0 _set_mention_flags over context.skipped
076fe5dc fix(chat): preserve skipped mention routing (#659) chat@4.33.0 debounce capacity max_queue_size, drain-all + skipped accumulation, has_mention and self._mention_handlers routing
6cb933eb fix(chat): isolate channel-scoped queue dispatch by thread (#832) chat@4.38.1 dispatch with latest.thread_id, thread-filtered skipped
5b538f6f fix(chat): keep thread locks alive during long handlers (#821) chat@4.39.0 _LockHeartbeat, _with_held_lock, max_lock_lifetime_ms, is_ownership_lost() checks, debounce loops after dispatch, extend_lock contract

Tests ported (tests/test_chat_faithful.py, from chat.test.ts)

tests/_fake_clock.py stands in for vi.useFakeTimers() and installTokenLockMock. It swaps chat_sdk.chat._sleep, _now_ms and _monotonic_ms for virtual time (with step_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):

  • Tick schedule (setInterval translation): 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 raises LockError, 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 before release_lock; cancelling the caller does not cancel the extend; an idle stop() does not yield (back-to-back queue messages are both handled); a handler error still stops the heartbeat and releases exactly once.
  • An extend that raises before held_until keeps ownership. Past held_until, the drain stops, the queued message stays for the next holder, and release_lock still runs. A False extend stops renewal after one call.
  • held_until seed: a database clock 31s behind the app does not make a fresh lock look lapsed.
  • One debounce loop: a message arriving mid-handler is dispatched afterwards and skipped resets.
  • Channel-scoped isolation with dict queue entries: to_json()-tagged, untagged snake_case, and untagged camelCase.
  • max_lock_lifetime_ms resolves 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 in stop(), no shield in the loop, no crash log, fixed-delay ticks, epoch-clock ticks); each change made a test fail.

Fidelity

--strict at the pin stays 733/733.

Regenerated after merging origin/main (bd5e45a). Against main's committed report:

missing 250 -> 242 (-8)
  packages/chat/src/chat.test.ts: 31 -> 23 (-8)

Divergences (recorded in docs/UPSTREAM_SYNC.md, CHANGELOG "Python-specific")

  1. held_until seed (the one behavioral divergence). Seeded from _now_ms() + DEFAULT_LOCK_TTL_MS when the heartbeat starts, not Lock.expires_at: the Python Postgres backend stamps expires_at with 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.
  2. asyncio translation notes (not behavioral). Ticks at a fixed rate on a monotonic clock, skipping ticks that come due while an extend is in flight (setInterval + if (inFlight) return). Each extend runs as its own task behind asyncio.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 that clearInterval does not). _run logs an unexpected exception at error (the interval callback's .catch).
  3. max_lock_lifetime_ms validation (config, not behavioral). Non-negative int at Chat init; JS coerces a numeric string where Python would raise TypeError on the first tick.

docs/UPSTREAM_SYNC.md item 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 the max_iterations removal. docs/ARCHITECTURE.md now describes the heartbeat and the new drain/debounce flow. The stale Python-only test_burst_preserves_skipped_unlike_debounce is removed: debounce now keeps skipped context too, and what the test still checked duplicated the faithful burst-collapse test.

Consumer impact (high)

  • Default drop strategy: a second message that arrives 30s or more into a long handler (for example a long Slack/Teams stream) now raises LockError (logged "Could not acquire lock on thread") instead of running concurrently. Downstream consumers such as chinchill should confirm their strategy (queue/debounce/burst, or on_lock_conflict).
  • Debounce handlers now always get a context whose skipped holds the superseded messages (it used to be None). Messages that arrive mid-handler are processed afterwards.
  • Channel-scoped locks (Telegram topics, custom lock_scope="channel"): queued and debounced messages are answered in their own thread, and skipped only contains messages from that thread. As upstream, a pending message from another thread that is not the latest is dropped rather than surfaced.
  • Skipped mentions route to on_mention.
  • Logs: message-queued / message-debounce-reset gain queue_depth; drain logs carry the message's own thread_id plus lock_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.py by +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/cancelling stop(), a max() on held_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:

  • Drain and debounce check heartbeat.is_ownership_lost() only at the top of each iteration (chat@4.41.1 chat.ts debounceLoop / 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 the StateAdapter protocol 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.
  • On a top-of-loop loss: warn and return; debounce drops its accumulated skipped, as upstream.
  • stop() is stopped = True; cancel; await shield(in_flight) with no bound, as upstream's await inFlight.
  • Only the Python-specific items under Divergences above are kept. chat.py diff: +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).

  • R1 (ebed4dd): 1 P2, fixed in b0ff50a. The heartbeat slept a full TTL/3 after each extend returned (fixed-delay); upstream's setInterval is 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).
  • R2 (b0ff50a): 1 P2, fixed in 3118274. The R1 fix put tick deadlines on the epoch clock; setInterval is monotonic, so a wall-clock step back could delay renewal past the TTL. Ticks now use _monotonic_ms(); the lifetime cap and held_until stay on epoch time, as upstream's Date.now(). Test: TestHeartbeatCadence::test_wall_clock_step_back_does_not_delay_the_next_tick (failed before the fix).
  • R3 (3118274): no actionable findings. The full suite passed in the reviewer's run.
  • No finding in R1-R3 re-raised an upstream-parity race, so no parity rebuttals were needed. Earlier findings about upstream-identical windows are closed as "won't fix: upstream parity" (cited in docs/UPSTREAM_SYNC.md item 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, --strict at chat@4.31.0 (all TS tests have Python equivalents), pyrefly 0 errors, pytest 6208 passed / 24 skipped.

Closes #190
Part of #184

…#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
@coderabbitai

coderabbitai Bot commented Sep 30, 2026 •

Copy link
Copy Markdown

Warning

Review limit reached

You'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.

Check out review usage here.

View limit details

Limit details: You’ve used the included review currently available.

Learn how review limits work.

Review configuration:

⚙️ Run configuration

Configuration used: Repository UI

Review profile: CHILL

Plan: Advanced

Run ID: 2fe79dfd-c55f-4940-9862-0c99bffe0fe3

📥 Commits

Reviewing files that changed from the base of the PR and between 807e47b and 3118274.

📒 Files selected for processing (9)
  • CHANGELOG.md
  • docs/ARCHITECTURE.md
  • docs/UPSTREAM_SYNC.md
  • scripts/fidelity_target.json
  • src/chat_sdk/chat.py
  • src/chat_sdk/types.py
  • tests/_fake_clock.py
  • tests/test_chat_faithful.py
  • tests/test_chat_lock_heartbeat.py

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.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@patrick-chinchill
patrick-chinchill marked this pull request as ready for review September 30, 2026 10:53
…e exit, monotonic held_until max, cap-gated ownership confirm (#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.
@patrick-chinchill

Copy link
Copy Markdown
Collaborator Author

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).

@patrick-chinchill
patrick-chinchill merged commit 34bf614 into main Sep 30, 2026
7 checks passed
@patrick-chinchill
patrick-chinchill deleted the sync/4.41-c1a branch September 30, 2026 12:29
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[4.41/C1a] Core concurrency: lock heartbeat, max_lock_lifetime_ms, per-thread drain isolation, debounce drain

1 participant