From 6862c7d0b85ffcdca2fa5221e155c036a0ad9ba0 Mon Sep 17 00:00:00 2001 From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com> Date: Sat, 25 Jul 2026 11:47:52 +0000 Subject: [PATCH] feat: add NORTH_STAR.md and MVP core loop (Telegram bot + LLM + proposal generator) --- .env.example | 23 ++ .gitignore | 37 +++ NORTH_STAR.md | 516 +++++++++++++++++++++++++++++++++ pyproject.toml | 52 ++++ steward/__init__.py | 1 + steward/bot/__init__.py | 0 steward/bot/telegram.py | 166 +++++++++++ steward/config.py | 35 +++ steward/llm/__init__.py | 0 steward/llm/client.py | 70 +++++ steward/main.py | 69 +++++ steward/proposals/__init__.py | 0 steward/proposals/generator.py | 106 +++++++ tests/__init__.py | 0 tests/test_bot.py | 119 ++++++++ tests/test_llm_client.py | 92 ++++++ tests/test_proposals.py | 106 +++++++ 17 files changed, 1392 insertions(+) create mode 100644 .env.example create mode 100644 .gitignore create mode 100644 NORTH_STAR.md create mode 100644 pyproject.toml create mode 100644 steward/__init__.py create mode 100644 steward/bot/__init__.py create mode 100644 steward/bot/telegram.py create mode 100644 steward/config.py create mode 100644 steward/llm/__init__.py create mode 100644 steward/llm/client.py create mode 100644 steward/main.py create mode 100644 steward/proposals/__init__.py create mode 100644 steward/proposals/generator.py create mode 100644 tests/__init__.py create mode 100644 tests/test_bot.py create mode 100644 tests/test_llm_client.py create mode 100644 tests/test_proposals.py diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..ff0520c --- /dev/null +++ b/.env.example @@ -0,0 +1,23 @@ +# Copy this file to .env and fill in your values. +# Never commit .env to version control. + +# Telegram bot token from @BotFather +TELEGRAM_BOT_TOKEN= + +# Comma-separated Telegram user IDs allowed to talk to Steward. +# Leave empty to allow everyone (not recommended for production). +TELEGRAM_ALLOWED_USER_IDS= + +# OpenAI (or compatible) credentials +OPENAI_API_KEY= +OPENAI_BASE_URL=https://api.openai.com/v1 +OPENAI_MODEL=gpt-4o + +# Optional: override the default system prompt +# OPENAI_SYSTEM_PROMPT= + +# Daily analysis target (optional) +# ANALYSIS_TARGET_URL=https://your-internal-api/endpoint +# ANALYSIS_TARGET_API_KEY= +# ANALYSIS_CRON_HOUR=8 +# ANALYSIS_CRON_MINUTE=0 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..2c88c73 --- /dev/null +++ b/.gitignore @@ -0,0 +1,37 @@ +# Python +__pycache__/ +*.py[cod] +*.pyo +*.pyd +*.so +*.egg +*.egg-info/ +dist/ +build/ +.eggs/ + +# Virtual environments +.venv/ +venv/ +env/ + +# Environment +.env + +# Pytest +.pytest_cache/ +.coverage +htmlcov/ + +# Mypy +.mypy_cache/ + +# Ruff +.ruff_cache/ + +# IDE +.vscode/ +.idea/ + +# OS +.DS_Store diff --git a/NORTH_STAR.md b/NORTH_STAR.md new file mode 100644 index 0000000..44c6c5e --- /dev/null +++ b/NORTH_STAR.md @@ -0,0 +1,516 @@ +# Steward +## Vision and Intent + +Version: 0.1 (Concept) +Status: Research +Author: Daniel Wagner + +--- + +# Executive Summary + +Steward is a long-running, AI-assisted personal operations platform designed to reduce cognitive load by acting as a persistent, trustworthy steward of both digital infrastructure and delegated personal objectives. + +Unlike traditional AI assistants, Steward is not intended to answer prompts in isolation. Instead, it maintains an ongoing understanding of its environment, remembers historical context, continuously evaluates events, and takes appropriate action within explicitly delegated authority. + +The primary research question underpinning Steward is: + +> Can an AI system safely become more useful over time without becoming less predictable? + +Steward seeks to answer this through carefully constrained autonomy, policy-driven execution and continuous learning from operational history. + +--- + +# Philosophy + +Steward is not *just* an AI chatbot. + +Steward is not an autonomous root user. + +Steward is not *just* an automation engine. + +Steward is a persistent operations platform. + +Its purpose is to quietly reduce cognitive load by: + +- observing +- remembering +- planning +- researching +- proposing +- executing safe actions +- continuously improving + +The system should feel less like ChatGPT and more like employing a careful, diligent junior platform engineer and executive assistant who never forgets anything. + +--- + +# Core Principles + +## Persistent + +Steward has persistent state and is never "reset" between interactions. + +Conceptually it is always running. In practice this means persistent state with event-driven wake, not a literal always-on compute loop, so persistence is not bought at the cost of continuous compute. + +It continuously evaluates changes in its environment as they arrive. + +Scheduled tasks are allowed but the initiator is an event. + +--- + +## Event Driven + +Everything is an event. + +Examples: + +- Prometheus alert +- Telegram conversation +- Calendar update +- Todoist completion +- GitHub Pull Request +- AWX job completion +- Weather change +- Home Assistant state change + +Events are evaluated against active goals. + +--- + +## Event Processing + +"Everything is an event" does not mean everything is evaluated. + +A homelab emits thousands of low-value events, and metrics sources can flap continuously. Raw events pass through a processing layer before reaching goal evaluation: + +- **Ingestion** normalises events from all sources into a common shape. +- **Filtering** drops noise below a relevance threshold. +- **Deduplication and debounce** collapse repeated or flapping events. +- **Correlation** groups related events into a single *situation* (e.g. three alerts from one failing disk). + +Steward evaluates *situations* against goals, not raw events. This bounds compute and prevents notification storms. + +Event content is untrusted input. See [Security Model](#security-model). + +--- + +## Goal Oriented + +Steward does not work from prompts. + +It works from goals. + +Examples: + +Maintain homelab reliability. + +Prepare for house move. + +Visit Queensland destinations before moving. + +Reduce operational effort. + +Research new technologies. + +Complete backlog projects. + +Prompts simply create, modify or clarify goals. + +--- + +## Policy First + +The AI never directly executes arbitrary commands. + +Every action must pass through a policy engine. + +Example policies: + +Allowed automatically + +- restart containers +- retry backups +- renew certificates +- apply security patches + +Requires approval + +- Kubernetes upgrades +- infrastructure changes +- firewall modifications +- snapshot deletion + +Never + +- delete backups +- modify Steward permissions +- disable audit logging + +--- + +## Transparency + +Every recommendation should explain: + +- why +- confidence +- evidence +- expected outcome +- rollback strategy + +No hidden reasoning. + +--- + +## Conservative + +Steward prefers: + +- reversible actions +- small changes +- incremental improvement + +rather than aggressive optimisation. + +--- + +# Domains + +Steward is intentionally domain-agnostic. + +Initially it will focus on Homelab Operations. + +Future domains include: + +- family planning +- travel +- finance +- home maintenance +- learning +- research +- project management + +Each domain has its own tools and permissions. + +--- + +# Cognitive Load Reduction + +The purpose is NOT task management. + +The purpose is reducing the need to remember. + +Traditional tools remember things once entered. + +Steward helps remember to remember. + +Examples + +"I noticed you mentioned replacing the UPS several times." + +"I noticed your dental check-up is overdue." + +"I've researched family camping locations." + +"I've identified a free weekend." + +--- + +# Communication Philosophy + +Steward communicates only when useful. + +Communication categories: + +Critical + +Immediate interruption. + +Decision Required + +Requests human approval. + +Background + +No notification required. + +Daily reports may exist but are not central. + +Real-time context-aware communication is the default. + +--- + +# Curiosity Budget + +One of Steward's defining concepts. + +If no operational work exists, Steward may spend limited compute researching or improving delegated objectives. + +Examples + +- Evaluate ArgoCD +- Benchmark local LLMs +- Research family holidays +- Compare UPS replacements +- Prototype monitoring improvements + +Curiosity is budgeted. + +The budget is a concrete, metered resource, not an aspiration. It is expressed in measurable units (for example: tokens/day, dollars/month, GPU-hours/week), set by the operator, and enforced by Steward refusing to exceed it. Consumption is observable so the operator can see what curiosity cost and what it produced. + +It is suspended immediately if operational work becomes necessary. + +--- + +# Confidence + +Steward maintains historical confidence scores. + +Example + +Restart Jellyfin + +Attempts: 34 + +Success: 33 + +Confidence: 97% + +Upgrade Kubernetes + +Attempts: 2 + +Success: 1 + +Rollback: 1 + +Confidence: 50% + +Confidence is earned. + +## Confidence and Autonomy + +Confidence influences autonomy, but must never override policy. + +Rules: + +- The policy tier of an action (automatic / approval / never) is human-set and immutable to Steward. High confidence never promotes an action into a more permissive tier. Doing so would be self-modification by the back door. +- Within a tier, confidence affects only *how* Steward acts: how strongly it recommends, how much evidence it attaches, whether it batches or surfaces individually. +- Confidence is scoped, not global. "Restart Jellyfin" confidence does not transfer to "restart Postgres". Scope is (action × target × context), and Steward must not generalise across scopes without evidence. +- Confidence decays. Environment changes (version upgrades, config changes, topology changes) invalidate historical success and should reset or discount the relevant scores. Stale confidence is treated as low confidence. + +--- + +# Security Model + +Safety is a property of the architecture, not just the model's good behaviour. + +## The policy engine is the trust boundary + +The policy engine is Steward's crown jewel and must be a separate, independently-audited component that the AI *calls*, not code the AI can read, modify, or bypass. The AI proposes actions; the policy engine decides. If the AI could edit the engine, every other guarantee collapses. + +## Event content is untrusted + +Events carry attacker-influenceable content: a GitHub PR title, a Telegram message, a webhook payload. Steward must treat all event content as untrusted input. + +- Event content can inform situations and goals but can never escalate authority or promote an action's policy tier. +- Instructions embedded in event content are data, not commands (defence against prompt injection). + +## Memory can rot + +"Nothing is forgotten" is both a feature and a liability. A learned preference can encode a mistake permanently, and learned lessons are themselves derived from potentially-untrusted history. + +- Learned memory is reviewable and retractable. A bad lesson must be able to be inspected and removed. +- New lessons that would change behaviour surface as *proposals*, consistent with [Reflection](#reflection), rather than silently altering future decisions. + +--- + +# Memory + +Different information belongs in different storage. The roles below are fixed; the named tools are current implementation choices, not commitments. + +Structured store (e.g. NocoDB) + +Examples + +- Tasks +- Goals +- Incidents +- Systems +- Services +- Confidence +- Policies + +Long-form store (e.g. vector database / RAG) + +Examples + +- documentation +- conversations +- runbooks +- release notes + +Metrics store (e.g. Prometheus) + +Examples + +- resource usage +- trends +- operational history + +Audit + +Immutable execution history. + +Nothing is forgotten. + +--- + +# Self Improvement + +Steward should become better over time. + +Important distinction + +Self-improving + +✅ Learn preferences + +✅ Learn successful remediations + +✅ Improve prompts + +✅ Improve playbooks + +✅ Improve planning + +Self-modifying + +❌ Rewrite policy engine + +❌ Change permissions + +❌ Remove safety controls + +Immutable system components remain protected. + +--- + +# Learning + +Steward continuously learns: + +Operator preferences + +Family preferences + +Infrastructure behaviour + +Recurring incidents + +Successful remediation + +Failed remediation + +Communication preferences + +Historical context + +This learning improves planning. + +--- + +# Reflection + +Periodic reflection sessions. + +Questions include: + +What repeated? + +What wasted time? + +What should be automated? + +Which playbooks failed? + +Which policies need review? + +Which research produced value? + +Reflection generates proposals rather than automatic change. + +--- + +# Research Workflow + +Research tasks behave like engineering work. + +Example + +Research destination --> Gather information --> Summarise findings --> Generate recommendations --> Present proposal --> Receive feedback --> Update understanding --> Re-plan + +--- + +# Engineering Workflow + +Operational improvements follow software engineering practices. + +Observe issue --> Generate proposal --> Approval --> Create Git branch --> Modify playbooks --> Run tests --> Open Pull Request --> Deploy through AWX --> Observe --> Learn + +--- + +# Long-term Vision + +Steward evolves from: + +Assistant --> Operator --> Engineer --> Trusted Steward + +Trust is earned through demonstrated reliability rather than increasing model capability. + +--- + +# Success Criteria + +Steward is successful when: + +I think less. + +I forget less. + +Routine work disappears. + +The homelab becomes increasingly self-maintaining. + +Projects continue progressing while I am busy. + +The system accumulates operational knowledge. + +The system proactively suggests valuable improvements. + +The system remains predictable. + +The system remains explainable. + +The system remains safe. + +## Measurable Proxies + +The criteria above are subjective. To tell whether v0.2 actually beat v0.1, track measurable proxies alongside them: + +- Operator interruptions per week (trend down). +- Ratio of actions auto-approved vs escalated for approval. +- Mean time to remediation for recurring incidents. +- False-alarm / unnecessary-notification rate. +- Curiosity spend vs value produced (proposals accepted). + +--- + +# Mission Statement + +Steward exists to maximise the reliability, maintainability and usefulness of its delegated domains while minimising operator cognitive load. + +It continuously observes, remembers, researches, plans and acts within explicitly delegated authority. + +It values safety over speed, explanation over opacity, and long-term trust over short-term autonomy. diff --git a/pyproject.toml b/pyproject.toml new file mode 100644 index 0000000..a40e0cb --- /dev/null +++ b/pyproject.toml @@ -0,0 +1,52 @@ +[build-system] +requires = ["setuptools>=68", "wheel"] +build-backend = "setuptools.build_meta" + +[project] +name = "steward" +version = "0.1.0" +description = "A long-running, AI-assisted personal operations platform" +readme = "README.md" +requires-python = ">=3.12" +dependencies = [ + "python-telegram-bot>=21.0", + "openai>=1.30", + "httpx>=0.27", + "python-dotenv>=1.0", + "apscheduler>=3.10", + "pydantic>=2.7", + "pydantic-settings>=2.3", +] + +[project.optional-dependencies] +dev = [ + "pytest>=8.2", + "pytest-asyncio>=0.23", + "pytest-mock>=3.14", + "ruff>=0.4", + "mypy>=1.10", + "respx>=0.21", +] + +[project.scripts] +steward = "steward.main:main" + +[tool.setuptools.packages.find] +where = ["."] +include = ["steward*"] + +[tool.ruff] +line-length = 100 +target-version = "py312" + +[tool.ruff.lint] +select = ["E", "F", "I", "UP"] + +[tool.mypy] +python_version = "3.12" +strict = true +ignore_missing_imports = true + +[tool.pytest.ini_options] +asyncio_mode = "auto" +testpaths = ["tests"] diff --git a/steward/__init__.py b/steward/__init__.py new file mode 100644 index 0000000..dcccb8d --- /dev/null +++ b/steward/__init__.py @@ -0,0 +1 @@ +"""Steward – AI-assisted personal operations platform.""" diff --git a/steward/bot/__init__.py b/steward/bot/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/steward/bot/telegram.py b/steward/bot/telegram.py new file mode 100644 index 0000000..32a3b3b --- /dev/null +++ b/steward/bot/telegram.py @@ -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) diff --git a/steward/config.py b/steward/config.py new file mode 100644 index 0000000..5c36303 --- /dev/null +++ b/steward/config.py @@ -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() diff --git a/steward/llm/__init__.py b/steward/llm/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/steward/llm/client.py b/steward/llm/client.py new file mode 100644 index 0000000..883cd4a --- /dev/null +++ b/steward/llm/client.py @@ -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 diff --git a/steward/main.py b/steward/main.py new file mode 100644 index 0000000..0d03b2d --- /dev/null +++ b/steward/main.py @@ -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() diff --git a/steward/proposals/__init__.py b/steward/proposals/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/steward/proposals/generator.py b/steward/proposals/generator.py new file mode 100644 index 0000000..68bd95c --- /dev/null +++ b/steward/proposals/generator.py @@ -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:** \n" + "**Why:** \n" + "**Evidence:** \n" + "**Expected outcome:** \n" + "**Rollback strategy:** \n" + "**Confidence:** - \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 diff --git a/tests/__init__.py b/tests/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/tests/test_bot.py b/tests/test_bot.py new file mode 100644 index 0000000..53a2a5a --- /dev/null +++ b/tests/test_bot.py @@ -0,0 +1,119 @@ +"""Tests for steward.bot.telegram (handler logic).""" + +from unittest.mock import AsyncMock, MagicMock + +import pytest +from telegram import Message, Update, User +from telegram.ext import CallbackContext + +from steward.bot.telegram import ( + _history, + _is_allowed, + clear_handler, + message_handler, + start_handler, +) +from steward.config import Settings +from steward.llm.client import LLMClient + + +def _make_settings(**kwargs) -> Settings: + defaults = dict( + telegram_bot_token="test-token", + openai_api_key="test-key", + ) + defaults.update(kwargs) + return Settings(**defaults) + + +def _make_update(user_id: int = 12345, text: str = "hello") -> Update: + user = MagicMock(spec=User) + user.id = user_id + + message = MagicMock(spec=Message) + message.text = text + message.reply_text = AsyncMock() + + update = MagicMock(spec=Update) + update.effective_user = user + update.message = message + return update + + +def _make_context(settings: Settings, llm: LLMClient | None = None) -> CallbackContext: # type: ignore[type-arg] + ctx = MagicMock(spec=CallbackContext) + ctx.bot_data = {"settings": settings, "llm": llm} + return ctx + + +class TestIsAllowed: + def test_no_allowlist_allows_everyone(self): + settings = _make_settings(telegram_allowed_user_ids=[]) + assert _is_allowed(999, settings) is True + + def test_allowlist_accepts_known_user(self): + settings = _make_settings(telegram_allowed_user_ids=[1, 2, 3]) + assert _is_allowed(2, settings) is True + + def test_allowlist_rejects_unknown_user(self): + settings = _make_settings(telegram_allowed_user_ids=[1, 2, 3]) + assert _is_allowed(999, settings) is False + + +@pytest.mark.asyncio +async def test_start_handler_replies(monkeypatch): + settings = _make_settings() + update = _make_update() + ctx = _make_context(settings) + + await start_handler(update, ctx) + + update.message.reply_text.assert_awaited_once() + text = update.message.reply_text.call_args.args[0] + assert "Steward" in text + + +@pytest.mark.asyncio +async def test_start_handler_ignores_disallowed_user(): + settings = _make_settings(telegram_allowed_user_ids=[9999]) + update = _make_update(user_id=1111) + ctx = _make_context(settings) + + await start_handler(update, ctx) + + update.message.reply_text.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_clear_handler_clears_history(): + settings = _make_settings() + user_id = 42 + _history[user_id] = [{"role": "user", "content": "old msg"}] + + update = _make_update(user_id=user_id) + ctx = _make_context(settings) + + await clear_handler(update, ctx) + + assert _history[user_id] == [] + update.message.reply_text.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_message_handler_calls_llm_and_replies(): + settings = _make_settings() + mock_llm = MagicMock(spec=LLMClient) + mock_llm.chat = AsyncMock(return_value="LLM response") + + update = _make_update(user_id=77, text="What is the weather?") + ctx = _make_context(settings, llm=mock_llm) + + _history[77].clear() + await message_handler(update, ctx) + + mock_llm.chat.assert_awaited_once() + update.message.reply_text.assert_awaited_once() + # History should now contain the user/assistant turn + assert len(_history[77]) == 2 + assert _history[77][0]["role"] == "user" + assert _history[77][1]["role"] == "assistant" diff --git a/tests/test_llm_client.py b/tests/test_llm_client.py new file mode 100644 index 0000000..2d16458 --- /dev/null +++ b/tests/test_llm_client.py @@ -0,0 +1,92 @@ +"""Tests for steward.llm.client.""" + +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +from steward.config import Settings +from steward.llm.client import LLMClient + + +def _make_settings(**kwargs) -> Settings: + defaults = dict( + telegram_bot_token="test-token", + openai_api_key="test-key", + openai_model="gpt-4o", + openai_system_prompt="You are Steward.", + ) + defaults.update(kwargs) + return Settings(**defaults) + + +@pytest.fixture +def settings(): + return _make_settings() + + +@pytest.fixture +def llm_client(settings): + return LLMClient(settings) + + +@pytest.mark.asyncio +async def test_chat_sends_correct_messages(llm_client): + """chat() should prepend the system prompt and append the user message.""" + mock_response = MagicMock() + mock_response.choices[0].message.content = "Hello from the LLM!" + + with patch.object( + llm_client._client.chat.completions, + "create", + new_callable=AsyncMock, + return_value=mock_response, + ) as mock_create: + result = await llm_client.chat("Hi there") + + assert result == "Hello from the LLM!" + call_kwargs = mock_create.call_args.kwargs + messages = call_kwargs["messages"] + assert messages[0]["role"] == "system" + assert messages[-1]["role"] == "user" + assert messages[-1]["content"] == "Hi there" + + +@pytest.mark.asyncio +async def test_chat_with_history(llm_client): + """chat() should include history between system and user messages.""" + history = [ + {"role": "user", "content": "previous question"}, + {"role": "assistant", "content": "previous answer"}, + ] + mock_response = MagicMock() + mock_response.choices[0].message.content = "reply" + + with patch.object( + llm_client._client.chat.completions, + "create", + new_callable=AsyncMock, + return_value=mock_response, + ) as mock_create: + await llm_client.chat("new question", history=history) + + messages = mock_create.call_args.kwargs["messages"] + roles = [m["role"] for m in messages] + assert roles == ["system", "user", "assistant", "user"] + + +@pytest.mark.asyncio +async def test_chat_custom_system_prompt(llm_client): + """chat() should use a custom system prompt when provided.""" + mock_response = MagicMock() + mock_response.choices[0].message.content = "ok" + + with patch.object( + llm_client._client.chat.completions, + "create", + new_callable=AsyncMock, + return_value=mock_response, + ) as mock_create: + await llm_client.chat("msg", system_prompt="Custom prompt") + + messages = mock_create.call_args.kwargs["messages"] + assert messages[0]["content"] == "Custom prompt" diff --git a/tests/test_proposals.py b/tests/test_proposals.py new file mode 100644 index 0000000..bb03e0d --- /dev/null +++ b/tests/test_proposals.py @@ -0,0 +1,106 @@ +"""Tests for steward.proposals.generator.""" + +from datetime import UTC +from unittest.mock import AsyncMock, MagicMock + +import httpx +import pytest +import respx + +from steward.config import Settings +from steward.llm.client import LLMClient +from steward.proposals.generator import Proposal, ProposalGenerator + + +def _make_settings(**kwargs) -> Settings: + defaults = dict( + openai_api_key="test-key", + analysis_target_url="https://example.com/api/status", + analysis_target_api_key="secret", + ) + defaults.update(kwargs) + return Settings(**defaults) + + +@pytest.fixture +def settings(): + return _make_settings() + + +@pytest.fixture +def mock_llm(): + llm = MagicMock(spec=LLMClient) + llm.chat = AsyncMock( + return_value="**Summary:** Everything is healthy.\n**Why:** No issues detected." + ) + return llm + + +@pytest.mark.asyncio +async def test_run_generates_proposal(settings, mock_llm): + """run() should fetch the URL and return a Proposal.""" + with respx.mock: + respx.get("https://example.com/api/status").mock( + return_value=httpx.Response(200, text='{"status": "ok"}') + ) + generator = ProposalGenerator(settings, mock_llm) + proposal = await generator.run() + + assert proposal is not None + assert isinstance(proposal, Proposal) + assert proposal.source_url == "https://example.com/api/status" + assert proposal.raw_data == '{"status": "ok"}' + mock_llm.chat.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_run_no_url_returns_none(mock_llm): + """run() should return None when analysis_target_url is not configured.""" + settings = _make_settings(analysis_target_url="") + generator = ProposalGenerator(settings, mock_llm) + result = await generator.run() + assert result is None + mock_llm.chat.assert_not_called() + + +@pytest.mark.asyncio +async def test_run_http_error_returns_none(settings, mock_llm): + """run() should return None when the HTTP request fails.""" + with respx.mock: + respx.get("https://example.com/api/status").mock( + return_value=httpx.Response(500, text="Internal Server Error") + ) + generator = ProposalGenerator(settings, mock_llm) + proposal = await generator.run() + + assert proposal is None + mock_llm.chat.assert_not_called() + + +@pytest.mark.asyncio +async def test_fetch_sends_bearer_token(settings, mock_llm): + """_fetch() should include ****** when api_key is configured.""" + with respx.mock: + route = respx.get("https://example.com/api/status").mock( + return_value=httpx.Response(200, text="data") + ) + generator = ProposalGenerator(settings, mock_llm) + await generator._fetch("https://example.com/api/status") + + request = route.calls.last.request + assert request.headers["Authorization"].startswith("Bearer ") + + +def test_proposal_format_for_telegram(): + """format_for_telegram() should contain the source URL and body.""" + from datetime import datetime + p = Proposal( + title="Test", + body="**Summary:** All good.", + source_url="https://example.com", + generated_at=datetime(2026, 1, 1, 8, 0, 0, tzinfo=UTC), + ) + text = p.format_for_telegram() + assert "https://example.com" in text + assert "All good." in text + assert "2026-01-01" in text