Skip to content

Messaging Module

Package: cogniverse_messaging (Application Layer) Location: libs/messaging/cogniverse_messaging/ Entry point: python -m cogniverse_messaging.gateway


Table of Contents

  1. Overview
  2. Package Structure
  3. MessagingGateway
  4. Command Routing
  5. Authentication
  6. Conversation History
  7. RuntimeClient
  8. Configuration
  9. Testing
  10. Architecture Position

Overview

The Messaging module runs a standalone gateway service that bridges Telegram to the Cogniverse runtime. It translates Telegram updates into runtime agent-dispatch calls, formats agent responses back into Telegram messages, and manages user registration (invite tokens). Agent dispatch, registration, tenant resolution, and per-chat conversation history all go through the runtime's HTTP API — the auth primitives (InviteTokenManager, UserTenantMapper) live in cogniverse_core.messaging_auth, and the runtime loads/saves conversation history itself around each agent call, so the gateway holds no backend connection.

Key responsibilities:

  • Telegram integration — polling (dev) or webhook (production) mode via python-telegram-bot
  • Command routing — maps /search, /summarize, /report, /research, /code, /wiki, /instructions, /memories, /jobs to runtime agent names and endpoints
  • User registration — invite-token based onboarding, mapping a Telegram user ID to a tenant ID
  • Conversation history — stores and retrieves per-chat turns via Mem0 so agents get context across messages
  • Runtime protocol adapter — a thin async HTTP client (RuntimeClient) wrapping the runtime's /agents/*/process, /wiki/*, and /admin/tenant/* endpoints

Package Structure

graph TD
    Root["<span style='color:#000'><b>cogniverse_messaging/</b></span>"]

    Root --> Gateway["<span style='color:#000'><b>gateway.py</b><br/>MessagingGateway, main() entry point</span>"]
    Root --> CommandRouter["<span style='color:#000'>command_router.py<br/>parse_message(), ParsedCommand</span>"]
    Root --> Auth["<span style='color:#000'>cogniverse_core.messaging_auth<br/>InviteTokenManager, UserTenantMapper</span>"]
    Root --> RuntimeClient["<span style='color:#000'>runtime_client.py<br/>RuntimeClient (async HTTP)</span>"]
    Root --> TelegramHandler["<span style='color:#000'>telegram_handler.py<br/>Response formatting, message chunking</span>"]

    Gateway --> CommandRouter
    Gateway --> Auth
    Gateway --> RuntimeClient
    Gateway --> TelegramHandler

    style Root fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Gateway fill:#ffcc80,stroke:#ef6c00,color:#000
    style CommandRouter fill:#81d4fa,stroke:#0288d1,color:#000
    style Auth fill:#81d4fa,stroke:#0288d1,color:#000
    style RuntimeClient fill:#81d4fa,stroke:#0288d1,color:#000
    style TelegramHandler fill:#81d4fa,stroke:#0288d1,color:#000

All modules are flat files directly under cogniverse_messaging/ (no subpackages).


MessagingGateway

Location: libs/messaging/cogniverse_messaging/gateway.py

MessagingGateway(
    bot_token: str,
    runtime_url: str,
    mode: str = "polling",       # "polling" (dev) or "webhook" (production)
    webhook_url: str = "",
    webhook_listen: str = "0.0.0.0",  # webhook mode: HTTP server bind address
    webhook_port: int = 8443,         # webhook mode: HTTP server bind port
    webhook_path: str = "",           # webhook mode: URL path Telegram POSTs updates to
    outbound_poll_seconds: float = 5.0,
)

Registration and tenant resolution go through the runtime's HTTP API: /start <token> calls POST /admin/messaging/register and each message resolves the sender via GET /admin/messaging/resolve (cached positively per user for 60s). The gateway therefore registers users out of the box in the deployed chart — it needs only RUNTIME_URL. A runtime/backend outage during either call replies "temporarily unavailable" (a failed registration keeps the token claimed for that user, who can retry; no other user can redeem it); tenant_id: null from resolve is the only thing that reads as "please register". Conversation history remains optional: without memory_manager the gateway skips history lookup/storage.

build_app() registers a CommandHandler for every slash command (start, help, and all nine command families) — the plain-text MessageHandler filters with ~filters.COMMAND, so a command without its own handler would be silently dropped by Telegram dispatch. Updates process concurrently (up to 32 in parallel) so one slow agent dispatch cannot stall other chats, and a registered error handler replies "Something went wrong handling that — please try again." when a handler raises (runtime unreachable, malformed response) instead of leaving the user in silence.

