Skip to content

Add asyncio support to the gRPC app extension (dapr.ext.grpc.aio) - #1206

Open
saikishore-p wants to merge 5 commits into
dapr:mainfrom
saikishore-p:feature/grpc-ext-aio
Open

saikishore-p wants to merge 5 commits into
dapr:mainfrom
saikishore-p:feature/grpc-ext-aio

Conversation

@saikishore-p

@saikishore-p saikishore-p commented Sep 12, 2026 •

Copy link
Copy Markdown

Description

Adds dapr.ext.grpc.aio, an asyncio-native App backed by a grpc.aio server, so Dapr callback handlers can be async def.

dapr.ext.grpc.App runs on grpc.server() with a 10-thread pool, so handlers must be synchronous. The async client (dapr.aio.clients.DaprClient) has been at parity with the sync client for a while, so an asyncio application can call Dapr asynchronously but cannot serve Dapr callbacks asynchronously. This PR closes that gap; there are no client-side changes.

This follows the implementation notes left on #695:

Just like we have dapr.ext.grpc which export App and Rule, we should have dapr.ext.grpc.aio which export the same. The only usage difference should be this single import.

Swapping the import is the only change an app needs, beyond async def handlers and awaiting run()/stop():

-from dapr.ext.grpc import App, InvokeMethodRequest, InvokeMethodResponse
+from dapr.ext.grpc.aio import App, InvokeMethodRequest, InvokeMethodResponse

 app = App()

 @app.method(name='my-method')
-def mymethod(request: InvokeMethodRequest) -> InvokeMethodResponse:
+async def mymethod(request: InvokeMethodRequest) -> InvokeMethodResponse:
     ...

-app.run(13551)
+asyncio.run(app.run(13551))

Reviewing this

Two commits, and the first is the one worth scrutiny:

  1. refactor: extract shared bases — touches only the three existing synchronous .py files, provably behaviour-preserving.
  2. feat: add dapr.ext.grpc.aio — purely additive apart from docs.

Review follow-up: db972b4 simplifies the closed-loop stop() path and adds the one-time warning for plain handlers, per the review below. The two merge commits only bring in main.

Why the refactor

Of the servicer's 456 lines, roughly 250 — register_topic's routing, _get_topic_callback, the bulk entry builders — contain no async content at all, and that is exactly the code that has churned most recently (SubscriptionMessage delivery, _route_map, bulk events). Copying it is what left #829 unmergeable after six months: it was written before all of those landed, and the copy silently fell behind.

So the two servicers are siblings over a shared base:

_CallbackServicerBase       # registries, topic routing, request → event translation
├── _CallbackServicer       # synchronous gRPC entry points
└── _AioCallbackServicer    # asyncio gRPC entry points

Neither subclasses the other, so there is no sync/async mixing in one MRO and no misleading isinstance. This differs from the repo's other aio modules (DaprGrpcClientAsync, dapr.ext.workflow.aio.DaprWorkflowClient), which are standalone copies — but their duplication is inherent, since every method body differs by an await. Most of this servicer's body doesn't.

Compatibility: nothing is removed. Attribute lookup walks the MRO, so every method still resolves on _CallbackServicer, instance state is untouched, and isinstance against both proto servicer types still holds — including for code reaching into internals like app._servicer._registered_topics, which examples/pubsub-simple does. The only observable delta is class-__dict__ reflection on a doubly-private class. Of the 19 methods, 12 move byte-identical; the 7 entry points delegate to 10 extracted helpers. The 48 pre-existing tests/ext/grpc tests pass unmodified against commit 1 alone.

tests/ext/grpc/aio/test_servicer.py::AsyncParityTests then locks the arrangement in: every RPC the sync servicer implements must be mirrored on the asyncio servicer as a coroutine function, and the registration helpers must stay shared rather than be reimplemented. A future RPC added to only one side fails the build.

_AioHealthCheckServicer is a standalone copy — at ~30 lines, a shared base would cost more than it saves.

If you'd rather not carry the refactor, dropping commit 1 and inlining the shared methods into a standalone _AioCallbackServicer is mechanical — happy to do that instead.

