Skip to content

Commit 1888e4f

Browse files
committed
feat: add SendLiveMessage support for group (multicast) channels
- Add A2AServiceGroupStub.SendLiveMessage using call_multicast_stream_stream, returning a MulticastBidiStreamHandler - Add SRPCMulticastTransport.send_live_message: concurrently sends requests and yields (source, StreamResponse) tuples from all group members - Add MulticastClient.send_live_message delegating to the transport Signed-off-by: Sam Betts <1769706+Tehsmash@users.noreply.github.com>
1 parent 0c35f76 commit 1888e4f

2 files changed

Lines changed: 52 additions & 0 deletions

File tree

slima2a/client_transport.py

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -473,6 +473,34 @@ async def get_extended_agent_card(
473473
async for source, response in self.stub.GetExtendedAgentCard(request):
474474
yield source, response
475475

476+
async def send_live_message(
477+
self,
478+
request_stream: "AsyncGenerator[Any, None]",
479+
*,
480+
context: ClientCallContext | None = None,
481+
) -> AsyncGenerator[tuple[Any, StreamResponse], None]:
482+
"""Sends a bidirectional streaming live message to all agents in the group.
483+
484+
Yields (source, StreamResponse) tuples as events arrive from any agent.
485+
"""
486+
bidi = self.stub.SendLiveMessage()
487+
async def _send():
488+
async for req in request_stream:
489+
await bidi.send_async(req.SerializeToString())
490+
await bidi.close_send_async()
491+
send_task = asyncio.ensure_future(_send())
492+
try:
493+
while True:
494+
msg = await bidi.recv_async()
495+
if msg.is_end():
496+
break
497+
if msg.is_error():
498+
raise msg.error
499+
if msg.is_data():
500+
yield msg.item.context, StreamResponse.FromString(msg.item.message)
501+
finally:
502+
await send_task
503+
476504
async def close(self) -> None:
477505
"""Closes the transport and releases any resources."""
478506
pass
@@ -648,6 +676,22 @@ async def get_extended_agent_card(
648676
):
649677
yield source, response
650678

679+
async def send_live_message(
680+
self,
681+
request_stream: "AsyncGenerator[Any, None]",
682+
*,
683+
context: ClientCallContext | None = None,
684+
) -> AsyncGenerator[tuple[Any, StreamResponse], None]:
685+
"""Sends a bidirectional streaming live message to all agents in the group.
686+
687+
Yields (source, StreamResponse) tuples as events arrive from any agent.
688+
Use ``source`` to demultiplex per-agent.
689+
"""
690+
async for source, response in self._transport.send_live_message(
691+
request_stream, context=context
692+
):
693+
yield source, response
694+
651695
async def close(self) -> None:
652696
await self._transport.close()
653697

slima2a/types/v1/a2a_pb2_slimrpc.py

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -382,6 +382,14 @@ async def DeleteTaskPushNotificationConfig(self, request: a2a__pb2.DeleteTaskPus
382382
if msg.is_data():
383383
yield msg.item.context, google__protobuf__empty_pb2.Empty.FromString(msg.item.message)
384384

385+
def SendLiveMessage(self, timeout: Optional[timedelta] = None, metadata: Optional[dict[str, str]] = None) -> slim_bindings.MulticastBidiStreamHandler:
386+
"""Open a bidirectional streaming SendLiveMessage call to all group members."""
387+
return self._channel.call_multicast_stream_stream(
388+
"lf.a2a.v1.A2AService",
389+
"SendLiveMessage",
390+
timeout,
391+
metadata,
392+
)
385393

386394

387395
class A2AServiceServicer:

0 commit comments

Comments
 (0)