import asyncio
import json
from urllib.request import Request, urlopen
from cogos import CogOS, CogOSConfig, UniversalSchemaDomain
from cogos.prompts import default_cm_prompt, DEFAULT_CHATBOT_PROMPT
class ExternalRAGPlugin:
name = "external_rag"
def __init__(self, endpoint: str):
self.endpoint = endpoint.rstrip("/")
async def chat_prompt_vars(self, message, *, session_id, state, registry, recalled):
payload = json.dumps({"query": message, "session_id": session_id}).encode()
req = Request(
url=f"{self.endpoint}/retrieve",
data=payload,
headers={"Content-Type": "application/json"},
method="POST",
)
def _do():
with urlopen(req, timeout=15) as resp:
data = json.loads(resp.read().decode())
return (data.get("context") or "").strip()
ctx = await asyncio.to_thread(_do)
if not ctx:
return {}
return {"rag_section": f"## Retrieved Context (RAG)\n{ctx}\n"}
config = CogOSConfig.from_file()
schema = UniversalSchemaDomain("user_profile")
schema.create_field("identity.name", "", "User's name")
cogos = (
CogOS(config)
.register(schema)
.cm_prompt(default_cm_prompt)
.chatbot_prompt(DEFAULT_CHATBOT_PROMPT) # supports {rag_section}
.use(ExternalRAGPlugin("http://localhost:9000"))
)