a2aproject / a2aproject/a2a-java

[Bug]: 0.3.3.Final - Agent terminates early when ADK tool returning RxJava Single hops threads

Aperta
#872 7 commenti 0 reazioni 0 assegnatari Vedi su GitHub
Lingua principale
Java
Stelle
490
Fork
172
Merge medio
1g 6h
PR unite (30g)
55

Descrizione

### What happened?

When an ADK tool returns an RxJava `Single` executes on a different thread to the executor, the agent endpoint returns right after the function call is made. I noticed this bug using Kotlin, where using `rxSingle` executes on the global dispatcher, rather than the executor, e.g:

```kt
@Schema(description = "Get the weather summary in a city")
fun getWeather(@Schema(name = "city", description = "City to query") city: String): Single> {
return rxSingle {
mapOf("result" to "It's 50 degrees and partly cloudy in $city")
}
}
```

When running an agent with this tool, e.g. with `message/send`, the result cuts out after the tool call.
```json
{
"jsonrpc": "2.0",
"id": "cli-check-2",
"result": {
"id": "565c05e1-8bfc-414c-950e-80d74e14350c",
"contextId": "cli-demo-context",
"status": {
"state": "working",
"timestamp": "2026-05-12T15:09:36.256177Z"
},
"artifacts": [
{
"artifactId": "80c03215-3225-4108-9f7f-040d5ad2ec82",
"parts": [
{
"data": {
"name": "getWeather",
"id": "adk-183771c0-58a0-4f1b-87e9-d406ffaef646",
"args": {
"city": "Edinburgh"
}
},
"metadata": {
"adk_type": "function_call"
},
"kind": "data"
}
],
"metadata": {}
}
],
"history": [
{
"role": "user",
"parts": [
{
"text": "What is the weather in Edinburgh?",
"kind": "text"
}
],
"messageId": "cli-check-2",
"contextId": "cli-demo-context",
"taskId": "565c05e1-8bfc-414c-950e-80d74e14350c",
"kind": "message"
}
],
"kind": "task"
}
}
```

Using `rxSingle(Dispatchers.Unconfined)` works only if the function never actually suspends, otherwise the issue persists. A Java equivalent would be to e.g. call `.subscribeOn(Schedulers.io())` on the `Single`, or any other change of Scheduler.

The issue appears to be in the `cleanupProducer` method of the `DefaultRequestHandler`, where the ChildQueue is closed prematurely. **It does appear to be fixed in 1.0.0 with AgentEmitter** - but I'm unsure if we'll be able to adopt it yet as the ADK doesn't support it (I've raised an issue for it [here](https://github.com/google/adk-java/issues/1194)). Is there any chance a fix could be backported to 0.3? I'm happy to provide help if needed.

### Code of Conduct

- [x] I agree to follow this project's Code of Conduct

Guida per i contributori

Apri la guida per i contributori

Valutazione

Questa issue non è ancora stata valutata.

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.