danielgtaylor / danielgtaylor/python-betterproto

[Feqture request] Add a way to produce server service implementation with using Stream directly

Aperta
#204 1 commento 0 reazioni 0 assegnatari Vedi su GitHub
Lingua principale
Python
Stelle
1.8k
Fork
234
Metriche di merge delle PR
Nessuna PR unita negli ultimi 30g

Descrizione

I'd love to use `grpclib.server.Stream` directly instead of `request_iterator/yield` combination to handle stream/steram gRPC server.

For example, the current way is a bit tricky to write consumer/producer style while `yield` must be placed directly under the function scope. So I need to use exception or some other tricky way to find if gRPC is closed or not.

```python
import asyncio
from contextlib import suppress

from .proto import ServerBase

# Assume that the followings are inbound/outbound for proxing gRPC messages
inbound: asyncio.Queue[str] = asyncio.Queue()
outbound: asyncio.Queue[str] = asyncio.Queue()

class ServerWithIteratorAndGenerator(ServerBase):
async def proxy(self, request_iterator: AsyncIterator[str]) -> AsyncIterator[str]:
class ConsumerClosedError(Exception):
pass

async def consumer_handler() -> None:
async for message in request_iterator:
await outbound.put(message)
raise ConsumerClosedError()

consumer = asyncio.create_task(consumer_handler())

with suppress(ConsumerClosedError):
while True:
producer = inbound.get()
done, _ = await asyncio.wait(
[consumer, producer],
return_when=asyncio.FIRST_COMPLETED,
)
# NOTE:
# If 'consumer' is closed (gRPC has closed), the code below raise ConsumerClosedError
message = await next(done)
if not message:
# inbound is closed.
break
yield message

# Make sure that consumer is closed
with suppress(ConsumerClosedError):
consumer.cancel()
```

If betterproto produce `ServerBase` with native `grpclib.server.Stream`, I could write above code as

```python
import asyncio

from .proto import ServerBase

# Assume that the followings are inbound/outbound for proxing gRPC messages
inbound: asyncio.Queue[str] = asyncio.Queue()
outbound: asyncio.Queue[str] = asyncio.Queue()

class ServerWithNativeStream(ServerBase):
async def proxy(self, stream: grpclib.server.Stream) -> None:
async def consumer_handler() -> None:
while True:
message = await stream.recv_message()
if message is None:
# gRPC is closed
break
await outbound.put(message)

async def producer_handler() -> None:
while True:
message = await inbound.get()
if message is None:
# inbound is closed
break
await stream.send_message(message)

consumer = asyncio.create_task(consumer_handler())
producer = asyncio.create_task(producer_handler())

done, pending = await asyncio.wait(
[consumer, producer],
return_when=asyncio.FIRST_COMPLETED,
)
for task in pending:
task.cancel()
for task in done:
task.result()
```

So above code does not have any tricky tech. to find if incoming stream is closed.

I know that I can use `_SeverBase__rpc_proxy()` method to overwrite it but I feel it's a bit redundant and ugly.

```python
import asyncio

from .proto import ServerBase

# Assume that the followings are inbound/outbound for proxing gRPC messages
inbound: asyncio.Queue[str] = asyncio.Queue()
outbound: asyncio.Queue[str] = asyncio.Queue()

class Server(ServerBase):
async def _ServerBase__rpc_proxy(self, stream: grpclib.server.Stream) -> None:
async def consumer_handler() -> None:
while True:
message = await stream.recv_message()
if message is None:
# gRPC is closed
break
await outbound.put(message)

async def producer_handler() -> None:
while True:
message = await inbound.get()
if message is None:
# inbound is closed
break
await stream.send_message(message)

consumer = asyncio.create_task(consumer_handler())
producer = asyncio.create_task(producer_handler())

done, pending = await asyncio.wait(
[consumer, producer],
return_when=asyncio.FIRST_COMPLETED,
)
for task in pending:
task.cancel()
for task in done:
task.result()
```

Guida per i contributori

Apri la guida per i contributori

Direzione di ricerca

Inizia leggendo l’implementazione generata di ServerBase e il suo punto di ingresso _ServerBase__rpc_proxy, quindi confronta l’interfaccia attuale request_iterator/yield con grpclib.server.Stream. Il lavoro è completato quando le implementazioni server generate possono ricevere direttamente lo Stream nativo per gli RPC di streaming senza sovrascrivere il metodo privato e il comportamento risultante è stato verificato rispetto ai controlli esistenti del progetto.

Scritto dal modello di indicizzazione a partire dal testo della issue.

Valutazione

Stack tecnologico
python
Ambito
api, backend-api-design
Tipo di issue
Funzionalità
Difficoltà
5/5
Tempo stimato
Più di una settimana
Stato di attività
Ferma
Chiarezza
Abbastanza chiara
Idoneità per principianti
25/100

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.