Skip to content

Latest commit

History

History
217 lines (158 loc) · 6.16 KB

File metadata and controls

217 lines (158 loc) · 6.16 KB

Python SDK

Package:robustmq.mq9
Requires: Python 3.10+, nats-py

Install

pip install robustmq

Quick start

importasynciofromrobustmq.mq9importClient, Priorityasyncdefmain():
asyncwithClient(server="nats://demo.robustmq.com:4222") asclient:
# Create mailboxmailbox=awaitclient.create(ttl=3600)
# Sendawaitclient.send(mailbox.mail_id, b"hello", priority=Priority.NORMAL)
# Subscribeasyncdefhandler(msg):
print(msg.payload)
sub=awaitclient.subscribe(mailbox.mail_id, handler)
# List metadata (no payload)metas=awaitclient.list(mailbox.mail_id)
# metas[i].msg_id, .priority, .ts# Deleteawaitclient.delete(mailbox.mail_id, metas[0].msg_id)
awaitsub.unsubscribe()
asyncio.run(main())

API

Client(server="nats://demo.robustmq.com:4222", *, max_reconnect_attempts=10,
reconnect_time_wait=2.0, request_timeout=5.0)
awaitclient.connect()
awaitclient.close()
# also: async with Client(...) as clientmailbox=awaitclient.create(ttl, *, public=False, name="", desc="")
# → Mailbox(mail_id, public, name, desc)awaitclient.send(mail_id, payload, *, priority=Priority.NORMAL)
# payload: bytes | str | dictsub=awaitclient.subscribe(mail_id, async_cb, *, priority="*", queue_group="")
# → nats Subscription; call sub.unsubscribe() to cancelmetas=awaitclient.list(mail_id)
# → list[MessageMeta(msg_id, priority, ts)]awaitclient.delete(mail_id, msg_id)

Priority

Priority.CRITICAL# "critical" — highest priorityPriority.URGENT# "urgent"Priority.NORMAL# "normal" — default, no suffix

Public mailbox

mailbox=awaitclient.create(ttl=3600, public=True, name="task.queue", desc="Tasks")
# mail_id == "task.queue"

Queue group (competitive consumption)

awaitclient.subscribe(mail_id, handler, queue_group="workers")
# Each message delivered to exactly one subscriber in the group

Error handling

importasynciofromrobustmq.mq9importClient, Priorityfromrobustmq.mq9.clientimportMq9Errorimportnats.errorsasyncdefmain():
# ── Connection failure ────────────────────────────────────────────────────# Raised when the NATS server is unreachable.try:
client=Client(server="nats://unreachable-host:4222")
awaitclient.connect()
exceptExceptionase:
print(f"Connection failed: {e}")
# ── Normal operation with mailbox / request errors ─────────────────────────asyncwithClient(server="nats://demo.robustmq.com:4222") asclient:
# Mailbox not found (404) — raised as Mq9Error with .code == 404try:
awaitclient.list("m-does-not-exist")
exceptMq9Errorase:
print(f"mq9 error {e.code}: {e}") # mq9 error 404: mailbox not found# Request timeout — raised as Mq9Error wrapping asyncio.TimeoutErrortry:
awaitclient.list("m-some-mailbox")
exceptMq9Errorase:
print(f"Timed out: {e}")
# No responders — server not running or subject invalidtry:
mailbox=awaitclient.create(ttl=3600)
exceptMq9Errorase:
ifisinstance(e.__cause__, nats.errors.NoRespondersError):
print("No mq9 server is listening on this subject")
asyncio.run(main())

Usage in AI Agent

Typical pattern: Agent creates its mailbox on startup, runs a message loop, and forwards results to another Agent.

importasynciofromrobustmq.mq9importClient, PrioritySERVER="nats://demo.robustmq.com:4222"DOWNSTREAM_MAIL_ID="agent-b.inbox"# well-known address of the next Agentasyncdefprocess(payload: bytes) ->bytes:
"""Business logic: transform the incoming task into a result."""task=payload.decode()
returnf"result:{task}".encode()
asyncdefrun_agent() ->None:
asyncwithClient(server=SERVER) asclient:
# 1. Create this Agent's own mailbox on startupmailbox=awaitclient.create(ttl=3600)
print(f"[agent] started, mail_id={mailbox.mail_id}")
# 2. Main loop: subscribe and process messagesasyncdefhandler(msg):
print(f"[agent] received [{msg.priority.value}] {msg.payload.decode()}")
result=awaitprocess(msg.payload)
# 3. Send result to the downstream Agentawaitclient.send(DOWNSTREAM_MAIL_ID, result, priority=Priority.NORMAL)
print(f"[agent] forwarded result to {DOWNSTREAM_MAIL_ID}")
sub=awaitclient.subscribe(mailbox.mail_id, handler)
# Keep the Agent runningtry:
awaitasyncio.Future() # run foreverexceptasyncio.CancelledError:
passfinally:
awaitsub.unsubscribe()
if__name__=="__main__":
asyncio.run(run_agent())
## LangGraph IntegrationThe`langchain-mq9`toolsworkdirectlyinsideLangGraphnodeseachnodecancreateamailbox, sendmessages, andreceivereplieswithoutanyextrawiring.
Seethefullworkingexamplein [demo/demo-langgraph/](../demo/demo-langgraph/) andthetoolreferencein [docs/langchain-mq9.md](langchain-mq9.md).
```pythonfromlangchain_mq9importCreateMailboxTool, SendMessageTool, GetMessagesToolfromlanggraph.graphimportStateGraph, ENDfromtypingimportTypedDictSERVER="nats://demo.robustmq.com:4222"classState(TypedDict):
mail_id: strmessages: strasyncdefnode_inbox(state: State) ->State:
create=CreateMailboxTool(server=SERVER)
state["mail_id"] =awaitcreate._arun(ttl=120)
returnstateasyncdefnode_read(state: State) ->State:
get=GetMessagesTool(server=SERVER)
state["messages"] =awaitget._arun(mail_id=state["mail_id"], limit=10)
returnstateg=StateGraph(State)
g.add_node("inbox", node_inbox)
g.add_node("read", node_read)
g.set_entry_point("inbox")
g.add_edge("inbox", "read")
g.add_edge("read", END)
graph=g.compile()