Skip to content

fix(flows): propagate downstream cancellation to the model stream - #1621

Open
sangkyoonnam wants to merge 1 commit into
google:mainfrom
sangkyoonnam:fix/1618-cancel-propagation
Open

sangkyoonnam wants to merge 1 commit into
google:mainfrom
sangkyoonnam:fix/1618-cancel-propagation

Conversation

@sangkyoonnam

@sangkyoonnam sangkyoonnam commented Oct 6, 2026 •

Copy link
Copy Markdown

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 Flowable from Runner.runAsync doesn't stop the model call. BaseLlmFlow.run shares each step as runOneStep(...).cache(), and FlowableCache never 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 a takeUntil in front of it: step.takeUntil(cancelled.toFlowable()).cache().doOnCancel(cancelled::onComplete), one CompletableSubject per step. A cancel completes the signal, takeUntil cancels runOneStep and completes the cached stream. What run() relies on from cache() is unchanged: unbounded upstream demand, terminal delivered as soon as a subscriber drains the buffer, buffer kept for the toList() 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 its Single disposed; whether its work stops is up to the tool. A model error arriving after the cancel goes to RxJavaPlugins.onError instead 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; Gemini builds its stream with fromFuture(...).flatMapIterable(...) and nothing there cancels the future, so this PR doesn't change that.

The ParallelAgent and LoopAgent escalation exits already cancel with takeUntil; those cancels now reach an active sibling's model stream too, with a regression test. #1434's cooperative CancellationToken cuts the stream inside callLlm and doesn't touch cache(), so the two compose.

Testing Plan

Unit Tests:

  • I have added or updated unit tests for my change.
  • All unit tests pass locally.

BaseLlmFlowCancelTest (new, 9 tests): cancel reaches the model through BaseLlmFlow.run and through InMemoryRunner.runAsync, a dispose before a function call leaves the tool unrun, a cancel inside onSubscribe doesn't leave the model streaming, a ParallelAgent sibling 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 the BaseLlmFlow change 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_sse disposes on onCompletion/onTimeout and the A2A executor's cancel() disposes its CompositeDisposable; with this change both reach the model. I read those paths, I haven't exercised them.

Checklist

  • I have read the CONTRIBUTING.md document.
  • My pull request contains a single commit.
  • I have performed a self-review of my own code.
  • I have commented my code, particularly in hard-to-understand areas.
  • I have added tests that prove my fix is effective or that my feature works.
  • New and existing unit tests pass locally with my changes.
  • I have manually tested my changes end-to-end.
  • Any dependent changes have been merged and published in downstream modules.

Additional context

runLive doesn't use cache() and isn't touched. Nothing in core subscribes to one run() 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.

@hemasekhar-p

Copy link
Copy Markdown
Contributor

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?

@hemasekhar-p hemasekhar-p added the waiting on reporter Waiting for reaction by reporter. Failing that, maintainers will eventually closed it as stale. label Oct 6, 2026
@sangkyoonnam
sangkyoonnam force-pushed the fix/1618-cancel-propagation branch from 26c291a to a63b9cd Compare October 6, 2026 14:16
@sangkyoonnam

Copy link
Copy Markdown
Author

Done. The branch is one commit again, a63b9cd, rebased onto main at 9debe96. The second commit was a merge from main I clicked through on GitHub right after opening the PR; the diff is unchanged.

@sangkyoonnam

Copy link
Copy Markdown
Author

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

@hemasekhar-p

Copy link
Copy Markdown
Contributor

@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?

@hemasekhar-p hemasekhar-p added waiting on reporter Waiting for reaction by reporter. Failing that, maintainers will eventually closed it as stale. and removed waiting on reporter Waiting for reaction by reporter. Failing that, maintainers will eventually closed it as stale. labels Oct 7, 2026
@sangkyoonnam
sangkyoonnam force-pushed the fix/1618-cancel-propagation branch from 0bb1da0 to a63b9cd Compare October 7, 2026 10:19
@sangkyoonnam

Copy link
Copy Markdown
Author

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.

@hemasekhar-p
hemasekhar-p force-pushed the fix/1618-cancel-propagation branch from a63b9cd to d217c29 Compare October 7, 2026 11:27
@hemasekhar-p

Copy link
Copy Markdown
Contributor

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

@hemasekhar-p hemasekhar-p added needs review and removed waiting on reporter Waiting for reaction by reporter. Failing that, maintainers will eventually closed it as stale. labels Oct 7, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants