commit 5a2c5df0415d53fe7064061039146303cf5ed0a2 Author: Wickedness Date: Mon Jun 8 22:19:08 2026 +0900 Initial commit diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..98ee43e --- /dev/null +++ b/.dockerignore @@ -0,0 +1,8 @@ +.env +.venv +__pycache__ +.pytest_cache +data +codex-home +workspaces +*.log diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..079f63c --- /dev/null +++ b/.env.example @@ -0,0 +1,25 @@ +# Telegram +TELEGRAM_BOT_TOKEN=put-your-token-here +ALLOWED_TELEGRAM_USER_IDS= + +# Bot storage and workspace +BOT_STATE_PATH=/app/data/state.json +BOT_WORKSPACE_ROOT=/workspaces +BOT_TIMEZONE=Asia/Seoul +SCHEDULE_POLL_SECONDS=30 +PERSISTENT_WS_IDLE_SECONDS=600 + +# Codex execution +CODEX_BIN=codex +CODEX_TIMEOUT_SECONDS=1800 +CODEX_BACKEND=auto +CODEX_APP_SERVER_URL=ws://127.0.0.1:4500 +CODEX_APP_SERVER_MODEL=gpt-5.5 +CODEX_MODEL_REASONING_EFFORT=high +CODEX_SPEED_MODE=standard +CODEX_READONLY_SANDBOX=read-only +CODEX_WRITE_SANDBOX=workspace-write +CODEX_RESUME_JSON_FLAG_STYLE=before_resume + +# Optional: set this only if you want plain messages to behave like /ask. +ALLOW_PLAIN_TEXT=true diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..96d0676 --- /dev/null +++ b/.gitignore @@ -0,0 +1,9 @@ +.env +.venv/ +__pycache__/ +*.py[cod] +.pytest_cache/ +data/ +deploy/.env +*.log + diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..a8d43db --- /dev/null +++ b/Dockerfile @@ -0,0 +1,36 @@ +FROM python:3.12-slim + +ARG CODEX_NPM_PACKAGE=@openai/codex + +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 \ + CODEX_HOME=/root/.codex \ + BOT_STATE_PATH=/app/data/state.json \ + BOT_WORKSPACE_ROOT=/workspaces \ + CODEX_APP_SERVER_URL=ws://127.0.0.1:4500 + +RUN apt-get update \ + && apt-get install -y --no-install-recommends \ + ca-certificates \ + curl \ + git \ + nodejs \ + npm \ + openssh-client \ + && npm install -g "${CODEX_NPM_PACKAGE}" \ + && apt-get clean \ + && rm -rf /var/lib/apt/lists/* + +WORKDIR /app + +COPY requirements.txt ./ +RUN pip install --no-cache-dir -r requirements.txt + +COPY app ./app +COPY scripts ./scripts +COPY README.md ./ + +RUN mkdir -p /app/data /workspaces /root/.codex \ + && chmod +x /app/scripts/start.sh + +CMD ["/app/scripts/start.sh"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..3fe4dd1 --- /dev/null +++ b/README.md @@ -0,0 +1,110 @@ +# NAS Telegram Codex Bridge + +Telegram bridge for a NAS-hosted Codex app-server. + +The bot does not interpret normal chat messages. It forwards them to Codex and +only keeps connection-management commands. + +## Commands + +```text +/start +/help +/whoami +/repo +/ask +/new +/status +/settings +/set_model +/set_effort +/set_speed +/fast +/models +/schedules +/records +/todos +/purchases +/birthdays +/sessions +/sessions archived +/use_session +/archive_session [thread_id] +/unarchive_session +``` + +`ALLOW_PLAIN_TEXT=true` makes regular Telegram messages behave like `/ask`. + +`/settings`, `/models`, `/schedules`, `/records`, and `/sessions` return inline +buttons, so you can change speed/model/reasoning settings, delete schedules, +manage personal records, and switch or archive sessions without copying ids. + +## Runtime settings + +Runtime settings are stored in `BOT_STATE_PATH` and applied to new turns: + +- `model`, for example `gpt-5.5` +- `reasoning_effort`: `minimal`, `low`, `medium`, `high`, or `xhigh` +- `speed_mode`: `standard` or `fast` + +Changing these settings closes existing persistent WebSocket clients. The next +Telegram message reconnects with the updated settings. Fast mode requires a +supported model/account; if the local Codex app-server rejects an override, the +bridge retries with narrower settings rather than failing the whole turn. + +The model picker first tries `codex debug models` inside the NAS container and +builds buttons from the model catalog Codex sees. If that fails, it falls back +to a short built-in list. + +## Scheduling + +The bridge registers a local MCP server named `telegram_bridge` in +`$CODEX_HOME/config.toml`. Codex can call these tools when the user asks for +scheduled Telegram notifications: + +- `telegram_create_schedule` +- `telegram_list_schedules` +- `telegram_delete_schedule` + +## Personal assistant records + +Codex can persist personal assistant records to the same NAS-backed state file. +Use natural language in Telegram, for example: + +```text +이건 업무 이력으로 저장해줘. 2025-09-15부터 2026-04-15까지 전국통합데이터 개방확대 사업에 참여했어. +이건 할 일로 기억해줘. 다음 주까지 신세계 TV쇼핑 채팅 연동 QA 정리. +구매 후보로 저장해줘. 모니터암 가격 8만원 이하가 되면 사고 싶어. +어머니 생일은 5월 12일이야. 기억해줘. +``` + +Codex uses these MCP tools: + +- `telegram_save_personal_record` +- `telegram_list_personal_records` +- `telegram_update_personal_record` +- `telegram_delete_personal_record` + +Record categories are `work_log`, `todo`, `purchase`, `birthday`, `contact`, +`note`, and `preference`. Telegram commands `/records`, `/todos`, `/purchases`, +and `/birthdays` show records with management buttons. + +The bridge scheduler runs inside the Telegram bot process. When a schedule is +due, it calls Codex app-server and sends the result back to the stored Telegram +chat id. + +## Deployment + +1. Copy `.env.example` to `.env`. +2. Fill `TELEGRAM_BOT_TOKEN`. +3. Use `/whoami` to find your Telegram user id and set + `ALLOWED_TELEGRAM_USER_IDS`. +4. Run `docker compose up -d --build`, or use `deploy/deploy.ps1` for the NAS. + +## Notes + +- Codex state and auth live under the mounted `codex-home` volume. +- The app-server listens on `ws://127.0.0.1:4500` inside the container. +- Normal chat messages use a persistent app-server WebSocket per Telegram chat. +- Idle persistent WebSockets are closed after `PERSISTENT_WS_IDLE_SECONDS`. +- Session commands call app-server thread APIs. diff --git a/app/__init__.py b/app/__init__.py new file mode 100644 index 0000000..09cee55 --- /dev/null +++ b/app/__init__.py @@ -0,0 +1,2 @@ +"""Telegram Codex bot package.""" + diff --git a/app/app_server_runner.py b/app/app_server_runner.py new file mode 100644 index 0000000..2881fca --- /dev/null +++ b/app/app_server_runner.py @@ -0,0 +1,525 @@ +from __future__ import annotations + +import asyncio +import json +import time +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any + +import websockets + +from app.codex_runner import CodexResult + + +class AppServerProtocolError(RuntimeError): + pass + + +@dataclass +class _TurnCollector: + agent_deltas: list[str] = field(default_factory=list) + final_messages: list[str] = field(default_factory=list) + commands: list[str] = field(default_factory=list) + errors: list[str] = field(default_factory=list) + done: bool = False + status: str | None = None + + def handle(self, message: dict[str, Any]) -> None: + method = message.get("method") + params = message.get("params") or {} + if not isinstance(params, dict): + return + + if method == "item/agentMessage/delta": + delta = params.get("delta") + if isinstance(delta, str): + self.agent_deltas.append(delta) + return + + if method == "item/completed": + item = params.get("item") + if isinstance(item, dict): + self._handle_completed_item(item) + return + + if method == "turn/completed": + self.done = True + turn = params.get("turn") + if isinstance(turn, dict): + self.status = turn.get("status") + error = turn.get("error") + if isinstance(error, dict) and error.get("message"): + self.errors.append(str(error["message"])) + return + + def _handle_completed_item(self, item: dict[str, Any]) -> None: + item_type = item.get("type") + if item_type in {"agentMessage", "agent_message"} and item.get("text"): + self.final_messages.append(str(item["text"])) + return + + if item_type in {"commandExecution", "command_execution"}: + command = item.get("command") + if isinstance(command, list): + self.commands.append(" ".join(str(part) for part in command)) + elif command: + self.commands.append(str(command)) + + def final_text(self) -> str: + if self.final_messages: + return self.final_messages[-1] + if self.agent_deltas: + return "".join(self.agent_deltas).strip() + if self.errors: + return "\n".join(self.errors) + return "(Codex finished without text output.)" + + +class AppServerRunner: + def __init__( + self, + url: str, + timeout_seconds: int, + model: str | None = None, + effort: str | None = None, + service_tier: str | None = None, + ) -> None: + self.url = url + self.timeout_seconds = timeout_seconds + self.model = model + self.effort = effort + self.service_tier = service_tier + self._next_id = 1 + + async def run( + self, + prompt: str, + cwd: Path, + sandbox: str, + session_id: str | None, + skip_git_repo_check: bool = False, + ) -> CodexResult: + del skip_git_repo_check + try: + return await asyncio.wait_for( + self._run_once(prompt, cwd, sandbox, session_id), + timeout=self.timeout_seconds, + ) + except asyncio.TimeoutError: + return CodexResult( + ok=False, + final_message=f"Codex app-server timed out after {self.timeout_seconds} seconds.", + session_id=session_id, + ) + except Exception as exc: + return CodexResult( + ok=False, + final_message=f"Codex app-server failed: {exc}", + session_id=session_id, + ) + + async def call(self, method: str, params: dict[str, Any] | None = None) -> Any: + collector = _TurnCollector() + async with websockets.connect( + self.url, + open_timeout=15, + ping_interval=20, + max_size=16 * 1024 * 1024, + ) as ws: + await self._initialize(ws, collector) + return await self._request(ws, method, params or {}, collector) + + async def list_threads( + self, + limit: int = 10, + archived: bool = False, + cursor: str | None = None, + ) -> dict[str, Any]: + params: dict[str, Any] = { + "limit": limit, + "archived": archived, + } + if cursor: + params["cursor"] = cursor + result = await self.call("thread/list", params) + return result if isinstance(result, dict) else {} + + async def loaded_threads(self) -> dict[str, Any]: + result = await self.call("thread/loaded/list", {}) + return result if isinstance(result, dict) else {} + + async def archive_thread(self, thread_id: str) -> None: + await self.call("thread/archive", {"threadId": thread_id}) + + async def unarchive_thread(self, thread_id: str) -> None: + await self.call("thread/unarchive", {"threadId": thread_id}) + + async def resume_thread(self, thread_id: str) -> str: + result = await self.call("thread/resume", {"threadId": thread_id}) + return _extract_thread_id(result) or thread_id + + async def _run_once( + self, + prompt: str, + cwd: Path, + sandbox: str, + session_id: str | None, + ) -> CodexResult: + collector = _TurnCollector() + async with websockets.connect( + self.url, + open_timeout=15, + ping_interval=20, + max_size=16 * 1024 * 1024, + ) as ws: + await self._initialize(ws, collector) + + thread_id = await self._ensure_thread(ws, cwd, sandbox, session_id, collector) + turn_params: dict[str, Any] = { + "threadId": thread_id, + "input": [{"type": "text", "text": prompt}], + } + + await self._request_with_settings_fallback( + ws, + "turn/start", + turn_params, + collector, + ) + while not collector.done: + message = await self._recv(ws) + await self._handle_non_response(ws, message, collector) + + ok = collector.status not in {"failed", "interrupted"} + return CodexResult( + ok=ok, + final_message=collector.final_text(), + session_id=thread_id, + commands=collector.commands, + stderr="\n".join(collector.errors), + returncode=0 if ok else 1, + ) + + async def _ensure_thread( + self, + ws: Any, + cwd: Path, + sandbox: str, + session_id: str | None, + collector: _TurnCollector, + ) -> str: + if session_id: + try: + result = await self._request_with_settings_fallback( + ws, + "thread/resume", + {"threadId": session_id}, + collector, + ) + return _extract_thread_id(result) or session_id + except AppServerProtocolError: + pass + + result = await self._request_with_settings_fallback( + ws, + "thread/start", + {}, + collector, + ) + thread_id = _extract_thread_id(result) + if not thread_id: + raise AppServerProtocolError("thread/start did not return a thread id") + return thread_id + + async def _initialize(self, ws: Any, collector: _TurnCollector) -> None: + await self._request( + ws, + "initialize", + { + "clientInfo": { + "name": "telegram_codex_bot", + "title": "Telegram Codex Bot", + "version": "0.1.0", + }, + "capabilities": {"experimentalApi": True}, + }, + collector, + ) + await self._notify(ws, "initialized", {}) + + def _settings_params(self) -> dict[str, Any]: + params: dict[str, Any] = {} + if self.model: + params["model"] = self.model + if self.effort: + params["effort"] = self.effort + if self.service_tier: + params["serviceTier"] = self.service_tier + return params + + def _settings_variants(self) -> list[tuple[str, dict[str, Any]]]: + full = self._settings_params() + variants: list[tuple[str, dict[str, Any]]] = [] + seen: set[str] = set() + + def add(label: str, params: dict[str, Any]) -> None: + key = json.dumps(params, sort_keys=True, ensure_ascii=True) + if key in seen: + return + seen.add(key) + variants.append((label, params)) + + add("model/effort/speed", full) + if "serviceTier" in full: + add("model/effort", {k: v for k, v in full.items() if k != "serviceTier"}) + if "effort" in full: + add("model/speed", {k: v for k, v in full.items() if k != "effort"}) + if "model" in full: + add("model", {"model": full["model"]}) + add("default settings", {}) + return variants + + async def _request_with_settings_fallback( + self, + ws: Any, + method: str, + base_params: dict[str, Any], + collector: _TurnCollector, + ) -> Any: + errors: list[str] = [] + last_error: AppServerProtocolError | None = None + for label, settings in self._settings_variants(): + params = dict(base_params) + params.update(settings) + try: + result = await self._request(ws, method, params, collector) + if errors: + collector.errors.append( + "Codex app-server rejected one or more runtime setting " + f"overrides; continued with {label}. Last error: {errors[-1]}" + ) + return result + except AppServerProtocolError as exc: + last_error = exc + errors.append(f"{label}: {exc}") + if not settings: + break + raise last_error or AppServerProtocolError(f"{method} failed") + + async def _request( + self, + ws: Any, + method: str, + params: dict[str, Any], + collector: _TurnCollector, + ) -> Any: + request_id = self._next_request_id() + await ws.send(json.dumps({"method": method, "id": request_id, "params": params})) + while True: + message = await self._recv(ws) + if message.get("id") == request_id: + if message.get("error"): + error = message["error"] + if isinstance(error, dict): + raise AppServerProtocolError(error.get("message") or str(error)) + raise AppServerProtocolError(str(error)) + return message.get("result") + await self._handle_non_response(ws, message, collector) + + async def _notify(self, ws: Any, method: str, params: dict[str, Any]) -> None: + await ws.send(json.dumps({"method": method, "params": params})) + + async def _recv(self, ws: Any) -> dict[str, Any]: + raw = await ws.recv() + if isinstance(raw, bytes): + raw = raw.decode("utf-8", errors="replace") + try: + message = json.loads(raw) + except json.JSONDecodeError as exc: + raise AppServerProtocolError(f"invalid JSON from app-server: {exc}") from exc + if not isinstance(message, dict): + raise AppServerProtocolError("app-server returned a non-object message") + return message + + async def _handle_non_response( + self, + ws: Any, + message: dict[str, Any], + collector: _TurnCollector, + ) -> None: + if "id" in message and "method" in message: + await ws.send( + json.dumps( + { + "id": message["id"], + "error": { + "code": -32601, + "message": "Telegram bot client cannot handle this app-server request.", + }, + } + ) + ) + return + collector.handle(message) + + def _next_request_id(self) -> int: + request_id = self._next_id + self._next_id += 1 + return request_id + + +def _extract_thread_id(result: Any) -> str | None: + if not isinstance(result, dict): + return None + thread = result.get("thread") + if isinstance(thread, dict) and thread.get("id"): + return str(thread["id"]) + if result.get("threadId"): + return str(result["threadId"]) + return None + + +class PersistentAppServerRunner(AppServerRunner): + def __init__( + self, + url: str, + timeout_seconds: int, + model: str | None = None, + effort: str | None = None, + service_tier: str | None = None, + ) -> None: + super().__init__( + url=url, + timeout_seconds=timeout_seconds, + model=model, + effort=effort, + service_tier=service_tier, + ) + self._ws: Any | None = None + self._lock = asyncio.Lock() + self.last_used_monotonic = time.monotonic() + + @property + def is_connected(self) -> bool: + return self._ws is not None + + def idle_seconds(self) -> float: + return time.monotonic() - self.last_used_monotonic + + async def close(self) -> None: + ws = self._ws + self._ws = None + if ws is not None: + await ws.close() + + async def close_if_idle(self, idle_seconds: int) -> bool: + if self._ws is None: + return False + if self.idle_seconds() < idle_seconds: + return False + async with self._lock: + if self._ws is None or self.idle_seconds() < idle_seconds: + return False + await self.close() + return True + + async def run( + self, + prompt: str, + cwd: Path, + sandbox: str, + session_id: str | None, + skip_git_repo_check: bool = False, + ) -> CodexResult: + del skip_git_repo_check + async with self._lock: + try: + return await asyncio.wait_for( + self._run_persistent(prompt, cwd, sandbox, session_id), + timeout=self.timeout_seconds, + ) + except asyncio.TimeoutError: + await self.close() + return CodexResult( + ok=False, + final_message=f"Codex app-server timed out after {self.timeout_seconds} seconds.", + session_id=session_id, + ) + except Exception as exc: + await self.close() + return CodexResult( + ok=False, + final_message=f"Codex app-server failed: {exc}", + session_id=session_id, + ) + + async def _run_persistent( + self, + prompt: str, + cwd: Path, + sandbox: str, + session_id: str | None, + ) -> CodexResult: + try: + return await self._run_on_connection(prompt, cwd, sandbox, session_id) + except Exception: + await self.close() + return await self._run_on_connection(prompt, cwd, sandbox, session_id) + + async def _run_on_connection( + self, + prompt: str, + cwd: Path, + sandbox: str, + session_id: str | None, + ) -> CodexResult: + self.last_used_monotonic = time.monotonic() + collector = _TurnCollector() + ws = await self._ensure_connection(collector) + thread_id = await self._ensure_thread(ws, cwd, sandbox, session_id, collector) + turn_params: dict[str, Any] = { + "threadId": thread_id, + "input": [{"type": "text", "text": prompt}], + } + + await self._request_with_settings_fallback( + ws, + "turn/start", + turn_params, + collector, + ) + while not collector.done: + message = await self._recv(ws) + await self._handle_non_response(ws, message, collector) + + ok = collector.status not in {"failed", "interrupted"} + self.last_used_monotonic = time.monotonic() + return CodexResult( + ok=ok, + final_message=collector.final_text(), + session_id=thread_id, + commands=collector.commands, + stderr="\n".join(collector.errors), + returncode=0 if ok else 1, + ) + + async def _ensure_connection(self, collector: _TurnCollector) -> Any: + if self._ws is not None: + return self._ws + self._ws = await websockets.connect( + self.url, + open_timeout=15, + ping_interval=20, + max_size=16 * 1024 * 1024, + ) + await self._initialize(self._ws, collector) + return self._ws + + +def _thread_sandbox(sandbox: str) -> str: + return { + "read-only": "readOnly", + "workspace-write": "workspaceWrite", + "danger-full-access": "dangerFullAccess", + }.get(sandbox, "readOnly") diff --git a/app/bot.py b/app/bot.py new file mode 100644 index 0000000..a9d778e --- /dev/null +++ b/app/bot.py @@ -0,0 +1,1421 @@ +from __future__ import annotations + +import asyncio +import json +import re +from pathlib import Path +from typing import Any + +from telegram import InlineKeyboardButton, InlineKeyboardMarkup, Update +from telegram.constants import ChatAction +from telegram.ext import ( + Application, + CallbackQueryHandler, + CommandHandler, + ContextTypes, + MessageHandler, + filters, +) + +from app.app_server_runner import AppServerRunner, PersistentAppServerRunner +from app.codex_runner import CodexResult, CodexRunner +from app.config import Config +from app.state import ChatState, StateStore, format_personal_record, format_schedule + + +HELP_TEXT = """Telegram-Codex 연결 브리지 + +일반 메시지는 그대로 Codex app-server로 전달됩니다. + +명령어: +/whoami - 내 Telegram user_id/chat_id 확인 +/repo <경로> - /workspaces 아래 작업 폴더 설정 +/ask <내용> - 내용을 Codex로 전달 +/new - 이 채팅의 Codex 세션 초기화 +/status - 브리지 상태 확인 +/settings - 모델, 추론강도, 속도 설정 확인 +/set_model - 모델 변경 +/set_effort - 추론강도 변경 +/set_speed - 속도 모드 변경 +/fast - 고속 모드 바로 변경/확인 +/models - 모델 선택 버튼 열기 +/schedules - 등록된 스케줄 목록/삭제 버튼 열기 +/records - 개인 기록 목록/삭제 버튼 열기 +/todos - 할 일 목록 열기 +/purchases - 구매 후보 목록 열기 +/birthdays - 생일 목록 열기 +/sessions - app-server 세션 목록 확인 +/sessions archived - 보관된 세션 목록 확인 +/use_session - 기존 세션으로 전환 +/archive_session [thread_id] - 세션 보관 +/unarchive_session - 보관된 세션 복원 +/help - 도움말 보기 +""" + +EFFORT_VALUES = {"minimal", "low", "medium", "high", "xhigh"} +RESET_WORDS = {"default", "reset", "unset", "기본", "초기화"} +MODEL_ALIASES = { + "spark": "gpt-5.3-codex-spark", + "스파크": "gpt-5.3-codex-spark", + "mini": "gpt-5.4-mini", + "미니": "gpt-5.4-mini", +} +SPEED_ALIASES = { + "standard": "standard", + "normal": "standard", + "off": "standard", + "보통": "standard", + "기본": "standard", + "fast": "fast", + "on": "fast", + "고속": "fast", +} +RECORD_CATEGORY_LABELS = { + "work_log": "업무이력", + "todo": "할 일", + "purchase": "구매", + "birthday": "생일", + "contact": "연락처", + "note": "메모", + "preference": "취향/설정", +} + + +class TelegramCodexBot: + def __init__(self, config: Config) -> None: + self.config = config + self.state = StateStore(config.state_path) + self.exec_runner = CodexRunner( + codex_bin=config.codex_bin, + timeout_seconds=config.codex_timeout_seconds, + resume_json_flag_style=config.resume_json_flag_style, + ) + self.management_runner = AppServerRunner( + url=config.app_server_url, + timeout_seconds=config.codex_timeout_seconds, + model=config.app_server_model, + effort=config.app_server_reasoning_effort, + service_tier=_service_tier_for_speed(config.app_server_speed_mode), + ) + self.chat_clients: dict[int, PersistentAppServerRunner] = {} + self.model_menus: dict[int, list[str]] = {} + self.scheduler_task: asyncio.Task | None = None + + def build_app(self) -> Application: + app = ( + Application.builder() + .token(self.config.telegram_bot_token) + .post_init(self.post_init) + .post_shutdown(self.post_shutdown) + .build() + ) + app.add_handler(CommandHandler("start", self.start)) + app.add_handler(CommandHandler("help", self.start)) + app.add_handler(CommandHandler("whoami", self.whoami)) + app.add_handler(CommandHandler("repo", self.repo)) + app.add_handler(CommandHandler("ask", self.ask)) + app.add_handler(CommandHandler("new", self.new)) + app.add_handler(CommandHandler("status", self.status)) + app.add_handler(CommandHandler("settings", self.settings)) + app.add_handler(CommandHandler("set_model", self.set_model)) + app.add_handler(CommandHandler("set_effort", self.set_effort)) + app.add_handler(CommandHandler("set_speed", self.set_speed)) + app.add_handler(CommandHandler("fast", self.fast)) + app.add_handler(CommandHandler("models", self.models)) + app.add_handler(CommandHandler("schedules", self.schedules)) + app.add_handler(CommandHandler("records", self.records)) + app.add_handler(CommandHandler("todos", self.todos)) + app.add_handler(CommandHandler("purchases", self.purchases)) + app.add_handler(CommandHandler("birthdays", self.birthdays)) + app.add_handler(CommandHandler("sessions", self.sessions)) + app.add_handler(CommandHandler("use_session", self.use_session)) + app.add_handler(CommandHandler("archive_session", self.archive_session)) + app.add_handler(CommandHandler("unarchive_session", self.unarchive_session)) + app.add_handler(CallbackQueryHandler(self.callback)) + if self.config.allow_plain_text: + app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, self.plain_text)) + return app + + async def post_init(self, app: Application) -> None: + self.scheduler_task = app.create_task(self.schedule_loop(app)) + + async def post_shutdown(self, app: Application) -> None: + if self.scheduler_task: + self.scheduler_task.cancel() + await asyncio.gather( + *(client.close() for client in self.chat_clients.values()), + return_exceptions=True, + ) + + def _is_allowed(self, update: Update) -> bool: + if not self.config.allowed_user_ids: + return False + user = update.effective_user + return bool(user and user.id in self.config.allowed_user_ids) + + async def _guard(self, update: Update) -> bool: + if self._is_allowed(update): + return True + await self._reply(update, "허용되지 않은 Telegram 사용자입니다. /whoami로 user_id를 확인해 주세요.") + return False + + async def _reply( + self, + update: Update, + text: str, + reply_markup: InlineKeyboardMarkup | None = None, + ) -> None: + if not update.effective_message: + return + chunks = _split_message(text) + for index, chunk in enumerate(chunks): + await update.effective_message.reply_text( + chunk, + reply_markup=reply_markup if index == len(chunks) - 1 else None, + ) + + async def _reply_or_edit( + self, + update: Update, + text: str, + reply_markup: InlineKeyboardMarkup | None = None, + ) -> None: + query = update.callback_query + if query and query.message: + try: + await query.edit_message_text(text, reply_markup=reply_markup) + return + except Exception: + await query.message.reply_text(text, reply_markup=reply_markup) + return + if update.effective_message: + await update.effective_message.reply_text(text, reply_markup=reply_markup) + + async def start(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + await self._reply_or_edit(update, HELP_TEXT, self._main_keyboard()) + + async def whoami(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + user = update.effective_user + chat = update.effective_chat + await self._reply( + update, + f"user_id={user.id if user else 'unknown'}\nchat_id={chat.id if chat else 'unknown'}", + ) + + async def repo(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + if not context.args: + await self._reply(update, "사용법: /repo <작업폴더 경로>") + return + + requested = " ".join(context.args).strip() + try: + repo_path = self._resolve_workspace_path(requested) + except ValueError as exc: + await self._reply(update, str(exc)) + return + + if not repo_path.exists(): + await self._reply(update, f"경로가 존재하지 않습니다: {repo_path}") + return + + self.state.set_repo(update.effective_chat.id, repo_path) + await self._close_client_for_chat(update.effective_chat.id) + await self._reply(update, f"작업 폴더를 설정했습니다:\n{repo_path}\nCodex 세션을 초기화했습니다.") + + async def ask(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + prompt = _args_text(context) + if not prompt: + await self._reply(update, "사용법: /ask <요청 내용>") + return + await self._run_codex_from_update(update, prompt) + + async def new(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + chat_id = update.effective_chat.id + self.state.clear_session(chat_id) + await self._close_client_for_chat(chat_id) + await self._reply(update, "이 Telegram 채팅의 Codex 세션을 초기화했습니다.") + + async def status(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + chat_id = update.effective_chat.id + chat_state = self.state.get_chat(chat_id) + client = self.chat_clients.get(chat_id) + schedules = self.state.list_schedules(chat_id=chat_id) + persistent_status = "연결 안 됨" + if client and client.is_connected: + persistent_status = f"연결됨, 유휴 {int(client.idle_seconds())}초" + await self._reply( + update, + "\n".join( + [ + f"작업폴더={chat_state.repo_path or '(설정 안 됨)'}", + f"백엔드={chat_state.session_backend or '(새 세션)'}", + f"세션ID={chat_state.session_id or '(새 세션)'}", + f"선호백엔드={self.config.codex_backend}", + f"persistent_ws={persistent_status}", + f"persistent_ws_유휴제한={self.config.persistent_ws_idle_seconds}초", + f"스케줄={len(schedules)}개", + self._format_runtime_settings(), + ] + ), + self._main_keyboard(), + ) + + async def settings(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + await self._reply_or_edit(update, self._format_runtime_settings(), self._settings_keyboard()) + + async def set_model(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + raw = _args_text(context) + if not raw: + await self._reply(update, "사용법: /set_model ") + return + + normalized = raw.strip() + lowered = normalized.lower() + if lowered in RESET_WORDS: + self.state.update_runtime_settings({"model": None}) + action = "모델을 환경 기본값으로 되돌렸습니다." + else: + model = await self._resolve_model_input(lowered, normalized) + self.state.update_runtime_settings({"model": model}) + action = f"모델을 {model}(으)로 설정했습니다." + await self._close_all_clients() + await self._reply(update, f"{action}\n\n{self._format_runtime_settings()}") + + async def set_effort(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + raw = _args_text(context).strip().lower() + if not raw: + await self._reply(update, "사용법: /set_effort ") + return + if raw in RESET_WORDS: + self.state.update_runtime_settings({"effort": None}) + action = "추론강도를 환경 기본값으로 되돌렸습니다." + elif raw in EFFORT_VALUES: + self.state.update_runtime_settings({"effort": raw}) + action = f"추론강도를 {raw}(으)로 설정했습니다." + else: + await self._reply(update, "사용 가능한 값: minimal, low, medium, high, xhigh, default") + return + await self._close_all_clients() + await self._reply(update, f"{action}\n\n{self._format_runtime_settings()}") + + async def set_speed(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + raw = _args_text(context).strip().lower() + if not raw: + await self._reply(update, "사용법: /set_speed ") + return + if raw in RESET_WORDS: + self.state.update_runtime_settings({"speed_mode": None}) + action = "속도 모드를 환경 기본값으로 되돌렸습니다." + else: + mode = SPEED_ALIASES.get(raw) + if not mode: + await self._reply(update, "사용 가능한 값: standard, fast, 보통, 고속, default") + return + self.state.update_runtime_settings({"speed_mode": mode}) + action = f"속도 모드를 {_display_speed(mode)}(으)로 설정했습니다." + await self._close_all_clients() + await self._reply(update, f"{action}\n\n{self._format_runtime_settings()}") + + async def fast(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + raw = _args_text(context).strip().lower() + if not raw or raw == "status": + await self._reply(update, self._format_runtime_settings()) + return + if raw not in {"on", "off"}: + await self._reply(update, "사용법: /fast ") + return + mode = "fast" if raw == "on" else "standard" + self.state.update_runtime_settings({"speed_mode": mode}) + await self._close_all_clients() + await self._reply(update, f"속도 모드를 {_display_speed(mode)}(으)로 설정했습니다.\n\n{self._format_runtime_settings()}") + + async def models(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + await self._show_model_menu(update) + + async def schedules(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + await self._show_schedule_list(update) + + async def records(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + arg = context.args[0] if context.args else None + category = arg if arg in RECORD_CATEGORY_LABELS else None + query = " ".join(context.args).strip() if arg and category is None else None + await self._show_record_list(update, category=category, query=query) + + async def todos(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + await self._show_record_list(update, category="todo", status="active") + + async def purchases(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + await self._show_record_list(update, category="purchase", status="active") + + async def birthdays(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + await self._show_record_list(update, category="birthday") + + async def sessions(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + archived = bool(context.args and context.args[0].lower() in {"archived", "archive"}) + try: + result = await self.management_runner.list_threads(limit=12, archived=archived) + except Exception as exc: + await self._reply(update, f"세션 목록을 가져오지 못했습니다: {exc}") + return + text, keyboard = _format_thread_list(result, archived=archived) + await self._reply_or_edit(update, text, keyboard) + + async def use_session(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + thread_id = _args_text(context) + if not thread_id: + await self._reply(update, "사용법: /use_session ") + return + try: + resumed_id = await self._use_session_for_chat(update.effective_chat.id, thread_id) + except Exception as exc: + await self._reply(update, f"세션을 전환하지 못했습니다: {exc}") + return + await self._reply(update, f"Codex 세션을 전환했습니다.\n세션ID={resumed_id}") + + async def callback(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + query = update.callback_query + if not query: + return + await query.answer() + if not await self._guard(update): + return + data = query.data or "" + chat_id = update.effective_chat.id + try: + if data == "sessions": + result = await self.management_runner.list_threads(limit=12, archived=False) + text, keyboard = _format_thread_list(result, archived=False) + await self._reply_or_edit(update, text, keyboard) + return + if data == "sessions_archived": + result = await self.management_runner.list_threads(limit=12, archived=True) + text, keyboard = _format_thread_list(result, archived=True) + await self._reply_or_edit(update, text, keyboard) + return + if data == "settings": + await self._reply_or_edit(update, self._format_runtime_settings(), self._settings_keyboard()) + return + if data == "settings_models": + await self._show_model_menu(update) + return + if data == "settings_effort": + await self._show_effort_menu(update) + return + if data == "schedules": + await self._show_schedule_list(update) + return + if data == "records": + await self._show_record_list(update) + return + if data.startswith("records_cat:"): + await self._show_record_list(update, category=data.split(":", 1)[1]) + return + if data.startswith("record_done:"): + await self._mark_record_done_from_button(update, data.split(":", 1)[1]) + return + if data.startswith("record_delete_confirm:"): + await self._confirm_record_delete(update, data.split(":", 1)[1]) + return + if data.startswith("record_delete:"): + await self._delete_record_from_button(update, data.split(":", 1)[1]) + return + if data.startswith("schedule_delete_confirm:"): + await self._confirm_schedule_delete(update, data.split(":", 1)[1]) + return + if data.startswith("schedule_delete:"): + await self._delete_schedule_from_button(update, data.split(":", 1)[1]) + return + if data.startswith("set_speed:"): + await self._set_speed_mode_from_button(update, data.split(":", 1)[1]) + return + if data.startswith("set_effort:"): + await self._set_effort_from_button(update, data.split(":", 1)[1]) + return + if data == "model_default": + self.state.update_runtime_settings({"model": None}) + await self._close_all_clients() + await self._reply_or_edit(update, "모델을 환경 기본값으로 되돌렸습니다.\n\n" + self._format_runtime_settings(), self._settings_keyboard()) + return + if data == "model_latest": + models = await self._model_candidates(chat_id) + model = _latest_model(models) + self.state.update_runtime_settings({"model": model}) + await self._close_all_clients() + await self._reply_or_edit(update, f"모델을 최신 추천값 {model}(으)로 설정했습니다.\n\n{self._format_runtime_settings()}", self._settings_keyboard()) + return + if data.startswith("model_pick:"): + await self._set_model_from_button(update, data.split(":", 1)[1]) + return + if data.startswith("use_session:"): + resumed_id = await self._use_session_for_chat(chat_id, data.split(":", 1)[1]) + await self._reply_or_edit(update, f"Codex 세션을 전환했습니다.\n세션ID={resumed_id}") + return + if data.startswith("archive_session:"): + thread_id = data.split(":", 1)[1] + await self._archive_session_for_chat(chat_id, thread_id) + await self._reply_or_edit(update, f"세션을 보관했습니다.\n세션ID={thread_id}") + return + if data.startswith("unarchive_session:"): + thread_id = data.split(":", 1)[1] + await self.management_runner.unarchive_thread(thread_id) + await self._reply_or_edit(update, f"세션을 복원했습니다.\n세션ID={thread_id}") + return + except Exception as exc: + await self._reply_or_edit(update, f"버튼 작업을 처리하지 못했습니다: {exc}") + return + await self._reply_or_edit(update, "알 수 없는 버튼 동작입니다.") + + def _main_keyboard(self) -> InlineKeyboardMarkup: + return InlineKeyboardMarkup( + [ + [ + InlineKeyboardButton("설정", callback_data="settings"), + InlineKeyboardButton("세션", callback_data="sessions"), + ], + [ + InlineKeyboardButton("스케줄", callback_data="schedules"), + InlineKeyboardButton("기록", callback_data="records"), + ], + [ + InlineKeyboardButton("모델", callback_data="settings_models"), + InlineKeyboardButton("할 일", callback_data="records_cat:todo"), + ], + ] + ) + + def _settings_keyboard(self) -> InlineKeyboardMarkup: + settings = self._runtime_settings() + speed_mode = str(settings.get("speed_mode") or "standard") + return InlineKeyboardMarkup( + [ + [ + InlineKeyboardButton(_checked_label("보통", speed_mode == "standard"), callback_data="set_speed:standard"), + InlineKeyboardButton(_checked_label("고속", speed_mode == "fast"), callback_data="set_speed:fast"), + ], + [ + InlineKeyboardButton("모델 선택", callback_data="settings_models"), + InlineKeyboardButton("추론강도", callback_data="settings_effort"), + ], + [ + InlineKeyboardButton("세션", callback_data="sessions"), + InlineKeyboardButton("스케줄", callback_data="schedules"), + ], + [ + InlineKeyboardButton("개인 기록", callback_data="records"), + ], + ] + ) + + async def _set_speed_mode_from_button(self, update: Update, mode: str) -> None: + if mode not in {"standard", "fast"}: + await self._reply_or_edit(update, "알 수 없는 속도 모드입니다.", self._settings_keyboard()) + return + self.state.update_runtime_settings({"speed_mode": mode}) + await self._close_all_clients() + await self._reply_or_edit( + update, + f"속도 모드를 {_display_speed(mode)}(으)로 설정했습니다.\n\n{self._format_runtime_settings()}", + self._settings_keyboard(), + ) + + async def _set_effort_from_button(self, update: Update, effort: str) -> None: + if effort == "default": + self.state.update_runtime_settings({"effort": None}) + action = "추론강도를 환경 기본값으로 되돌렸습니다." + elif effort in EFFORT_VALUES: + self.state.update_runtime_settings({"effort": effort}) + action = f"추론강도를 {effort}(으)로 설정했습니다." + else: + await self._reply_or_edit(update, "알 수 없는 추론강도입니다.", self._settings_keyboard()) + return + await self._close_all_clients() + await self._reply_or_edit(update, f"{action}\n\n{self._format_runtime_settings()}", self._settings_keyboard()) + + async def _set_model_from_button(self, update: Update, index_text: str) -> None: + chat_id = update.effective_chat.id + try: + index = int(index_text) + except ValueError: + await self._reply_or_edit(update, "알 수 없는 모델 선택입니다.", self._settings_keyboard()) + return + models = self.model_menus.get(chat_id) + if not models: + models = await self._model_candidates(chat_id) + if index < 0 or index >= len(models): + await self._reply_or_edit(update, "모델 목록이 오래되었습니다. /models로 다시 열어주세요.", self._settings_keyboard()) + return + model = models[index] + self.state.update_runtime_settings({"model": model}) + await self._close_all_clients() + await self._reply_or_edit(update, f"모델을 {model}(으)로 설정했습니다.\n\n{self._format_runtime_settings()}", self._settings_keyboard()) + + async def _show_effort_menu(self, update: Update) -> None: + settings = self._runtime_settings() + current = settings.get("effort") + rows: list[list[InlineKeyboardButton]] = [] + rows.append([InlineKeyboardButton(_checked_label("기본값", current is None), callback_data="set_effort:default")]) + for effort in ["minimal", "low", "medium", "high", "xhigh"]: + rows.append([InlineKeyboardButton(_checked_label(effort, current == effort), callback_data=f"set_effort:{effort}")]) + rows.append([InlineKeyboardButton("설정으로 돌아가기", callback_data="settings")]) + await self._reply_or_edit(update, "추론강도를 선택해 주세요.", InlineKeyboardMarkup(rows)) + + async def _show_model_menu(self, update: Update) -> None: + chat_id = update.effective_chat.id + models = await self._model_candidates(chat_id) + current = self._runtime_settings().get("model") + latest = _latest_model(models) + rows: list[list[InlineKeyboardButton]] = [ + [InlineKeyboardButton(_checked_label("기본값", current is None), callback_data="model_default")], + [InlineKeyboardButton(f"최신 추천: {latest}", callback_data="model_latest")], + ] + for index, model in enumerate(models[:12]): + rows.append( + [InlineKeyboardButton(_checked_label(model, current == model), callback_data=f"model_pick:{index}")] + ) + rows.append([InlineKeyboardButton("목록 새로고침", callback_data="settings_models")]) + rows.append([InlineKeyboardButton("설정으로 돌아가기", callback_data="settings")]) + self.model_menus[chat_id] = models + await self._reply_or_edit( + update, + "모델을 선택해 주세요.\n목록은 NAS의 Codex 모델 카탈로그를 우선 사용합니다.", + InlineKeyboardMarkup(rows), + ) + + async def _model_candidates(self, chat_id: int) -> list[str]: + models = await self._discover_codex_models() + if not models: + models = self._fallback_models() + models = _sort_models(models, preferred=self.config.app_server_model) + self.model_menus[chat_id] = models + return models + + async def _resolve_model_input(self, lowered: str, original: str) -> str: + if lowered in {"latest", "최신"}: + return _latest_model(await self._model_candidates(chat_id=0)) + return MODEL_ALIASES.get(lowered, original) + + async def _discover_codex_models(self) -> list[str]: + for args in (["debug", "models"], ["debug", "models", "--bundled"]): + try: + proc = await asyncio.create_subprocess_exec( + self.config.codex_bin, + *args, + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + stdout, _stderr = await asyncio.wait_for(proc.communicate(), timeout=20) + except Exception: + continue + if proc.returncode != 0: + continue + try: + payload = json.loads(stdout.decode("utf-8", errors="replace")) + except json.JSONDecodeError: + continue + models = _extract_model_ids(payload) + if models: + return models + return [] + + def _fallback_models(self) -> list[str]: + return _dedupe( + [ + self.config.app_server_model or "", + "gpt-5.5", + "gpt-5.4", + "gpt-5.4-mini", + "gpt-5.3-codex-spark", + ] + ) + + async def _show_schedule_list(self, update: Update) -> None: + chat_id = update.effective_chat.id + schedules = self.state.list_schedules(chat_id=chat_id) + if not schedules: + await self._reply_or_edit( + update, + "등록된 스케줄이 없습니다.", + InlineKeyboardMarkup([[InlineKeyboardButton("설정으로 돌아가기", callback_data="settings")]]), + ) + return + + lines = ["등록된 스케줄:"] + rows: list[list[InlineKeyboardButton]] = [] + for schedule in schedules[:20]: + schedule_id = int(schedule.get("id") or 0) + lines.append(_format_schedule_display(schedule)) + rows.append( + [ + InlineKeyboardButton( + f"삭제 #{schedule_id}", + callback_data=f"schedule_delete_confirm:{schedule_id}", + ) + ] + ) + rows.append([InlineKeyboardButton("새로고침", callback_data="schedules")]) + rows.append([InlineKeyboardButton("설정으로 돌아가기", callback_data="settings")]) + await self._reply_or_edit(update, "\n".join(lines), InlineKeyboardMarkup(rows)) + + async def _confirm_schedule_delete(self, update: Update, schedule_id_text: str) -> None: + try: + schedule_id = int(schedule_id_text) + except ValueError: + await self._reply_or_edit(update, "알 수 없는 스케줄입니다.") + return + schedule = self.state.get_schedule(schedule_id) + if not schedule or int(schedule.get("chat_id", 0)) != update.effective_chat.id: + await self._reply_or_edit(update, "해당 스케줄을 찾을 수 없습니다.") + return + rows = [ + [ + InlineKeyboardButton("삭제 확정", callback_data=f"schedule_delete:{schedule_id}"), + InlineKeyboardButton("취소", callback_data="schedules"), + ] + ] + await self._reply_or_edit( + update, + f"이 스케줄을 삭제할까요?\n{_format_schedule_display(schedule)}", + InlineKeyboardMarkup(rows), + ) + + async def _delete_schedule_from_button(self, update: Update, schedule_id_text: str) -> None: + try: + schedule_id = int(schedule_id_text) + except ValueError: + await self._reply_or_edit(update, "알 수 없는 스케줄입니다.") + return + removed = self.state.remove_schedule(schedule_id, chat_id=update.effective_chat.id) + if not removed: + await self._reply_or_edit(update, "해당 스케줄을 찾을 수 없습니다.") + return + await self._reply_or_edit( + update, + f"스케줄을 삭제했습니다.\n{_format_schedule_display(removed)}", + InlineKeyboardMarkup([[InlineKeyboardButton("스케줄 목록", callback_data="schedules")]]), + ) + + async def _show_record_list( + self, + update: Update, + category: str | None = None, + status: str | None = None, + query: str | None = None, + ) -> None: + chat_id = update.effective_chat.id + records = self.state.list_personal_records( + chat_id=chat_id, + category=category, + status=status, + query=query, + limit=20, + ) + if not records: + rows = self._record_category_rows() + rows.append([InlineKeyboardButton("설정으로 돌아가기", callback_data="settings")]) + title = _record_list_title(category, status, query) + await self._reply_or_edit(update, f"{title}: 없음", InlineKeyboardMarkup(rows)) + return + + lines = [_record_list_title(category, status, query) + ":"] + rows: list[list[InlineKeyboardButton]] = [] + for record in records: + record_id = int(record.get("id") or 0) + lines.append(_format_record_display(record)) + row: list[InlineKeyboardButton] = [] + if record.get("status", "active") == "active" and record.get("category") in {"todo", "purchase"}: + row.append(InlineKeyboardButton(f"완료 #{record_id}", callback_data=f"record_done:{record_id}")) + row.append(InlineKeyboardButton(f"삭제 #{record_id}", callback_data=f"record_delete_confirm:{record_id}")) + rows.append(row) + rows.extend(self._record_category_rows()) + rows.append([InlineKeyboardButton("새로고침", callback_data=f"records_cat:{category}" if category else "records")]) + rows.append([InlineKeyboardButton("설정으로 돌아가기", callback_data="settings")]) + await self._reply_or_edit(update, "\n".join(lines), InlineKeyboardMarkup(rows)) + + def _record_category_rows(self) -> list[list[InlineKeyboardButton]]: + return [ + [ + InlineKeyboardButton("업무이력", callback_data="records_cat:work_log"), + InlineKeyboardButton("할 일", callback_data="records_cat:todo"), + ], + [ + InlineKeyboardButton("구매", callback_data="records_cat:purchase"), + InlineKeyboardButton("생일", callback_data="records_cat:birthday"), + ], + [ + InlineKeyboardButton("메모", callback_data="records_cat:note"), + InlineKeyboardButton("연락처", callback_data="records_cat:contact"), + ], + ] + + async def _mark_record_done_from_button(self, update: Update, record_id_text: str) -> None: + try: + record_id = int(record_id_text) + except ValueError: + await self._reply_or_edit(update, "알 수 없는 기록입니다.") + return + updated = self.state.update_personal_record( + record_id, + {"status": "done"}, + chat_id=update.effective_chat.id, + ) + if not updated: + await self._reply_or_edit(update, "해당 기록을 찾을 수 없습니다.") + return + await self._reply_or_edit( + update, + f"완료로 표시했습니다.\n{_format_record_display(updated)}", + InlineKeyboardMarkup([[InlineKeyboardButton("기록 목록", callback_data="records")]]), + ) + + async def _confirm_record_delete(self, update: Update, record_id_text: str) -> None: + try: + record_id = int(record_id_text) + except ValueError: + await self._reply_or_edit(update, "알 수 없는 기록입니다.") + return + record = self.state.get_personal_record(record_id) + if not record or int(record.get("chat_id", 0)) != update.effective_chat.id: + await self._reply_or_edit(update, "해당 기록을 찾을 수 없습니다.") + return + rows = [ + [ + InlineKeyboardButton("삭제 확정", callback_data=f"record_delete:{record_id}"), + InlineKeyboardButton("취소", callback_data="records"), + ] + ] + await self._reply_or_edit( + update, + f"이 기록을 삭제할까요?\n{_format_record_display(record)}", + InlineKeyboardMarkup(rows), + ) + + async def _delete_record_from_button(self, update: Update, record_id_text: str) -> None: + try: + record_id = int(record_id_text) + except ValueError: + await self._reply_or_edit(update, "알 수 없는 기록입니다.") + return + removed = self.state.remove_personal_record(record_id, chat_id=update.effective_chat.id) + if not removed: + await self._reply_or_edit(update, "해당 기록을 찾을 수 없습니다.") + return + await self._reply_or_edit( + update, + f"기록을 삭제했습니다.\n{_format_record_display(removed)}", + InlineKeyboardMarkup([[InlineKeyboardButton("기록 목록", callback_data="records")]]), + ) + + async def archive_session(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + chat_state = self.state.get_chat(update.effective_chat.id) + thread_id = _args_text(context) or chat_state.session_id + if not thread_id: + await self._reply(update, "사용법: /archive_session [thread_id]") + return + try: + await self._archive_session_for_chat(update.effective_chat.id, thread_id) + except Exception as exc: + await self._reply(update, f"세션을 보관하지 못했습니다: {exc}") + return + await self._reply(update, f"세션을 보관했습니다.\n세션ID={thread_id}") + + async def unarchive_session(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + thread_id = _args_text(context) + if not thread_id: + await self._reply(update, "사용법: /unarchive_session ") + return + try: + await self.management_runner.unarchive_thread(thread_id) + except Exception as exc: + await self._reply(update, f"세션을 복원하지 못했습니다: {exc}") + return + await self._reply(update, f"세션을 복원했습니다.\n세션ID={thread_id}") + + async def plain_text(self, update: Update, context: ContextTypes.DEFAULT_TYPE) -> None: + if not await self._guard(update): + return + text = update.effective_message.text if update.effective_message else "" + text = text.strip() + if not text: + return + await self._run_codex_from_update(update, text) + + async def _run_codex_from_update(self, update: Update, prompt: str) -> None: + chat_id = update.effective_chat.id + if update.effective_chat: + await update.effective_chat.send_action(ChatAction.TYPING) + result, fallback_note = await self._run_codex_for_chat( + chat_id=chat_id, + prompt=prompt, + sandbox=self.config.readonly_sandbox, + update_chat_state=True, + ) + await self._reply_result(update, result, fallback_note) + + async def _reply_result(self, update: Update, result: CodexResult, fallback_note: str = "") -> None: + if result.ok: + text = result.final_message + if result.stderr and "runtime setting overrides" in result.stderr: + text = f"브리지 참고: {result.stderr}\n\n{text}" + if fallback_note: + text = f"{fallback_note}\n\n{text}" + await self._reply(update, text) + return + + stderr_text = f"\n\nstderr:\n{result.stderr[-1000:]}" if result.stderr else "" + note_text = f"{fallback_note}\n\n" if fallback_note else "" + await self._reply(update, f"실패(returncode={result.returncode})\n\n{note_text}{result.final_message}{stderr_text}") + + async def _run_codex_for_chat( + self, + chat_id: int, + prompt: str, + sandbox: str, + update_chat_state: bool, + session_id_override: str | None = None, + ) -> tuple[CodexResult, str]: + chat_state = self.state.get_chat(chat_id) + cwd = Path(chat_state.repo_path) if chat_state.repo_path else self.config.workspace_root + if not cwd.exists(): + return ( + CodexResult( + ok=False, + final_message=f"작업 폴더가 존재하지 않습니다: {cwd}", + session_id=chat_state.session_id, + returncode=1, + ), + "", + ) + + bridge_prompt = self._with_bridge_context( + user_prompt=prompt, + chat_id=chat_id, + chat_state=chat_state, + cwd=cwd, + ) + temp_state = ChatState( + repo_path=chat_state.repo_path, + session_id=session_id_override or chat_state.session_id, + session_backend="app-server" if session_id_override else chat_state.session_backend, + ) + skip_git_repo_check = chat_state.repo_path is None and sandbox == self.config.readonly_sandbox + result, backend_used, fallback_note = await self._run_with_backend( + chat_id=chat_id, + prompt=bridge_prompt, + cwd=cwd, + sandbox=sandbox, + chat_state=temp_state, + skip_git_repo_check=skip_git_repo_check, + ) + + if update_chat_state: + chat_state.session_id = result.session_id or chat_state.session_id + if result.session_id: + chat_state.session_backend = backend_used + self.state.update_chat(chat_id, chat_state) + + return result, fallback_note + + async def _run_with_backend( + self, + chat_id: int, + prompt: str, + cwd: Path, + sandbox: str, + chat_state: ChatState, + skip_git_repo_check: bool, + ) -> tuple[CodexResult, str, str]: + if self._should_use_app_server(sandbox): + app_session = _session_id_for_backend(chat_state, "app-server") + app_result = await self._client_for_chat(chat_id).run( + prompt=prompt, + cwd=cwd, + sandbox=sandbox, + session_id=app_session, + skip_git_repo_check=skip_git_repo_check, + ) + if app_result.ok or self.config.codex_backend == "app-server": + return app_result, "app-server", "" + + exec_session = _session_id_for_backend(chat_state, "exec") + exec_result = await self.exec_runner.run( + prompt=prompt, + cwd=cwd, + sandbox=sandbox, + session_id=exec_session, + skip_git_repo_check=skip_git_repo_check, + ) + note = f"app-server 실행이 실패해서 이번 요청은 exec fallback으로 처리했습니다.\napp-server 오류: {app_result.final_message}" + if not exec_result.ok: + exec_result.final_message = ( + "app-server와 exec fallback이 모두 실패했습니다.\n\n" + f"app-server: {app_result.final_message}\n\n" + f"exec: {exec_result.final_message}" + ) + return exec_result, "exec", note + + exec_session = _session_id_for_backend(chat_state, "exec") + exec_result = await self.exec_runner.run( + prompt=prompt, + cwd=cwd, + sandbox=sandbox, + session_id=exec_session, + skip_git_repo_check=skip_git_repo_check, + ) + return exec_result, "exec", "" + + async def schedule_loop(self, app: Application) -> None: + await asyncio.sleep(5) + while True: + try: + await self._close_idle_clients() + for schedule in self.state.claim_due_schedules(): + app.create_task(self._run_schedule(app, schedule)) + except Exception as exc: + print(f"schedule loop error: {exc}", flush=True) + await asyncio.sleep(max(5, self.config.schedule_poll_seconds)) + + async def _close_idle_clients(self) -> None: + closed_chat_ids: list[int] = [] + for chat_id, client in list(self.chat_clients.items()): + if await client.close_if_idle(self.config.persistent_ws_idle_seconds): + closed_chat_ids.append(chat_id) + for chat_id in closed_chat_ids: + print(f"closed idle app-server websocket for chat_id={chat_id}", flush=True) + + async def _close_client_for_chat(self, chat_id: int) -> None: + client = self.chat_clients.pop(chat_id, None) + if client is not None: + await client.close() + + async def _close_all_clients(self) -> None: + clients = list(self.chat_clients.values()) + self.chat_clients.clear() + await asyncio.gather(*(client.close() for client in clients), return_exceptions=True) + + async def _run_schedule(self, app: Application, schedule: dict) -> None: + chat_id = int(schedule["chat_id"]) + name = schedule.get("name") or f"스케줄 #{schedule.get('id')}" + prompt = _scheduled_prompt(schedule) + try: + await app.bot.send_chat_action(chat_id=chat_id, action=ChatAction.TYPING) + result, fallback_note = await self._run_codex_for_chat( + chat_id=chat_id, + prompt=prompt, + sandbox=self.config.readonly_sandbox, + update_chat_state=False, + session_id_override=schedule.get("thread_id") or None, + ) + if result.session_id and not schedule.get("thread_id"): + schedule["thread_id"] = result.session_id + self.state.update_schedule(schedule) + if result.ok: + text = f"스케줄 실행 결과: {name}\n\n{result.final_message}" + if fallback_note: + text = f"{fallback_note}\n\n{text}" + else: + text = f"스케줄 실행 실패: {name}\n\n{result.final_message}" + for chunk in _split_message(text): + await app.bot.send_message(chat_id=chat_id, text=chunk) + except Exception as exc: + await app.bot.send_message(chat_id=chat_id, text=f"스케줄 실행 실패: {name}\n\n{exc}") + + def _client_for_chat(self, chat_id: int) -> PersistentAppServerRunner: + client = self.chat_clients.get(chat_id) + if client is None: + client = PersistentAppServerRunner( + url=self.config.app_server_url, + timeout_seconds=self.config.codex_timeout_seconds, + ) + self.chat_clients[chat_id] = client + self._apply_runtime_settings(client) + return client + + async def _use_session_for_chat(self, chat_id: int, thread_id: str) -> str: + self._apply_runtime_settings(self.management_runner) + resumed_id = await self.management_runner.resume_thread(thread_id) + chat_state = self.state.get_chat(chat_id) + chat_state.session_id = resumed_id + chat_state.session_backend = "app-server" + self.state.update_chat(chat_id, chat_state) + await self._close_client_for_chat(chat_id) + return resumed_id + + async def _archive_session_for_chat(self, chat_id: int, thread_id: str) -> None: + await self.management_runner.archive_thread(thread_id) + chat_state = self.state.get_chat(chat_id) + if chat_state.session_id == thread_id: + chat_state.session_id = None + chat_state.session_backend = None + self.state.update_chat(chat_id, chat_state) + await self._close_client_for_chat(chat_id) + + def _runtime_defaults(self) -> dict[str, str | None]: + return { + "model": self.config.app_server_model, + "effort": self.config.app_server_reasoning_effort, + "speed_mode": self.config.app_server_speed_mode, + } + + def _runtime_settings(self) -> dict[str, str | None]: + settings = self.state.get_runtime_settings(self._runtime_defaults()) + speed_mode = str(settings.get("speed_mode") or "standard").lower() + if speed_mode not in {"standard", "fast"}: + speed_mode = "standard" + settings["speed_mode"] = speed_mode + return settings + + def _apply_runtime_settings(self, runner: AppServerRunner) -> None: + settings = self._runtime_settings() + runner.model = settings.get("model") + runner.effort = settings.get("effort") + runner.service_tier = _service_tier_for_speed(str(settings.get("speed_mode") or "standard")) + + def _format_runtime_settings(self) -> str: + settings = self._runtime_settings() + speed_mode = str(settings.get("speed_mode") or "standard") + service_tier = _service_tier_for_speed(speed_mode) or "standard" + return "\n".join( + [ + f"모델={settings.get('model') or '(기본값)'}", + f"추론강도={settings.get('effort') or '(기본값)'}", + f"속도={_display_speed(speed_mode)}", + f"서비스티어={service_tier}", + ] + ) + + def _with_bridge_context( + self, + user_prompt: str, + chat_id: int, + chat_state: ChatState, + cwd: Path, + ) -> str: + thread_id = chat_state.session_id or "(not assigned yet)" + schedules = self.state.list_schedules(chat_id=chat_id) + schedule_lines = "\n".join(f"- {format_schedule(schedule)}" for schedule in schedules[:12]) + if not schedule_lines: + schedule_lines = "- none" + records = self.state.list_personal_records(chat_id=chat_id, limit=8) + record_lines = "\n".join(f"- {format_personal_record(record)}" for record in records) + if not record_lines: + record_lines = "- none" + runtime_settings = self._format_runtime_settings() + return ( + "[Telegram bridge context]\n" + "You are connected to the user through a private Telegram bridge.\n" + f"Default Telegram chat_id for replies, schedules, and notifications: {chat_id}\n" + f"Current Codex app-server thread_id: {thread_id}\n" + f"Current working directory: {cwd}\n" + "Current runtime settings:\n" + f"{runtime_settings}\n" + "If the user asks to schedule, remind, monitor, or notify later, use the " + "telegram_create_schedule MCP tool. Use the default chat_id above and do not ask " + "for a destination unless the user explicitly wants a different chat. The schedule " + "tool's prompt field should describe only the task to perform at run time.\n" + "If the user asks to list or delete schedules, use telegram_list_schedules or " + "telegram_delete_schedule.\n" + "If the user asks you to remember, save, log, manage, or later recall personal " + "assistant information, use the personal record MCP tools. Categories are: " + "work_log, todo, purchase, birthday, contact, note, preference. Store durable " + "facts such as project history, todos, purchase candidates, birthdays, contacts, " + "and personal preferences when the user explicitly provides them or asks to save them. " + "Use telegram_list_personal_records before answering questions that depend on saved " + "personal history or personal assistant records.\n" + "Currently registered bridge schedules for this chat:\n" + f"{schedule_lines}\n" + "Recent personal assistant records for this chat:\n" + f"{record_lines}\n" + "[/Telegram bridge context]\n\n" + "User message:\n" + f"{user_prompt}" + ) + + def _should_use_app_server(self, sandbox: str) -> bool: + if self.config.codex_backend == "exec": + return False + return sandbox == self.config.readonly_sandbox + + def _resolve_workspace_path(self, requested: str) -> Path: + root = self.config.workspace_root.resolve() + path = Path(requested) + if not path.is_absolute(): + path = root / requested + resolved = path.resolve() + if resolved != root and root not in resolved.parents: + raise ValueError(f"작업 경로는 workspace root 안에 있어야 합니다: {root}") + return resolved + + +def _session_id_for_backend(chat_state: ChatState, backend: str) -> str | None: + if not chat_state.session_id: + return None + if chat_state.session_backend == backend: + return chat_state.session_id + if chat_state.session_backend is None and backend == "exec": + return chat_state.session_id + return None + + +def _service_tier_for_speed(speed_mode: str) -> str | None: + return "fast" if speed_mode == "fast" else None + + +def _display_speed(speed_mode: str) -> str: + return "fast(고속)" if speed_mode == "fast" else "standard(보통)" + + +def _args_text(context: ContextTypes.DEFAULT_TYPE) -> str: + return " ".join(context.args).strip() + + +def _checked_label(label: str, checked: bool) -> str: + return f"✓ {label}" if checked else label + + +def _extract_model_ids(payload: Any) -> list[str]: + candidates: list[str] = [] + + def visit(value: Any, key: str = "") -> None: + if isinstance(value, dict): + for nested_key, nested_value in value.items(): + visit(nested_value, str(nested_key).lower()) + return + if isinstance(value, list): + for item in value: + visit(item, key) + return + if not isinstance(value, str): + return + if key in {"id", "model", "model_id", "name", "slug"} or _looks_like_model_id(value): + if _looks_like_model_id(value): + candidates.append(value.strip()) + + visit(payload) + return _dedupe(candidates) + + +def _looks_like_model_id(value: str) -> bool: + value = value.strip() + if not value or len(value) > 80 or " " in value: + return False + lowered = value.lower() + if lowered.startswith(("gpt-", "o1", "o3", "o4")): + return True + return "codex" in lowered and re.match(r"^[a-z0-9._:/-]+$", lowered) is not None + + +def _dedupe(values: list[str]) -> list[str]: + seen: set[str] = set() + result: list[str] = [] + for value in values: + value = value.strip() + if not value or value in seen: + continue + seen.add(value) + result.append(value) + return result + + +def _sort_models(models: list[str], preferred: str | None) -> list[str]: + latest = _latest_model(models) + + def key(model: str) -> tuple[int, tuple[int, ...], str]: + if preferred and model == preferred: + return (0, (), model) + if model == latest: + return (1, (), model) + if "mini" not in model and "spark" not in model and model.startswith("gpt-"): + return (2, tuple(-part for part in _model_version_parts(model)), model) + if "mini" in model: + return (3, tuple(-part for part in _model_version_parts(model)), model) + if "spark" in model or "codex" in model: + return (4, tuple(-part for part in _model_version_parts(model)), model) + return (5, tuple(-part for part in _model_version_parts(model)), model) + + return sorted(_dedupe(models), key=key) + + +def _latest_model(models: list[str]) -> str: + candidates = [ + model + for model in models + if model.startswith("gpt-") and "mini" not in model and "spark" not in model + ] + if not candidates: + candidates = [model for model in models if model.startswith("gpt-")] + if not candidates: + return models[0] if models else "gpt-5.5" + return max(candidates, key=lambda model: (_model_version_parts(model), model)) + + +def _model_version_parts(model: str) -> list[int]: + return [int(part) for part in re.findall(r"\d+", model)] + + +def _format_schedule_display(schedule: dict) -> str: + schedule_type = schedule.get("schedule_type", "daily") + if schedule_type == "interval": + cadence = f"{int(schedule.get('interval_minutes') or 60)}분마다" + else: + cadence = ( + f"매일 {int(schedule.get('hour') or 8):02}:" + f"{int(schedule.get('minute') or 0):02} {schedule.get('timezone') or 'Asia/Seoul'}" + ) + status_map = { + "active": "활성", + "paused": "일시정지", + "deleted": "삭제됨", + } + status = status_map.get(str(schedule.get("status", "active")), str(schedule.get("status", "active"))) + name = schedule.get("name") or "스케줄 작업" + return f"#{schedule.get('id')} {name} - {cadence} - {status}" + + +def _format_record_display(record: dict) -> str: + category = str(record.get("category") or "note") + category_label = RECORD_CATEGORY_LABELS.get(category, category) + status_map = { + "active": "진행중", + "done": "완료", + "cancelled": "취소", + "archived": "보관", + } + status = status_map.get(str(record.get("status") or "active"), str(record.get("status") or "active")) + title = record.get("title") or "(제목 없음)" + date_bits = [] + for label, key in (("일자", "date"), ("시작", "start_date"), ("종료", "end_date"), ("기한", "due_date")): + if record.get(key): + date_bits.append(f"{label}:{record[key]}") + tags = record.get("tags") or [] + tag_text = f" #{' #'.join(str(tag) for tag in tags)}" if tags else "" + date_text = f" ({', '.join(date_bits)})" if date_bits else "" + return f"#{record.get('id')} [{category_label}] {title} - {status}{date_text}{tag_text}" + + +def _record_list_title(category: str | None, status: str | None, query: str | None) -> str: + if query: + return f"개인 기록 검색: {query}" + if category: + label = RECORD_CATEGORY_LABELS.get(category, category) + if status == "active": + return f"{label} 목록" + return f"{label} 기록" + return "개인 기록" + + +def _format_thread_list(result: dict, archived: bool) -> tuple[str, InlineKeyboardMarkup | None]: + threads = result.get("data") or result.get("threads") or result.get("items") or [] + if not isinstance(threads, list) or not threads: + label = "보관된 세션" if archived else "세션" + keyboard = InlineKeyboardMarkup( + [[InlineKeyboardButton("새로고침", callback_data="sessions_archived" if archived else "sessions")]] + ) + return f"{label}: 없음", keyboard + + lines = ["보관된 세션:" if archived else "세션:"] + buttons: list[list[InlineKeyboardButton]] = [] + for index, thread in enumerate(threads, start=1): + if not isinstance(thread, dict): + continue + thread_id = str(thread.get("id") or thread.get("threadId") or "(unknown)") + name = _readable_thread_field(thread.get("name") or thread.get("title"), "(untitled)") + status = _readable_thread_field(thread.get("status"), "unknown") + cwd = thread.get("cwd") or thread.get("workingDirectory") or "" + updated = thread.get("updatedAt") or thread.get("updated_at") or "" + lines.append(f"{index}. {name} [{status}]") + lines.append(f" id={thread_id}") + if cwd: + lines.append(f" cwd={cwd}") + if updated: + lines.append(f" updated={updated}") + if thread_id != "(unknown)": + action = "복원" if archived else "보관" + action_callback = "unarchive_session" if archived else "archive_session" + buttons.append( + [ + InlineKeyboardButton(f"사용 {index}", callback_data=f"use_session:{thread_id}"), + InlineKeyboardButton(f"{action} {index}", callback_data=f"{action_callback}:{thread_id}"), + ] + ) + + next_cursor = result.get("nextCursor") or result.get("next_cursor") + if next_cursor: + lines.append(f"\nnext_cursor={next_cursor}") + buttons.append( + [InlineKeyboardButton("새로고침", callback_data="sessions_archived" if archived else "sessions")] + ) + if not archived: + buttons.append([InlineKeyboardButton("보관된 세션", callback_data="sessions_archived")]) + else: + buttons.append([InlineKeyboardButton("활성 세션", callback_data="sessions")]) + return "\n".join(lines), InlineKeyboardMarkup(buttons) + + +def _readable_thread_field(value: object, fallback: str) -> str: + if isinstance(value, str) and value.strip(): + return value.strip() + if isinstance(value, dict): + for key in ("text", "name", "title", "status", "type"): + nested = value.get(key) + if isinstance(nested, str) and nested.strip(): + return nested.strip() + return fallback + + +def _scheduled_prompt(schedule: dict) -> str: + return ( + "This is a scheduled Telegram bridge task. Run the task now and produce a concise " + "Telegram-ready response. Do not create another schedule unless the task explicitly " + "asks you to change scheduling.\n\n" + f"Schedule: {format_schedule(schedule)}\n\n" + f"Task:\n{schedule.get('prompt') or ''}" + ) + + +def _split_message(text: str, limit: int = 3500) -> list[str]: + if len(text) <= limit: + return [text] + chunks: list[str] = [] + remaining = text + while remaining: + chunks.append(remaining[:limit]) + remaining = remaining[limit:] + return chunks diff --git a/app/bridge_mcp.py b/app/bridge_mcp.py new file mode 100644 index 0000000..cb98bf3 --- /dev/null +++ b/app/bridge_mcp.py @@ -0,0 +1,476 @@ +from __future__ import annotations + +import json +import os +import sys +from pathlib import Path +from typing import Any + +from app.state import StateStore, format_personal_record, format_schedule + + +PROTOCOL_VERSION = "2024-11-05" + + +def main() -> None: + server = BridgeMcpServer( + StateStore(Path(os.getenv("BOT_STATE_PATH", "/app/data/state.json"))), + default_timezone=os.getenv("BOT_TIMEZONE", "Asia/Seoul"), + ) + server.run() + + +class BridgeMcpServer: + def __init__(self, state: StateStore, default_timezone: str) -> None: + self.state = state + self.default_timezone = default_timezone + + def run(self) -> None: + for line in sys.stdin: + line = line.strip() + if not line: + continue + try: + message = json.loads(line) + response = self.handle(message) + except Exception as exc: + response = _error(None, -32603, str(exc)) + if response is not None: + sys.stdout.write(json.dumps(response, ensure_ascii=False) + "\n") + sys.stdout.flush() + + def handle(self, message: dict[str, Any]) -> dict[str, Any] | None: + method = message.get("method") + request_id = message.get("id") + + if method == "initialize": + params = message.get("params") or {} + return _result( + request_id, + { + "protocolVersion": params.get("protocolVersion") or PROTOCOL_VERSION, + "capabilities": {"tools": {}}, + "serverInfo": { + "name": "telegram-codex-bridge", + "version": "0.1.0", + }, + "instructions": ( + "Use these tools when the user asks to schedule recurring Telegram " + "notifications, manage bridge schedules, or persist personal assistant " + "records such as work history, todos, purchase candidates, birthdays, " + "contacts, and notes. The Telegram chat_id and " + "current thread_id are provided in the bridge context inside each user turn. " + "Use that chat_id as the default notification target and do not ask for a " + "destination unless the user explicitly wants a different chat." + ), + }, + ) + + if method == "notifications/initialized": + return None + + if method == "tools/list": + return _result(request_id, {"tools": _tools()}) + + if method == "tools/call": + params = message.get("params") or {} + name = params.get("name") + arguments = params.get("arguments") or {} + try: + return _result(request_id, self.call_tool(str(name), arguments)) + except Exception as exc: + return _result( + request_id, + { + "content": [{"type": "text", "text": f"Tool failed: {exc}"}], + "isError": True, + }, + ) + + if request_id is None: + return None + return _error(request_id, -32601, f"Unknown method: {method}") + + def call_tool(self, name: str, arguments: dict[str, Any]) -> dict[str, Any]: + if name == "telegram_create_schedule": + return self.create_schedule(arguments) + if name == "telegram_list_schedules": + return self.list_schedules(arguments) + if name == "telegram_delete_schedule": + return self.delete_schedule(arguments) + if name == "telegram_save_personal_record": + return self.save_personal_record(arguments) + if name == "telegram_list_personal_records": + return self.list_personal_records(arguments) + if name == "telegram_update_personal_record": + return self.update_personal_record(arguments) + if name == "telegram_delete_personal_record": + return self.delete_personal_record(arguments) + raise ValueError(f"Unknown tool: {name}") + + def create_schedule(self, arguments: dict[str, Any]) -> dict[str, Any]: + chat_id = int(arguments["chat_id"]) + schedule_type = str(arguments.get("schedule_type") or "daily") + if schedule_type not in {"daily", "interval"}: + raise ValueError("schedule_type must be daily or interval") + + schedule: dict[str, Any] = { + "chat_id": chat_id, + "thread_id": arguments.get("thread_id") or None, + "name": str(arguments.get("name") or "Scheduled Codex task"), + "prompt": str(arguments["prompt"]), + "schedule_type": schedule_type, + "timezone": str(arguments.get("timezone") or self.default_timezone), + "status": "active", + } + if schedule_type == "interval": + schedule["interval_minutes"] = max(1, int(arguments.get("interval_minutes") or 60)) + else: + schedule["hour"] = int(arguments.get("hour") if arguments.get("hour") is not None else 8) + schedule["minute"] = int(arguments.get("minute") if arguments.get("minute") is not None else 0) + + created = self.state.add_schedule(schedule) + text = f"Created Telegram schedule: {format_schedule(created)}" + return { + "content": [{"type": "text", "text": text}], + "structuredContent": {"schedule": created}, + } + + def list_schedules(self, arguments: dict[str, Any]) -> dict[str, Any]: + chat_id = arguments.get("chat_id") + schedules = self.state.list_schedules( + chat_id=int(chat_id) if chat_id is not None else None, + include_paused=bool(arguments.get("include_paused", True)), + ) + if not schedules: + text = "No Telegram bridge schedules are registered." + else: + text = "\n".join(format_schedule(schedule) for schedule in schedules) + return { + "content": [{"type": "text", "text": text}], + "structuredContent": {"schedules": schedules}, + } + + def delete_schedule(self, arguments: dict[str, Any]) -> dict[str, Any]: + schedule_id = int(arguments["schedule_id"]) + chat_id = arguments.get("chat_id") + removed = self.state.remove_schedule( + schedule_id, + chat_id=int(chat_id) if chat_id is not None else None, + ) + if removed is None: + raise ValueError(f"schedule not found: {schedule_id}") + text = f"Deleted Telegram schedule: {format_schedule(removed)}" + return { + "content": [{"type": "text", "text": text}], + "structuredContent": {"deleted": removed}, + } + + def save_personal_record(self, arguments: dict[str, Any]) -> dict[str, Any]: + chat_id = int(arguments["chat_id"]) + category = str(arguments.get("category") or "note") + if category not in _record_categories(): + raise ValueError(f"category must be one of: {', '.join(_record_categories())}") + tags = arguments.get("tags") or [] + if not isinstance(tags, list): + tags = [str(tags)] + metadata = arguments.get("metadata") or {} + if not isinstance(metadata, dict): + metadata = {} + record = { + "chat_id": chat_id, + "category": category, + "title": str(arguments.get("title") or "Untitled"), + "content": str(arguments.get("content") or ""), + "status": str(arguments.get("status") or "active"), + "date": arguments.get("date") or None, + "start_date": arguments.get("start_date") or None, + "end_date": arguments.get("end_date") or None, + "due_date": arguments.get("due_date") or None, + "tags": [str(tag) for tag in tags], + "metadata": metadata, + } + created = self.state.add_personal_record(record) + text = f"Saved personal record: {format_personal_record(created)}" + return { + "content": [{"type": "text", "text": text}], + "structuredContent": {"record": created}, + } + + def list_personal_records(self, arguments: dict[str, Any]) -> dict[str, Any]: + chat_id = arguments.get("chat_id") + records = self.state.list_personal_records( + chat_id=int(chat_id) if chat_id is not None else None, + category=arguments.get("category") or None, + status=arguments.get("status") or None, + query=arguments.get("query") or None, + limit=int(arguments.get("limit") or 30), + ) + if not records: + text = "No personal records found." + else: + text = "\n".join(format_personal_record(record) for record in records) + return { + "content": [{"type": "text", "text": text}], + "structuredContent": {"records": records}, + } + + def update_personal_record(self, arguments: dict[str, Any]) -> dict[str, Any]: + record_id = int(arguments["record_id"]) + chat_id = arguments.get("chat_id") + updates = dict(arguments.get("updates") or {}) + if "category" in updates and updates["category"] not in _record_categories(): + raise ValueError(f"category must be one of: {', '.join(_record_categories())}") + if "tags" in updates and not isinstance(updates["tags"], list): + updates["tags"] = [str(updates["tags"])] + updated = self.state.update_personal_record( + record_id, + updates, + chat_id=int(chat_id) if chat_id is not None else None, + ) + if updated is None: + raise ValueError(f"personal record not found: {record_id}") + text = f"Updated personal record: {format_personal_record(updated)}" + return { + "content": [{"type": "text", "text": text}], + "structuredContent": {"record": updated}, + } + + def delete_personal_record(self, arguments: dict[str, Any]) -> dict[str, Any]: + record_id = int(arguments["record_id"]) + chat_id = arguments.get("chat_id") + removed = self.state.remove_personal_record( + record_id, + chat_id=int(chat_id) if chat_id is not None else None, + ) + if removed is None: + raise ValueError(f"personal record not found: {record_id}") + text = f"Deleted personal record: {format_personal_record(removed)}" + return { + "content": [{"type": "text", "text": text}], + "structuredContent": {"deleted": removed}, + } + + +def _tools() -> list[dict[str, Any]]: + return [ + { + "name": "telegram_create_schedule", + "description": ( + "Create a recurring Telegram notification schedule. Use this when the user asks " + "Codex to remind, notify, monitor, or run a recurring task through Telegram. " + "The prompt should describe only the work to run at schedule time, not the cadence." + ), + "inputSchema": { + "type": "object", + "properties": { + "chat_id": { + "type": "integer", + "description": "Telegram chat id that should receive the scheduled result.", + }, + "thread_id": { + "type": "string", + "description": "Current Codex app-server thread id, if available.", + }, + "name": {"type": "string", "description": "Short schedule name."}, + "prompt": { + "type": "string", + "description": "Task prompt to run at the scheduled time.", + }, + "schedule_type": { + "type": "string", + "enum": ["daily", "interval"], + "description": "daily for wall-clock time, interval for every N minutes.", + }, + "hour": { + "type": "integer", + "minimum": 0, + "maximum": 23, + "description": "Daily run hour in the selected timezone.", + }, + "minute": { + "type": "integer", + "minimum": 0, + "maximum": 59, + "description": "Daily run minute in the selected timezone.", + }, + "interval_minutes": { + "type": "integer", + "minimum": 1, + "description": "Interval cadence in minutes.", + }, + "timezone": { + "type": "string", + "description": "IANA timezone, usually Asia/Seoul.", + }, + }, + "required": ["chat_id", "name", "prompt", "schedule_type"], + }, + "annotations": { + "title": "Create Telegram schedule", + "readOnlyHint": False, + "destructiveHint": False, + "idempotentHint": False, + "openWorldHint": False, + }, + }, + { + "name": "telegram_list_schedules", + "description": "List Telegram bridge schedules, optionally scoped to a chat_id.", + "inputSchema": { + "type": "object", + "properties": { + "chat_id": {"type": "integer"}, + "include_paused": {"type": "boolean"}, + }, + }, + "annotations": { + "title": "List Telegram schedules", + "readOnlyHint": True, + "destructiveHint": False, + "idempotentHint": True, + "openWorldHint": False, + }, + }, + { + "name": "telegram_delete_schedule", + "description": "Delete a Telegram bridge schedule by id.", + "inputSchema": { + "type": "object", + "properties": { + "schedule_id": {"type": "integer"}, + "chat_id": { + "type": "integer", + "description": "Optional Telegram chat id safety check.", + }, + }, + "required": ["schedule_id"], + }, + "annotations": { + "title": "Delete Telegram schedule", + "readOnlyHint": False, + "destructiveHint": True, + "idempotentHint": False, + "openWorldHint": False, + }, + }, + { + "name": "telegram_save_personal_record", + "description": ( + "Save a durable personal assistant record to NAS-backed Telegram bridge storage. " + "Use for work history, todos, purchase plans, birthdays, contacts, personal notes, " + "preferences, and other facts the user explicitly wants remembered later." + ), + "inputSchema": { + "type": "object", + "properties": { + "chat_id": {"type": "integer"}, + "category": { + "type": "string", + "enum": _record_categories(), + "description": "Record category.", + }, + "title": {"type": "string", "description": "Short title."}, + "content": {"type": "string", "description": "Full details to remember."}, + "status": { + "type": "string", + "description": "active, done, cancelled, archived, or another short status.", + }, + "date": {"type": "string", "description": "Primary date, ISO YYYY-MM-DD if possible."}, + "start_date": {"type": "string", "description": "Start date, ISO YYYY-MM-DD if possible."}, + "end_date": {"type": "string", "description": "End date, ISO YYYY-MM-DD if possible."}, + "due_date": {"type": "string", "description": "Due date, ISO YYYY-MM-DD if possible."}, + "tags": {"type": "array", "items": {"type": "string"}}, + "metadata": {"type": "object", "additionalProperties": True}, + }, + "required": ["chat_id", "category", "title", "content"], + }, + "annotations": { + "title": "Save personal record", + "readOnlyHint": False, + "destructiveHint": False, + "idempotentHint": False, + "openWorldHint": False, + }, + }, + { + "name": "telegram_list_personal_records", + "description": "List or search personal assistant records stored by the Telegram bridge.", + "inputSchema": { + "type": "object", + "properties": { + "chat_id": {"type": "integer"}, + "category": {"type": "string", "enum": _record_categories()}, + "status": {"type": "string"}, + "query": {"type": "string"}, + "limit": {"type": "integer", "minimum": 1, "maximum": 100}, + }, + }, + "annotations": { + "title": "List personal records", + "readOnlyHint": True, + "destructiveHint": False, + "idempotentHint": True, + "openWorldHint": False, + }, + }, + { + "name": "telegram_update_personal_record", + "description": "Update a personal assistant record by id.", + "inputSchema": { + "type": "object", + "properties": { + "record_id": {"type": "integer"}, + "chat_id": {"type": "integer"}, + "updates": {"type": "object", "additionalProperties": True}, + }, + "required": ["record_id", "updates"], + }, + "annotations": { + "title": "Update personal record", + "readOnlyHint": False, + "destructiveHint": False, + "idempotentHint": False, + "openWorldHint": False, + }, + }, + { + "name": "telegram_delete_personal_record", + "description": "Delete a personal assistant record by id.", + "inputSchema": { + "type": "object", + "properties": { + "record_id": {"type": "integer"}, + "chat_id": {"type": "integer"}, + }, + "required": ["record_id"], + }, + "annotations": { + "title": "Delete personal record", + "readOnlyHint": False, + "destructiveHint": True, + "idempotentHint": False, + "openWorldHint": False, + }, + }, + ] + + +def _record_categories() -> list[str]: + return ["work_log", "todo", "purchase", "birthday", "contact", "note", "preference"] + + +def _result(request_id: Any, result: dict[str, Any]) -> dict[str, Any]: + return {"jsonrpc": "2.0", "id": request_id, "result": result} + + +def _error(request_id: Any, code: int, message: str) -> dict[str, Any]: + return { + "jsonrpc": "2.0", + "id": request_id, + "error": {"code": code, "message": message}, + } + + +if __name__ == "__main__": + main() diff --git a/app/codex_runner.py b/app/codex_runner.py new file mode 100644 index 0000000..8e55528 --- /dev/null +++ b/app/codex_runner.py @@ -0,0 +1,148 @@ +from __future__ import annotations + +import asyncio +import json +from dataclasses import dataclass, field +from pathlib import Path + + +@dataclass +class CodexResult: + ok: bool + final_message: str + session_id: str | None + commands: list[str] = field(default_factory=list) + stderr: str = "" + returncode: int | None = None + + +class CodexRunner: + def __init__( + self, + codex_bin: str, + timeout_seconds: int, + resume_json_flag_style: str, + ) -> None: + self.codex_bin = codex_bin + self.timeout_seconds = timeout_seconds + self.resume_json_flag_style = resume_json_flag_style + + def _command( + self, + prompt: str, + sandbox: str, + session_id: str | None, + skip_git_repo_check: bool, + ) -> list[str]: + base = [self.codex_bin, "exec", "--json", "--sandbox", sandbox] + if skip_git_repo_check: + base.append("--skip-git-repo-check") + if not session_id: + return [*base, prompt] + + if self.resume_json_flag_style == "after_resume": + return [ + self.codex_bin, + "exec", + "resume", + session_id, + "--json", + "--sandbox", + sandbox, + *(["--skip-git-repo-check"] if skip_git_repo_check else []), + prompt, + ] + return [*base, "resume", session_id, prompt] + + async def run( + self, + prompt: str, + cwd: Path, + sandbox: str, + session_id: str | None, + skip_git_repo_check: bool = False, + ) -> CodexResult: + cmd = self._command( + prompt=prompt, + sandbox=sandbox, + session_id=session_id, + skip_git_repo_check=skip_git_repo_check, + ) + proc = await asyncio.create_subprocess_exec( + *cmd, + cwd=str(cwd), + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + + try: + stdout, stderr = await asyncio.wait_for( + proc.communicate(), + timeout=self.timeout_seconds, + ) + except asyncio.TimeoutError: + proc.kill() + await proc.wait() + return CodexResult( + ok=False, + final_message=f"Codex timed out after {self.timeout_seconds} seconds.", + session_id=session_id, + ) + + parsed = self._parse_jsonl(stdout.decode("utf-8", errors="replace")) + stderr_text = stderr.decode("utf-8", errors="replace").strip() + if proc.returncode != 0: + message = parsed.final_message or stderr_text or "Codex command failed." + return CodexResult( + ok=False, + final_message=message, + session_id=parsed.session_id or session_id, + commands=parsed.commands, + stderr=stderr_text, + returncode=proc.returncode, + ) + + return CodexResult( + ok=True, + final_message=parsed.final_message or "(Codex finished without text output.)", + session_id=parsed.session_id or session_id, + commands=parsed.commands, + stderr=stderr_text, + returncode=proc.returncode, + ) + + @staticmethod + def _parse_jsonl(text: str) -> CodexResult: + session_id: str | None = None + final_messages: list[str] = [] + commands: list[str] = [] + fallback_lines: list[str] = [] + + for line in text.splitlines(): + line = line.strip() + if not line: + continue + try: + event = json.loads(line) + except json.JSONDecodeError: + fallback_lines.append(line) + continue + + if event.get("type") == "thread.started": + session_id = event.get("thread_id") or session_id + continue + + item = event.get("item") + if isinstance(item, dict): + if item.get("type") == "agent_message" and item.get("text"): + final_messages.append(str(item["text"])) + if item.get("type") == "command_execution" and item.get("command"): + commands.append(str(item["command"])) + + final_message = final_messages[-1] if final_messages else "\n".join(fallback_lines) + return CodexResult( + ok=True, + final_message=final_message, + session_id=session_id, + commands=commands, + ) diff --git a/app/config.py b/app/config.py new file mode 100644 index 0000000..d49845e --- /dev/null +++ b/app/config.py @@ -0,0 +1,96 @@ +from __future__ import annotations + +import os +from dataclasses import dataclass +from pathlib import Path + +from dotenv import load_dotenv + + +def _bool_env(name: str, default: bool) -> bool: + raw = os.getenv(name) + if raw is None: + return default + return raw.strip().lower() in {"1", "true", "yes", "y", "on"} + + +def _int_env(name: str, default: int) -> int: + raw = os.getenv(name) + if not raw: + return default + return int(raw) + + +def _id_set(name: str) -> set[int]: + raw = os.getenv(name, "") + ids: set[int] = set() + for part in raw.replace(";", ",").split(","): + part = part.strip() + if not part: + continue + ids.add(int(part)) + return ids + + +@dataclass(frozen=True) +class Config: + telegram_bot_token: str + allowed_user_ids: set[int] + state_path: Path + workspace_root: Path + codex_bin: str + codex_timeout_seconds: int + codex_backend: str + app_server_url: str + app_server_model: str | None + app_server_reasoning_effort: str | None + app_server_speed_mode: str + readonly_sandbox: str + write_sandbox: str + resume_json_flag_style: str + allow_plain_text: bool + timezone: str + schedule_poll_seconds: int + persistent_ws_idle_seconds: int + + @classmethod + def load(cls) -> "Config": + load_dotenv() + token = os.getenv("TELEGRAM_BOT_TOKEN", "").strip() + if not token: + raise RuntimeError("TELEGRAM_BOT_TOKEN is required") + + backend = os.getenv("CODEX_BACKEND", "auto").strip().lower() + if backend not in {"auto", "app-server", "exec"}: + raise RuntimeError("CODEX_BACKEND must be one of: auto, app-server, exec") + + app_server_model = os.getenv("CODEX_APP_SERVER_MODEL", "").strip() or "gpt-5.5" + app_server_reasoning_effort = ( + os.getenv("CODEX_MODEL_REASONING_EFFORT", "").strip() or "high" + ) + app_server_speed_mode = os.getenv("CODEX_SPEED_MODE", "standard").strip().lower() + if app_server_speed_mode not in {"standard", "fast"}: + raise RuntimeError("CODEX_SPEED_MODE must be one of: standard, fast") + + return cls( + telegram_bot_token=token, + allowed_user_ids=_id_set("ALLOWED_TELEGRAM_USER_IDS"), + state_path=Path(os.getenv("BOT_STATE_PATH", "/app/data/state.json")), + workspace_root=Path(os.getenv("BOT_WORKSPACE_ROOT", "/workspaces")), + codex_bin=os.getenv("CODEX_BIN", "codex"), + codex_timeout_seconds=_int_env("CODEX_TIMEOUT_SECONDS", 1800), + codex_backend=backend, + app_server_url=os.getenv("CODEX_APP_SERVER_URL", "ws://127.0.0.1:4500"), + app_server_model=app_server_model, + app_server_reasoning_effort=app_server_reasoning_effort, + app_server_speed_mode=app_server_speed_mode, + readonly_sandbox=os.getenv("CODEX_READONLY_SANDBOX", "read-only"), + write_sandbox=os.getenv("CODEX_WRITE_SANDBOX", "workspace-write"), + resume_json_flag_style=os.getenv( + "CODEX_RESUME_JSON_FLAG_STYLE", "before_resume" + ), + allow_plain_text=_bool_env("ALLOW_PLAIN_TEXT", True), + timezone=os.getenv("BOT_TIMEZONE", "Asia/Seoul"), + schedule_poll_seconds=_int_env("SCHEDULE_POLL_SECONDS", 30), + persistent_ws_idle_seconds=_int_env("PERSISTENT_WS_IDLE_SECONDS", 600), + ) diff --git a/app/ensure_codex_config.py b/app/ensure_codex_config.py new file mode 100644 index 0000000..69a1fbe --- /dev/null +++ b/app/ensure_codex_config.py @@ -0,0 +1,50 @@ +from __future__ import annotations + +import os +import re +from pathlib import Path + + +BEGIN = "# BEGIN telegram-codex-bridge managed MCP" +END = "# END telegram-codex-bridge managed MCP" + + +def main() -> None: + codex_home = Path(os.getenv("CODEX_HOME", "/root/.codex")) + codex_home.mkdir(parents=True, exist_ok=True) + config_path = codex_home / "config.toml" + state_path = os.getenv("BOT_STATE_PATH", "/app/data/state.json") + timezone = os.getenv("BOT_TIMEZONE", "Asia/Seoul") + + block = f"""{BEGIN} +[mcp_servers.telegram_bridge] +command = "python" +args = ["-m", "app.bridge_mcp"] +tool_timeout_sec = 30.0 +default_tools_approval_mode = "auto" + +[mcp_servers.telegram_bridge.env] +PYTHONPATH = "/app" +BOT_STATE_PATH = "{_toml_string(state_path)}" +BOT_TIMEZONE = "{_toml_string(timezone)}" +{END} +""" + + existing = config_path.read_text(encoding="utf-8") if config_path.exists() else "" + pattern = re.compile( + rf"{re.escape(BEGIN)}.*?{re.escape(END)}\s*", + flags=re.DOTALL, + ) + if pattern.search(existing): + updated = pattern.sub(block, existing) + else: + updated = existing.rstrip() + "\n\n" + block if existing.strip() else block + config_path.write_text(updated, encoding="utf-8") + + +def _toml_string(value: str) -> str: + return value.replace("\\", "\\\\").replace('"', '\\"') + + +if __name__ == "__main__": + main() diff --git a/app/main.py b/app/main.py new file mode 100644 index 0000000..aef24d3 --- /dev/null +++ b/app/main.py @@ -0,0 +1,15 @@ +from __future__ import annotations + +from app.bot import TelegramCodexBot +from app.config import Config + + +def main() -> None: + config = Config.load() + bot = TelegramCodexBot(config) + app = bot.build_app() + app.run_polling(allowed_updates=["message", "callback_query"]) + + +if __name__ == "__main__": + main() diff --git a/app/state.py b/app/state.py new file mode 100644 index 0000000..a75fd1f --- /dev/null +++ b/app/state.py @@ -0,0 +1,337 @@ +from __future__ import annotations + +import json +from dataclasses import dataclass, field +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any +from zoneinfo import ZoneInfo + + +@dataclass +class ChatState: + repo_path: str | None = None + session_id: str | None = None + session_backend: str | None = None + + +class StateStore: + def __init__(self, path: Path) -> None: + self.path = path + self.path.parent.mkdir(parents=True, exist_ok=True) + self._data: dict[str, Any] = { + "chats": {}, + "schedules": [], + "next_schedule_id": 1, + "personal_records": [], + "next_record_id": 1, + "runtime_settings": {}, + } + self.load() + + def load(self) -> None: + if not self.path.exists(): + self._ensure_defaults() + return + self._data = json.loads(self.path.read_text(encoding="utf-8")) + self._ensure_defaults() + + def _ensure_defaults(self) -> None: + self._data.setdefault("chats", {}) + self._data.setdefault("schedules", []) + self._data.setdefault("next_schedule_id", 1) + self._data.setdefault("personal_records", []) + self._data.setdefault("next_record_id", 1) + self._data.setdefault("runtime_settings", {}) + + def save(self) -> None: + tmp_path = self.path.with_suffix(self.path.suffix + ".tmp") + tmp_path.write_text( + json.dumps(self._data, ensure_ascii=False, indent=2), + encoding="utf-8", + ) + tmp_path.replace(self.path) + + def get_chat(self, chat_id: int) -> ChatState: + self.load() + raw = self._data.setdefault("chats", {}).setdefault(str(chat_id), {}) + return ChatState( + repo_path=raw.get("repo_path"), + session_id=raw.get("session_id"), + session_backend=raw.get("session_backend"), + ) + + def update_chat(self, chat_id: int, chat_state: ChatState) -> None: + self.load() + self._data.setdefault("chats", {})[str(chat_id)] = { + "repo_path": chat_state.repo_path, + "session_id": chat_state.session_id, + "session_backend": chat_state.session_backend, + } + self.save() + + def clear_session(self, chat_id: int) -> None: + chat = self.get_chat(chat_id) + chat.session_id = None + chat.session_backend = None + self.update_chat(chat_id, chat) + + def set_repo(self, chat_id: int, repo_path: Path) -> None: + chat = self.get_chat(chat_id) + chat.repo_path = str(repo_path) + chat.session_id = None + chat.session_backend = None + self.update_chat(chat_id, chat) + + def get_runtime_settings(self, defaults: dict[str, Any]) -> dict[str, Any]: + self.load() + settings = dict(defaults) + settings.update(self._data.setdefault("runtime_settings", {})) + return settings + + def update_runtime_settings(self, updates: dict[str, Any]) -> dict[str, Any]: + self.load() + settings = self._data.setdefault("runtime_settings", {}) + for key, value in updates.items(): + if value is None: + settings.pop(key, None) + else: + settings[key] = value + self.save() + return dict(settings) + + def add_schedule(self, schedule: dict[str, Any]) -> dict[str, Any]: + self.load() + schedule = dict(schedule) + schedule["id"] = int(self._data.setdefault("next_schedule_id", 1)) + self._data["next_schedule_id"] = schedule["id"] + 1 + schedule.setdefault("status", "active") + schedule.setdefault("created_at", _utc_now().isoformat()) + schedule.setdefault("next_run_at", compute_next_run(schedule, _utc_now()).isoformat()) + self._data.setdefault("schedules", []).append(schedule) + self.save() + return schedule + + def list_schedules(self, chat_id: int | None = None, include_paused: bool = True) -> list[dict[str, Any]]: + self.load() + schedules = list(self._data.setdefault("schedules", [])) + if chat_id is not None: + schedules = [s for s in schedules if int(s.get("chat_id", 0)) == chat_id] + if not include_paused: + schedules = [s for s in schedules if s.get("status", "active") == "active"] + return schedules + + def get_schedule(self, schedule_id: int) -> dict[str, Any] | None: + self.load() + for schedule in self._data.setdefault("schedules", []): + if int(schedule.get("id", 0)) == schedule_id: + return dict(schedule) + return None + + def update_schedule(self, schedule: dict[str, Any]) -> None: + self.load() + schedule_id = int(schedule["id"]) + schedules = self._data.setdefault("schedules", []) + for index, existing in enumerate(schedules): + if int(existing.get("id", 0)) == schedule_id: + schedules[index] = dict(schedule) + self.save() + return + raise KeyError(f"schedule not found: {schedule_id}") + + def remove_schedule(self, schedule_id: int, chat_id: int | None = None) -> dict[str, Any] | None: + self.load() + schedules = self._data.setdefault("schedules", []) + for index, schedule in enumerate(schedules): + if int(schedule.get("id", 0)) != schedule_id: + continue + if chat_id is not None and int(schedule.get("chat_id", 0)) != chat_id: + return None + removed = schedules.pop(index) + self.save() + return dict(removed) + return None + + def claim_due_schedules(self, now: datetime | None = None) -> list[dict[str, Any]]: + now = now or _utc_now() + self.load() + claimed: list[dict[str, Any]] = [] + schedules = self._data.setdefault("schedules", []) + for schedule in schedules: + if schedule.get("status", "active") != "active": + continue + next_run_at = _parse_datetime(schedule.get("next_run_at")) + if next_run_at is None or next_run_at > now: + continue + claimed.append(dict(schedule)) + schedule["last_started_at"] = now.isoformat() + schedule["next_run_at"] = compute_next_run(schedule, now).isoformat() + if claimed: + self.save() + return claimed + + def add_personal_record(self, record: dict[str, Any]) -> dict[str, Any]: + self.load() + record = dict(record) + record["id"] = int(self._data.setdefault("next_record_id", 1)) + self._data["next_record_id"] = record["id"] + 1 + record.setdefault("status", "active") + record.setdefault("created_at", _utc_now().isoformat()) + record.setdefault("updated_at", record["created_at"]) + record.setdefault("tags", []) + self._data.setdefault("personal_records", []).append(record) + self.save() + return record + + def list_personal_records( + self, + chat_id: int | None = None, + category: str | None = None, + status: str | None = None, + query: str | None = None, + limit: int = 50, + ) -> list[dict[str, Any]]: + self.load() + records = list(self._data.setdefault("personal_records", [])) + if chat_id is not None: + records = [r for r in records if int(r.get("chat_id", 0)) == chat_id] + if category: + records = [r for r in records if str(r.get("category") or "") == category] + if status: + records = [r for r in records if str(r.get("status") or "") == status] + if query: + needle = query.casefold() + records = [r for r in records if needle in _record_search_text(r).casefold()] + records.sort(key=lambda r: str(r.get("updated_at") or r.get("created_at") or ""), reverse=True) + return [dict(r) for r in records[: max(1, limit)]] + + def get_personal_record(self, record_id: int) -> dict[str, Any] | None: + self.load() + for record in self._data.setdefault("personal_records", []): + if int(record.get("id", 0)) == record_id: + return dict(record) + return None + + def update_personal_record(self, record_id: int, updates: dict[str, Any], chat_id: int | None = None) -> dict[str, Any] | None: + self.load() + records = self._data.setdefault("personal_records", []) + for index, record in enumerate(records): + if int(record.get("id", 0)) != record_id: + continue + if chat_id is not None and int(record.get("chat_id", 0)) != chat_id: + return None + updated = dict(record) + for key, value in updates.items(): + if value is None: + updated.pop(key, None) + else: + updated[key] = value + updated["updated_at"] = _utc_now().isoformat() + records[index] = updated + self.save() + return dict(updated) + return None + + def remove_personal_record(self, record_id: int, chat_id: int | None = None) -> dict[str, Any] | None: + self.load() + records = self._data.setdefault("personal_records", []) + for index, record in enumerate(records): + if int(record.get("id", 0)) != record_id: + continue + if chat_id is not None and int(record.get("chat_id", 0)) != chat_id: + return None + removed = records.pop(index) + self.save() + return dict(removed) + return None + + +def compute_next_run(schedule: dict[str, Any], after: datetime | None = None) -> datetime: + after = after or _utc_now() + if after.tzinfo is None: + after = after.replace(tzinfo=timezone.utc) + + schedule_type = schedule.get("schedule_type", "daily") + if schedule_type == "interval": + interval = max(1, int(schedule.get("interval_minutes") or 60)) + return after + timedelta(minutes=interval) + + tz_name = str(schedule.get("timezone") or "Asia/Seoul") + try: + tz = ZoneInfo(tz_name) + except Exception: + tz = ZoneInfo("Asia/Seoul") + + local_after = after.astimezone(tz) + hour = int(schedule.get("hour") or 8) + minute = int(schedule.get("minute") or 0) + candidate = local_after.replace(hour=hour, minute=minute, second=0, microsecond=0) + if candidate <= local_after: + candidate += timedelta(days=1) + return candidate.astimezone(timezone.utc) + + +def format_schedule(schedule: dict[str, Any]) -> str: + schedule_type = schedule.get("schedule_type", "daily") + if schedule_type == "interval": + cadence = f"every {int(schedule.get('interval_minutes') or 60)} minutes" + else: + cadence = ( + f"daily {int(schedule.get('hour') or 8):02}:" + f"{int(schedule.get('minute') or 0):02} {schedule.get('timezone') or 'Asia/Seoul'}" + ) + name = schedule.get("name") or "scheduled task" + return f"#{schedule.get('id')} {name} - {cadence} - {schedule.get('status', 'active')}" + + +def format_personal_record(record: dict[str, Any]) -> str: + category = record.get("category") or "note" + title = record.get("title") or "(untitled)" + status = record.get("status") or "active" + date_bits = [] + for key in ("date", "start_date", "end_date", "due_date"): + if record.get(key): + date_bits.append(f"{key}={record[key]}") + tags = record.get("tags") or [] + tag_text = f" tags={','.join(str(tag) for tag in tags)}" if tags else "" + date_text = f" {' '.join(date_bits)}" if date_bits else "" + return f"#{record.get('id')} [{category}] {title} - {status}{date_text}{tag_text}" + + +def _record_search_text(record: dict[str, Any]) -> str: + parts: list[str] = [] + for key in ( + "category", + "title", + "content", + "status", + "date", + "start_date", + "end_date", + "due_date", + ): + if record.get(key): + parts.append(str(record[key])) + tags = record.get("tags") + if isinstance(tags, list): + parts.extend(str(tag) for tag in tags) + metadata = record.get("metadata") + if isinstance(metadata, dict): + parts.extend(str(value) for value in metadata.values()) + return "\n".join(parts) + + +def _parse_datetime(value: Any) -> datetime | None: + if not value: + return None + try: + parsed = datetime.fromisoformat(str(value)) + except ValueError: + return None + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc) + + +def _utc_now() -> datetime: + return datetime.now(timezone.utc) diff --git a/deploy/README.md b/deploy/README.md new file mode 100644 index 0000000..a7e868f --- /dev/null +++ b/deploy/README.md @@ -0,0 +1,31 @@ +# Deployment Notes + +Default NAS target: + +```text +host: comtropy.synology.me +port: 50022 +remote dir: /volume1/docker/telegram-codex-bot +``` + +PowerShell deployment: + +```powershell +.\deploy\make-env.ps1 +.\deploy\deploy.ps1 -NasUser your-nas-admin-user +``` + +The script uploads files over legacy SCP, then runs: + +```bash +docker compose up -d --build +``` + +If the first Codex command fails inside Telegram, open an SSH session and run: + +```bash +cd /volume1/docker/telegram-codex-bot +docker compose exec telegram-codex-bot codex login --device-auth +``` + +The `codex-home` directory is mounted to `/root/.codex` so login state persists. diff --git a/deploy/check-status.ps1 b/deploy/check-status.ps1 new file mode 100644 index 0000000..a129396 --- /dev/null +++ b/deploy/check-status.ps1 @@ -0,0 +1,24 @@ +param( + [string]$NasUser = "y2keui", + [string]$NasHost = "comtropy.synology.me", + [int]$NasPort = 50022, + [string]$RemoteDir = "/volume1/docker/telegram-codex-bot", + [string]$SshKeyPath = "$env:USERPROFILE\.ssh\nas_codex_ed25519" +) + +$ErrorActionPreference = "Stop" + +$remote = "${NasUser}@${NasHost}" +$sshAuthArgs = @() +if ($SshKeyPath -and (Test-Path $SshKeyPath)) { + $sshAuthArgs = @("-i", $SshKeyPath, "-o", "IdentitiesOnly=yes") +} +$remoteCommand = @" +export PATH=/usr/local/bin:/var/packages/ContainerManager/target/usr/bin:/usr/bin:/bin:/usr/sbin:/sbin +cd '$RemoteDir' +sudo /usr/local/bin/docker-compose ps +echo '--- logs ---' +sudo /usr/local/bin/docker-compose logs --tail=160 telegram-codex-bot +"@ + +ssh @sshAuthArgs -tt -p $NasPort $remote $remoteCommand diff --git a/deploy/deploy.ps1 b/deploy/deploy.ps1 new file mode 100644 index 0000000..e6efc7f --- /dev/null +++ b/deploy/deploy.ps1 @@ -0,0 +1,69 @@ +param( + [Parameter(Mandatory = $true)] + [string]$NasUser, + + [string]$NasHost = "comtropy.synology.me", + [int]$NasPort = 50022, + [string]$RemoteDir = "/volume1/docker/telegram-codex-bot", + [string]$SshKeyPath = "$env:USERPROFILE\.ssh\nas_codex_ed25519" +) + +$ErrorActionPreference = "Stop" + +$root = Split-Path -Parent (Split-Path -Parent $MyInvocation.MyCommand.Path) +$remote = "${NasUser}@${NasHost}" +$remotePath = "/usr/local/bin:/var/packages/ContainerManager/target/usr/bin:/usr/bin:/bin:/usr/sbin:/sbin" +$sshAuthArgs = @() +if ($SshKeyPath -and (Test-Path $SshKeyPath)) { + $sshAuthArgs = @("-i", $SshKeyPath, "-o", "IdentitiesOnly=yes") +} + +if (-not (Test-Path (Join-Path $root ".env"))) { + Write-Host "Missing .env. Copy .env.example to .env and fill secrets first." -ForegroundColor Yellow + exit 1 +} + +function Invoke-Checked { + param( + [Parameter(Mandatory = $true)] + [scriptblock]$Script, + + [Parameter(Mandatory = $true)] + [string]$Description + ) + + & $Script + if ($LASTEXITCODE -ne 0) { + throw "$Description failed with exit code $LASTEXITCODE" + } +} + +Invoke-Checked -Description "Creating remote directory" -Script { + ssh @sshAuthArgs -p $NasPort $remote "export PATH='$remotePath'; mkdir -p '$RemoteDir'" +} + +Invoke-Checked -Description "Cleaning old uploaded source files" -Script { + ssh @sshAuthArgs -p $NasPort $remote "export PATH='$remotePath'; rm -rf '$RemoteDir/app' '$RemoteDir/scripts'" +} + +$uploadPaths = @( + (Join-Path $root "app"), + (Join-Path $root "scripts"), + (Join-Path $root ".dockerignore"), + (Join-Path $root ".env"), + (Join-Path $root "Dockerfile"), + (Join-Path $root "README.md"), + (Join-Path $root "docker-compose.yml"), + (Join-Path $root "requirements.txt") +) + +Invoke-Checked -Description "Uploading project files" -Script { + $scpArgs = @("-O") + $sshAuthArgs + @("-P", "$NasPort", "-r") + $uploadPaths + @("${remote}:$RemoteDir/") + & scp @scpArgs +} + +Invoke-Checked -Description "Building and starting container" -Script { + ssh @sshAuthArgs -tt -p $NasPort $remote "export PATH='$remotePath'; cd '$RemoteDir' && mkdir -p data workspaces codex-home && (sudo /usr/local/bin/docker-compose up -d --build || sudo /var/packages/ContainerManager/target/usr/bin/docker-compose up -d --build || sudo docker compose up -d --build)" +} + +Write-Host "Deployment finished: $RemoteDir" diff --git a/deploy/make-env.ps1 b/deploy/make-env.ps1 new file mode 100644 index 0000000..e19b537 --- /dev/null +++ b/deploy/make-env.ps1 @@ -0,0 +1,36 @@ +param( + [string]$WorkspaceRoot = "/workspaces", + [string]$StatePath = "/app/data/state.json" +) + +$ErrorActionPreference = "Stop" + +$root = Split-Path -Parent (Split-Path -Parent $MyInvocation.MyCommand.Path) +$envPath = Join-Path $root ".env" + +$secureToken = Read-Host "Telegram bot token" -AsSecureString +$tokenPtr = [Runtime.InteropServices.Marshal]::SecureStringToBSTR($secureToken) +try { + $token = [Runtime.InteropServices.Marshal]::PtrToStringBSTR($tokenPtr) +} +finally { + if ($tokenPtr -ne [IntPtr]::Zero) { + [Runtime.InteropServices.Marshal]::ZeroFreeBSTR($tokenPtr) + } +} +$allowedIds = Read-Host "Allowed Telegram user ids (comma-separated; leave blank for /whoami only)" + +@" +TELEGRAM_BOT_TOKEN=$token +ALLOWED_TELEGRAM_USER_IDS=$allowedIds +BOT_STATE_PATH=$StatePath +BOT_WORKSPACE_ROOT=$WorkspaceRoot +CODEX_BIN=codex +CODEX_TIMEOUT_SECONDS=1800 +CODEX_READONLY_SANDBOX=read-only +CODEX_WRITE_SANDBOX=workspace-write +CODEX_RESUME_JSON_FLAG_STYLE=before_resume +ALLOW_PLAIN_TEXT=true +"@ | Set-Content -Path $envPath -Encoding UTF8 + +Write-Host "Wrote $envPath" diff --git a/deploy/register-ssh-key.ps1 b/deploy/register-ssh-key.ps1 new file mode 100644 index 0000000..e04dc3d --- /dev/null +++ b/deploy/register-ssh-key.ps1 @@ -0,0 +1,31 @@ +param( + [string]$NasUser = "y2keui", + [string]$NasHost = "comtropy.synology.me", + [int]$NasPort = 50022, + [string]$SshKeyPath = "$env:USERPROFILE\.ssh\nas_codex_ed25519" +) + +$ErrorActionPreference = "Stop" + +$publicKeyPath = "$SshKeyPath.pub" +if (-not (Test-Path $publicKeyPath)) { + Write-Host "Missing public key: $publicKeyPath" -ForegroundColor Yellow + exit 1 +} + +$remote = "${NasUser}@${NasHost}" +$remoteCommand = @' +export PATH=/usr/local/bin:/var/packages/ContainerManager/target/usr/bin:/usr/bin:/bin:/usr/sbin:/sbin +umask 077 +mkdir -p ~/.ssh +cat >> ~/.ssh/authorized_keys +chmod 700 ~/.ssh +chmod 600 ~/.ssh/authorized_keys +'@ + +Get-Content -Raw $publicKeyPath | ssh -p $NasPort $remote $remoteCommand +if ($LASTEXITCODE -ne 0) { + throw "Registering SSH public key failed with exit code $LASTEXITCODE" +} + +Write-Host "Public key registered: $publicKeyPath" diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..bc0e5ea --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,14 @@ +services: + telegram-codex-bot: + build: + context: . + image: telegram-codex-bot:local + container_name: telegram-codex-bot + restart: unless-stopped + env_file: + - .env + volumes: + - ./data:/app/data + - ./workspaces:/workspaces + - ./codex-home:/root/.codex + diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..99a33f5 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,3 @@ +python-dotenv>=1.0,<2 +python-telegram-bot>=21,<23 +websockets>=13,<16 diff --git a/scripts/start.sh b/scripts/start.sh new file mode 100644 index 0000000..c5779d3 --- /dev/null +++ b/scripts/start.sh @@ -0,0 +1,54 @@ +#!/bin/sh +set -eu + +CODEX_BIN="${CODEX_BIN:-codex}" +BOT_WORKSPACE_ROOT="${BOT_WORKSPACE_ROOT:-/workspaces}" +CODEX_APP_SERVER_URL="${CODEX_APP_SERVER_URL:-ws://127.0.0.1:4500}" +CODEX_APP_SERVER_LISTEN="${CODEX_APP_SERVER_LISTEN:-$CODEX_APP_SERVER_URL}" +CODEX_APP_SERVER_ENABLED="${CODEX_APP_SERVER_ENABLED:-true}" +CODEX_APP_SERVER_CWD="${CODEX_APP_SERVER_CWD:-$BOT_WORKSPACE_ROOT}" +CODEX_APP_SERVER_SANDBOX="${CODEX_APP_SERVER_SANDBOX:-${CODEX_READONLY_SANDBOX:-read-only}}" +CODEX_APP_SERVER_APPROVAL_POLICY="${CODEX_APP_SERVER_APPROVAL_POLICY:-never}" + +APP_SERVER_PID="" + +cleanup() { + if [ -n "$APP_SERVER_PID" ] && kill -0 "$APP_SERVER_PID" 2>/dev/null; then + kill "$APP_SERVER_PID" 2>/dev/null || true + wait "$APP_SERVER_PID" 2>/dev/null || true + fi +} + +trap cleanup INT TERM EXIT + +mkdir -p /app/data "$BOT_WORKSPACE_ROOT" "${CODEX_HOME:-/root/.codex}" +python -m app.ensure_codex_config + +if [ "$CODEX_APP_SERVER_ENABLED" = "true" ]; then + echo "Starting Codex app-server at $CODEX_APP_SERVER_LISTEN" + "$CODEX_BIN" \ + --cd "$CODEX_APP_SERVER_CWD" \ + --sandbox "$CODEX_APP_SERVER_SANDBOX" \ + --ask-for-approval "$CODEX_APP_SERVER_APPROVAL_POLICY" \ + app-server \ + --listen "$CODEX_APP_SERVER_LISTEN" \ + & + APP_SERVER_PID="$!" + + if echo "$CODEX_APP_SERVER_LISTEN" | grep -q '^ws://'; then + HEALTH_URL="$(echo "$CODEX_APP_SERVER_LISTEN" | sed 's#^ws://#http://#')/readyz" + for _ in $(seq 1 30); do + if curl -fsS "$HEALTH_URL" >/dev/null 2>&1; then + echo "Codex app-server is ready" + break + fi + if ! kill -0 "$APP_SERVER_PID" 2>/dev/null; then + echo "Codex app-server exited before becoming ready" + break + fi + sleep 1 + done + fi +fi + +python -m app.main