feat: add NORTH_STAR.md and MVP core loop (Telegram bot + LLM + proposal generator)
This commit is contained in:
committed by
GitHub
parent
c071fe6bf7
commit
6862c7d0b8
@@ -0,0 +1 @@
|
||||
"""Steward – AI-assisted personal operations platform."""
|
||||
@@ -0,0 +1,166 @@
|
||||
"""Telegram bot interface for Steward."""
|
||||
|
||||
import logging
|
||||
from collections import defaultdict
|
||||
|
||||
from telegram import Update
|
||||
from telegram.constants import ParseMode
|
||||
from telegram.ext import (
|
||||
Application,
|
||||
CommandHandler,
|
||||
ContextTypes,
|
||||
MessageHandler,
|
||||
filters,
|
||||
)
|
||||
|
||||
from steward.config import Settings
|
||||
from steward.llm.client import LLMClient
|
||||
from steward.proposals.generator import Proposal, ProposalGenerator
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Per-user conversation history (in-memory for the MVP)
|
||||
_history: dict[int, list[dict[str, str]]] = defaultdict(list)
|
||||
_MAX_HISTORY = 20 # keep the last N turns per user
|
||||
|
||||
|
||||
def _is_allowed(user_id: int, settings: Settings) -> bool:
|
||||
"""Return True if the user is in the allow-list (or no list is configured)."""
|
||||
if not settings.telegram_allowed_user_ids:
|
||||
return True
|
||||
return user_id in settings.telegram_allowed_user_ids
|
||||
|
||||
|
||||
async def _send_long(update: Update, text: str) -> None:
|
||||
"""Send text, splitting if it exceeds Telegram's 4096-char limit."""
|
||||
limit = 4096
|
||||
for i in range(0, len(text), limit):
|
||||
await update.message.reply_text( # type: ignore[union-attr]
|
||||
text[i : i + limit],
|
||||
parse_mode=ParseMode.MARKDOWN,
|
||||
)
|
||||
|
||||
|
||||
async def start_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
|
||||
"""Handle /start."""
|
||||
settings: Settings = context.bot_data["settings"]
|
||||
user = update.effective_user
|
||||
if user is None or not _is_allowed(user.id, settings):
|
||||
return
|
||||
|
||||
await update.message.reply_text( # type: ignore[union-attr]
|
||||
"Hello, I'm *Steward* \U0001f916\n\n"
|
||||
"I'm your AI-assisted personal operations platform.\n"
|
||||
"Talk to me naturally, or use:\n"
|
||||
"/help – show available commands\n"
|
||||
"/clear – reset our conversation history\n"
|
||||
"/analyse – run a manual API analysis right now",
|
||||
parse_mode=ParseMode.MARKDOWN,
|
||||
)
|
||||
|
||||
|
||||
async def help_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
|
||||
"""Handle /help."""
|
||||
settings: Settings = context.bot_data["settings"]
|
||||
user = update.effective_user
|
||||
if user is None or not _is_allowed(user.id, settings):
|
||||
return
|
||||
|
||||
await update.message.reply_text( # type: ignore[union-attr]
|
||||
"*Steward commands*\n\n"
|
||||
"/start – greeting\n"
|
||||
"/help – this message\n"
|
||||
"/clear – reset conversation history\n"
|
||||
"/analyse – trigger an immediate API analysis and proposal",
|
||||
parse_mode=ParseMode.MARKDOWN,
|
||||
)
|
||||
|
||||
|
||||
async def clear_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
|
||||
"""Handle /clear – wipe conversation history for this user."""
|
||||
settings: Settings = context.bot_data["settings"]
|
||||
user = update.effective_user
|
||||
if user is None or not _is_allowed(user.id, settings):
|
||||
return
|
||||
|
||||
_history[user.id].clear()
|
||||
await update.message.reply_text("Conversation history cleared.") # type: ignore[union-attr]
|
||||
|
||||
|
||||
async def analyse_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
|
||||
"""Handle /analyse – run the proposal generator on demand."""
|
||||
settings: Settings = context.bot_data["settings"]
|
||||
llm: LLMClient = context.bot_data["llm"]
|
||||
user = update.effective_user
|
||||
if user is None or not _is_allowed(user.id, settings):
|
||||
return
|
||||
|
||||
await update.message.reply_text("Running analysis, please wait...") # type: ignore[union-attr]
|
||||
generator = ProposalGenerator(settings, llm)
|
||||
proposal = await generator.run()
|
||||
if proposal is None:
|
||||
await update.message.reply_text( # type: ignore[union-attr]
|
||||
"Analysis could not be completed. "
|
||||
"Check that `ANALYSIS_TARGET_URL` is configured."
|
||||
)
|
||||
return
|
||||
|
||||
await _send_long(update, proposal.format_for_telegram())
|
||||
|
||||
|
||||
async def message_handler(update: Update, context: ContextTypes.DEFAULT_TYPE) -> None:
|
||||
"""Handle plain text messages – forward to LLM and reply."""
|
||||
settings: Settings = context.bot_data["settings"]
|
||||
llm: LLMClient = context.bot_data["llm"]
|
||||
user = update.effective_user
|
||||
if user is None or not _is_allowed(user.id, settings):
|
||||
return
|
||||
|
||||
text = update.message.text # type: ignore[union-attr]
|
||||
if not text:
|
||||
return
|
||||
|
||||
history = _history[user.id]
|
||||
reply = await llm.chat(text, history=history)
|
||||
|
||||
# Update history
|
||||
history.append({"role": "user", "content": text})
|
||||
history.append({"role": "assistant", "content": reply})
|
||||
# Trim to keep only the most recent turns (2 messages per turn)
|
||||
if len(history) > _MAX_HISTORY * 2:
|
||||
_history[user.id] = history[-( _MAX_HISTORY * 2):]
|
||||
|
||||
await _send_long(update, reply)
|
||||
|
||||
|
||||
def build_application(settings: Settings, llm: LLMClient) -> Application: # type: ignore[type-arg]
|
||||
"""Build and return the Telegram Application."""
|
||||
app = (
|
||||
Application.builder()
|
||||
.token(settings.telegram_bot_token)
|
||||
.build()
|
||||
)
|
||||
app.bot_data["settings"] = settings
|
||||
app.bot_data["llm"] = llm
|
||||
|
||||
app.add_handler(CommandHandler("start", start_handler))
|
||||
app.add_handler(CommandHandler("help", help_handler))
|
||||
app.add_handler(CommandHandler("clear", clear_handler))
|
||||
app.add_handler(CommandHandler("analyse", analyse_handler))
|
||||
app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, message_handler))
|
||||
|
||||
return app
|
||||
|
||||
|
||||
async def send_proposal(app: Application, proposal: Proposal, user_ids: list[int]) -> None: # type: ignore[type-arg]
|
||||
"""Send a proposal message to all configured user IDs."""
|
||||
text = proposal.format_for_telegram()
|
||||
for uid in user_ids:
|
||||
try:
|
||||
await app.bot.send_message(
|
||||
chat_id=uid,
|
||||
text=text,
|
||||
parse_mode=ParseMode.MARKDOWN,
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("Failed to send proposal to user %d", uid)
|
||||
@@ -0,0 +1,35 @@
|
||||
"""Configuration management for Steward."""
|
||||
|
||||
from pydantic_settings import BaseSettings, SettingsConfigDict
|
||||
|
||||
|
||||
class Settings(BaseSettings):
|
||||
"""Application settings loaded from environment variables or .env file."""
|
||||
|
||||
model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8", extra="ignore")
|
||||
|
||||
# Telegram
|
||||
telegram_bot_token: str = ""
|
||||
telegram_allowed_user_ids: list[int] = []
|
||||
|
||||
# LLM (OpenAI-compatible)
|
||||
openai_api_key: str = ""
|
||||
openai_base_url: str = "https://api.openai.com/v1"
|
||||
openai_model: str = "gpt-4o"
|
||||
openai_system_prompt: str = (
|
||||
"You are Steward, a persistent, trustworthy AI-assisted personal operations platform. "
|
||||
"You reduce cognitive load by observing, remembering, planning, and proposing actions. "
|
||||
"You are conservative, transparent, and policy-aware. "
|
||||
"Always explain your reasoning."
|
||||
)
|
||||
|
||||
# Analysis / proposal worker
|
||||
analysis_target_url: str = ""
|
||||
analysis_target_api_key: str = ""
|
||||
analysis_cron_hour: int = 8
|
||||
analysis_cron_minute: int = 0
|
||||
|
||||
|
||||
def get_settings() -> Settings:
|
||||
"""Return application settings singleton."""
|
||||
return Settings()
|
||||
@@ -0,0 +1,70 @@
|
||||
"""LLM client wrapper (OpenAI-compatible)."""
|
||||
|
||||
import logging
|
||||
from collections.abc import AsyncIterator
|
||||
|
||||
from openai import AsyncOpenAI
|
||||
|
||||
from steward.config import Settings
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class LLMClient:
|
||||
"""Thin async wrapper around the OpenAI chat-completions API."""
|
||||
|
||||
def __init__(self, settings: Settings) -> None:
|
||||
self._settings = settings
|
||||
self._client = AsyncOpenAI(
|
||||
api_key=settings.openai_api_key,
|
||||
base_url=settings.openai_base_url,
|
||||
)
|
||||
|
||||
async def chat(
|
||||
self,
|
||||
user_message: str,
|
||||
*,
|
||||
history: list[dict[str, str]] | None = None,
|
||||
system_prompt: str | None = None,
|
||||
) -> str:
|
||||
"""Send a user message (with optional history) and return the assistant reply."""
|
||||
system = system_prompt or self._settings.openai_system_prompt
|
||||
messages: list[dict[str, str]] = [{"role": "system", "content": system}]
|
||||
if history:
|
||||
messages.extend(history)
|
||||
messages.append({"role": "user", "content": user_message})
|
||||
|
||||
logger.debug(
|
||||
"LLM request: model=%s messages=%d", self._settings.openai_model, len(messages)
|
||||
)
|
||||
response = await self._client.chat.completions.create(
|
||||
model=self._settings.openai_model,
|
||||
messages=messages, # type: ignore[arg-type]
|
||||
)
|
||||
reply = response.choices[0].message.content or ""
|
||||
logger.debug("LLM reply: %d chars", len(reply))
|
||||
return reply
|
||||
|
||||
async def stream(
|
||||
self,
|
||||
user_message: str,
|
||||
*,
|
||||
history: list[dict[str, str]] | None = None,
|
||||
system_prompt: str | None = None,
|
||||
) -> AsyncIterator[str]:
|
||||
"""Stream the assistant reply token by token."""
|
||||
system = system_prompt or self._settings.openai_system_prompt
|
||||
messages: list[dict[str, str]] = [{"role": "system", "content": system}]
|
||||
if history:
|
||||
messages.extend(history)
|
||||
messages.append({"role": "user", "content": user_message})
|
||||
|
||||
stream = await self._client.chat.completions.create(
|
||||
model=self._settings.openai_model,
|
||||
messages=messages, # type: ignore[arg-type]
|
||||
stream=True,
|
||||
)
|
||||
async for chunk in stream:
|
||||
delta = chunk.choices[0].delta.content
|
||||
if delta:
|
||||
yield delta
|
||||
@@ -0,0 +1,69 @@
|
||||
"""Steward application entry point."""
|
||||
|
||||
import logging
|
||||
import sys
|
||||
|
||||
from apscheduler.schedulers.asyncio import AsyncIOScheduler
|
||||
|
||||
from steward.bot.telegram import build_application, send_proposal
|
||||
from steward.config import get_settings
|
||||
from steward.llm.client import LLMClient
|
||||
from steward.proposals.generator import ProposalGenerator
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
format="%(asctime)s %(levelname)s %(name)s: %(message)s",
|
||||
)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
async def _run_scheduled_analysis(
|
||||
generator: ProposalGenerator,
|
||||
app, # type: ignore[type-arg]
|
||||
user_ids: list[int],
|
||||
) -> None:
|
||||
"""Scheduled job: run analysis and push proposal via Telegram."""
|
||||
logger.info("Running scheduled analysis")
|
||||
proposal = await generator.run()
|
||||
if proposal:
|
||||
await send_proposal(app, proposal, user_ids)
|
||||
else:
|
||||
logger.info("No proposal generated (generator returned None)")
|
||||
|
||||
|
||||
def main() -> None:
|
||||
"""Start Steward."""
|
||||
settings = get_settings()
|
||||
|
||||
if not settings.telegram_bot_token:
|
||||
logger.error("TELEGRAM_BOT_TOKEN is not set – cannot start")
|
||||
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)
|
||||
app = build_application(settings, llm)
|
||||
generator = ProposalGenerator(settings, llm)
|
||||
|
||||
scheduler = AsyncIOScheduler()
|
||||
scheduler.add_job(
|
||||
_run_scheduled_analysis,
|
||||
"cron",
|
||||
hour=settings.analysis_cron_hour,
|
||||
minute=settings.analysis_cron_minute,
|
||||
args=[generator, app, settings.telegram_allowed_user_ids],
|
||||
)
|
||||
scheduler.start()
|
||||
|
||||
logger.info(
|
||||
"Steward starting: model=%s analysis_url=%s",
|
||||
settings.openai_model,
|
||||
settings.analysis_target_url or "(none)",
|
||||
)
|
||||
app.run_polling(allowed_updates=["message"])
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,106 @@
|
||||
"""Proposal data model and generator.
|
||||
|
||||
This module implements a minimal "generate proposal" workflow:
|
||||
|
||||
1. Fetch data from a target API endpoint.
|
||||
2. Ask the LLM to analyse the data and produce a structured proposal.
|
||||
3. Return the proposal for delivery (e.g. via Telegram).
|
||||
|
||||
The Proposal dataclass is intentionally simple for the MVP – it captures the
|
||||
fields described in the manifesto (why, confidence, evidence, expected outcome,
|
||||
rollback strategy) without any persistence layer yet.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import UTC, datetime
|
||||
|
||||
import httpx
|
||||
|
||||
from steward.config import Settings
|
||||
from steward.llm.client import LLMClient
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
_ANALYSIS_SYSTEM_PROMPT = (
|
||||
"You are Steward, an AI operations platform. "
|
||||
"You have been given raw data from an internal API. "
|
||||
"Analyse the data and produce a concise proposal in the following format:\n\n"
|
||||
"**Summary:** <one-sentence summary>\n"
|
||||
"**Why:** <reason this proposal matters>\n"
|
||||
"**Evidence:** <key data points from the API response>\n"
|
||||
"**Expected outcome:** <what will improve if adopted>\n"
|
||||
"**Rollback strategy:** <how to undo if things go wrong>\n"
|
||||
"**Confidence:** <Low | Medium | High> - <brief justification>\n\n"
|
||||
"Be conservative. If the data is healthy and no action is needed, say so explicitly."
|
||||
)
|
||||
|
||||
|
||||
@dataclass
|
||||
class Proposal:
|
||||
"""A structured action proposal generated by Steward."""
|
||||
|
||||
title: str
|
||||
body: str
|
||||
source_url: str
|
||||
generated_at: datetime = field(default_factory=lambda: datetime.now(UTC))
|
||||
raw_data: str = ""
|
||||
|
||||
def format_for_telegram(self) -> str:
|
||||
"""Return a Markdown-formatted string suitable for a Telegram message."""
|
||||
ts = self.generated_at.strftime("%Y-%m-%d %H:%M UTC")
|
||||
return (
|
||||
f"\U0001f50d *Steward Proposal*\n"
|
||||
f"_{ts}_\n\n"
|
||||
f"*Source:* `{self.source_url}`\n\n"
|
||||
f"{self.body}"
|
||||
)
|
||||
|
||||
|
||||
class ProposalGenerator:
|
||||
"""Fetches data from a target API and generates a proposal via the LLM."""
|
||||
|
||||
def __init__(self, settings: Settings, llm: LLMClient) -> None:
|
||||
self._settings = settings
|
||||
self._llm = llm
|
||||
|
||||
async def run(self) -> Proposal | None:
|
||||
"""Fetch the target API and return a Proposal, or None on error."""
|
||||
url = self._settings.analysis_target_url
|
||||
if not url:
|
||||
logger.warning("analysis_target_url is not configured - skipping proposal generation")
|
||||
return None
|
||||
|
||||
raw = await self._fetch(url)
|
||||
if raw is None:
|
||||
return None
|
||||
|
||||
body = await self._llm.chat(
|
||||
f"Here is the API response from {url}:\n\n{raw}",
|
||||
system_prompt=_ANALYSIS_SYSTEM_PROMPT,
|
||||
)
|
||||
|
||||
return Proposal(
|
||||
title="Daily Analysis",
|
||||
body=body,
|
||||
source_url=url,
|
||||
raw_data=raw,
|
||||
)
|
||||
|
||||
async def _fetch(self, url: str) -> str | None:
|
||||
"""Fetch the URL and return the response body as text."""
|
||||
headers: dict[str, str] = {}
|
||||
api_key = self._settings.analysis_target_api_key
|
||||
if api_key:
|
||||
headers["Authorization"] = "Bearer " + api_key
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
response = await client.get(url, headers=headers)
|
||||
response.raise_for_status()
|
||||
return response.text
|
||||
except httpx.HTTPError as exc:
|
||||
logger.error("Failed to fetch %s: %s", url, exc)
|
||||
return None
|
||||
Reference in New Issue
Block a user