Usage:

from cogniverse_messaging.gateway import MessagingGateway

gateway = MessagingGateway(
    bot_token="123456:ABC-token",
    runtime_url="http://localhost:28000",
    mode="polling",
)
await gateway.run()  # dispatches to run_polling() or run_webhook() based on mode

Running the module directly (python -m cogniverse_messaging.gateway) reads TELEGRAM_BOT_TOKEN (required), RUNTIME_URL (default http://localhost:28000), GATEWAY_MODE (default polling), TELEGRAM_WEBHOOK_URL (required when GATEWAY_MODE=webhook), GATEWAY_WEBHOOK_LISTEN (default 0.0.0.0), GATEWAY_WEBHOOK_PORT (default 8443), GATEWAY_WEBHOOK_PATH (default ""), and GATEWAY_OUTBOUND_POLL_SECONDS (default 5) from the environment.

Outbound delivery

Alongside inbound handling, the gateway runs a background _outbound_drain_loop — started in both run_polling and run_webhook, cancelled in their shutdown finally. Every GATEWAY_OUTBOUND_POLL_SECONDS (default 5) it calls RuntimeClient.drain_outbound() (GET /admin/messaging/outbound/drain) and sends each returned {chat_id, text} via the bot. A drain failure (runtime blip) is logged and retried next tick. Because the runtime clears a message on the drain that returns it, a failed send keeps the message in the gateway's in-memory retry buffer and re-attempts it on later ticks — after OUTBOUND_SEND_MAX_ATTEMPTS (3) failed sends it is dropped with an error, so one dead chat never stops the others or grows the buffer without bound. A message missing chat_id/text is dropped immediately (it can never send). This is the delivery side of the runtime's POST /messaging/send — the path job-completion notifications reach a tenant's linked chats.

Known limitation — the retry buffer is in-memory only. self._outbound_retry is a plain process-local list, never written to disk or a store. If the gateway process restarts (deploy, crash, OOM, pod reschedule) while a message is between a failed send and its next retry tick, the message is lost: the runtime's durable outbound:pending queue already forgot it the moment drain_outbound() returned it, and the retry buffer holding it disappears with the process. The more common restart case — a process restart before a message is ever drained — is unaffected, since Redis still has it and the next gateway instance drains it normally; the gap is specifically the narrow window after a drain and before the retry either succeeds or exhausts its OUTBOUND_SEND_MAX_ATTEMPTS. The messages on this path are routine Telegram job-completion notifications, not approval gates or alerts anything blocks on, and their underlying content (job result, wiki save) is already persisted/logged independently of the Telegram push, so a drop here loses a convenience notification, not data. run()'s SIGTERM/shutdown finally calls _log_dropped_outbound_retry(), which logs and clears any non-empty retry buffer so a planned restart at least leaves an error log line per dropped message — it cannot resend at that point (the bot and runtime client are already torn down earlier in the same shutdown), and it does not run at all on an unclean kill (OOM, SIGKILL), same as any other Python shutdown hook.


Command Routing

Location: libs/messaging/cogniverse_messaging/command_router.py

parse_message(text=None, has_photo=False, has_video=False, photo_file_id=None, video_file_id=None) -> ParsedCommand classifies an incoming Telegram message. Agent slash commands map directly to runtime agent names:

Command Agent
/search <query> search_agent
/summarize <query> summarizer_agent
/report <query> detailed_report_agent
/research <query> deep_research_agent
/code <query> coding_agent

/wiki, /instructions, /memories, and /jobs are parsed into their own ParsedCommand fields (is_wiki/wiki_subcommand, etc.) and dispatched by MessagingGateway._handle_*_command to the matching RuntimeClient method. A message with no recognized command and no media falls through to gateway_agent.

Media messages set has_media=True and route by type:

Media Agent Query Why
Photo image_search_agent caption, else "Find visually similar images" The photo itself is the query — its bytes are embedded and matched against stored image embeddings
Video search_agent caption, else "Find similar video content" A video file has no single query embedding, so it stays caption-text search over the video index

For a photo, MessagingGateway._handle_message downloads the file through the bot API (context.bot.get_file(file_id) → download_as_bytearray()) and base64-encodes it into the dispatch context as media_content_b64 (with media_mime), alongside media_type / media_file_id. The bot token lives only in the gateway, so the runtime cannot fetch from Telegram itself — the bytes have to travel with the request. MessagingGateway._download_photo_b64 refuses anything over MAX_PHOTO_BYTES (5 MB), checking both the size Telegram advertises and the bytes actually received, and replies with a plain message rather than dispatching when the download fails.


Authentication

Location: libs/core/cogniverse_core/messaging_auth.py (shared with the runtime, which serves registration routes over HTTP)

  • InviteTokenManager(config_manager) — generates, claims and consumes invite tokens, stored in the _system tenant's config store (ConfigScope.SYSTEM, service "messaging_gateway"). generate_token(tenant_id, expires_in_hours=24) returns a UUID hex token. claim_token(token, platform, external_user_id) binds an unused, unexpired token to one external user with a compare-and-set on its record and returns the tenant_id; it returns the tenant_id again for the user already holding it and None for an unknown, expired or used token or one another user holds, and raises on a store outage. mark_token_used(token, platform, external_user_id) consumes a token the user holds; it returns False on a failed consume write (logged; the token stays bound to that user). validate_token(token) returns the tenant_id of a token nobody has claimed or used. The runtime's POST /admin/messaging/register route drives claim → register → consume in that order: of any number of concurrent registers of one token on any process or replica exactly one user claims it and every other user gets 404, a failed registration keeps the token for the user who claimed it, and that user's retry resumes; the gateway only speaks HTTP.
  • UserTenantMapper(memory_manager) — maps a Telegram user ID to a tenant ID via Mem0, storing the mapping under the system tenant partition (SYSTEM_TENANT_ID) with agent_name="_messaging_gateway" and infer=False so the raw mapping text isn't rewritten by LLM extraction.

Conversation History

Conversation history is server-side: the gateway sends only context_id (the Telegram chat id) with each dispatch, and the runtime's agent dispatcher loads that context's recent turns before the agent runs and appends the two new turns after (cogniverse_core.conversation.ConversationStore, keyed by (tenant_id, context_id)). The gateway therefore holds no Mem0 connection and multi-turn memory works in the deployed chart. History is enrichment: a Mem0 outage degrades to no-history (the agent still answers) rather than failing the reply. The two new turns persist after the reply is sent, so the write's cost never lands on the user's latency; their order and the saves still landing are kept in the runtime's Redis, so the chat's next message reads them whichever runtime worker or replica receives it, and a Redis outage answers the dispatch with 503 rather than a reply that would never be stored. A save the runtime could not complete is reported through the dispatcher's conversation_persist_status(). See Core → ConversationStore and the dispatcher.


RuntimeClient

Location: libs/messaging/cogniverse_messaging/runtime_client.py

Thin async httpx wrapper around the runtime's HTTP API — the gateway never imports agent or core code directly.

RuntimeClient(runtime_url: str, timeout: float = 300.0, dispatch_timeout: float = 120.0)

timeout bounds SSE event-stream reads; dispatch_timeout is the interactive-chat read budget for dispatch_agent (a hung runtime must not hold a gateway worker for the full stream timeout). Every other call (CRUD, drain, health) uses the shared client's 30s read default, and all connects fail within 5s so an unreachable runtime surfaces in seconds. dispatch_agent degrades a dead / hung / non-JSON runtime to a status dict ({"status": "unavailable"|"warming"|"error"}) rather than raising — the gateway renders a dead runtime as "Service temporarily unavailable" and a timed-out cold start as "Runtime is warming up, try again shortly.", matching how register_user / resolve_tenant already behave.

Key methods:

Method Endpoint
health() GET /health (returns True/False, never raises)
dispatch_agent(agent_name, query, tenant_id, context_id=None, conversation_history=None, top_k=10, context=None) POST /agents/{agent_name}/process
stream_events(task_id) GET /events/workflows/{task_id} (SSE)
create_invite_token(tenant_id, expires_in_hours=24) POST /admin/messaging/invite
register_user(platform, external_user_id, token) POST /admin/messaging/register — returns {"status": "registered"\|"invalid_token"\|"unavailable", ...}, never raises
resolve_tenant(platform, external_user_id) GET /admin/messaging/resolve — {"status": "ok", "tenant_id": str\|None} or {"status": "unavailable"}
save_wiki_session(tenant_id, query, response, agent_name="gateway_agent", entities=None) POST /wiki/save
search_wiki(tenant_id, query, top_k=5) POST /wiki/search
get_wiki_topic(tenant_id, slug) GET /wiki/topic/{slug}
get_wiki_index(tenant_id) GET /wiki/index
lint_wiki(tenant_id) GET /wiki/lint
delete_wiki_topic(tenant_id, slug) DELETE /wiki/topic/{slug}
set_instructions(tenant_id, text) PUT /admin/tenant/{tenant}/instructions
get_instructions(tenant_id) GET /admin/tenant/{tenant}/instructions
list_memories(tenant_id, agent_name=None) GET /admin/tenant/{tenant}/memories
clear_memories(tenant_id, agent_name=None) DELETE /admin/tenant/{tenant}/memories
list_jobs(tenant_id) GET /admin/tenant/{tenant}/jobs
create_job(tenant_id, name, schedule, query, post_actions=None) POST /admin/tenant/{tenant}/jobs
delete_job(tenant_id, job_id) DELETE /admin/tenant/{tenant}/jobs/{job_id}
close() closes the underlying httpx.AsyncClient

save_wiki_session is defined but not currently called by the gateway — /wiki save replies that sessions auto-save in the background instead of invoking it (see Command Routing). Non-2xx responses from the CRUD methods are normalized to {"status": "error", "status_code": ..., "message": ...} by _json_or_error, so callers never need to branch on httpx exceptions; dispatch_agent and create_invite_token normalize errors inline instead.


Configuration

Deployed via the messaging section of the Helm chart (charts/cogniverse/values.yaml), disabled by default:

messaging:
  enabled: false
  mode: polling   # polling for dev, webhook for production
  replicaCount: 1

Enable it at deploy time with cogniverse up --messaging (requires TELEGRAM_BOT_TOKEN in the environment).

Environment Variable Purpose
TELEGRAM_BOT_TOKEN Required. Telegram bot token.
RUNTIME_URL Runtime base URL (default http://localhost:28000).
GATEWAY_MODE polling or webhook (default polling).
TELEGRAM_WEBHOOK_URL Required when GATEWAY_MODE=webhook.
GATEWAY_WEBHOOK_LISTEN Webhook server bind address (default 0.0.0.0).
GATEWAY_WEBHOOK_PORT Webhook server bind port (default 8443).
GATEWAY_WEBHOOK_PATH URL path Telegram POSTs updates to (default "").
GATEWAY_OUTBOUND_POLL_SECONDS Seconds between outbound-queue drains for delivery (default 5).

Testing

uv run pytest tests/messaging/unit/ -v --tb=long

Covers command parsing, invite-token auth, gateway command dispatch, and the RuntimeClient CRUD wrappers. tests/messaging/integration/test_gateway_webhook_serves.py verifies run_webhook() actually binds and serves an HTTP server that Telegram can POST updates to, not just registers the webhook URL. Round-trip coverage against a real runtime lives in tests/runtime/integration/test_inbound_messaging_primitive.py and tests/runtime/integration/test_inbound_messaging_redis.py; end-to-end coverage lives in tests/e2e/test_messaging_e2e.py and tests/e2e/test_messaging_gateway_e2e.py.

Photo-download coverage has two boundary levels. The regular integration test uses a real python-telegram-bot client against a local Bot API stub. test_gateway_photo_download_real_telegram.py is an opt-in local_only test that calls Telegram's real getUpdates, getFile, and file-download APIs, then verifies the exact downloaded bytes are forwarded to image_search_agent. Set TELEGRAM_BOT_TOKEN through the normal secret resolution path, send that bot a photo, and optionally set TELEGRAM_CHAT_ID to select a chat. It skips with an actionable reason when credentials, a pending photo, or getUpdates availability are absent; Telegram disables getUpdates while a webhook is configured. Never store either value in a tracked file.


Architecture Position

flowchart TB
    subgraph AppLayer["<span style='color:#000'>Application Layer</span>"]
        Messaging["<span style='color:#000'>cogniverse-messaging ◄─ YOU ARE HERE<br/>Telegram gateway</span>"]
        Runtime["<span style='color:#000'>cogniverse-runtime</span>"]
    end

    Telegram(("<span style='color:#000'>Telegram</span>")) --> Messaging
    Messaging -->|HTTP| Runtime

    style AppLayer fill:#90caf9,stroke:#1565c0,color:#000
    style Messaging fill:#64b5f6,stroke:#1565c0,color:#000
    style Runtime fill:#64b5f6,stroke:#1565c0,color:#000

cogniverse-messaging is not imported by any other libs/* package, and is not declared as a workspace dependency of any package — it talks to the runtime over HTTP and only reaches into cogniverse_core/cogniverse_sdk directly for invite-token storage in auth.py.

Dependencies: python-telegram-bot[webhooks], httpx (declared); cogniverse_core, cogniverse_sdk (imported directly, not declared in pyproject.toml)

Dependents: none (standalone service)