Skip to content

Commit e2e4ccc

Browse files
committed
feat: add asyncio support to the gRPC app extension (dapr.ext.grpc.aio)
dapr.ext.grpc.App runs on grpc.server() with a thread pool, so handlers must be synchronous. Applications built on asyncio have had to vendor their own callback server to use `async def` handlers. Add dapr.ext.grpc.aio, backed by grpc.aio.server(). It exports the same names as dapr.ext.grpc; swapping the import is the only change an app needs, beyond `async def` handlers and awaiting run()/stop(). _AioCallbackServicer is a sibling of _CallbackServicer over the base extracted in the previous commit - neither subclasses the other, so there is no sync/async mixing in one MRO and no misleading isinstance relationship. It supplies only the gRPC entry points, which await the handler result and the grpc.aio context coroutines. The asyncio servicer covers the full current surface: service invocation, pub/sub including bulk events, input bindings, job events on both the stable and alpha services, and health checks. AppCallbackAlphaServicer is registered on the server, and send_initial_metadata is awaited, since it is a coroutine on the aio context rather than a plain method. Handlers may be plain functions - results are awaited only when awaitable - so register_health_check(lambda: None) keeps working. The server is built on the first run()/start() call rather than in __init__, because grpc.aio.server() binds to whichever event loop is current when it is called; add_external_service() therefore queues its registration and replays it at creation. Adds unit tests including an end-to-end suite over a real grpc.aio server and parity guards asserting every sync RPC is mirrored as a coroutine, plus the invoke-simple-async and pubsub-simple-async examples. The async client (dapr.aio.clients.DaprClient) already exists; this closes the remaining server-side gap. Closes #695 Signed-off-by: Sai Kishore Punagani <63619246+saikishore-p@users.noreply.github.com>
1 parent d1f0f96 commit e2e4ccc

24 files changed

Lines changed: 2620 additions & 9 deletions

‎AGENTS.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ Extensions are bundled into the core `dapr` wheel and exposed as installable ext
5757
| Extra | Import path | Purpose | Active development |
5858
|-------|-------------|---------|--------------------|
5959
| `dapr[workflow]` | `dapr.ext.workflow` | Durable workflow orchestration (durabletask vendored internally) | **High**, major focus area |
60-
| `dapr[grpc]` | `dapr.ext.grpc` | gRPC server for Dapr callbacks (methods, pub/sub, bindings, jobs) | Moderate |
60+
| `dapr[grpc]` | `dapr.ext.grpc`, `dapr.ext.grpc.aio` | gRPC server for Dapr callbacks (methods, pub/sub, bindings, jobs), sync and asyncio | Moderate |
6161
| `dapr[fastapi]` | `dapr.ext.fastapi` | FastAPI integration for pub/sub and actors | Moderate |
6262
| `dapr[flask]` | `dapr.ext.flask` | Flask integration for pub/sub and actors (legacy `flask_dapr` import path is a deprecated shim) | Low |
6363
| `dapr[langgraph]` | `dapr.ext.langgraph` | LangGraph checkpoint persistence to Dapr state store | Moderate |

‎dapr/ext/grpc/AGENTS.md‎

Lines changed: 61 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -7,16 +7,26 @@ The gRPC extension provides a **server-side callback framework** for Dapr applic
77
```
88
dapr/ext/grpc/
99
├── __init__.py # Public API exports
10-
├── app.py # App class — main entry point
11-
├── _servicer.py # _CallbackServicer — internal routing
12-
├── _health_servicer.py # _HealthCheckServicer
10+
├── app.py # App class — main entry point (sync)
11+
├── _servicer.py # _CallbackServicerBase (shared) + _CallbackServicer (sync)
12+
├── _health_servicer.py # _HealthCheckServicerBase (shared) + _HealthCheckServicer
13+
├── aio/ # asyncio-native variant, same public surface
14+
│ ├── __init__.py # Public API exports
15+
│ ├── app.py # App class backed by grpc.aio
16+
│ ├── _servicer.py # _AioCallbackServicer — async gRPC entry points
17+
│ └── _health_servicer.py # _AioHealthCheckServicer
1318
└── py.typed
1419
1520
tests/ext/grpc/
1621
├── test_app.py # Decorator registration tests
1722
├── test_servicer.py # Routing, handlers, bulk events
1823
├── test_health_servicer.py # Health check tests
19-
└── test_topic_event_response.py # Response status tests
24+
├── test_topic_event_response.py # Response status tests
25+
└── aio/
26+
├── test_app.py # Decorator registration, lazy server creation, lifecycle
27+
├── test_servicer.py # Async handlers, sync/aio parity guards
28+
├── test_health_servicer.py # Async and sync health check callbacks
29+
└── test_server.py # End-to-end over a real grpc.aio server
2030
```
2131

