Introduce a shared ConversationService (steward/bot/core.py) that owns the LLM call, history, knowledge-base search, and thread-memory keying behind a normalized ThreadKey, so both Telegram and Matrix drive the same pipeline. - Add steward/bot/matrix.py: a mautrix-python appservice bot that receives Synapse transactions and replies via the client-server API. - Refactor telegram.py handlers into thin wrappers over ConversationService. - Generalize ThreadMemoryStore/ThreadSummary to platform-scoped keys with legacy chat_id:thread_id migration. - Add a matrix config section (homeserver, tokens, room/user allowlists). - Rewrite main.py as async, starting Telegram and/or Matrix on one event loop. - Add mautrix>=0.21.0 dependency. Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
144 lines
4.7 KiB
Python
144 lines
4.7 KiB
Python
"""Steward application entry point."""
|
||
|
||
import asyncio
|
||
import logging
|
||
import signal
|
||
import sys
|
||
from datetime import time
|
||
from typing import Any
|
||
|
||
from telegram.ext import Application, ContextTypes
|
||
|
||
from steward.bot.core import ConversationService
|
||
from steward.bot.matrix import StewardMatrixBot
|
||
from steward.bot.telegram import build_application, send_proposal
|
||
from steward.config import get_settings
|
||
from steward.llm.client import LLMClient
|
||
from steward.memory.thread_store import ThreadMemoryStore
|
||
from steward.proposals.generator import ProposalGenerator
|
||
from steward.tools.client import ToolClient
|
||
|
||
logging.basicConfig(
|
||
level=logging.INFO,
|
||
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
|
||
)
|
||
logging.getLogger("httpx").setLevel(logging.WARNING)
|
||
logging.getLogger("httpcore").setLevel(logging.WARNING)
|
||
logger = logging.getLogger(__name__)
|
||
|
||
|
||
async def _run_scheduled_analysis(context: ContextTypes.DEFAULT_TYPE) -> None:
|
||
"""Scheduled job: run analysis and push proposal via Telegram."""
|
||
logger.info("Running scheduled analysis")
|
||
generator: ProposalGenerator = context.bot_data["generator"]
|
||
user_ids: list[int] = context.bot_data["user_ids"]
|
||
proposal = await generator.run()
|
||
if proposal:
|
||
await send_proposal(context.application, proposal, user_ids)
|
||
else:
|
||
logger.info("No proposal generated (generator returned None)")
|
||
|
||
|
||
async def _run_telegram(app: Application, shutdown_event: asyncio.Event) -> None: # type: ignore[type-arg]
|
||
"""Start the Telegram bot and wait for shutdown."""
|
||
await app.initialize()
|
||
await app.start()
|
||
if app.updater is not None:
|
||
await app.updater.start_polling(allowed_updates=["message"])
|
||
logger.info("Telegram bot started")
|
||
try:
|
||
await shutdown_event.wait()
|
||
finally:
|
||
if app.updater is not None:
|
||
await app.updater.stop()
|
||
await app.stop()
|
||
await app.shutdown()
|
||
|
||
|
||
async def _run_matrix(bot: StewardMatrixBot, shutdown_event: asyncio.Event) -> None:
|
||
"""Start the Matrix appservice and wait for shutdown."""
|
||
await bot.run()
|
||
try:
|
||
await shutdown_event.wait()
|
||
finally:
|
||
await bot.stop()
|
||
|
||
|
||
async def main() -> None:
|
||
"""Start Steward (Telegram and/or Matrix)."""
|
||
settings = get_settings()
|
||
|
||
if not settings.telegram_bot_token and not settings.matrix_enabled:
|
||
logger.error("No platform configured – set TELEGRAM_BOT_TOKEN or Matrix appservice config")
|
||
sys.exit(1)
|
||
|
||
if not settings.openai_api_key:
|
||
logger.error("OPENAI_API_KEY is not set – cannot start")
|
||
sys.exit(1)
|
||
|
||
llm = LLMClient(settings)
|
||
thread_store = ThreadMemoryStore(settings.thread_memory_path)
|
||
|
||
tool_client: ToolClient | None = None
|
||
if settings.mcp_server_url:
|
||
tool_client = ToolClient(
|
||
base_url=settings.mcp_server_url,
|
||
api_key=settings.mcp_server_api_key,
|
||
)
|
||
logger.info("MCP tool server configured: %s", settings.mcp_server_url)
|
||
else:
|
||
logger.info("No MCP_SERVER_URL configured – tool calling disabled")
|
||
|
||
service = ConversationService(settings, llm, thread_store, tool_client)
|
||
|
||
shutdown_event = asyncio.Event()
|
||
loop = asyncio.get_running_loop()
|
||
for sig in (signal.SIGINT, signal.SIGTERM):
|
||
loop.add_signal_handler(sig, shutdown_event.set)
|
||
|
||
tasks: list[asyncio.Task[Any]] = []
|
||
|
||
if settings.telegram_bot_token:
|
||
app = build_application(settings, service)
|
||
generator = ProposalGenerator(settings, llm)
|
||
app.bot_data["generator"] = generator
|
||
app.bot_data["user_ids"] = settings.telegram_allowed_user_ids
|
||
|
||
if app.job_queue is not None:
|
||
app.job_queue.run_daily(
|
||
_run_scheduled_analysis,
|
||
time=time(hour=settings.analysis_cron_hour, minute=settings.analysis_cron_minute),
|
||
)
|
||
else:
|
||
logger.warning("JobQueue not available – scheduled analysis disabled")
|
||
|
||
logger.info(
|
||
"Steward starting: model=%s analysis_url=%s tools=%s",
|
||
settings.openai_model,
|
||
settings.analysis_target_url or "(none)",
|
||
settings.mcp_server_url or "(none)",
|
||
)
|
||
tasks.append(asyncio.create_task(_run_telegram(app, shutdown_event)))
|
||
|
||
if settings.matrix_enabled:
|
||
matrix_bot = StewardMatrixBot(settings, service)
|
||
tasks.append(asyncio.create_task(_run_matrix(matrix_bot, shutdown_event)))
|
||
|
||
if not tasks:
|
||
logger.error("No platform started")
|
||
sys.exit(1)
|
||
|
||
try:
|
||
await asyncio.gather(*tasks)
|
||
except asyncio.CancelledError:
|
||
pass
|
||
|
||
|
||
if __name__ == "__main__":
|
||
asyncio.run(main())
|
||
|
||
|
||
def run() -> None:
|
||
"""Synchronous entry point for the ``steward`` console script."""
|
||
asyncio.run(main())
|