a2aproject / a2aproject/a2a-java

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

Abierto
#872 7 comentarios 0 reacciones 0 asignados Ver en GitHub
Lenguaje dominante
Java
Estrellas
490
Forks
172
Merge medio
1 d 6 h
PR fusionados (30 d)
55

Descripción

### 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

Guía de contribución

Abrir la guía de contribución

Evaluación

Este issue todavía no se ha evaluado.

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.