class AsyncApiGenerator:
def __init__(self, title: str = "Istos Network", version: str = "1.0.0"):
self.doc: Dict[str, Any] = {
"asyncapi": "3.0.0",
"info": {
"title": title,
"version": version
},
"channels": {},
"operations": {},
"components": {
"messages": {},
"schemas": {}
}
}
self._message_count = 0
def _register_message(self, name_hint: str, schema: dict) -> str:
self._message_count += 1
msg_id = f"{name_hint}Message_{self._message_count}"
self.doc["components"]["messages"][msg_id] = {
"payload": schema
}
return f"#/components/messages/{msg_id}"
def generate(self, istos_instance: Any) -> str:
"""Generates the AsyncAPI YAML specification."""
self.doc["asyncapi"] = "3.0.0"
if "operations" not in self.doc:
self.doc["operations"] = {}
# Helper to setup channel
def _ensure_channel(ch_name, address, title, description):
if ch_name not in self.doc["channels"]:
self.doc["channels"][ch_name] = {
"address": address,
"title": title,
"description": description,
"messages": {}
}
# 1. Handlers (RPC - Receive Query and Reply)
for handler in istos_instance._handlers:
schemas = get_function_schemas(handler.func)
ch_name = handler.prefix.replace('/', '_').replace('*', 'star')
_ensure_channel(ch_name, handler.prefix, f"Handler: {handler.func.__name__}", inspect.getdoc(handler.func) or "")
op_id = f"handle_{handler.func.__name__}"
op: Dict[str, Any] = {
"action": "receive",
"channel": { "$ref": f"#/channels/{ch_name}" },
"tags": [{"name": "@handle"}, {"name": "RPC"}]
}
if schemas["payload_schema"]:
msg_key = handler.func.__name__ + "_req"
msg_ref = self._register_message(msg_key, schemas["payload_schema"])
self.doc["channels"][ch_name]["messages"][msg_key] = { "$ref": msg_ref }
op["messages"] = [{ "$ref": f"#/channels/{ch_name}/messages/{msg_key}" }]
if schemas["return_schema"]:
rep_msg_key = handler.func.__name__ + "_rep"
rep_msg_ref = self._register_message(rep_msg_key, schemas["return_schema"])
self.doc["channels"][ch_name]["messages"][rep_msg_key] = { "$ref": rep_msg_ref }
op["reply"] = {
"messages": [{ "$ref": f"#/channels/{ch_name}/messages/{rep_msg_key}" }]
}
self.doc["operations"][op_id] = op
# 2. Subscribers (Pub/Sub - Receive Events)
for sub in istos_instance._subscribers:
schemas = get_function_schemas(sub.func)
ch_name = sub.prefix.replace('/', '_').replace('*', 'star')
_ensure_channel(ch_name, sub.prefix, f"Topic: {sub.prefix}", inspect.getdoc(sub.func) or "")
op_id = f"subscribe_{sub.func.__name__}"
op = {
"action": "receive",
"channel": { "$ref": f"#/channels/{ch_name}" },
"tags": [{"name": "@subscribe"}, {"name": "Pub/Sub"}]
}
if schemas["payload_schema"]:
msg_key = sub.func.__name__ + "_event"
msg_ref = self._register_message(msg_key, schemas["payload_schema"])
self.doc["channels"][ch_name]["messages"][msg_key] = { "$ref": msg_ref }
op["messages"] = [{ "$ref": f"#/channels/{ch_name}/messages/{msg_key}" }]
self.doc["operations"][op_id] = op
# 3. Publishers (Pub/Sub - Send Events)
for pub in istos_instance._publishers:
schemas = get_function_schemas(pub.func)
ch_name = pub.prefix.replace('/', '_').replace('*', 'star')
_ensure_channel(ch_name, pub.prefix, f"Topic: {pub.prefix}", inspect.getdoc(pub.func) or "")
op_id = f"publish_{pub.func.__name__}"
op = {
"action": "send",
"channel": { "$ref": f"#/channels/{ch_name}" },
"tags": [{"name": "@publish"}, {"name": "Pub/Sub"}]
}
if schemas["return_schema"]:
msg_key = pub.func.__name__ + "_event"
msg_ref = self._register_message(msg_key, schemas["return_schema"])
self.doc["channels"][ch_name]["messages"][msg_key] = { "$ref": msg_ref }
op["messages"] = [{ "$ref": f"#/channels/{ch_name}/messages/{msg_key}" }]
self.doc["operations"][op_id] = op
# 4. Queries (RPC - Send Query and Expect Reply)
for query in getattr(istos_instance, "_queries", []):
schemas = get_function_schemas(query.func)
ch_name = query.prefix.replace('/', '_').replace('*', 'star')
_ensure_channel(ch_name, query.prefix, f"Query: {query.func.__name__}", inspect.getdoc(query.func) or "")
op_id = f"query_{query.func.__name__}"
op = {
"action": "send",
"channel": { "$ref": f"#/channels/{ch_name}" },
"tags": [{"name": "@query"}, {"name": "RPC"}]
}
if schemas["payload_schema"]:
msg_key = query.func.__name__ + "_req"
msg_ref = self._register_message(msg_key, schemas["payload_schema"])
self.doc["channels"][ch_name]["messages"][msg_key] = { "$ref": msg_ref }
op["messages"] = [{ "$ref": f"#/channels/{ch_name}/messages/{msg_key}" }]
if schemas["return_schema"]:
rep_msg_key = query.func.__name__ + "_rep"
rep_msg_ref = self._register_message(rep_msg_key, schemas["return_schema"])
self.doc["channels"][ch_name]["messages"][rep_msg_key] = { "$ref": rep_msg_ref }
op["reply"] = {
"messages": [{ "$ref": f"#/channels/{ch_name}/messages/{rep_msg_key}" }]
}
self.doc["operations"][op_id] = op
def _safe_schemas(func: Callable) -> Dict[str, Any]:
# A @channel handler takes a ChannelSession, which has no JSON Schema;
# never let an unrepresentable parameter sink the whole document.
try:
return get_function_schemas(func)
except Exception:
return {"payload_schema": None, "return_schema": None}
# 5. Streams (RPC - one request, many reply chunks)
for stream in getattr(istos_instance, "_streams", []):
schemas = _safe_schemas(stream.func)
ch_name = stream.prefix.replace('/', '_').replace('*', 'star')
_ensure_channel(ch_name, stream.prefix, f"Stream: {stream.func.__name__}", inspect.getdoc(stream.func) or "")
op_id = f"stream_{stream.func.__name__}"
op = {
"action": "receive",
"channel": { "$ref": f"#/channels/{ch_name}" },
"tags": [{"name": "@stream"}, {"name": "Streaming"}]
}
if schemas["payload_schema"]:
msg_key = stream.func.__name__ + "_req"
msg_ref = self._register_message(msg_key, schemas["payload_schema"])
self.doc["channels"][ch_name]["messages"][msg_key] = { "$ref": msg_ref }
op["messages"] = [{ "$ref": f"#/channels/{ch_name}/messages/{msg_key}" }]
self.doc["operations"][op_id] = op
# 6. Channels (bidirectional - full-duplex session)
for channel in getattr(istos_instance, "_channels", []):
schemas = _safe_schemas(channel.func)
ch_name = channel.prefix.replace('/', '_').replace('*', 'star')
_ensure_channel(ch_name, channel.prefix, f"Channel: {channel.func.__name__}", inspect.getdoc(channel.func) or "")
op_id = f"channel_{channel.func.__name__}"
op = {
"action": "receive",
"channel": { "$ref": f"#/channels/{ch_name}" },
"tags": [{"name": "@channel"}, {"name": "Bidirectional"}]
}
if schemas["payload_schema"]:
msg_key = channel.func.__name__ + "_open"
msg_ref = self._register_message(msg_key, schemas["payload_schema"])
self.doc["channels"][ch_name]["messages"][msg_key] = { "$ref": msg_ref }
op["messages"] = [{ "$ref": f"#/channels/{ch_name}/messages/{msg_key}" }]
self.doc["operations"][op_id] = op
return str(yaml.dump(self.doc, sort_keys=False))