fix(langgraph): wire node stream on_next handlers to registered executors #10795
Merged
hurui200320
merged 2 commits from 2026-04-22 07:53:54 +00:00
bugfix/m3-node-stream-on-next-noop into master
Labels
Clear labels
auto/needs-reevaluation
controller-managed
overdue
auto/blocked-by-deps
auto/ci-timeout
auto/claimed-implementer
auto/claimed-merge
auto/claimed-reviewer
auto/driver-down
auto/invariant-violation
auto/last-attempt-tier-0
auto/last-attempt-tier-1
auto/last-attempt-tier-2
auto/last-attempt-tier-min
Automation Tracking
auto/needs-conflict-resolution
auto/needs-implementer
auto/postmortem
auto/ready-to-merge
auto/restart-throttled
auto/revert
auto/sentinel
auto/stale-inactivity
auto/unstable
Blocked
Needs Feedback
Signed-off: Owner
Signed-off: Scrum Master
Signed-off: Tech Lead
Spike
Controller deferred this PR; awaiting Phase 6+ scope-evaluator or operator re-enablement.
Auto-agents controller manages this PR/issue (see tools/controller/deploy/RUNBOOK.md). Remove this label to abandon controller management.
PR blocked by an open issue dependency. Operator must close the dep (or remove the dependency link) before the merge driver can act. Auto-cleared by merge_drive when no open deps remain.
Most recent merge cycle hit CI timeout. Driver excludes this PR while last merge_cycle row is < 30 min old; label persists thereafter as visible history.
Currently being processed by an implementer worker.
Currently being processed by the merge driver.
Currently being processed by a reviewer worker.
Merge driver heartbeat stale; pipeline halted. Closed automatically on next clean tick.
Detected master commit violating the strict merge invariant. Tracked as an issue (not a PR label); kept here for label completeness.
In-cycle escalation: most recent attempt ran at the Tier 0 slot (`tier-0`). Slot's model defined in .opencode/models/tiers.yaml.
In-cycle escalation: most recent attempt ran at the Tier 1 slot (`tier-1`). Slot's model defined in .opencode/models/tiers.yaml.
In-cycle escalation: most recent attempt ran at the Tier 2 slot (`tier-2`). Slot's model defined in .opencode/models/tiers.yaml. Gated behind IMPLEMENTER_ESCALATION_TIER2_ENABLED.
In-cycle escalation: most recent attempt ran at the Tier -1 slot (`tier-min`). Slot's model defined in .opencode/models/tiers.yaml. Suffix is ``-min`` (not ``--1``) so the Forgejo UI reads naturally.
Tracking issues used by the AI Automation system for agents to communicate and report.
Rebase conflict needs LLM conflict-resolver.
Failing CI needs implementer attention.
Documenting a driver incident or rollback.
Reviewer has APPROVED this PR and no later REQUEST_CHANGES is outstanding. The merge driver requires this label to even consider a PR for merging. Set by the reviewer worker on APPROVE; cleared on REQUEST_CHANGES.
Train repeatedly lost master-tempo races. Driver excludes via merge_cycle until cooldown elapses; label persists as visible history.
Revert PR backing out an invariant violation. Fast-tracked through the merge driver.
Sentinel PR duplicated from upstream into a personal fork by tools/duplicate_prs_to_fork.py for pipeline testing. Lives only in the fork; the canonical pipeline never sees it.
No implementer activity for N days. Flagged for human review. Auto-cleared on next push to head branch.
Repeatedly fails on current master (>= 3 ci-fail-on-rebased-sha releases in 12 h). Excluded from driver until human triage.
A ticket in a blocked state and unable to complete until some other task is completed first.
Bounty
$100
A bounty of $100 for any open-source contributor who provides a MR that solves this issue
Bounty
$1000
A bounty of $1000 for any open-source contributor who provides a MR that solves this issue
Bounty
$10000
A bounty of $10000 for any open-source contributor who provides a MR that solves this issue
Bounty
$20
A bounty of $20 for any open-source contributor who provides a MR that solves this issue
Bounty
$2000
A bounty of $2000 for any open-source contributor who provides a MR that solves this issue
Bounty
$250
A bounty of $250 for any open-source contributor who provides a MR that solves this issue
Bounty
$50
A bounty of $50 for any open-source contributor who provides a MR that solves this issue
Bounty
$500
A bounty of $500 for any open-source contributor who provides a MR that solves this issue
Bounty
$5000
A bounty of $5000 for any open-source contributor who provides a MR that solves this issue
Bounty
$750
A bounty of $750 for any open-source contributor who provides a MR that solves this issue
MoSCoW
Could have
Could have feature in order to satisfy the epic/legendary.
MoSCoW
Must have
Must have feature in order to satisfy the epic/legendary.
MoSCoW
Should have
Should have feature in order to satisfy the epic/legendary.
There are questions in the ticket that can not be completed until the project owner provides clarity.
Points
1
1 man-hours worth of work for an expert with no learning curve.
Points
13
13 man-hours worth of work for an expert with no learning curve.
Points
2
2 man-hours worth of work for an expert with no learning curve.
Points
21
21 man-hours worth of work for an expert with no learning curve.
Points
3
3 man-hours worth of work for an expert with no learning curve.
Points
34
34 man-hours worth of work for an expert with no learning curve.
Points
5
5 man-hours worth of work for an expert with no learning curve.
Points
55
55 man-hours worth of work for an expert with no learning curve.
Points
8
8 man-hours worth of work for an expert with no learning curve.
Points
88
88 man-hours worth of work for an expert with no learning curve.
Priority
Backlog
This ticket has backlogged priority and is not to be worked on yet
Priority
CI Blocker
Critical priority issue that blocks CI/CD pipeline and prevents PR merges
Priority
Critical
The priority is critical
Priority
High
The priority is high
Priority
Low
The priority is low
Priority
Medium
The priority is medium
When an epic or legendary is in review it must be signed off by owner, tech lead, and scrum master before being marked as completed.
When an epic or legendary is in review it must be signed off by owner, tech lead, and scrum master before being marked as completed.
When an epic or legendary is in review it must be signed off by owner, tech lead, and scrum master before being marked as completed.
A ticket for learning a tool or technology that is needed to be able to do future planning and design.
State
Completed
The ticket has been fully implemented, completed, and merged with the source code. This label should only be applied once a ticket is closed.
State
Duplicate
A ticket that represents the same content as an existing ticket.
State
In Progress
A ticket that is actively being developed.
State
In Review
A ticket that has had some code completed to implement but is waiting to pass peer review and is not yet merged in.
State
Paused
This ticket's work started but wasn't finished. It's on hold (likely in a feature branch) and will be resumed later, either due to a blocker or a delay.
State
Unverified
All new tickets start in this state. A developer may set it to show the ticket is unverified. This means we haven't agreed to work on it. It will either move to a verified state or be closed as wontdo.
State
Verified
The issue has been verified by a developer as legitimate. It will be worked on and verified tickets are now considered part of the backlog.
State
Wont Do
This ticket has been decided it wont be done. This may mean the bug has been determined to not be real (cant verify) or the feature is one we have decided we dont want to adopt.
Type
Automation
Any edits or discussion about the AI automated coding system.
Type
Bug
Something that doesnt work as intended.
Type
Discussion
Anytime a ticket represents a discussion about a subject and doesnt fall into one of the other categories.
Type
Documentation
An error or improvement needed in the documentation.
Type
Epic
Any first tier epic. That is, an epic which contains only issues as children and will not have sub-epics.
Type
Feature
Some new functionality not present.
Type
Legendary
A type of Epic which will contain other Epics.
Type
Refactor
A code change that restructures existing code without changing its external behavior.
Type
Support
Someone needs help using the project.
Type
Task
A generic task that doesnt fit into the other type categories.
Type
Testing
Work exclusively focusing on fixing or expanding testing.
No Label
Type
Bug
Projects
Clear projects
No project
Assignees
aditya (Aditya Chhabra)
aleenaumair (Aleena Umair)
brent.edwards (Brent Edwards)
CoreRasurae (Luis Mendes)
drew (Drew Morris)
eugen.thaci (Eugen Thaci)
freemo (Jeffrey Phillips Freeman)
HAL9000 (HAL 9000)
HAL9001 (HAL9001)
hamza.khyari (Hamza Khyari)
hurui200320 (Rui Hu)
justin.morris
khird (Kyle Hird)
org.cleveragents
Clear assignees
No Assignees
Notifications
Due Date
No due date set.
Blocks
Reference: cleveragents/cleveragents-core#10795
Reference in New Issue
Block a user
Blocking a user prevents them from interacting with repositories, such as opening or commenting on pull requests or issues. Learn more about blocking a user.
Delete Branch "bugfix/m3-node-stream-on-next-noop"
Deleting a branch is permanent. Although the deleted branch may continue to exist for a short time before it actually gets removed, it CANNOT be undone in most cases. Continue?
Summary
Fixes the disconnected node execution pipeline in
LangGraph._setup_node_stream_subscriptionswhere allon_nexthandlers were no-op (pass), silently dropping every message delivered to node streams. Additionally addresses critical deadlock, thread-safety, resource management, and error-handling issues in the now-activated executor path.Closes #6511
Changes
Problem
_setup_node_stream_subscriptionssubscribed observers with emptyon_nexthandlerssync_executorper node (set via_register_node_executor) was only reachable through an implicitmapoperator in the observable pipeline — not the subscription handlernode_namevariable_register_node_executor(per-messageThreadPoolExecutor, per-messageasyncio.run(), non-thread-safeStateManager.update_state) were dormant but became reachable onceon_nextwas wiredFix
on_nextto the executor via_make_on_next_handler(name)— a proper method onLangGraphthat creates a closure looking up the executor fromself._node_executorsat invocation timeexecutor(msg)is wrapped intry/exceptso failures are logged vialogger.exceptionrather than propagating unhandled through the RxPy Subject call chain.TimeoutErroris distinguished from other exceptions with awarning-level log for clearer diagnostics. Shutdown-relatedRuntimeError(containing "stopping" or "not running") is also logged atwarninglevel to reduce noise during normal graph shutdownon_nextclosure captures its ownnameparameter via the factory methodself._node_executors: dict[str, Callable]replacessetattron the stream router, avoiding namespace pollutionThreadPoolExecutor—self._executor_poolis created once in__init__and shut down instop(), replacing per-message pool creationsync_executorusesasyncio.run_coroutine_threadsafewhen the scheduler's loop is running and the caller is on a different thread; falls back to thread pool when on the same thread to avoid deadlockStateManager.update_stateacquires athreading.RLock, captures amodel_copy(deep=True)under the lock, then emits and returns the copy outside the lock. This prevents re-entrant deadlock from subscriber callbacks while ensuring subscribers receive immutable point-in-time snapshots. The same deep-copy-under-lock pattern is applied to all five mutation methods:update_state,load_checkpoint,time_travel,reset, andreplace_stateStateManager.replace_state()deep-copies the input inside the lock before storing it, preventing external holders of the original reference from mutating internal state without acquiring the lock. Both the stored copy and the emission copy are created independently from the original input to avoid redundant double-copying.LangGraph.execute()uses this instead of directly assigning tostate_manager.state, closing a race condition whereexecute()bypassed the locking disciplineStateManager.reset()deep-copies the providedinitial_statebefore storing it, applying the same defensive copy pattern asreplace_state()execute()passes a deep copy of state tosend_message()so stream subscribers cannot hold a reference to the internal state objectsync_executorreads timeout fromnode.config.timeout(whereNodeConfigdefines it) instead ofgetattr(node, "timeout", None)which always returnedNoneexecution_historyusescollections.deque(maxlen=1000)for atomic bounded growth, eliminating TOCTOU race conditionsexecute()now raisesValueError(matchingstart()) instead of silently returning unmodified statefuture.result()calls use a configurable timeout (default 300s) with descriptiveTimeoutErroron expiryTimeoutError,future.cancel()is called before re-raising. A warning is logged whencancel()returnsFalse(task already running), indicating thread pool pressure. Inline comments document that cancellation is best-effortsync_executorcatchesCancelledErrorexplicitly in both therun_coroutine_threadsafepath and the thread pool path to produce a clear "graph stopping" message instead of a misleading generic errorstop()andstart()useshutdown(wait=True, cancel_futures=True)to properly drain in-flight tasks and prevent old-cycle tasks from corrupting new-cycle stateLangGraph.__del__shuts down_executor_poolifstop()was never called, with broadExceptionsuppression for interpreter shutdown resilience;start()shuts down the previous pool before creating a new one to prevent thread leaks on repeated callsstart()re-creates_executor_poolso the graph can be restarted afterstop()sync_executorandexecute()raise a descriptiveRuntimeErrorwhen the graph is not running, preventing misleading errors from submitting to a shut-down poolstart()setsis_running = Falsebefore shutting down the old pool and restores it after creating the new pool, preventing new submissions during the pool swap windowsync_executorwraps_executor_pool.submit()intry/except RuntimeErrorto detect the shutdown condition whenstop()is called between theis_runningguard and thesubmit()call, re-raising with the same descriptive message used by theis_runningguardclear_history()now acquires_lockfor consistency with other mutation methods_save_checkpointfallback path acquires the lock when called without pre-built data to ensure consistent snapshotsStateManager._lockupgraded fromthreading.Locktothreading.RLock, eliminating the deadlock risk when_save_checkpoint(None)is called while the lock is already held by the same threados.getpid()in addition toupdate_countand timestamp, preventing silent data loss when multiple processes share the samecheckpoint_dir_history_lock(declared but never used) removed_make_on_next_handlerand_make_on_error_handlerare proper methods with consistent type aliases (_OnNextCallable,_OnErrorCallable)get_state()— Returnsmodel_copy(deep=True)under the lock to prevent observing partially-mutated stateMAX_EXECUTION_HISTORYrenamed from_MAX_EXECUTION_HISTORYsince it is part of the testable contractrun_asyncannotated as-> StreamMessageinstead of-> Any;replace_state()docstring reworded to describe general purpose_make_on_next_handlerreferences follow-up ticket #10799 for routing executor return values to successor node streamsupdate_state()andreplace_state()docstrings now explicitly warn that emission order is not guaranteed to match mutation order under concurrent accessThreadPoolExecutor()documents the defaultmin(32, cpu_count+4)sizingcast(dict[str, Any], input_data)removed fromexecute()since Pyright already narrows the type after theisinstancecheckstep_execute_graphnow usestry/finallymatching the pattern instep_execute_expecting_errorcontext.add_cleanup(context.graph.stop)registered in the@whenstep so the thread pool is cleaned up even if assertions fail_cleanup_bg_loopwrapscall_soon_threadsafeintry/except RuntimeErrorto ensure thread join and loop close always execute even if the loop is already closedstart()docstring warns thatshutdown(wait=True)blocks indefinitely if a running task is stuck in blocking I/Odeque.append()is thread-safe under CPython GILa2a-sdk1.0.0 removed the legacyA2AClientclass froma2a.client, breaking the TDD test atfeatures/tdd_a2a_sdk_dependency.feature:21. Pinned to>=0.3.0,<1.0.0inpyproject.tomlto prevent the breaking version from being installed in CI. Migration to 1.0.0 is separate workFiles Changed
src/cleveragents/langgraph/graph.py— Core fix + deadlock prevention, timeout, future cancellation, CancelledError handling in both paths, resource management, deque, type annotations, is_running guard on execute(), dead code removal, public constant rename, best-effort cancellation documentation, TimeoutError vs Exception log level distinction, RuntimeError shutdown detection in on_next handler, TOCTOU-safe submit() with try/except, cancel() failure warning log, TODO(#10799) for downstream propagation, version-pinned comment for scheduler._loop, TOCTOU-safe start(), start() hang documentation, redundant cast removed, default pool sizing documented, deque thread-safety commentsrc/cleveragents/langgraph/state.py— Thread-safety RLock on all mutation methods, deep-copy emission inside lock for all five mutation methods, replace_state() creates both copies from original input (no redundant double-copy), reset() deep-copy input, file I/O outside lock with lock-safe fallback, clear_history locking, process-safe checkpoint filenames with PID, emission ordering documentationfeatures/consolidated_langgraph.feature— 20 new BDD scenarios (21 total including pre-existing)features/steps/langgraph_graph_coverage_steps.py— Step definitions with strengthened assertions, cleanup safety, timeout tests, is_running guard test, replace_state direct test with internal state identity assertion, CancelledError test for both paths, del test, execute() not-running guard test, TimeoutError-through-stream test, exception-safe event loop cleanup, restart cleanup registration, background loop cleanup guard, barrier timeout and thread completion assertion, correct parameter names in slow_execute functions, plain assignment for Behave context attributespyproject.toml— Pinneda2a-sdk>=0.3.0,<1.0.0to prevent breakingA2AClientremoval in 1.0.0Testing
run_coroutine_threadsafepath with mock verification that the correct code path was takenQuality Gates
nox -e lintnox -e typechecknox -e unit_testsnox -e integration_testsnox -e e2e_testsnox -e coverage_reportReview Issues Addressed (Cycle 8)
Major (M1 — CI Blocker, Fixed)
a2a-sdk1.0.0 released upstream, removing the legacyA2AClientclass froma2a.client. The existing dependency constraint>=0.3.0allowed CI to install 1.0.0, breaking the TDD testfeatures/tdd_a2a_sdk_dependency.feature:21("a2a SDK provides the A2AClient class"). Pinned toa2a-sdk>=0.3.0,<1.0.0inpyproject.toml. Migration to 1.0.0 is out of scope for this ticketReview Issues Addressed (Cycle 7)
Major (M1 — CI Blocker, Fixed)
features/steps/langgraph_graph_coverage_steps.py: Two@given/@thendecorator strings were wrapped across multiple lines, causingruff format --checkto reject them. Collapsed both to single lines (E501 is suppressed forfeatures/steps/*.py)Minor (All Fixed)
sync_executorbetweenis_runningguard andsubmit(): Wrapped_executor_pool.submit()intry/except RuntimeErrorthat detects the shutdown condition and re-raises with the same descriptive message used by theis_runningguardCancelledErrorduring shutdown logged atexceptionlevel: Added dedicatedexcept RuntimeErrorclause in_make_on_next_handlerthat checks for "stopping"/"not running" patterns and logs atwarninglevel instead ofexceptionleveltimeout=10.0tobarrier.wait()and addedassert all(not t.is_alive() for t in threads)after the join loopreplace_state(): Changed emission copy to be created from the originalnew_stateinput rather than from the already-copied stored state, eliminating redundant serializationos.getpid()to checkpoint filename format:checkpoint_{timestamp}_{update_count}_{pid}.jsonslow_executetest functions: Renamed_statetostatein both functionscontext.stream_emissions: list[GraphState] = []to plain assignment with comment:context.stream_emissions = [] # list[GraphState]Nits (All Fixed)
_save_checkpoint(None)deadlock risk: UpgradedStateManager._lockfromthreading.Locktothreading.RLock(reentrant lock), making the deadlock impossible. Updated docstring to reflect the changedeque.append()thread safety undocumented: Added inline comment:# Thread-safe under CPython GIL; deque.append is atomic.future_tp.cancel()returnsFalse, indicating the timed-out task could not be cancelled and the thread slot remains occupiedreplace_statetest accesses internalstateattribute without comment: Added explanatory comment:# Direct access intentional: identity check requires the actual internal object, not a get_state() copy.Previous Review Cycles
Cycle 6 (All Fixed)
CancelledErrornot caught inrun_coroutine_threadsafepath: Addedexcept concurrent.futures.CancelledErrorcatch mirroring the thread pool pathTimeoutErrorwarning-level log through stream path: Added BDD scenariostart():start()now setsis_running = Falsebefore shutting down the old pooltry/finallycontext.add_cleanup()start()can hang indefinitely: Added docstring warningCycle 5 (All Fixed)
replace_state()stores caller's mutable reference directly: Fixed by deep-copying inputCancelledErrorpath: Added BDD scenarioCycles 1–4
See git history for earlier review cycle fixes (all addressed in prior amendments).
Deferred (Pre-existing, Out of Scope)
close()does not acquire_lockbefore settingis_closed = True(pre-existing)Node.execution_countincrement is not thread-safe (pre-existing, newly reachable)asyncio.run()creates a new event loop per thread-pool invocation (acceptable for current use case)a2a-sdk0.3.x to 1.0.0 (separate ticket scope — 1.0.0 removes the legacyA2AClientclass and introducesClientas the replacement)41b5e8633ftoebd47fc5f8ebd47fc5f8toed93e19f45ed93e19f45tofb2680266efb2680266eto08b7ba627b08b7ba627bto68a2e06f3268a2e06f32to00a3630fc000a3630fc0to7455ca6a6d7455ca6a6dtoc96ba34f14c96ba34f14to92dfd2f6f392dfd2f6f3to0d267934a7@HAL9000 This PR is ready to review: rebased onto the latest master, all CI checks passed.
hurui200320 referenced this pull request2026-04-22 06:01:13 +00:00
Review Summary
All required checks have passed and this PR successfully addresses bug #6511 by wiring the on_next handlers, improving thread-safety, resource management, and error handling in LangGraph. Tests have been added and updated, coverage remains ≥97%, and no prohibited patterns (e.g., # type: ignore) were introduced. The code adheres to project specifications and conventions. No blocking issues found. Approving.
Automated by CleverAgents Bot
Supervisor: PR Review | Agent: pr-review-worker