Behaviour notes

  • Full current surface, not a 2025 snapshot: service invocation, pub/sub with SubscriptionMessage delivery and the cloudevents deprecation path, topic rules, dead letter topics, disable_topic_validation, bulk topic events, job events on both the stable and alpha services, input bindings, health checks, and DAPR_GRPC_MAX_INBOUND_MESSAGE_SIZE_BYTES.
  • AppCallbackAlphaServicer is registered on the aio server, so OnJobEventAlpha1 and OnBulkTopicEventAlpha1 are actually reachable.
  • send_initial_metadata is awaited. It is a coroutine on grpc.aio's context but a plain method on the sync one; calling it without await silently drops response headers.
  • The server is built on the first run()/start(), not 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; calling it once the app is running raises, since a serving grpc.aio server cannot take new services.
  • start() is added as a non-blocking alternative to run(), for serving the app alongside other work on the same loop (e.g. an ASGI lifespan handler). run() is start() plus wait_for_termination().
  • run() stops the server on the way out, including when its task is cancelled, so an interrupted app doesn't leave its port bound. The synchronous App relies on __del__ for that, which can't await a coroutine.
  • stop(grace=None) takes a grace period. An App belongs to the event loop that started it: if that loop has closed, stop() raises, because grpc.aio needs the loop to drain the server and its listener stays bound until the process exits.
  • Plain functions are still accepted as handlers — results are awaited only when awaitable — but they run inline on the event loop, so method, subscribe, binding and job_event emit a one-time UserWarning when given one. register_health_check is exempt, since lambda: None is its documented usage.

Testing

  • 81 unit tests under tests/ext/grpc/aio/, mirroring the sync suite, plus test_server.py, which runs a real grpc.aio server on an ephemeral port and drives it through generated stubs — covering invoke (data, content type, response headers, invocation metadata), topic events, bindings, job events, health checks, UNIMPLEMENTED routing, and lifecycle: cancellation during and after startup, concurrent starts, event-loop binding, bind failure, a cancelled graceful drain, a slow handler not blocking a concurrent one, and the one-time warning for plain handlers.
  • Two new examples with tests/examples/ coverage: invoke-simple-async and pubsub-simple-async.

Run locally on Python 3.10 (the lint/type floor and the bottom of the CI matrix) against a real Dapr 1.18.4 runtime — sidecar, placement, scheduler and Redis:

Suite Command Result
Unit pytest -m "not e2e" ./tests --ignore=tests/integration --ignore=tests/examples 2153 passed
gRPC extension pytest tests/ext/grpc/ 129 passed
New asyncio tests pytest tests/ext/grpc/aio/ 81 passed
Integration pytest tests/integration/ -m "not dapr_head" 94 passed
Examples pytest tests/examples/ 49 of 50 passed, see below
Types mypy clean, 225 source files
Lint / format ruff check / ruff format --check clean
unittest runner python -m unittest discover ./tests/ext/grpc OK

The one example that didn't pass is test_configuration, which is order-dependent when the whole suite runs serially: it passes on its own, does the same on unmodified main, doesn't import dapr.ext.grpc, and is untouched by this PR. dapr.ext.grpc.* is not in mypy's ignore list, so the new code is fully type-checked.

The first commit was also verified in isolation: with only refactor: extract shared bases... applied, the 48 pre-existing tests/ext/grpc tests pass unmodified and the full unit suite reports 1640 — exactly the count on its parent commit.

pubsub-simple-async is the server-side callback model; the existing pubsub-streaming-async is the client-side streaming API. Both READMEs now cross-reference each other so the distinction is clear.

Documentation

dapr/docs#5350 adds an asyncio section to the Python gRPC extension page, on the v1.19 branch.

Issue reference

Please reference the issue this PR will close: #695

Some history, since it isn't obvious from the issue: #695 was closed as completed on 2026-04-01, but nothing in the repo implements it. The prior attempt, #829, was auto-closed for inactivity after 60 days. That PR was written against the old ext/dapr-ext-grpc/ layout and predates SubscriptionMessage, the _route_map routing rewrite, bulk topic events and the stable OnJobEvent, so it could not simply be rebased — this is a fresh implementation against current main.

We hit this in production and ended up vendoring our own asyncio callback server to work around it, which is what prompted picking the work back up.

Checklist

Please make sure you've completed the relevant tasks for this PR, out of the following list:

  • Code compiles correctly
  • Created/updated tests
  • Extended the documentation

@saikishore-p
saikishore-p requested review from a team as code owners September 12, 2026 16:05
@saikishore-p
saikishore-p force-pushed the feature/grpc-ext-aio branch 2 times, most recently from 04eb3b9 to e2e4ccc Compare September 13, 2026 04:05
…icers

Splits _CallbackServicer into a base holding everything that does not invoke a
user handler - the handler registries, topic routing, and the translation of
incoming requests into the SDK types handlers receive - and a thin subclass
supplying the synchronous gRPC entry points. _HealthCheckServicer is split the
same way, with callback registration moving to the base.

