Skip to content

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

Description

@patrick-chinchill

Summary

Port upstream's Chat concurrency rework: lock heartbeat capped by max_lock_lifetime_ms; drains dispatch each message under its own thread id (lock_scope="channel"); debounce keeps skipped context and drains mid-handler arrivals; skipped mentions route to on_mention. Today a handler running >30s loses its lock and a second message can run concurrently on the thread.

Upstream changes

  • eccc6b91 fix(chat): detect mentions in skipped queued messages (#656) — chat@4.32.0 — setMentionFlags() runs detection on the dispatched and every context.skipped message; any mention routes to mention handlers.
  • 076fe5dc fix(chat): preserve skipped mention routing (#659) — chat@4.33.0 — debounce capacity maxQueueSize (was 1); debounce drains all pending, accumulates superseded into skipped, dispatches with MessageContext; routing message.isMention || (hasMention && mentionHandlers.length > 0).
  • 6cb933eb fix(chat): isolate channel-scoped queue dispatch by thread (#832) — chat@4.38.1 — security-relevant (cross-thread dispatch): drains dispatch with latest.message.threadId, skipped filtered to it, totalSinceLastHandler = skipped.length + 1.
  • 5b538f6f fix(chat): keep thread locks alive during long handlers (#821) — chat@4.39.0 — withHeldLock()/startLockHeartbeat() extend every DEFAULT_LOCK_TTL_MS / 3 without stacking, stop after maxLockLifetimeMs (600000); ownership lost on a false extend or backend unreachable past heldUntil; stop() awaits the in-flight extend before releaseLock; drains check isOwnershipLost(); debounce loops after dispatch; adds ConcurrencyConfig.maxLockLifetimeMs + extendLock token-compare contract.

Current Python behavior

  • src/chat_sdk/chat.py:82 DEFAULT_LOCK_TTL_MS = 30_000; grep -rn 'heartbeat\|max_lock_lifetime' src/chat_sdk is empty. chat.py:1996-2019 _handle_drop never extends the lock.
  • chat.py:2064-2163 _handle_queue_or_debounce: :2081 effective_max = 1 if strategy == "debounce" else max_queue_size; :2118 self-enqueue with size 1; burst (:2149-2157) sleeps then does a one-off extend_lock.
  • chat.py:2167-2212 _debounce_loop: Python-only max_iterations = 20 (:2175); dequeues one entry, dispatches with no MessageContext and the holder's thread_id (:2211), then breaks.
  • chat.py:2216-2269 _drain_queue: dispatches with the holder's thread_id (:2263); skipped not thread-filtered (:2247); inline "Lock lost ... aborting" extend_lock checks (:2239-2244, :2266-2269).
  • chat.py:2316 sets is_mention only on the dispatched message; :2372 has no skipped-mention path.
  • types.py:221-231 ConcurrencyConfig lacks max_lock_lifetime_ms; types.py:1205 extend_lock has no contract docstring. All backends already token-compare: state/memory.py:137, state/redis.py:45 Lua script, state/postgres.py:274 (AND token = $4 AND expires_at > now()).
  • docs/UPSTREAM_SYNC.md:257 claims "debounce loop iteration limits (20) ... must match"; upstream has no cap at 4.31.0 or 4.41.1.

Scope

  • types.py: ConcurrencyConfig.max_lock_lifetime_ms: int = 600_000 (+ max_queue_size comment covers debounce); token-compare contract in the StateAdapter.extend_lock docstring. chat.py: DEFAULT_MAX_LOCK_LIFETIME_MS, resolved in all three concurrency branches (chat.py:310-332).
  • Private _LockHeartbeat (is_ownership_lost(), async stop()) and _with_held_lock(lock, thread_id, lock_key, fn): start heartbeat, run fn(heartbeat), finally await stop() then release_lock. Route _handle_drop (incl. the on_lock_conflict force path) and the lock-holding branch of _handle_queue_or_debounce through it.
  • Heartbeat: tick every DEFAULT_LOCK_TTL_MS / 3; stop renewing (warn) once elapsed ≥ max_lock_lifetime_ms; False extend → lost, exit; extend exception → lost if wall-clock ≥ held_until, else warn; held_until = lock.expires_at, then now + TTL per successful extend; is_ownership_lost() = lost or now_ms >= held_until.
  • _drain_queue(heartbeat, adapter, lock_key) / _debounce_loop(heartbeat, adapter, lock_key): drop thread_id; check is_ownership_lost() per drain iteration / after each debounce sleep; remove inline extend_lock/"Lock lost" branches; dispatch with latest.thread_id, skipped filtered to it, total_since_last_handler=len(skipped)+1.
  • Debounce: capacity max_queue_size for busy-path and self-enqueue; drop-newest early return still excludes debounce (as upstream). Drain all pending per tick, accumulate superseded into skipped, dispatch with MessageContext, reset, loop again instead of break.
  • Remove max_iterations: max_lock_lifetime_ms now bounds the loop (renewal stops → lock lapses one TTL later → is_ownership_lost() exits). Correct UPSTREAM_SYNC.md:257.
  • Burst: drop the post-sleep extend_lock; call _drain_queue(heartbeat, ...).
  • _set_mention_flags(adapter, message, context) -> bool: fill is_mention on the message and each context.skipped in place, return has_mention. Keep 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 to is None. Route with if message.is_mention or (has_mention and self._mention_handlers):.

Out of scope

Porting notes

  • Cancellation (recommended design): each extend is its own task awaited via asyncio.shield; stop() sets stopped, cancels the loop task, then awaits the in-flight extend (gather(..., return_exceptions=True)) before release_lock. Never swallow a CancelledError aimed at the heartbeat; suppress post-stop warns (upstream if (!stopped)). Awaiting each extend before sleeping again equals upstream's no-stacking if (inFlight) return.
  • Clocks: Lock.expires_at is epoch ms (types.py:1186-1191) → held_until is wall-clock ms; the lifetime cap may use time.monotonic() (divergence: upstream uses Date.now()). Put both behind helpers next to _sleep (chat.py:205) for a fake clock.
  • Ownership lost: drains return early, leaving the queue to the new holder; release_lock stays a no-op for a foreign token.
  • Rehydration: dequeued thread_id must survive JSON (Message.from_json reads threadId/thread_id, types.py:704; dict fallback chat.py:2518); cover a dict queue entry in the isolation test.

Tests

From packages/chat/src/chat.test.ts (fidelity-mapped → tests/test_chat_faithful.py):

Python-specific:

  • Token-checked, expiry-aware fake lock on a fake clock (port of installTokenLockMock; MockStateAdapter.extend_lock always returns True, shared/mock_adapter.py:273): peak_in_flight == 1; heartbeat task done() after release.
  • stop() during an extend blocked on an asyncio.Event waits for it; release_lock strictly after (AsyncMock call order).
  • Extend raising before held_until keeps ownership; after it, drain stops and release_lock still runs.
  • Debounce message arriving mid-handler is dispatched after the handler returns, no new webhook. No real sleeps beyond the debounce window.

Acceptance criteria

  • Full validation command from CLAUDE.md passes.
  • chat.test.ts fidelity misses against chat@4.41.1 (non-strict, per [4.41/P0] Fidelity tooling for the 4.41 wave: single pin constant, SHA pin, it.each expansion, map new core test files #185) drop by the tests above.
  • docs/UPSTREAM_SYNC.md: :257 corrected, max_iterations removal recorded, row for the asyncio heartbeat design (monotonic lifetime clock).
  • CHANGELOG under "Unreleased (4.41 wave)", consumer-visible: under drop (default) a second message ≥30s into a long handler raises LockError instead of running concurrently; debounce handlers get a context with superseded skipped; channel-scoped queued messages answer in their own thread; skipped mentions route to on_mention.
  • SELF_REVIEW.md checks applied (rebind/state coherence on the heartbeat, no double release).

Dependencies

Blocked by #185. Blocks #191, #195, #203.

Metadata

  • Effort: L (~250 LOC source, ~450 LOC tests)
  • Consumer impact: high — Slack/Teams streaming users on default drop will see a second message during a long stream rejected (LockError, logged) instead of run concurrently; downstream consumers (e.g. chinchill) should confirm their strategy. Security-relevant: 6cb933eb.
  • Suggested branch: sync/4.41-c1a

Part of #184.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions