|
18 | 18 | import asyncio |
19 | 19 | import functools |
20 | 20 | import inspect |
| 21 | +import json |
21 | 22 | import datetime |
22 | 23 | import enum |
23 | 24 | import os |
24 | 25 | import mimetypes |
25 | 26 | import weakref |
26 | | -from typing import List, Union, Callable, Dict, Awaitable, Optional, Mapping, cast, TypeVar |
| 27 | +from typing import Any, List, Union, Callable, Dict, Awaitable, Optional, Mapping, cast, TypeVar |
27 | 28 | from abc import abstractmethod, ABC |
28 | 29 |
|
29 | 30 | from ._ffi_client import FfiClient, FfiHandle |
|
59 | 60 | _chain_outgoing, |
60 | 61 | ) |
61 | 62 | from ._proto.rpc_pb2 import RpcMethodInvocationResponseRequest |
| 63 | +from .actions import ( |
| 64 | + ACTION_DECLINED_CODE, |
| 65 | + ACTION_METHOD_PREFIX, |
| 66 | + ACTIONS_ATTRIBUTE, |
| 67 | + ActionDeclinedError, |
| 68 | + ActionEntry, |
| 69 | + ActionHandler, |
| 70 | + ActionRegistration, |
| 71 | + invoke_handler, |
| 72 | + parse_actions, |
| 73 | + serialize_actions, |
| 74 | +) |
62 | 75 | from .log import logger |
63 | 76 |
|
64 | 77 | from .data_stream import ( |
@@ -173,6 +186,11 @@ def attributes(self) -> dict[str, str]: |
173 | 186 | """Custom attributes associated with the participant.""" |
174 | 187 | return dict(self._info.attributes) |
175 | 188 |
|
| 189 | + @property |
| 190 | + def actions(self) -> List[ActionEntry]: |
| 191 | + """Snapshot of the actions this participant exposes.""" |
| 192 | + return parse_actions(self._info.attributes.get(ACTIONS_ATTRIBUTE)) |
| 193 | + |
176 | 194 | @property |
177 | 195 | def kind(self) -> proto_participant.ParticipantKind.ValueType: |
178 | 196 | """Participant's kind (e.g., regular participant, ingress, egress, sip, agent).""" |
@@ -276,6 +294,7 @@ def __init__( |
276 | 294 | self._track_publications: dict[str, LocalTrackPublication] = {} |
277 | 295 | self._rpc_handlers: Dict[str, RpcHandler] = {} |
278 | 296 | self._rpc_interceptors: List[RpcInterceptor] = [] |
| 297 | + self._action_catalog: Dict[str, ActionEntry] = {} |
279 | 298 | # Handles of data stream writers that have been opened but not yet |
280 | 299 | # closed, so the room can drop them at disconnect. The FFI close |
281 | 300 | # request consumes the handle (take_handle), so an entry is removed as |
@@ -584,6 +603,74 @@ def unregister_rpc_method(self, method: str) -> None: |
584 | 603 |
|
585 | 604 | FfiClient.instance.request(req) |
586 | 605 |
|
| 606 | + async def register_action( |
| 607 | + self, |
| 608 | + name: str, |
| 609 | + description: str, |
| 610 | + parameters: Dict[str, Any], |
| 611 | + handler: ActionHandler, |
| 612 | + *, |
| 613 | + consent: str = "none", |
| 614 | + ) -> ActionRegistration: |
| 615 | + """ |
| 616 | + Expose an action other participants can discover and call. |
| 617 | +
|
| 618 | + The handler receives the parsed arguments and an `ActionContext`, may be sync or async, |
| 619 | + and returns a JSON-serializable result. Raise `ActionDeclinedError` to decline. |
| 620 | + """ |
| 621 | + method = ACTION_METHOD_PREFIX + name |
| 622 | + |
| 623 | + async def rpc_handler(data: RpcInvocationData) -> str: |
| 624 | + return await invoke_handler(handler, data.payload, data.caller_identity) |
| 625 | + |
| 626 | + self.register_rpc_method(method, rpc_handler) |
| 627 | + self._action_catalog[name] = ActionEntry(name, description, parameters, consent) |
| 628 | + await self._publish_actions() |
| 629 | + |
| 630 | + async def unregister() -> None: |
| 631 | + self.unregister_rpc_method(method) |
| 632 | + self._action_catalog.pop(name, None) |
| 633 | + await self._publish_actions() |
| 634 | + |
| 635 | + return ActionRegistration(unregister) |
| 636 | + |
| 637 | + async def call_action( |
| 638 | + self, |
| 639 | + destination_identity: str, |
| 640 | + name: str, |
| 641 | + args: Optional[Dict[str, Any]] = None, |
| 642 | + *, |
| 643 | + response_timeout: Optional[float] = None, |
| 644 | + ) -> Any: |
| 645 | + """ |
| 646 | + Call an action on another participant and return its parsed result. |
| 647 | +
|
| 648 | + Raises: |
| 649 | + ActionDeclinedError: The participant declined the call. |
| 650 | + RpcError: Unknown action, timeout, or other transport failure. |
| 651 | + """ |
| 652 | + try: |
| 653 | + response = await self.perform_rpc( |
| 654 | + destination_identity=destination_identity, |
| 655 | + method=ACTION_METHOD_PREFIX + name, |
| 656 | + payload=json.dumps(args or {}), |
| 657 | + response_timeout=response_timeout, |
| 658 | + ) |
| 659 | + except RpcError as e: |
| 660 | + if e.code == ACTION_DECLINED_CODE: |
| 661 | + raise ActionDeclinedError(e.message) from e |
| 662 | + raise |
| 663 | + return json.loads(response) if response else None |
| 664 | + |
| 665 | + async def _publish_actions(self) -> None: |
| 666 | + await self.set_attributes( |
| 667 | + {ACTIONS_ATTRIBUTE: serialize_actions(list(self._action_catalog.values()))} |
| 668 | + ) |
| 669 | + |
| 670 | + def _republish_actions(self) -> None: |
| 671 | + if self._action_catalog: |
| 672 | + asyncio.ensure_future(self._publish_actions()) |
| 673 | + |
587 | 674 | def set_track_subscription_permissions( |
588 | 675 | self, |
589 | 676 | *, |
|
0 commit comments