This is preparation for asyncio servicers, which need all of the former and none
of the latter. Keeping that logic in one place avoids the duplication that left
the previous attempt at async support (dapr#829) unmergeable once SubscriptionMessage
delivery, the _route_map rewrite and bulk topic events landed on the synchronous
side only.

Behaviour is unchanged. Nothing is removed from either class's reachable surface
- attribute lookup walks the MRO, so callers reaching into internals such as
app._servicer, which examples/pubsub-simple does, are unaffected. The existing
tests pass unmodified.

Also extracts _resolve_topic_event_type in app.py so the deprecation warning for
unannotated topic handlers has a single definition to share.

Signed-off-by: Sai Kishore Punagani <63619246+saikishore-p@users.noreply.github.com>
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 dapr#695

Signed-off-by: Sai Kishore Punagani <63619246+saikishore-p@users.noreply.github.com>

@CasperGN CasperGN left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nice work, and thanks for the clear write-up. I'm happy with the shared-base approach: it's the right trade-off given how much of the servicer has no async content, and the parity test keeps the two sides from drifting.

What I checked:

  • I read the refactor commit line by line. It's behaviour-preserving: the only reordering is that InvokeMethodResponse is now checked before bytes/str/GrpcMessage, and those types don't overlap. The bulk "no handler" path now raises inside the helper rather than returning None to the caller, with the same result on the wire.
  • tests/ext/grpc passes (126), and so do the unittest runner, mypy (186 files) and ruff.
  • I ran a live smoke test against a real grpc.aio server. 20 concurrent calls to a 50 ms async def handler finished in about 0.05 s, so they really ran concurrently. Unknown methods return UNIMPLEMENTED and handler errors return UNKNOWN, the same as the sync app. Start, stop, then start again on the same port in a second asyncio.run() works.

Three things I'd like addressed before merging:

  1. Docs. This adds a public module (dapr.ext.grpc.aio), so it needs a section on the Python SDK gRPC extension page in dapr/docs, with a link to the docs PR here. The README alone won't reach most users.

  2. Simplify the closed-loop recovery. _abandon_if_owning_loop_is_closed (app.py L135-L165), _abandoned_ports (L79, L195-L200) and the SO_REUSEPORT warning are a lot of state to maintain for a case where the caller has already misused the app. Please replace them with a clear error from stop() (L271) when the owning loop is closed, saying that the listener stays bound until the process exits.

  3. Warn when a sync handler is registered on the aio app. Plain functions run inline on the event loop (_needs_await), and one blocking handler stalls every other RPC. Please emit a one-time warning from the method/subscribe/binding/job_event decorators (L362-L364, L401-L413, L430-L432, L457-L459) when the handler isn't a coroutine function. Leave register_health_check (L315) out, since lambda: None is its documented usage.

stop() now raises when the event loop the app was started on has closed,
instead of abandoning the server and tracking its port. grpc.aio needs the
owning loop to drain the server, so the listener stays bound until the
process exits; closing the loop with the app running is a caller error, and
the error says so. Removes _abandon_if_owning_loop_is_closed, its bookkeeping
and the SO_REUSEPORT warning.

method, subscribe, binding and job_event now emit a one-time UserWarning when
the handler is not a coroutine function, since a plain handler runs inline on
the event loop and a blocking one stalls every other RPC.
register_health_check is exempt: lambda: None is its documented usage.

Signed-off-by: Sai Kishore Punagani <63619246+saikishore-p@users.noreply.github.com>
@saikishore-p

Copy link
Copy Markdown
Author

Thanks for the thorough review, and for checking it against a live server.

I've merged current main and addressed all three in db972b4:

  1. Docs — Python SDK: document the asyncio gRPC app (dapr.ext.grpc.aio) docs#5350 adds an asyncio section to the gRPC extension page, on the v1.19 branch since the module ships in 1.19.

  2. Closed-loop recovery — _abandon_if_owning_loop_is_closed, _abandoned_servers, _abandoned_ports and the SO_REUSEPORT warning are gone. stop() now raises a RuntimeError when the owning loop is closed, saying the listener stays bound until the process exits. It keeps no state for that case, so the App stays unusable after the misuse rather than half-recovering. Net −14 lines in aio/app.py.

  3. Sync handler warning — method, subscribe, binding and job_event each call a small helper that emits a UserWarning once per app when the handler isn't a coroutine function. register_health_check is left out. Tests cover the once-per-app behaviour, async-only apps staying silent, and the health check exemption.

One question on (3). I check for "coroutine function" with inspect.iscoroutinefunction, and it gives a false warning for some handler shapes that do work correctly on the aio app, because the servicer awaits whatever they return:

  • an object with async def __call__
  • a sync decorator wrapping an async def handler, with or without functools.wraps

If you'd rather avoid warning on those, I can unwrap __wrapped__ and check __call__ before deciding, along the lines of is_async_callable in the workflow extension, about 10 lines, copied rather than imported so dapr[grpc] doesn't depend on the workflow extra.

One case neither approach handles: unittest.mock.AsyncMock is only recognised as a coroutine function from Python 3.12, so registering one as a handler warns on 3.10/3.11 and not on 3.12+. It only matters for test suites that treat warnings as errors, so I've left it alone, but let me know if you'd like it special-cased.

@saikishore-p
saikishore-p requested a review from CasperGN October 3, 2026 04:32
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Feature Request] Async GRPC Subscriber (dapr-ext-grpc)

2 participants