KubeMQ
Client SDKsPythonReference

RPC

Commands and queries — KubeMQ Python SDK reference.

AsyncCQClient exposes send_command / send_query and subscribe_to_commands / subscribe_to_queries. Responses reuse CommandResponse / QueryResponse wrappers.

Commands

send_command

await client.send_command(message: CommandMessage) -> CommandResponse

CommandMessage fields:

FieldTypeRequiredDescription
channelstrYesTarget channel
bodybytesConditionalPayload
metadatastrConditionalMetadata
timeout_in_secondsintYesRPC timeout
tagsdict[str, str]NoTags

subscribe_to_commands

subscribe_to_commands is an async generator — iterate it with async for:

async for command in client.subscribe_to_commands(
    subscription=CommandsSubscription(...),
    cancellation_token=token,
):
    # handle command and send response
    await client.send_response(CommandResponse(command_received=command, is_executed=True))

send_response (command)

await client.send_response(response: CommandResponse) -> None

Queries

send_query

await client.send_query(message: QueryMessage) -> QueryResponse

QueryMessage extends CommandMessage with:

FieldTypeDescription
cache_keystrCache key
cache_ttl_in_secondsintTTL

subscribe_to_queries

subscribe_to_queries is an async generator — iterate it with async for:

async for query in client.subscribe_to_queries(
    subscription=QueriesSubscription(...),
    cancellation_token=token,
):
    # handle query and send response
    await client.send_response(QueryResponse(query_received=query, is_executed=True, body=b"result"))

Quick Usage

rpc.py
cmd = CommandMessage(
    channel="svc.commands",
    body=b"ping",
    timeout_in_seconds=5,
)
resp = await client.send_command(cmd)

See Also

Was this page helpful?

On this page