2232
Installed via the `grpc` extra on core dapr: `pip install "dapr[grpc]"`.
@@ -42,6 +52,8 @@ from dapr.ext.grpc import (
4252

4353
Note: `InvokeMethodRequest`, `InvokeMethodResponse`, `BindingRequest`, `TopicEventResponse`, `Job`, `JobEvent`, and failure policies are actually defined in the core SDK (`dapr/clients/grpc/`) and re-exported here.
4454

55+
`dapr.ext.grpc.aio` exports the same names. Swapping the import is the only change an app needs, beyond `async def` handlers and awaiting `run()`/`stop()`.
56+
4557
## App class (`app.py`)
4658

4759
The central entry point. Creates a gRPC server and provides decorators for handler registration.
@@ -88,9 +100,44 @@ app.register_health_check(lambda: None) # Not a decorator — direct registrati
88100
- `TopicEventResponse('success'|'retry'|'drop')` → explicit status
89101
- `None` → defaults to SUCCESS
90102

103+
## Asyncio app (`aio/`)
104+
105+
`dapr.ext.grpc.aio.App` is the asyncio counterpart, backed by `grpc.aio.server()`. Same decorators, same registration semantics, same wire behavior.
106+
107+
```python
108+
import asyncio
109+
from dapr.ext.grpc.aio import App
110+
111+
app = App()
112+
113+
@app.subscribe(pubsub_name='pubsub', topic='orders')
114+
async def handle_event(event: SubscriptionMessage) -> TopicEventResponse:
115+
...
116+
117+
asyncio.run(app.run(3010))
118+
```
119+
120+
Differences from the synchronous `App`, all of them forced by the async runtime:
121+
122+
- **`run()` and `stop()` are coroutines.** `stop(grace=None)` takes a grace period (the sync `stop()` is always immediate). There is no `__del__` hook, because a coroutine cannot be awaited from one — `run()` instead stops the server in a `finally`, so a cancelled app does not leave its port bound.
123+
- **`start()` exists** as a non-blocking alternative to `run()`, for serving the app alongside other work on the same loop (e.g. from an ASGI lifespan handler). `run()` is `start()` plus `wait_for_termination()`.
124+
- **The server is built lazily**, on the first `run()`/`start()` call, not in `__init__`. `grpc.aio.server()` binds to whichever event loop is current when it is called, so building it in `__init__` would attach it to the wrong loop. `add_external_service()` therefore queues its registration and replays it when the server is created, and raises if called once the app is running.
125+
- **`start()` and `stop()` are serialised** by an `asyncio.Lock`. grpc.aio segfaults if a `stop()` call is *concurrently in flight* with a `start()` call, so the two must never overlap; the lock also means a restart waits for an in-progress drain rather than binding a second server to the same port. (Stopping a server whose own `start()` has already unwound — the cleanup path in `_start`'s `except` — is sequential, not concurrent, and is safe. `stop()` on a server that never started returns cleanly; on one whose `start()` was *cancelled part-way* it raises `InvalidStateError`, which that path suppresses.)
126+
- **The lock is built per running loop**, not in `__init__`: an `asyncio.Lock` binds to the loop of its first *contended* acquire, so one built at construction raises "bound to a different event loop" on a second `asyncio.run()` — and silently stops excluding anything before that, because the uncontended path returns before the loop check.
127+
- **Decorators return the handler**, so the decorated name stays bound. The sync decorators return `None`.
128+
- **Handlers may be plain functions.** Results are awaited only when awaitable, so `register_health_check(lambda: None)` still works. A plain handler runs inline on the event loop and must not block.
129+
130+
### Sharing with the sync implementation
131+
132+
`_CallbackServicerBase` (in `_servicer.py`) holds everything that does not invoke a user handler: the handler registries, topic routing, and the request→event translation. `_CallbackServicer` and `_AioCallbackServicer` are **siblings** on top of it — neither subclasses the other — and each supplies only the gRPC entry points.
133+
134+
This keeps the churn-prone routing logic (`_get_topic_callback`, `register_topic`, the bulk entry builders) in one place while leaving the two servicers free to differ where they must. `tests/ext/grpc/aio/test_servicer.py::AsyncParityTests` enforces the arrangement: every RPC the sync servicer implements must be mirrored on the aio servicer as a coroutine function, and the registration helpers must stay shared rather than be reimplemented.
135+
136+
`_HealthCheckServicerBase` splits the health servicer the same way: registration in the base, the gRPC entry point in each sibling.
137+
91138
## Internal routing (`_servicer.py`)
92139

93-
`_CallbackServicer` implements `AppCallbackServicer` + `AppCallbackAlphaServicer` gRPC service interfaces. It maintains internal registries:
140+
`_CallbackServicerBase` implements `AppCallbackServicer` + `AppCallbackAlphaServicer` gRPC service interfaces; `_CallbackServicer` adds the synchronous entry points. It maintains internal registries:
94141

95142
- `_invoke_method_map` — method name → handler
96143
- `_topic_map` — topic key → handler
@@ -124,6 +171,13 @@ app.register_health_check(lambda: None) # Not a decorator — direct registrati
124171
uv run python -m unittest discover -v ./tests/ext/grpc
125172
```
126173

174+
`unittest discover` covers the whole tree including `aio/` — `IsolatedAsyncioTestCase` and
175+
`subTest` are both native unittest. pytest runs the same tests:
176+
177+
```bash
178+
uv run pytest ./tests/ext/grpc
179+
```
180+
127181
Test patterns:
128182
- `test_app.py` — decorator registration, health check registration
129183
- `test_servicer.py` — handler invocation with mock gRPC context, return type handling (str, bytes, proto, response object), topic subscriptions, bulk events, bindings, duplicate registration errors
@@ -132,8 +186,8 @@ Test patterns:
132186

133187
## Key details
134188

135-
- **Synchronous only**: Uses `grpc.server()` with `ThreadPoolExecutor(10)`. No async handler support.
189+
- **Sync app threading**: `dapr.ext.grpc.App` uses `grpc.server()` with `ThreadPoolExecutor(10)`. For `async def` handlers use `dapr.ext.grpc.aio.App`, which serves on a `grpc.aio` event loop instead.
136190
- **Default port**: 3010 (from `dapr.conf.global_settings.GRPC_APP_PORT`)
137-
- **Topic handler event type**: inferred from the handler annotation. Annotating the event parameter with `dapr.ext.grpc.SubscriptionMessage` — the same SDK-owned type the streaming subscription API (`DaprClient.subscribe`) delivers, with `metadata()` populated from the gRPC invocation metadata — delivers that type. Unannotated or otherwise-annotated handlers receive the DEPRECATED `cloudevents.sdk.event.v1.Event` and `subscribe()` emits a `DeprecationWarning` at registration. Deprecation timeline: 1.20 delivers `SubscriptionMessage` to unannotated handlers (legacy only via explicit `v1.Event` annotation), 1.21 drops `cloudevents` from the `grpc` extra (import becomes conditional), 1.22 removes the legacy path entirely (same release the `flask_dapr` shim goes away). New code must annotate with `SubscriptionMessage`. (Internally the choice is plumbed through `_CallbackServicer.register_topic(legacy_cloudevent=...)`.)
191+
- **Topic handler event type**: inferred from the handler annotation. Annotating the event parameter with `dapr.ext.grpc.SubscriptionMessage` — the same SDK-owned type the streaming subscription API (`DaprClient.subscribe`) delivers, with `metadata()` populated from the gRPC invocation metadata — delivers that type. Unannotated or otherwise-annotated handlers receive the DEPRECATED `cloudevents.sdk.event.v1.Event` and `subscribe()` emits a `DeprecationWarning` at registration. Deprecation timeline: 1.20 delivers `SubscriptionMessage` to unannotated handlers (legacy only via explicit `v1.Event` annotation), 1.21 drops `cloudevents` from the `grpc` extra (import becomes conditional), 1.22 removes the legacy path entirely (same release the `flask_dapr` shim goes away). New code must annotate with `SubscriptionMessage`. (Internally the choice is plumbed through `_CallbackServicerBase.register_topic(legacy_cloudevent=...)`.)
138192
- **Duplicate registration**: Registering the same method/topic/binding name twice raises `ValueError`
139193
- **Missing handlers**: Calling an unregistered method/topic/binding raises `NotImplementedError` (gRPC UNIMPLEMENTED)

‎dapr/ext/grpc/README.md‎

Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,5 +12,32 @@ pip install "dapr[grpc]"
1212
from dapr.ext.grpc import App
1313
```
1414

15+
An asyncio-native app backed by a `grpc.aio` server is available under
16+
`dapr.ext.grpc.aio`. It exposes the same names and decorators; handlers may be
17+
`async def` and are awaited, and `run()`/`stop()` are coroutines:
18+
19+
```python
20+
import asyncio
21+
22+
from dapr.ext.grpc.aio import App, InvokeMethodRequest, InvokeMethodResponse
23+
24+
app = App()
25+
26+
27+
@app.method(name='my-method')
28+
async def my_method(request: InvokeMethodRequest) -> InvokeMethodResponse:
29+
...
30+
31+
32+
asyncio.run(app.run(50051))
33+
```
34+
35+
Plain (non-async) handlers are still accepted, but they run inline on the event
36+
loop, so they must not block. See the [`invoke-simple-async`][invoke-async] and
37+
[`pubsub-simple-async`][pubsub-async] examples.
38+
1539
See the root [README](../../../README.md) for migration steps from the legacy
16-
`dapr-ext-grpc` distribution.
40+
`dapr-ext-grpc` distribution.
41+
42+
[invoke-async]: ../../../examples/invoke-simple-async
43+
[pubsub-async]: ../../../examples/pubsub-simple-async

‎dapr/ext/grpc/aio/__init__.py‎

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
# -*- coding: utf-8 -*-
2+
3+
"""
4+
Copyright 2025 The Dapr Authors
5+
Licensed under the Apache License, Version 2.0 (the "License");
6+
you may not use this file except in compliance with the License.
7+
You may obtain a copy of the License at
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
Unless required by applicable law or agreed to in writing, software
10+
distributed under the License is distributed on an "AS IS" BASIS,
11+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
See the License for the specific language governing permissions and
13+
limitations under the License.
14+
"""
15+
16+
from dapr.clients.grpc._jobs import ConstantFailurePolicy, DropFailurePolicy, FailurePolicy, Job
17+
from dapr.clients.grpc._request import BindingRequest, InvokeMethodRequest, JobEvent
18+
from dapr.clients.grpc._response import InvokeMethodResponse, TopicEventResponse
19+
from dapr.common.pubsub.subscription import SubscriptionMessage
20+
21+
# No cloudevents ImportError guard here: importing this subpackage imports the parent
22+
# dapr.ext.grpc first, which already raises the actionable error.
23+
from dapr.ext.grpc.aio.app import App, Rule # type:ignore
24+
25+
__all__ = [
26+
'App',
27+
'Rule',
28+
'SubscriptionMessage',
29+
'InvokeMethodRequest',
30+
'InvokeMethodResponse',
31+
'BindingRequest',
32+
'TopicEventResponse',
33+
'Job',
34+
'JobEvent',
35+
'FailurePolicy',
36+
'DropFailurePolicy',
37+
'ConstantFailurePolicy',
38+
]
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
# -*- coding: utf-8 -*-
2+
3+
"""
4+
Copyright 2025 The Dapr Authors
5+
Licensed under the Apache License, Version 2.0 (the "License");
6+
you may not use this file except in compliance with the License.
7+
You may obtain a copy of the License at
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
Unless required by applicable law or agreed to in writing, software
10+
distributed under the License is distributed on an "AS IS" BASIS,
11+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
See the License for the specific language governing permissions and
13+
limitations under the License.
14+
"""
15+
16+
from dapr.ext.grpc._health_servicer import _HealthCheckServicerBase
17+
from dapr.ext.grpc.aio._servicer import _needs_await
18+
from dapr.proto.runtime.v1.appcallback_pb2 import HealthCheckResponse
19+
20+
21+
class _AioHealthCheckServicer(_HealthCheckServicerBase):
22+
"""The asyncio-native implementation of HealthCheck Server.
23+
24+
Shares registration with the synchronous servicer via their common base; only the gRPC
25+
entry point differs, awaiting the callback result.
26+
"""
27+
28+
async def HealthCheck(self, request, context):
29+
"""Health check."""
30+
health_check = self._require_health_check_cb(context)
31+
result = health_check()
32+
if _needs_await(result):
33+
await result
34+
return HealthCheckResponse()

0 commit comments

Comments
 (0)