Source code for soulfire.example_plugin
from __future__ import annotations
from effect_py import EffectGen, fn
from .errors import SoulFireOperationError
from .plugin.example.v1.example_connect import ExamplePluginServiceClient
from .plugin.example.v1.example_pb2 import EchoRequest, EchoResponse, Tick, WatchTicksRequest
from .plugin_api_pb2 import PluginApiDescriptor
from .plugins import PluginCatalog
from .streams import Stream
from .transport import rpc, rpc_stream
_PLUGIN_ID = "example"
_SERVICE_NAME = "soulfire.plugin.example.v1.ExamplePluginService"
[docs]
class ExamplePluginClient:
__slots__ = ("_client",)
def __init__(self, client: ExamplePluginServiceClient) -> None:
self._client = client
[docs]
@fn("ExamplePluginClient.echo")
def echo(
self, instance_id: str, message: str, *, timeout_ms: int | None = None
) -> EffectGen[EchoResponse, SoulFireOperationError]:
return (
yield from rpc(
"ExamplePluginClient.echo",
lambda: self._client.echo(
EchoRequest(instance_id=instance_id, message=message), timeout_ms=timeout_ms
),
)
)
[docs]
def watch_ticks(
self, instance_id: str, count: int, *, timeout_ms: int | None = None
) -> Stream[Tick, SoulFireOperationError]:
return rpc_stream(
"ExamplePluginClient.watch_ticks",
lambda: self._client.watch_ticks(
WatchTicksRequest(instance_id=instance_id, count=count), timeout_ms=timeout_ms
),
)
class _ExamplePluginModule:
plugin_id = _PLUGIN_ID
@staticmethod
def is_compatible(descriptor: PluginApiDescriptor) -> bool:
return _is_compatible(descriptor)
@staticmethod
def create(catalog: PluginCatalog, _: PluginApiDescriptor) -> ExamplePluginClient:
return ExamplePluginClient(catalog.service(ExamplePluginServiceClient))
example_plugin = _ExamplePluginModule()
def _is_compatible(descriptor: PluginApiDescriptor) -> bool:
return descriptor.api_major_version == 1 and any(
service.full_name == _SERVICE_NAME for service in descriptor.services
)