Repository navigation
fix(flows): propagate downstream cancellation to the model stream - #1621
sangkyoonnam wants to merge 1 commit into
Conversation
|
Hi @sangkyoonnam, thank you for your contribution. We appreciate you taking the time to submit this pull request. As per contribution policy, PR should consists of a single commit. Could you please change your commits accordingly? |
26c291a to
a63b9cd
Compare
|
The reporter validated this head (a63b9cd) in production: 311 regression tests green, a mid-stream client disconnect left no further model or tool activity, and their cooperative cancel flag still works alongside. Details in #1618 (comment). |
|
@sangkyoonnam,thank you for the update and I noticed there is nothing in the merge commit. To proceed with the review, could you please squash them into a single commit? |
0bb1da0 to
a63b9cd
Compare
|
Squashed again: the branch is back to the single commit a63b9cd. The extra commit was an empty merge from main made with the Update branch button; I won't use it on this PR again. |
a63b9cd to
d217c29
Compare
|
@sangkyoonnam, thanks for the quick update. Currently this PR is under review by our team, we will keep you posted if any additional information is required. thank you. |
Link to Issue or Description of Change
1. Link to an existing issue (if applicable):
This covers the propagation half of #1618; a public cancel API (#1341) stays separate. #1565 (approved) moves the changed line into
runStep; I'll rebase onto it once it lands, the swap is the same there.Problem:
Disposing the
FlowablefromRunner.runAsyncdoesn't stop the model call.BaseLlmFlow.runshares each step asrunOneStep(...).cache(), andFlowableCachenever cancels its upstream:CacheSubscription.cancel()only drops the subscriber from the list. Once the consumer disposes, the step runs on with zero subscribers, including the after-model callbacks and the step's tool calls. I measured the last part: a function call the model emits after the dispose still runs the tool.Solution:
Keep
cache()and put atakeUntilin front of it:step.takeUntil(cancelled.toFlowable()).cache().doOnCancel(cancelled::onComplete), oneCompletableSubjectper step. A cancel completes the signal,takeUntilcancelsrunOneStepand completes the cached stream. Whatrun()relies on fromcache()is unchanged: unbounded upstream demand, terminal delivered as soon as a subscriber drains the buffer, buffer kept for thetoList()subscriber that arrives after the step finished.After the change a disconnect cancels the step's pending reactive work: the in-flight model
Flowable, after-model callbacks not yet run, tool calls not yet started. A tool already executing has itsSingledisposed; whether its work stops is up to the tool. A model error arriving after the cancel goes toRxJavaPlugins.onErrorinstead of sitting in the cache for nobody, which is RxJava's normal post-cancel routing. Whether the HTTP request is aborted is up to the model client;Geminibuilds its stream withfromFuture(...).flatMapIterable(...)and nothing there cancels the future, so this PR doesn't change that.The
ParallelAgentandLoopAgentescalation exits already cancel withtakeUntil; those cancels now reach an active sibling's model stream too, with a regression test. #1434's cooperativeCancellationTokencuts the stream insidecallLlmand doesn't touchcache(), so the two compose.Testing Plan
Unit Tests:
BaseLlmFlowCancelTest(new, 9 tests): cancel reaches the model throughBaseLlmFlow.runand throughInMemoryRunner.runAsync, a dispose before a function call leaves the tool unrun, a cancel insideonSubscribedoesn't leave the model streaming, aParallelAgentsibling escalation cancels the model;cache()semantics kept under bounded demand (run(ctx).test(1)gets the final event and completion, and a model error after a buffered event) and for a second subscriber after a cancel; a control run still runs the tool. Without theBaseLlmFlowchange 6 of the 9 fail../mvnw -pl core test: 2079 tests, 0 failures, 24 skipped.Manual End-to-End (E2E) Tests:
Not run. The dev server's
/run_ssedisposes ononCompletion/onTimeoutand the A2A executor'scancel()disposes itsCompositeDisposable; with this change both reach the model. I read those paths, I haven't exercised them.Checklist
Additional context
runLivedoesn't usecache()and isn't touched. Nothing in core subscribes to onerun()result twice at the same time, so "a subscriber cancels" and "the last subscriber cancels" are the same event here.I wrote this with Claude Code and reviewed and ran it myself.