From 6de70e1910b0effa8b5201382d64ae8cb9b97efe Mon Sep 17 00:00:00 2001 From: Your Name Date: Sat, 18 Jul 2026 17:03:30 +0800 Subject: [PATCH] feat(host188): forward controlled callback correlation --- .../services/awooop_ansible_post_verifier.py | 12 +- ...claw_callback_forwarder_ansible_catalog.py | 22 +- ...law_callback_forwarder_payload_contract.py | 135 ++ ...callback_overlay_post_verifier_contract.py | 6 +- .../app/bot/telegram.py | 1775 +++++++++++++++++ .../app/core/config.py | 252 +++ .../services/telegram_callback_forwarder.py | 408 ++++ .../manifest.json | 32 + .../188-openclaw-callback-forwarder.yml | 108 +- 9 files changed, 2725 insertions(+), 25 deletions(-) create mode 100644 apps/api/tests/test_host188_openclaw_callback_forwarder_payload_contract.py create mode 100644 infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/bot/telegram.py create mode 100644 infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/core/config.py create mode 100644 infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/services/telegram_callback_forwarder.py create mode 100644 infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/manifest.json diff --git a/apps/api/src/services/awooop_ansible_post_verifier.py b/apps/api/src/services/awooop_ansible_post_verifier.py index cb50e3297..5c3a4289b 100644 --- a/apps/api/src/services/awooop_ansible_post_verifier.py +++ b/apps/api/src/services/awooop_ansible_post_verifier.py @@ -251,14 +251,14 @@ _POSTCONDITION_REGISTRY: dict[str, tuple[AssetPostcondition, ...]] = { "set -euo pipefail; " "root=/home/ollama/.local/share/awoooi-controlled-apply/" "openclaw-callback-forwarder/" - "payloads/ffc45311b6403b30cb6ae6e1aa83dae1a047e9c9 && " + "payloads/8f367f8e7d6104c44b29bd5fa2d894db09ddf365 && " "manifest=\"$root/manifest.json\" && " "test -d \"$root\" && test ! -L \"$root\" && " "test \"$(stat -c '%U:%G:%a' \"$root\")\" = 'ollama:ollama:755' && " "test -f \"$manifest\" && test ! -L \"$manifest\" && " "test \"$(stat -c '%U:%G:%a' \"$manifest\")\" = 'ollama:ollama:444' && " "test \"$(sha256sum \"$manifest\" | cut -d' ' -f1)\" = " - "dafeb467b744ee7e80dc043d51f9e5d7e6905a58103f1710d519b553b56c6e48", + "0e5edaf700d264b44cc0eae3b9429af642be0b1563a0bffa633942d30d2c2e1c", ), AssetPostcondition( "host_188_openclaw_callback_overlay_enabled", @@ -273,7 +273,7 @@ _POSTCONDITION_REGISTRY: dict[str, tuple[AssetPostcondition, ...]] = { "test \"$(stat -c '%U:%G:%a' \"$override\")\" = " "'ollama:ollama:444' && " "test \"$(sha256sum \"$override\" | cut -d' ' -f1)\" = " - "8f6b7e8a0d6746c1331e99910656e64deebf67c9582bab3e1d36f71bb7a253c0", + "33c4594d4e4092b2a2367583a67f4c9a82ead0d28ca5078eca4896276084da0d", ), AssetPostcondition( "host_188_openclaw_callback_payload_runtime", @@ -282,7 +282,7 @@ _POSTCONDITION_REGISTRY: dict[str, tuple[AssetPostcondition, ...]] = { "set -euo pipefail; " "root=/home/ollama/.local/share/awoooi-controlled-apply/" "openclaw-callback-forwarder/" - "payloads/ffc45311b6403b30cb6ae6e1aa83dae1a047e9c9 && " + "payloads/8f367f8e7d6104c44b29bd5fa2d894db09ddf365 && " "telegram=\"$root/app/bot/telegram.py\" && " "config=\"$root/app/core/config.py\" && " "forwarder=\"$root/app/services/telegram_callback_forwarder.py\" && " @@ -297,7 +297,7 @@ _POSTCONDITION_REGISTRY: dict[str, tuple[AssetPostcondition, ...]] = { "test -f \"$forwarder\" && test ! -L \"$forwarder\" && " "test \"$(stat -c '%U:%G:%a' \"$forwarder\")\" = 'ollama:ollama:444' && " "test \"$(sha256sum \"$forwarder\" | cut -d' ' -f1)\" = " - "0d7693e34f63db4491d974f9f76a3a1c793d32cda193cb55ed3d8c86297b4fe3 && " + "8e8e0faf6322abb21d670ad9094f1805d7a4b3643da8fe88d74a8c9591066bf3 && " "ids=$(docker ps -q " "--filter label=com.docker.compose.project=clawbot " "--filter label=com.docker.compose.service=clawbot) && " @@ -314,7 +314,7 @@ _POSTCONDITION_REGISTRY: dict[str, tuple[AssetPostcondition, ...]] = { "e1b8b633f7b4780d3078d6c42a55e49d200255d3597a34e0cd2fd649b84eaea1 && " "test \"$(docker exec \"$ids\" sha256sum " "/app/app/services/telegram_callback_forwarder.py | cut -d' ' -f1)\" = " - "0d7693e34f63db4491d974f9f76a3a1c793d32cda193cb55ed3d8c86297b4fe3 && " + "8e8e0faf6322abb21d670ad9094f1805d7a4b3643da8fe88d74a8c9591066bf3 && " "test \"$(docker inspect --format='{{range .Mounts}}{{if eq " ".Destination \"/app/app/bot/telegram.py\"}}{{.Source}}|" "{{.Destination}}|{{.Type}}|{{.RW}}{{end}}{{end}}' \"$id\")\" = " diff --git a/apps/api/tests/test_host188_openclaw_callback_forwarder_ansible_catalog.py b/apps/api/tests/test_host188_openclaw_callback_forwarder_ansible_catalog.py index ac064686f..53a02561a 100644 --- a/apps/api/tests/test_host188_openclaw_callback_forwarder_ansible_catalog.py +++ b/apps/api/tests/test_host188_openclaw_callback_forwarder_ansible_catalog.py @@ -31,12 +31,13 @@ SERVICE_REGISTRY = ROOT / "ops" / "config" / "service-registry.yaml" MONITORING_REGISTRY = ROOT / "ops" / "monitoring" / "service-registry.yaml" CATALOG_ID = "ansible:188-openclaw-callback-forwarder" CANONICAL_ASSET_ID = "service:openclaw:host188" -EXPECTED_SHA = "ffc45311b6403b30cb6ae6e1aa83dae1a047e9c9" +EXPECTED_SHA = "8f367f8e7d6104c44b29bd5fa2d894db09ddf365" +UPSTREAM_SOURCE_COMMIT = "ffc45311b6403b30cb6ae6e1aa83dae1a047e9c9" EXPECTED_MANIFEST_SHA256 = ( - "dafeb467b744ee7e80dc043d51f9e5d7e6905a58103f1710d519b553b56c6e48" + "0e5edaf700d264b44cc0eae3b9429af642be0b1563a0bffa633942d30d2c2e1c" ) EXPECTED_OVERRIDE_SHA256 = ( - "8f6b7e8a0d6746c1331e99910656e64deebf67c9582bab3e1d36f71bb7a253c0" + "33c4594d4e4092b2a2367583a67f4c9a82ead0d28ca5078eca4896276084da0d" ) PAYLOAD_ROOT = ( ROOT @@ -54,7 +55,7 @@ EXPECTED_PAYLOAD_SHA256 = { "e1b8b633f7b4780d3078d6c42a55e49d200255d3597a34e0cd2fd649b84eaea1" ), "app/services/telegram_callback_forwarder.py": ( - "0d7693e34f63db4491d974f9f76a3a1c793d32cda193cb55ed3d8c86297b4fe3" + "8e8e0faf6322abb21d670ad9094f1805d7a4b3643da8fe88d74a8c9591066bf3" ), } @@ -353,6 +354,15 @@ def test_openclaw_playbook_is_check_first_bounded_and_rollback_capable() -> None "/home/ollama/.local/share/awoooi-controlled-apply/" f"openclaw-callback-forwarder/payloads/{EXPECTED_SHA}" ), + ).replace( + "{{ callback_identity_path }}", + ( + "/home/ollama/.local/share/awoooi-controlled-apply/" + "openclaw-callback-forwarder/receipts/current-deployment.json" + ), + ).replace( + "{{ callback_identity_container_path }}", + "/run/awoooi/openclaw-callback-forwarder/current-deployment.json", ) assert hashlib.sha256(rendered_override.encode()).hexdigest() == ( EXPECTED_OVERRIDE_SHA256 @@ -411,7 +421,9 @@ def test_openclaw_payload_manifest_is_hash_pinned_and_overlay_only() -> None: assert hashlib.sha256(manifest_path.read_bytes()).hexdigest() == ( EXPECTED_MANIFEST_SHA256 ) - assert manifest["source_commit"] == EXPECTED_SHA + assert manifest["source_commit"] == UPSTREAM_SOURCE_COMMIT + assert manifest["overlay_revision"] == EXPECTED_SHA + assert manifest["overlay_revision_kind"] == "git_blob_of_callback_forwarder" assert manifest["deployment_mode"] == ( "host_user_owned_read_only_compose_override_no_build_no_pull" ) diff --git a/apps/api/tests/test_host188_openclaw_callback_forwarder_payload_contract.py b/apps/api/tests/test_host188_openclaw_callback_forwarder_payload_contract.py new file mode 100644 index 000000000..a817a69ce --- /dev/null +++ b/apps/api/tests/test_host188_openclaw_callback_forwarder_payload_contract.py @@ -0,0 +1,135 @@ +from __future__ import annotations + +import importlib.util +import json +import sys +from pathlib import Path +from types import ModuleType, SimpleNamespace + +import pytest + +ROOT = Path(__file__).resolve().parents[3] +REVISION = "8f367f8e7d6104c44b29bd5fa2d894db09ddf365" +PAYLOAD = ( + ROOT + / "infra" + / "ansible" + / "files" + / "openclaw-callback-forwarder" + / REVISION + / "app" + / "services" + / "telegram_callback_forwarder.py" +) + + +def _load_payload() -> ModuleType: + spec = importlib.util.spec_from_file_location( + "host188_controlled_callback_forwarder", + PAYLOAD, + ) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + sys.modules[spec.name] = module + spec.loader.exec_module(module) + return module + + +def _identity() -> dict[str, str]: + return { + "schema_version": "openclaw_callback_deployment_identity_v1", + "catalog_id": "ansible:188-openclaw-callback-forwarder", + "trace_id": "trace:host188:001", + "run_id": "run:host188:001", + "work_item_id": "AIA-SRE-017", + } + + +def _update() -> dict: + return { + "update_id": 77, + "callback_query": { + "data": "detail:INC-20260716-TEST", + "from": {"id": 7, "username": "operator"}, + "id": "callback-forward-77", + "message": { + "chat": {"id": -1001}, + "message_id": 9, + "text": "incident", + }, + }, + } + + +class _CaptureClient: + def __init__(self) -> None: + self.body = b"" + + async def post(self, _url: str, *, content: bytes, headers: dict) -> object: + self.body = content + assert headers["X-Telegram-Forward-Schema"] == ( + "telegram_callback_forward_v1" + ) + return SimpleNamespace(status_code=200) + + +def test_identity_loader_accepts_exact_regular_file_and_rejects_unsafe_shapes( + tmp_path: Path, +) -> None: + module = _load_payload() + identity_path = tmp_path / "current-deployment.json" + identity_path.write_text(json.dumps(_identity()), encoding="utf-8") + + assert module.load_controlled_correlation(str(identity_path)) == { + "schema_version": "awoooi_controlled_callback_correlation_v1", + "trace_id": "trace:host188:001", + "run_id": "run:host188:001", + "work_item_id": "AIA-SRE-017", + } + + identity_path.write_text(json.dumps({**_identity(), "extra": "no"}), encoding="utf-8") + assert module.load_controlled_correlation(str(identity_path)) is None + identity_path.write_text("x" * 4097, encoding="utf-8") + assert module.load_controlled_correlation(str(identity_path)) is None + + target = tmp_path / "target.json" + target.write_text(json.dumps(_identity()), encoding="utf-8") + identity_path.unlink() + identity_path.symlink_to(target) + assert module.load_controlled_correlation(str(identity_path)) is None + + +@pytest.mark.asyncio +async def test_forwarder_adds_exact_correlation_and_preserves_legacy_fallback( + tmp_path: Path, +) -> None: + module = _load_payload() + identity_path = tmp_path / "current-deployment.json" + identity_path.write_text(json.dumps(_identity()), encoding="utf-8") + client = _CaptureClient() + forwarder = module.TelegramCallbackForwarder( + enabled=True, + url="https://awoooi.wooo.work/api/v1/telegram/callback-forward", + signing_key="test-key", + controlled_identity_path=str(identity_path), + ) + + result = await forwarder.forward(_update(), now=1, client=client) + + assert result.acknowledged is True + assert json.loads(client.body)["controlled_correlation"] == { + "schema_version": "awoooi_controlled_callback_correlation_v1", + "trace_id": "trace:host188:001", + "run_id": "run:host188:001", + "work_item_id": "AIA-SRE-017", + } + + identity_path.unlink() + legacy_client = _CaptureClient() + legacy_result = await forwarder.forward( + _update(), + now=2, + client=legacy_client, + ) + assert legacy_result.acknowledged is True + assert "controlled_correlation" not in json.loads(legacy_client.body) diff --git a/apps/api/tests/test_host188_openclaw_callback_overlay_post_verifier_contract.py b/apps/api/tests/test_host188_openclaw_callback_overlay_post_verifier_contract.py index 00b01532a..9e0711c99 100644 --- a/apps/api/tests/test_host188_openclaw_callback_overlay_post_verifier_contract.py +++ b/apps/api/tests/test_host188_openclaw_callback_overlay_post_verifier_contract.py @@ -15,7 +15,7 @@ CATALOG_ID = "ansible:188-openclaw-callback-forwarder" EXPECTED_PAYLOAD_SHA256 = { "af9600f32305ae73184651a8d3437063d7c8d35e3c3eb0910a2e27ac874c9302", "e1b8b633f7b4780d3078d6c42a55e49d200255d3597a34e0cd2fd649b84eaea1", - "0d7693e34f63db4491d974f9f76a3a1c793d32cda193cb55ed3d8c86297b4fe3", + "8e8e0faf6322abb21d670ad9094f1805d7a4b3643da8fe88d74a8c9591066bf3", } FRESH_RECEIPT_BLOCKER = ( "fresh_callback_receipt_run_scope_unavailable_in_static_postcondition_registry" @@ -72,8 +72,8 @@ def test_host188_overlay_probes_are_exact_and_secret_blind() -> None: assert "stat -c '%U:%G:%a'" in joined assert "ollama:ollama:444" in joined assert "ollama:ollama:755" in joined - assert "dafeb467b744ee7e80dc043d51f9e5d7e6905a58103f1710d519b553b56c6e48" in joined - assert "8f6b7e8a0d6746c1331e99910656e64deebf67c9582bab3e1d36f71bb7a253c0" in joined + assert "0e5edaf700d264b44cc0eae3b9429af642be0b1563a0bffa633942d30d2c2e1c" in joined + assert "33c4594d4e4092b2a2367583a67f4c9a82ead0d28ca5078eca4896276084da0d" in joined assert "docker-compose.override.yml" in joined assert "COMPOSE_FILE" not in joined assert all(digest in joined for digest in EXPECTED_PAYLOAD_SHA256) diff --git a/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/bot/telegram.py b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/bot/telegram.py new file mode 100644 index 000000000..f9426b4b0 --- /dev/null +++ b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/bot/telegram.py @@ -0,0 +1,1775 @@ +""" +OpenClaw v5.4 - Telegram Bot +ChatOps 雙向通訊、審批按鈕、多主機指令路由、安全控制 + +更新: +- v5.4: Meta Watchdog 監控監控系統 (Phase M) +- v5.3: 整合安全機制 (操作分級、權限控制、審計日誌) +- v5.2: 支援多主機指令路由 (kubectl -> 120, docker -> 188) +- v5.1: 快速回應機制 +- v5.0: 初始版本 +""" + +import asyncio +import logging +from typing import Callable + +from telegram import Update, InlineKeyboardButton, InlineKeyboardMarkup +from telegram.ext import ( + Application, + CommandHandler, + CallbackQueryHandler, + MessageHandler, + ContextTypes, + filters, +) + +from app.core.config import settings +from app.core.auth import ( + is_authorized, + can_execute, + get_user_role, + get_role_name, + get_role_emoji, + UserRole, +) +from app.core.operation_levels import ( + OperationLevel, + get_operation, + get_operation_level, + requires_confirmation, + is_forbidden, + get_level_emoji, +) +from app.core.audit import audit_logger, AuditAction +from app.core.user_tracker import track_user, get_all_users, get_user_count +from app.core.user_manager import add_user as add_auth_user, remove_user as remove_auth_user, list_users, get_user_role_dynamic +from app.services.grafana_service import get_grafana_service +from app.services.meta_watchdog import init_meta_watchdog, shutdown_meta_watchdog, get_meta_watchdog +from app.services.telegram_callback_forwarder import ( + AWOOOI_READ_CALLBACK_PATTERN, + CLAWBOT_APPROVAL_CALLBACK_PATTERN, + TelegramCallbackForwarder, + callback_action, +) +from app.bot.confirm_handler import confirmation_handler +from app.bot.incident_handler import incident_handler +from app.bot.commands import ( + StatusCommand, + LogsCommand, + ResourcesCommand, + RestartCommand, + AlertsCommand, + RepairCommand, + DiagnoseCommand, +) + +logger = logging.getLogger(__name__) + + +class TelegramBot: + """Telegram Bot 管理器""" + + def __init__(self): + self.token = settings.TELEGRAM_BOT_TOKEN + self.admin_chat_id = settings.TELEGRAM_ADMIN_CHAT_ID + self.allowed_users = settings.get_allowed_user_ids() + self.app: Application | None = None + self.is_running = False + self._pending_approvals: dict[str, dict] = {} + self.callback_forwarder = TelegramCallbackForwarder( + enabled=settings.TELEGRAM_CALLBACK_FORWARD_ENABLED, + url=settings.TELEGRAM_CALLBACK_FORWARD_URL, + signing_key=self.token, + allowed_hosts=settings.get_callback_forward_allowed_hosts(), + timeout_seconds=settings.TELEGRAM_CALLBACK_FORWARD_TIMEOUT_SECONDS, + max_body_bytes=settings.TELEGRAM_CALLBACK_FORWARD_MAX_BODY_BYTES, + ) + + # 初始化模組化命令處理器 + self.status_cmd = StatusCommand() + self.logs_cmd = LogsCommand() + self.resources_cmd = ResourcesCommand() + self.restart_cmd = RestartCommand() + self.alerts_cmd = AlertsCommand() + self.repair_cmd = RepairCommand() + self.diagnose_cmd = DiagnoseCommand() + + async def start(self): + """啟動 Telegram Bot""" + if not self.token: + logger.warning("Telegram Bot Token not configured, skipping...") + return + + try: + self.app = Application.builder().token(self.token).build() + + # 註冊命令處理器 - 基本 + self.app.add_handler(CommandHandler("start", self._cmd_start)) + self.app.add_handler(CommandHandler("help", self._cmd_help)) + self.app.add_handler(CommandHandler("users", self._cmd_users)) + self.app.add_handler(CommandHandler("adduser", self._cmd_adduser)) + self.app.add_handler(CommandHandler("removeuser", self._cmd_removeuser)) + self.app.add_handler(CommandHandler("authlist", self._cmd_authlist)) + self.app.add_handler(CommandHandler("whoami", self._cmd_whoami)) + + # 註冊模組化命令處理器 - Phase 1 核心命令 + self.app.add_handler(CommandHandler("status", self.status_cmd.handle)) + self.app.add_handler(CommandHandler("logs", self.logs_cmd.handle)) + self.app.add_handler(CommandHandler("disk", self.resources_cmd.handle_disk)) + self.app.add_handler(CommandHandler("memory", self.resources_cmd.handle_memory)) + self.app.add_handler(CommandHandler("cpu", self.resources_cmd.handle_cpu)) + self.app.add_handler(CommandHandler("resources", self.resources_cmd.handle_resources)) + self.app.add_handler(CommandHandler("restart", self.restart_cmd.handle)) + + # 註冊模組化命令處理器 - Phase 2 告警與修復 + self.app.add_handler(CommandHandler("alerts", self.alerts_cmd.handle)) + self.app.add_handler(CommandHandler("repair", self.repair_cmd.handle)) + + # 註冊模組化命令處理器 - Phase 3 AI 診斷 + self.app.add_handler(CommandHandler("diagnose", self.diagnose_cmd.handle)) + self.app.add_handler(CommandHandler("ask", self.diagnose_cmd.handle_ask)) + + # 註冊命令處理器 - K8s 操作 + # 註冊命令處理器 - Phase 4 Grafana Dashboard + self.app.add_handler(CommandHandler("dashboard", self._cmd_dashboard)) + self.app.add_handler(CommandHandler("pods", self._cmd_pods)) + self.app.add_handler(CommandHandler("services", self._cmd_services)) + self.app.add_handler(CommandHandler("scale", self._cmd_scale)) + self.app.add_handler(CommandHandler("kubectl", self._cmd_kubectl)) + + # 註冊命令處理器 - Docker 操作 + self.app.add_handler(CommandHandler("containers", self._cmd_containers)) + self.app.add_handler(CommandHandler("docker", self._cmd_docker)) + + # 註冊確認流程回調處理器 (InlineKeyboard) + confirm_handler = confirmation_handler.get_callback_handler() + self.app.add_handler(confirm_handler) + logger.info(f"[INIT] confirm_handler registered: {confirm_handler}") + + # 註冊 Incident 生命週期處理器 (接手/完成按鈕) + inc_handler = incident_handler.get_callback_handler() + self.app.add_handler(inc_handler) + logger.info(f"[INIT] incident_handler registered: {inc_handler}") + + # ClawBot 自有 callback 保持原 handler;只讀 AWOOOI callback 走 HMAC bridge。 + approval_handler = CallbackQueryHandler( + self._handle_approval_callback, + pattern=CLAWBOT_APPROVAL_CALLBACK_PATTERN, + ) + self.app.add_handler(approval_handler) + logger.info(f"[INIT] approval_handler registered: {approval_handler}") + self.app.add_handler( + CallbackQueryHandler( + self._handle_external_callback, + pattern=AWOOOI_READ_CALLBACK_PATTERN, + ) + ) + # 未知或寫入型 callback 一律 fail closed,避免 ghost button。 + self.app.add_handler(CallbackQueryHandler(self._handle_blocked_callback)) + + # 啟動確認超時清理任務 + await confirmation_handler.start_cleanup_task() + + # 註冊 Meta Watchdog 命令 (Phase M) + self.app.add_handler(CommandHandler("watchdog", self._cmd_watchdog)) + + # 註冊一般訊息處理器 (AI 對話) + self.app.add_handler(MessageHandler(filters.TEXT & ~filters.COMMAND, self._handle_message)) + + # 啟動 Bot (non-blocking) + await self.app.initialize() + await self.app.start() + await self.app.updater.start_polling(drop_pending_updates=True) + + self.is_running = True + logger.info("Telegram Bot started successfully") + + # 啟動 Meta Watchdog (Phase M: 監控監控系統) + try: + from app.services.meta_watchdog import MetaWatchdogConfig + watchdog_config = MetaWatchdogConfig( + healthchecks_ping_url=settings.HEALTHCHECKS_PING_URL, + ) + await init_meta_watchdog( + notify_callback=self.send_message, + config=watchdog_config, + ) + logger.info(f"Meta Watchdog started (Healthchecks.io: {'enabled' if settings.HEALTHCHECKS_PING_URL else 'disabled'})") + except Exception as e: + logger.warning(f"Meta Watchdog start failed (non-critical): {e}") + + except Exception as e: + logger.error(f"Failed to start Telegram Bot: {e}") + raise + + async def stop(self): + """停止 Telegram Bot""" + if self.app and self.is_running: + try: + # 停止 Meta Watchdog + await shutdown_meta_watchdog() + logger.info("Meta Watchdog stopped") + + await self.app.updater.stop() + await self.app.stop() + await self.app.shutdown() + self.is_running = False + logger.info("Telegram Bot stopped") + except Exception as e: + logger.error(f"Error stopping Telegram Bot: {e}") + + def _is_allowed_user(self, user_id: int, operation_id: str = None) -> bool: + """ + 檢查用戶是否有權限 + + Args: + user_id: Telegram 用戶 ID + operation_id: 操作 ID (用於更精確的權限檢查) + + Returns: + bool 是否有權限 + """ + # 使用新的 auth 模組 + if not is_authorized(user_id): + audit_logger.log_auth_failed(user_id, reason="未授權用戶") + return False + + # 如果指定了操作,檢查是否可執行 + if operation_id: + can_exec, reason = can_execute(user_id, operation_id) + if not can_exec: + audit_logger.log_auth_denied(user_id, operation_id, reason) + return False + + return True + + def _get_user_info(self, user) -> str: + """取得用戶資訊字串""" + role = get_user_role(user.id) + emoji = get_role_emoji(role) + role_name = get_role_name(role) + return f"{emoji} {user.username or user.id} ({role_name})" + + # ===== 命令處理器 ===== + + async def _cmd_start(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """處理 /start 命令""" + user = update.effective_user + # 追蹤用戶互動 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + role = get_user_role(user.id) + role_emoji = get_role_emoji(role) + role_name = get_role_name(role) + + logger.info(f"User {user.id} ({user.username}) started the bot, role={role.name}") + audit_logger.log_auth_success(user.id, user.username or str(user.id), update.effective_chat.id) + + if role == UserRole.UNKNOWN: + await update.message.reply_text( + f"🤖 *OpenClaw v{settings.APP_VERSION}*\n\n" + f"⚠️ 您尚未被授權使用此系統。\n" + f"請聯繫管理員申請權限。\n\n" + f"您的用戶 ID: `{user.id}`", + parse_mode="Markdown" + ) + return + + await update.message.reply_text( + f"🤖 *OpenClaw v{settings.APP_VERSION}*\n\n" + f"歡迎使用 WOOO AIOps 智能運維助手!\n\n" + f"*您的身份:*\n" + f"{role_emoji} {role_name}\n\n" + f"*核心功能:*\n" + f"• 🔧 主機操作 (kubectl, docker)\n" + f"• 🤖 AI 智能對話\n" + f"• 🚨 系統告警通知\n" + f"• ✅ 審批自動修復\n\n" + f"輸入 /help 查看可用命令\n" + f"輸入 /whoami 查看您的權限", + parse_mode="Markdown" + ) + + async def _cmd_status(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """處理 /status 命令""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + pending_count = len(self._pending_approvals) + + await update.message.reply_text( + f"📊 *OpenClaw 狀態*\n\n" + f"🟢 服務狀態: 運行中\n" + f"📍 環境: {settings.APP_ENV}\n" + f"⏳ 待審批: {pending_count} 項\n", + parse_mode="Markdown" + ) + + async def _cmd_watchdog(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """處理 /watchdog 命令 - Meta Watchdog 狀態檢查""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + watchdog = get_meta_watchdog() + if not watchdog: + await update.message.reply_text( + "⚠️ *Meta Watchdog 未啟動*\n\n" + "Meta Watchdog 服務尚未初始化。", + parse_mode="Markdown" + ) + return + + # 檢查參數 + args = context.args + if args and args[0].lower() == "check": + # 立即執行檢查 + await update.message.reply_text("🔄 正在執行 Meta Watchdog 檢查...") + report = await watchdog.check_now() + await update.message.reply_text(report, parse_mode="Markdown") + else: + # 顯示當前狀態 + report = watchdog.get_status_report() + await update.message.reply_text(report, parse_mode="Markdown") + + async def _cmd_help(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """處理 /help 命令""" + await update.message.reply_text( + "🤖 *OpenClaw v5.4 命令列表*\n\n" + "*基本指令:*\n" + "/start - 開始使用\n" + "/help - 顯示此說明\n" + "/whoami - 查看您的權限\n\n" + "*系統監控:*\n" + "/status [host|k8s] - 系統狀態\n" + "/disk /memory /cpu - 資源監控\n" + "/resources - 綜合資源\n\n" + "*日誌查詢:*\n" + "/logs \n\n" + "*告警與修復 (Phase 2):*\n" + "/alerts - 查看當前告警\n" + "/alerts ack - 確認告警\n" + "/repair strategies - 修復策略\n" + "/repair <策略> <目標> - 執行修復\n\n" + "*AI 智能診斷 (Phase 3):*\n" + "/diagnose - 診斷服務\n" + "/diagnose alert - 分析告警\n" + "/diagnose logs - 分析日誌\n" + "/ask <問題> - AI 問答\n\n" + "*Grafana Dashboard (Phase 4):*\n" + "/dashboard - 列出儀表板\n" + "/dashboard - 詳情與連結\n\n" + "*Meta Watchdog (Phase M):*\n" + "/watchdog - 監控系統狀態\n" + "/watchdog check - 立即檢查\n\n" + "*K8s 操作:*\n" + "/pods /services - 查詢資源\n" + "/restart - 重啟 (需確認)\n" + "/scale - 擴縮容\n\n" + "*Docker 操作:*\n" + "/containers - 列出容器\n" + "/docker - 執行命令\n\n" + "🟢 自動 | 🟡 確認 | 🟠 審批", + parse_mode="Markdown" + ) + + + async def _cmd_users(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """處理 /users 命令 - 查看所有用戶 (僅 SUPERADMIN)""" + user = update.effective_user + role = get_user_role(user.id) + + if role != UserRole.SUPERADMIN: + await update.message.reply_text("⛔ 此命令僅限超級管理員使用") + return + + users = get_all_users() + if not users: + await update.message.reply_text("📭 尚無用戶記錄") + return + + users.sort(key=lambda x: x.get("last_seen", ""), reverse=True) + + msg = f"👥 *已知用戶列表* ({len(users)} 人)\n\n" + + for u in users[:20]: + uid = u.get("user_id") + uname = u.get("username", "-") + fname = u.get("first_name", "-") + count = u.get("interaction_count", 0) + last_seen = u.get("last_seen", "-")[:10] + + urole = get_user_role(uid) + status = "✅" if urole != UserRole.UNKNOWN else "❌" + + msg += f"{status} `{uid}` @{uname}\n" + msg += f" 名稱: {fname} | 互動: {count}次\n" + msg += f" 最後活動: {last_seen}\n\n" + + if len(users) > 20: + msg += f"... 還有 {len(users) - 20} 個用戶" + + await update.message.reply_text(msg, parse_mode="Markdown") + + async def _cmd_adduser(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """添加授權用戶: /adduser """ + user = update.effective_user + + # 追蹤用戶 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + + if len(context.args) < 2: + await update.message.reply_text( + "⚠️ *用法:* /adduser \n\n" + "*可用角色:*\n" + "• `viewer` - 只讀權限\n" + "• `operator` - 運維權限\n" + "• `admin` - 管理權限\n\n" + "*範例:*\n" + "`/adduser 123456789 viewer`", + parse_mode="Markdown" + ) + return + + try: + target_id = int(context.args[0]) + except ValueError: + await update.message.reply_text("❌ 無效的用戶 ID") + return + + role = context.args[1].lower() + success, message = add_auth_user(target_id, role, user.id) + + await update.message.reply_text(message) + + async def _cmd_removeuser(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """移除授權用戶: /removeuser """ + user = update.effective_user + + # 追蹤用戶 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + + if len(context.args) < 1: + await update.message.reply_text( + "⚠️ *用法:* /removeuser \n\n" + "*範例:*\n" + "`/removeuser 123456789`", + parse_mode="Markdown" + ) + return + + try: + target_id = int(context.args[0]) + except ValueError: + await update.message.reply_text("❌ 無效的用戶 ID") + return + + success, message = remove_auth_user(target_id, user.id) + + await update.message.reply_text(message) + + async def _cmd_authlist(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """列出所有授權用戶""" + user = update.effective_user + role = get_user_role(user.id) + + # 追蹤用戶 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + + if role != UserRole.SUPERADMIN: + await update.message.reply_text("⛔ 此命令僅限超級管理員使用") + return + + msg = list_users() + await update.message.reply_text(msg, parse_mode="Markdown") + + + async def _cmd_whoami(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """顯示用戶身份與權限資訊""" + user = update.effective_user + # 追蹤用戶互動 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + role = get_user_role(user.id) + role_emoji = get_role_emoji(role) + role_name = get_role_name(role) + + # 根據角色顯示可用操作 + if role == UserRole.UNKNOWN: + permissions = "❌ 無權限" + elif role == UserRole.VIEWER: + permissions = "🟢 只讀操作 (查詢 Pods、Services、日誌等)" + elif role == UserRole.OPERATOR: + permissions = ( + "🟢 只讀操作\n" + "🟡 需確認操作 (重啟、擴縮容等)" + ) + elif role == UserRole.ADMIN: + permissions = ( + "🟢 只讀操作\n" + "🟡 需確認操作\n" + "🟠 需審批操作 (回滾、部署等)" + ) + else: # SUPERADMIN + permissions = ( + "🟢 只讀操作\n" + "🟡 需確認操作\n" + "🟠 需審批操作\n" + "🔴 危險操作 (需額外確認)" + ) + + await update.message.reply_text( + f"👤 *您的身份資訊*\n\n" + f"*用戶名:* @{user.username or 'N/A'}\n" + f"*用戶 ID:* `{user.id}`\n" + f"*角色:* {role_emoji} {role_name}\n\n" + f"*可用操作:*\n{permissions}", + parse_mode="Markdown" + ) + + async def _cmd_dashboard(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """ + Grafana Dashboard 命令 + 用法: + /dashboard - 列出可用的 Dashboard + /dashboard - 顯示 Dashboard 詳情與連結 + /dashboard render - 渲染 Dashboard 圖片 + /dashboard - 渲染特定面板 + """ + user = update.effective_user + # 追蹤用戶互動 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + + # 檢查權限 + if not is_authorized(user.id): + await update.message.reply_text( + "❌ 您沒有權限使用此 Bot。\n" + f"您的 Chat ID: `{user.id}`\n" + "請聯繫管理員添加授權。", + parse_mode="Markdown" + ) + return + + args = context.args + grafana = get_grafana_service() + + # 無參數 - 列出所有 Dashboard + if not args: + dashboard_list = grafana.list_dashboards() + await update.message.reply_text(dashboard_list, parse_mode="Markdown") + return + + dashboard_name = args[0].lower() + + # 檢查 Dashboard 是否存在 + dashboard = grafana.resolve_dashboard(dashboard_name) + if not dashboard: + await update.message.reply_text( + f"❌ Dashboard `{dashboard_name}` 不存在。\n\n" + f"使用 `/dashboard` 查看可用的 Dashboard 列表。", + parse_mode="Markdown" + ) + return + + # 第二個參數 - render 或 panel_id + if len(args) >= 2: + second_arg = args[1].lower() + + if second_arg == "render": + # 渲染整個 Dashboard + await update.message.reply_text( + f"🔄 正在渲染 Dashboard `{dashboard_name}`...\n" + f"這可能需要幾秒鐘...", + parse_mode="Markdown" + ) + + try: + image_bytes = await grafana.render_dashboard(dashboard_name) + if image_bytes: + from io import BytesIO + await update.message.reply_photo( + photo=BytesIO(image_bytes), + caption=f"📊 **{dashboard['title']}**\n" + f"⏰ 時間範圍: 最近 1 小時", + parse_mode="Markdown" + ) + else: + # 如果無法渲染,提供 URL + url = await grafana.get_dashboard_url(dashboard_name) + await update.message.reply_text( + f"⚠️ 無法渲染 Dashboard 圖片。\n" + f"這可能是因為 Grafana Image Renderer 插件未安裝。\n\n" + f"🔗 請直接訪問: {url}", + parse_mode="Markdown" + ) + except Exception as e: + logger.error(f"Dashboard render error: {e}") + url = await grafana.get_dashboard_url(dashboard_name) + await update.message.reply_text( + f"❌ 渲染失敗: {str(e)[:100]}\n\n" + f"🔗 請直接訪問: {url}", + parse_mode="Markdown" + ) + return + + elif second_arg.isdigit(): + # 渲染特定面板 + panel_id = int(second_arg) + await update.message.reply_text( + f"🔄 正在渲染 Panel {panel_id}...", + parse_mode="Markdown" + ) + + try: + image_bytes = await grafana.render_dashboard( + dashboard_name, + panel_id=panel_id + ) + if image_bytes: + from io import BytesIO + await update.message.reply_photo( + photo=BytesIO(image_bytes), + caption=f"📊 **{dashboard['title']}** - Panel {panel_id}", + parse_mode="Markdown" + ) + else: + await update.message.reply_text( + f"⚠️ 無法渲染 Panel {panel_id}。\n" + f"請確認 Panel ID 是否正確。", + parse_mode="Markdown" + ) + except Exception as e: + logger.error(f"Panel render error: {e}") + await update.message.reply_text( + f"❌ 渲染失敗: {str(e)[:100]}", + parse_mode="Markdown" + ) + return + + # 顯示 Dashboard 詳情 + info = await grafana.get_dashboard_info(dashboard_name) + await update.message.reply_text(info, parse_mode="Markdown") + + + async def _cmd_pods(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """列出 K8s Pods (自動路由到 K3s Master)""" + user = update.effective_user + # 追蹤用戶互動 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + operation_id = "kubectl_get_pods" + + # 權限檢查 (使用操作 ID) + if not self._is_allowed_user(user.id, operation_id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + # 記錄操作 + audit_logger.log_operation_requested( + user.id, user.username or str(user.id), + operation_id, "查詢 Pods" + ) + + # 取得 namespace 參數 + namespace = context.args[0] if context.args else "wooo-aiops-uat" + await update.message.reply_text(f"⏳ 正在查詢 Pods ({namespace})...") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法,自動 SSH 到 K3s Master (120) + result = await ssh_agent.kubectl_routed( + ["get", "pods", "-o", "wide"], + namespace=namespace + ) + + if result["success"]: + output = result["output"][:3500] # Telegram 訊息限制 + await update.message.reply_text(f"📦 *Pods ({namespace}):*\n```\n{output}\n```", parse_mode="Markdown") + audit_logger.log_operation_executed( + user.id, user.username or str(user.id), + operation_id, "查詢 Pods", {"success": True} + ) + else: + await update.message.reply_text(f"❌ 查詢失敗: {result.get('error', '未知錯誤')}") + audit_logger.log_operation_failed( + user.id, user.username or str(user.id), + operation_id, "查詢 Pods", result.get('error', '未知錯誤') + ) + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + audit_logger.log_operation_failed( + user.id, user.username or str(user.id), + operation_id, "查詢 Pods", str(e) + ) + + async def _cmd_services(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """列出 K8s Services (自動路由到 K3s Master)""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + namespace = context.args[0] if context.args else "wooo-aiops-uat" + await update.message.reply_text(f"⏳ 正在查詢 Services ({namespace})...") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法,自動 SSH 到 K3s Master (120) + result = await ssh_agent.kubectl_routed( + ["get", "svc"], + namespace=namespace + ) + + if result["success"]: + output = result["output"][:3500] + await update.message.reply_text(f"🌐 *Services ({namespace}):*\n```\n{output}\n```", parse_mode="Markdown") + else: + await update.message.reply_text(f"❌ 查詢失敗: {result.get('error', '未知錯誤')}") + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + + async def _cmd_restart(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """重啟 Deployment (自動路由到 K3s Master)""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + if not context.args: + await update.message.reply_text("❓ 用法: /restart \n例如: /restart wooo-api") + return + + deployment = context.args[0] + namespace = context.args[1] if len(context.args) > 1 else "wooo-aiops-uat" + await update.message.reply_text(f"⏳ 正在重啟 {deployment}...") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法 + result = await ssh_agent.kubectl_routed( + ["rollout", "restart", f"deployment/{deployment}"], + namespace=namespace + ) + + if result["success"]: + await update.message.reply_text(f"✅ {deployment} 重啟成功!") + else: + await update.message.reply_text(f"❌ 重啟失敗: {result.get('error', '未知錯誤')}") + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + + async def _cmd_scale(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """擴縮容 Deployment (自動路由到 K3s Master)""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + if len(context.args) < 2: + await update.message.reply_text("❓ 用法: /scale <副本數>\n例如: /scale wooo-api 3") + return + + deployment = context.args[0] + replicas = context.args[1] + namespace = context.args[2] if len(context.args) > 2 else "wooo-aiops-uat" + await update.message.reply_text(f"⏳ 正在將 {deployment} 擴縮至 {replicas} 副本...") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法 + result = await ssh_agent.kubectl_routed( + ["scale", f"deployment/{deployment}", f"--replicas={replicas}"], + namespace=namespace + ) + + if result["success"]: + await update.message.reply_text(f"✅ {deployment} 已擴縮至 {replicas} 副本!") + else: + await update.message.reply_text(f"❌ 擴縮失敗: {result.get('error', '未知錯誤')}") + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + + async def _cmd_logs(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """查看 Pod 日誌 (自動路由到 K3s Master)""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + if not context.args: + await update.message.reply_text("❓ 用法: /logs [行數]\n例如: /logs wooo-api-xxx 100") + return + + pod = context.args[0] + tail_lines = context.args[1] if len(context.args) > 1 else "50" + await update.message.reply_text(f"⏳ 正在獲取 {pod} 日誌...") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法 + result = await ssh_agent.kubectl_routed( + ["logs", pod, f"--tail={tail_lines}"], + namespace="wooo-aiops-uat" + ) + + if result["success"]: + output = result["output"][:3500] + await update.message.reply_text(f"📜 *{pod} 日誌 (最後{tail_lines}行):*\n```\n{output}\n```", parse_mode="Markdown") + else: + await update.message.reply_text(f"❌ 獲取日誌失敗: {result.get('error', '未知錯誤')}") + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + + async def _cmd_kubectl(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """執行 kubectl 命令 (自動路由到 K3s Master)""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + if not context.args: + await update.message.reply_text("❓ 用法: /kubectl <命令>\n例如: /kubectl get nodes") + return + + cmd_args = " ".join(context.args) + await update.message.reply_text(f"⏳ 執行: kubectl {cmd_args} (→ K3s Master)") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法,自動 SSH 到 K3s Master + result = await ssh_agent.kubectl_routed(list(context.args)) + + if result["success"]: + output = result["output"][:3500] + await update.message.reply_text(f"```\n{output}\n```", parse_mode="Markdown") + else: + await update.message.reply_text(f"❌ 執行失敗: {result.get('error', '未知錯誤')}") + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + + async def _cmd_containers(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """列出 Docker 容器 (本地執行)""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + # 可選指定主機 + host_id = context.args[0] if context.args else None + host_info = f" (主機 {host_id})" if host_id else " (本地)" + await update.message.reply_text(f"⏳ 正在查詢容器{host_info}...") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法 + result = await ssh_agent.docker_routed( + ["ps", "--format", "table {{.Names}}\t{{.Status}}\t{{.Ports}}"], + host_id=host_id + ) + + if result["success"]: + output = result["output"][:3500] + await update.message.reply_text(f"🐳 *Docker 容器{host_info}:*\n```\n{output}\n```", parse_mode="Markdown") + else: + await update.message.reply_text(f"❌ 查詢失敗: {result.get('error', '未知錯誤')}") + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + + async def _cmd_docker(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """執行 docker 命令 (本地執行)""" + if not self._is_allowed_user(update.effective_user.id): + await update.message.reply_text("⛔ 您沒有權限執行此操作") + return + + if not context.args: + await update.message.reply_text("❓ 用法: /docker <命令>\n例如: /docker ps") + return + + cmd_args = " ".join(context.args) + await update.message.reply_text(f"⏳ 執行: docker {cmd_args}") + + try: + from app.services.ssh_agent import SSHAgent + ssh_agent = SSHAgent() + # 使用路由方法 + result = await ssh_agent.docker_routed(list(context.args)) + + if result["success"]: + output = result["output"][:3500] + await update.message.reply_text(f"```\n{output}\n```", parse_mode="Markdown") + else: + await update.message.reply_text(f"❌ 執行失敗: {result.get('error', '未知錯誤')}") + except Exception as e: + await update.message.reply_text(f"❌ 執行錯誤: {e}") + + async def _handle_message(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """處理一般訊息 - 快速回應優先,複雜問題才用 AI""" + user = update.effective_user + # 追蹤用戶互動 + track_user(user.id, user.username, user.first_name, update.effective_chat.id) + message_text = update.message.text + + logger.info(f"Message from {user.id} ({user.username}): {message_text}") + + lower_text = message_text.lower().strip() + + # ===== 1. 精確匹配快速回應 ===== + quick_responses = { + # 問候 + "你好": "👋 您好!我是 OpenClaw,WOOO AIOps 智能運維助手。\n\n輸入 /help 查看可用指令!", + "嗨": "👋 嗨!有什麼可以幫您的嗎?輸入 /help 查看指令。", + "hello": "👋 Hello! I'm OpenClaw. Type /help for commands.", + "hi": "👋 Hi! Type /help to see available commands.", + "在嗎": "🤖 我在!24/7 全天候待命。有什麼需要幫忙的嗎?", + "在嗎?": "🤖 我在!24/7 全天候待命。有什麼需要幫忙的嗎?", + # 狀態查詢 + "狀態": "📊 查看系統狀態:\n• /status - OpenClaw 狀態\n• /pods - K8s Pods\n• /containers - Docker 容器", + "系統狀態": "📊 查看系統狀態:\n• /status - OpenClaw 狀態\n• /pods - K8s Pods\n• /containers - Docker 容器", + # 感謝 + "謝謝": "🙏 不客氣!有需要隨時呼叫我。", + "感謝": "🙏 很高興能幫上忙!", + "thanks": "🙏 You're welcome!", + "thank you": "🙏 You're welcome! Happy to help.", + # 其他 + "ok": "👍", + "好": "👍", + "收到": "👍 收到!", + } + + if lower_text in quick_responses: + await update.message.reply_text(quick_responses[lower_text]) + return + + # ===== 2. 關鍵字智能路由 ===== + keyword_responses = self._get_keyword_response(lower_text) + if keyword_responses: + await update.message.reply_text(keyword_responses, parse_mode="Markdown") + return + + # ===== 3. 權限檢查 ===== + if not self._is_allowed_user(user.id): + await update.message.reply_text("⛔ 您沒有權限使用此功能") + return + + # ===== 4. 呼叫 AI (Nemotron / OpenClaw) ===== + # 2026-03-31 ogt: Phase 22 ADR-044 - 啟用 AI 對話 (NVIDIA Nemotron) + try: + thinking_msg = await update.message.reply_text("🤔 思考中...") + response = await self._get_ai_response(message_text) + await thinking_msg.edit_text(response) + except Exception as e: + logger.error(f"AI response failed: {e}") + await update.message.reply_text( + "⚠️ AI 暫時無法使用,請稍後再試或輸入 /help 查看指令。" + ) + + def _get_keyword_response(self, text: str) -> str | None: + """根據關鍵字返回智能回應""" + # Pod 相關 + if any(kw in text for kw in ["pod", "pods", "容器狀態", "k8s狀態"]): + return "📦 *查看 Pods 狀態:*\n\n使用 `/pods` 指令查看所有 K8s Pods。\n\n範例:\n• `/pods` - 列出所有 Pods\n• `/logs ` - 查看 Pod 日誌" + + # 重啟相關 + if any(kw in text for kw in ["重啟", "restart", "重新啟動", "reboot"]): + return "🔄 *重啟服務:*\n\n使用 `/restart ` 指令。\n\n範例:\n• `/restart wooo-api` - 重啟 API 服務\n• `/restart wooo-frontend` - 重啟前端" + + # 日誌相關 + if any(kw in text for kw in ["日誌", "log", "logs", "紀錄", "錯誤"]): + return "📜 *查看日誌:*\n\n使用 `/logs ` 指令。\n\n先用 `/pods` 找到 Pod 名稱,再查看日誌。\n\n範例:\n• `/logs wooo-api-xxx`" + + # Docker 相關 + if any(kw in text for kw in ["docker", "容器", "container"]): + return "🐳 *Docker 操作:*\n\n• `/containers` - 列出所有容器\n• `/docker ps` - 容器狀態\n• `/docker logs ` - 容器日誌\n• `/docker restart ` - 重啟容器" + + # 擴縮容相關 + if any(kw in text for kw in ["擴容", "縮容", "scale", "副本", "replica"]): + return "📈 *擴縮容:*\n\n使用 `/scale <數量>` 指令。\n\n範例:\n• `/scale wooo-api 3` - 擴展到 3 個副本\n• `/scale wooo-api 1` - 縮減到 1 個副本" + + # 部署相關 + if any(kw in text for kw in ["部署", "deploy", "發布", "上線"]): + return "🚀 *部署操作:*\n\n目前部署透過 GitLab CI/CD 自動執行。\n\n查看部署狀態:\n• `/pods` - 查看 Pod 狀態\n• `/kubectl get deployments` - 部署列表" + + # 監控相關 + if any(kw in text for kw in ["監控", "monitor", "grafana", "prometheus", "指標"]): + return "📊 *監控系統:*\n\n• Grafana: https://grafana.wooo.work\n• Prometheus: 內部監控\n\n查看服務狀態:\n• `/pods` - K8s Pods\n• `/services` - K8s Services" + + # 幫助相關 + if any(kw in text for kw in ["幫助", "help", "怎麼用", "使用說明", "指令"]): + return "❓ *需要幫助?*\n\n輸入 `/help` 查看完整指令列表。\n\n常用指令:\n• `/pods` - 查看 Pods\n• `/containers` - 查看 Docker\n• `/restart` - 重啟服務" + + # 網站相關 + if any(kw in text for kw in ["網站", "website", "頁面", "url", "連結"]): + return "🌐 *相關網站:*\n\n• AIOps: https://aiops.wooo.work\n• n8n: https://n8n.wooo.work\n• MOMO: https://momo.wooo.work\n• 岑洋: https://tsenyang.wooo.work" + + # 告警相關 + if any(kw in text for kw in ["告警", "alert", "警報", "通知"]): + return "🔔 *告警管理:*\n\n告警會自動推送到此頻道。\n\n收到告警時,您可以:\n• 點擊「批准」執行自動修復\n• 點擊「拒絕」手動處理\n• 使用 `/pods` 查看狀態" + + return None + + async def _get_ai_response(self, user_message: str) -> str: + """使用 LLM 生成回應""" + try: + from app.services.llm_optimizer import LLMOptimizer + optimizer = LLMOptimizer() + + # 構建 AIOps 對話 Prompt + prompt = f"""你是 OpenClaw,WOOO AIOps 智能運維助手。 + +你的能力包括: +1. 監控系統狀態 (Prometheus, Grafana) +2. K8s 操作 (kubectl get/restart/scale/logs) +3. Docker 容器管理 +4. 自動修復 (restart, scale_up, scale_down, rollback) +5. 告警處理和根本原因分析 + +用戶訊息: {user_message} + +請用繁體中文簡潔回覆,如果用戶詢問操作相關問題,提供相關的 /命令 指引。 +限制在 500 字以內。""" + + if optimizer.provider == "anthropic": + response = await optimizer._call_anthropic(prompt) + elif optimizer.provider == "openai": + response = await optimizer._call_openai(prompt) + elif optimizer.provider == "ollama": + response = await optimizer._call_ollama(prompt) + else: + return "⚠️ LLM 服務未配置" + + return response[:2000] # Telegram 訊息限制 + + except Exception as e: + logger.error(f"LLM call failed: {e}") + raise + + # ===== 審批回調處理 ===== + + async def _handle_external_callback( + self, + update: Update, + context: ContextTypes.DEFAULT_TYPE, + ): + """Forward an allowlisted AWOOOI read callback without exposing secrets.""" + del context + result = await self.callback_forwarder.forward(update.to_dict()) + if result.acknowledged: + logger.info( + "AWOOOI callback forward acknowledged", + extra={ + "callback_action": callback_action(update.callback_query.data), + "forward_status": result.reason, + "upstream_status": result.status_code, + }, + ) + return + + logger.warning( + "AWOOOI callback forward failed closed", + extra={ + "callback_action": callback_action(update.callback_query.data), + "forward_status": result.reason, + "upstream_status": result.status_code, + }, + ) + try: + await update.callback_query.answer( + "⚠️ 此按鈕目前無法安全處理,未執行任何動作。", + show_alert=True, + ) + except Exception: + logger.warning("Telegram sanitized callback failure reply unavailable") + + async def _handle_blocked_callback( + self, + update: Update, + context: ContextTypes.DEFAULT_TYPE, + ): + """Reject unknown/write callbacks without forwarding or executing them.""" + del context + logger.warning( + "Telegram callback blocked by action allowlist", + extra={"callback_action": callback_action(update.callback_query.data)}, + ) + try: + await update.callback_query.answer( + "⛔ 此按鈕未通過受控操作合約,未執行任何動作。", + show_alert=True, + ) + except Exception: + logger.warning("Telegram sanitized blocked callback reply unavailable") + + async def _handle_approval_callback(self, update: Update, context: ContextTypes.DEFAULT_TYPE): + """處理審批按鈕回調 (包含修復按鈕)""" + import sys + print(f"[CALLBACK] _handle_approval_callback ENTERED at {__file__}", flush=True, file=sys.stderr) + try: + query = update.callback_query + print(f"[CALLBACK] query.data = {query.data}", flush=True, file=sys.stderr) + await query.answer() + except Exception as e: + print(f"[CALLBACK] ERROR: {e}", flush=True, file=sys.stderr) + import traceback + traceback.print_exc() + raise + + user = query.from_user + if not self._is_allowed_user(user.id): + await query.edit_message_text("⛔ 您沒有權限執行此操作") + return + + data = query.data + logger.info(f"[DEBUG] Callback data received: {data}") + + # === 處理修復按鈕回調 === + if data.startswith("repair:"): + logger.info(f"[DEBUG] Routing to _handle_repair_callback") + await self._handle_repair_callback(query, user) + return + + # === Phase 6.1: 處理 RAG 反饋按鈕回調 === + if data.startswith("rag:"): + logger.info(f"[DEBUG] Routing to _handle_rag_feedback_callback") + await self._handle_rag_feedback_callback(query, user) + return + + # === 處理審批按鈕回調 === + if not data.startswith(("approve:", "reject:")): + return + + action, action_id = data.split(":", 1) + is_approved = action == "approve" + + logger.info(f"Approval callback: {action_id} -> {action} by {user.username}") + + # 調用 API 進行審批 + import httpx + try: + async with httpx.AsyncClient() as client: + endpoint = "approve" if is_approved else "reject" + response = await client.post( + f"http://localhost:{settings.PORT}/api/v1/actions/{action_id}/{endpoint}", + timeout=30.0 + ) + result = response.json() + + status_emoji = "✅" if is_approved else "❌" + status_text = "已批准" if is_approved else "已拒絕" + + # ========================================================================= + # 🔴 統帥鐵律:簽核後必須保留原始告警內容 + # Memory: feedback_approval_preserve_content.md + # ========================================================================= + original_text = query.message.text or "" + separator = "──────────────" + stamp = f"{status_emoji} 已由 @{user.username or user.id} {status_text}" + + # 組合: 原始內容 + 分隔線 + 鋼印 + updated_text = f"{original_text}\n{separator}\n{stamp}" + + await query.edit_message_text( + updated_text, + parse_mode="Markdown", + reply_markup=None, # 移除按鈕 + ) + + # 清理暫存 + self._pending_approvals.pop(action_id, None) + + except Exception as e: + logger.error(f"Approval callback failed: {e}") + await query.edit_message_text(f"⚠️ 操作失敗: {e}") + + async def _handle_repair_callback(self, query, user): + """ + 處理修復按鈕回調 (OPS.74 v2) + + callback_data 格式: repair:{strategy}:{service} + 例: repair:restart:wooo-api + + 架構決策: + - 保留原始告警訊息內容,只移除按鈕 + - Redis 分散式鎖防止並發/連點 + - alert_acknowledge/ignore 直接確認,不調用修復命令 + - 10s SSH timeout + - Audit Log JSON stdout + """ + import asyncio + import json + from datetime import datetime + from app.bot.commands.repair import repair_command + from app.core.repair_lock import repair_lock, RedisUnavailableError + from app.core.service_catalog import service_catalog + + data = query.data + parts = data.split(":", 2) + + if len(parts) < 3: + await query.answer("⚠️ 無效的修復請求格式") + return + + _, strategy, service = parts + chat_id = query.message.chat_id + message_id = query.message.message_id + original_text = query.message.text or query.message.caption or "" + start_time = datetime.now() + + logger.info(f"[REPAIR] Callback: {strategy} -> {service} by {user.username}") + + # 策略名稱對照表 + strategy_names = { + "restart": "🔄 重啟服務", + "scale_up": "📈 擴展副本", + "scale_down": "📉 縮減副本", + "clear_cache": "🗑️ 清除快取", + "rollback": "⏪ 回滾版本", + "ignore": "🙈 忽略告警", + "alert_acknowledge": "✅ 確認收到", + } + strategy_display = strategy_names.get(strategy, strategy) + + # === 1. 檢查服務是否允許此策略 === + if not service_catalog.can_execute(service, strategy): + warning = service_catalog.get_warning(service) + if warning: + await query.answer(warning, show_alert=True) + else: + allowed = service_catalog.get_strategies(service) + allowed_str = ", ".join(allowed) + await query.answer(f"⚠️ 不支援 {strategy},允許: {allowed_str}", show_alert=True) + return + + # === 2. 特殊處理:確認/忽略(不需要 Redis 鎖或修復操作)=== + if strategy in ("alert_acknowledge", "ignore"): + # 移除按鈕但保留原始訊息 + await query.edit_message_reply_markup(reply_markup=None) + # 發送確認通知 + ack_text = ( + f"{strategy_display}\n\n" + f"🎯 服務: `{service}`\n" + f"👤 操作者: @{user.username or user.id}\n" + f"⏰ 時間: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}" + ) + await self.app.bot.send_message(chat_id=chat_id, text=ack_text, parse_mode="Markdown") + # Audit Log + audit_log = { + "timestamp": datetime.now().isoformat(), + "level": "AUDIT", + "event": "alert_acknowledged", + "actor_telegram_id": user.id, + "actor_username": user.username, + "action": strategy, + "service": service, + "result": "success" + } + logger.info(json.dumps(audit_log, ensure_ascii=False)) + return + + # === 3. 嘗試獲取 Redis 分散式鎖 === + try: + lock_acquired = await repair_lock.acquire(service, ttl=60) + if not lock_acquired: + await query.answer(f"⏳ {service} 正在處理中,請稍候", show_alert=True) + return + except RedisUnavailableError as e: + await query.answer(str(e), show_alert=True) + logger.error(f"[REPAIR] Redis unavailable: {e}") + return + + try: + # === 4. 移除按鈕(保留原始訊息)=== + await query.edit_message_reply_markup(reply_markup=None) + + # === 5. 發送處理中通知 === + processing_msg = await self.app.bot.send_message( + chat_id=chat_id, + text=f"⏳ 正在執行 {strategy_display}...\n🎯 目標: `{service}`\n👤 操作者: @{user.username or user.id}", + parse_mode="Markdown" + ) + + # === 6. 執行修復 (帶 60s timeout) === + try: + result = await asyncio.wait_for( + repair_command._do_repair( + strategy=strategy, + target=service, + namespace="wooo-aiops-uat", + user_id=user.id, + username=user.username + ), + timeout=60.0 # Docker restart 可能需要較長時間 + ) + except asyncio.TimeoutError: + result = {"success": False, "error": "修復操作逾時 (60s)"} + + # === 7. 計算執行時間 === + duration_ms = int((datetime.now() - start_time).total_seconds() * 1000) + + # === 8. 更新處理中訊息為結果 === + if result.get("success"): + output = result.get("output", "")[:300] + result_text = ( + f"✅ 修復成功\n\n" + f"🎯 服務: `{service}`\n" + f"📋 策略: {strategy_display}\n" + f"👤 操作者: @{user.username or user.id}\n" + f"⏱️ 耗時: {duration_ms}ms" + ) + if output: + result_text += f"\n\n```\n{output}\n```" + audit_result = "success" + else: + error_msg = result.get("error", "未知錯誤") + result_text = ( + f"❌ 修復失敗\n\n" + f"🎯 服務: `{service}`\n" + f"📋 策略: {strategy_display}\n" + f"👤 操作者: @{user.username or user.id}\n" + f"⏱️ 耗時: {duration_ms}ms\n\n" + f"錯誤: `{error_msg}`" + ) + audit_result = "failed" + + # 編輯處理中訊息 + await self.app.bot.edit_message_text( + chat_id=chat_id, + message_id=processing_msg.message_id, + text=result_text, + parse_mode="Markdown" + ) + + # === 9. Audit Log (JSON stdout) === + audit_log = { + "timestamp": datetime.now().isoformat(), + "level": "AUDIT", + "event": "repair_executed", + "actor_telegram_id": user.id, + "actor_username": user.username, + "action": strategy, + "service": service, + "service_type": service_catalog.get_type(service), + "result": audit_result, + "duration_ms": duration_ms, + "error_message": result.get("error") if not result.get("success") else None + } + logger.info(json.dumps(audit_log, ensure_ascii=False)) + + except Exception as e: + logger.error(f"[REPAIR] Callback failed: {e}") + await self.app.bot.send_message( + chat_id=chat_id, + text=f"⚠️ 執行錯誤\n\n🎯 服務: `{service}`\n錯誤: `{str(e)}`", + parse_mode="Markdown" + ) + + finally: + # === 10. 釋放鎖 === + await repair_lock.release(service) + + # ========================================================================= + # Phase 6.1: RAG 反饋機制 (Feedback Loop) + # ========================================================================= + + async def _handle_rag_feedback_callback(self, query, user): + """ + 處理 RAG 反饋按鈕回調 (Phase 6.1) + + callback_data 格式: + - rag:up:{service}:{error_type} (讚成: 有幫助) + - rag:down:{service}:{error_type} (反對: 無效/過時) + + 功能: + - 防洗票: 每個使用者只能投一次票 + - 自動淘汰: downvote >= 3 時標記為 deprecated + """ + from app.core.knowledge_base import get_knowledge_base + + data = query.data + parts = data.split(":", 3) + + if len(parts) < 4: + await query.answer("無效的反饋格式") + return + + _, vote_type, service, error_type = parts + user_id = user.id + + logger.info( + f"[RAG Feedback] {vote_type} vote from {user.username or user_id} " + f"for {service}/{error_type}" + ) + + try: + kb = await get_knowledge_base() + + if vote_type == "up": + success, message = await kb.upvote(service, error_type, user_id) + elif vote_type == "down": + success, message = await kb.downvote(service, error_type, user_id) + else: + await query.answer("無效的投票類型") + return + + # 顯示結果 (toast notification) + emoji = "👍" if vote_type == "up" else "👎" + await query.answer(f"{emoji} {message}", show_alert=True) + + # 更新訊息: 移除反饋按鈕 (已投票) + original_text = query.message.text or "" + new_keyboard = {"inline_keyboard": []} + + # 保留非 RAG 反饋按鈕 + if query.message.reply_markup: + for row in query.message.reply_markup.inline_keyboard: + new_row = [] + for btn in row: + if btn.callback_data and not btn.callback_data.startswith("rag:"): + new_row.append({ + "text": btn.text, + "callback_data": btn.callback_data + } if btn.callback_data else { + "text": btn.text, + "url": btn.url + }) + elif btn.url: # URL 按鈕 + new_row.append({"text": btn.text, "url": btn.url}) + if new_row: + new_keyboard["inline_keyboard"].append(new_row) + + # 更新訊息文字 (標記已投票) + vote_status = "✅ 感謝回饋,此知識已被標記為有用" if vote_type == "up" else "⚠️ 感謝回饋,已記錄此知識可能過時" + if "已找到歷史修復經驗" in original_text: + updated_text = original_text.replace( + "🧠 *已找到歷史修復經驗*\n└ 請評估是否適用:", + f"🧠 *已找到歷史修復經驗*\n└ {vote_status}" + ) + else: + updated_text = original_text + + await query.edit_message_text( + text=updated_text, + parse_mode="Markdown", + reply_markup=new_keyboard if new_keyboard["inline_keyboard"] else None, + disable_web_page_preview=True, + ) + + logger.info(f"[RAG Feedback] Processed: {vote_type} for {service}/{error_type}") + + except Exception as e: + logger.error(f"[RAG Feedback] Error: {e}") + await query.answer(f"反饋處理失敗: {e}", show_alert=True) + + # ===== 發送訊息方法 ===== + + async def send_approval_request( + self, + action_id: str, + action_type: str, + target: str, + action: str, + context: dict | None = None, + ): + """發送審批請求""" + if not self.app: + logger.warning("Telegram Bot not initialized") + return + + # 構建訊息 + message = ( + f"🚨 *需要審批*\n\n" + f"*動作 ID:* `{action_id}`\n" + f"*類型:* {action_type}\n" + f"*目標:* `{target}`\n" + f"*操作:* {action}\n" + ) + + if context: + if context.get("error_type"): + message += f"*錯誤類型:* {context['error_type']}\n" + if context.get("count"): + message += f"*發生次數:* {context['count']}\n" + if context.get("signoz_trace_url"): + message += f"[🔍 查看 Trace]({context['signoz_trace_url']})\n" + + message += "\n請選擇操作:" + + # 審批按鈕 + keyboard = InlineKeyboardMarkup([ + [ + InlineKeyboardButton("✅ 批准", callback_data=f"approve:{action_id}"), + InlineKeyboardButton("❌ 拒絕", callback_data=f"reject:{action_id}"), + ] + ]) + + try: + await self.app.bot.send_message( + chat_id=self.admin_chat_id, + text=message, + parse_mode="Markdown", + reply_markup=keyboard, + ) + self._pending_approvals[action_id] = { + "type": action_type, + "target": target, + "action": action, + } + logger.info(f"Approval request sent for {action_id}") + + except Exception as e: + logger.error(f"Failed to send approval request: {e}") + + async def send_action_result(self, action_id: str, success: bool, message: str): + """發送動作執行結果""" + if not self.app: + return + + emoji = "✅" if success else "❌" + status = "成功" if success else "失敗" + + text = ( + f"{emoji} *動作執行{status}*\n\n" + f"*動作 ID:* `{action_id}`\n" + f"*結果:* {message}" + ) + + try: + await self.app.bot.send_message( + chat_id=self.admin_chat_id, + text=text, + parse_mode="Markdown", + ) + except Exception as e: + logger.error(f"Failed to send action result: {e}") + + async def send_inspection_alert( + self, + url: str, + errors: list[str], + screenshot_path: str | None = None, + ): + """發送巡檢告警""" + if not self.app: + return + + text = ( + f"🔍 *網站巡檢告警*\n\n" + f"*URL:* {url}\n" + f"*發現問題:*\n" + ) + for error in errors[:5]: # 最多顯示 5 個 + text += f"• {error}\n" + + try: + # 如果有截圖,發送圖片 + if screenshot_path: + with open(screenshot_path, "rb") as photo: + await self.app.bot.send_photo( + chat_id=self.admin_chat_id, + photo=photo, + caption=text, + parse_mode="Markdown", + ) + else: + await self.app.bot.send_message( + chat_id=self.admin_chat_id, + text=text, + parse_mode="Markdown", + ) + except Exception as e: + logger.error(f"Failed to send inspection alert: {e}") + + async def send_signoz_alert( + self, + alert_name: str, + labels: dict, + annotations: dict, + ): + """發送 SigNoz 告警""" + if not self.app: + return + + severity = labels.get("severity", "warning") + emoji = "🔴" if severity == "critical" else "🟡" if severity == "warning" else "🔵" + + text = ( + f"{emoji} *SigNoz 告警*\n\n" + f"*名稱:* {alert_name}\n" + f"*嚴重度:* {severity}\n" + ) + + if annotations.get("description"): + text += f"*描述:* {annotations['description']}\n" + + if labels.get("service"): + text += f"*服務:* {labels['service']}\n" + + try: + await self.app.bot.send_message( + chat_id=self.admin_chat_id, + text=text, + parse_mode="Markdown", + ) + except Exception as e: + logger.error(f"Failed to send SigNoz alert: {e}") + + async def send_message(self, text: str, parse_mode: str = "Markdown"): + """發送一般訊息""" + if not self.app: + return + + try: + await self.app.bot.send_message( + chat_id=self.admin_chat_id, + text=text, + parse_mode=parse_mode, + ) + except Exception as e: + logger.error(f"Failed to send message: {e}") + + +# ===== 全域單例 ===== +telegram = TelegramBot() + + +# ===== 模組層級包裝函數 (供其他模組 import) ===== + +async def send_alert_notification(message: str): + """ + 發送告警通知 (HTML 格式) + + Args: + message: HTML 格式的告警訊息 + """ + if not telegram.app: + logger.warning("Telegram Bot not initialized, cannot send alert") + return + + try: + await telegram.app.bot.send_message( + chat_id=telegram.admin_chat_id, + text=message, + parse_mode="HTML", + ) + except Exception as e: + logger.error(f"Failed to send alert notification: {e}") + + +async def send_healing_approval(service: str, action: str, alert_context: str): + """ + 發送修復審批請求 + + Args: + service: 要修復的服務名稱 + action: 修復動作 (restart, scale_up, etc.) + alert_context: HTML 格式的告警內容 + """ + from telegram import InlineKeyboardButton, InlineKeyboardMarkup + + if not telegram.app: + logger.warning("Telegram Bot not initialized, cannot send approval") + return + + # 構建審批訊息 + msg = ( + f"{alert_context}\n\n" + f"━━━━━━━━━━━━━━━━━━━━━\n" + f"🔧 建議修復動作\n" + f"🎯 服務: {service}\n" + f"📋 動作: {action}\n\n" + f"請選擇操作:" + ) + + # 修復按鈕 + keyboard = InlineKeyboardMarkup([ + [ + InlineKeyboardButton("✅ 批准修復", callback_data=f"repair:{action}:{service}"), + InlineKeyboardButton("❌ 忽略", callback_data=f"reject:alert:{service}"), + ] + ]) + + try: + await telegram.app.bot.send_message( + chat_id=telegram.admin_chat_id, + text=msg, + parse_mode="HTML", + reply_markup=keyboard, + ) + logger.info(f"Healing approval sent for {service} ({action})") + except Exception as e: + logger.error(f"Failed to send healing approval: {e}") + + +# ===== main.py 相容介面 ===== + +_external_bot_app = None + + +def set_bot_app(app): + """ + 設定外部 Telegram Bot Application (供 main.py 呼叫) + + Args: + app: telegram.ext.Application 實例 + """ + global _external_bot_app + _external_bot_app = app + # 同時更新 telegram 單例 + telegram.app = app + + +def setup_bot(app): + """ + 設定 Bot 命令處理器 (供 main.py 呼叫) + + Args: + app: telegram.ext.Application 實例 + """ + from telegram.ext import CommandHandler, CallbackQueryHandler, MessageHandler, filters + + # 註冊命令處理器 - 基本 + app.add_handler(CommandHandler("start", telegram._cmd_start)) + app.add_handler(CommandHandler("help", telegram._cmd_help)) + + # 註冊模組化命令處理器 - Phase 1 核心命令 + app.add_handler(CommandHandler("status", telegram.status_cmd.handle)) + app.add_handler(CommandHandler("logs", telegram.logs_cmd.handle)) + app.add_handler(CommandHandler("disk", telegram.resources_cmd.handle_disk)) + app.add_handler(CommandHandler("memory", telegram.resources_cmd.handle_memory)) + app.add_handler(CommandHandler("cpu", telegram.resources_cmd.handle_cpu)) + app.add_handler(CommandHandler("resources", telegram.resources_cmd.handle_resources)) + app.add_handler(CommandHandler("restart", telegram.restart_cmd.handle)) + + # 註冊模組化命令處理器 - Phase 2 告警與修復 + app.add_handler(CommandHandler("alerts", telegram.alerts_cmd.handle)) + app.add_handler(CommandHandler("repair", telegram.repair_cmd.handle)) + + # 註冊模組化命令處理器 - Phase 3 AI 診斷 + app.add_handler(CommandHandler("diagnose", telegram.diagnose_cmd.handle)) + app.add_handler(CommandHandler("ask", telegram.diagnose_cmd.handle_ask)) + + # K8s / Docker 操作 + app.add_handler(CommandHandler("pods", telegram._cmd_pods)) + app.add_handler(CommandHandler("services", telegram._cmd_services)) + app.add_handler(CommandHandler("scale", telegram._cmd_scale)) + app.add_handler(CommandHandler("kubectl", telegram._cmd_kubectl)) + app.add_handler(CommandHandler("containers", telegram._cmd_containers)) + app.add_handler(CommandHandler("docker", telegram._cmd_docker)) + + # ClawBot 原有 callback handlers 必須優先於 AWOOOI bridge。 + app.add_handler(confirmation_handler.get_callback_handler()) + app.add_handler(incident_handler.get_callback_handler()) + app.add_handler( + CallbackQueryHandler( + telegram._handle_approval_callback, + pattern=CLAWBOT_APPROVAL_CALLBACK_PATTERN, + ) + ) + app.add_handler( + CallbackQueryHandler( + telegram._handle_external_callback, + pattern=AWOOOI_READ_CALLBACK_PATTERN, + ) + ) + app.add_handler(CallbackQueryHandler(telegram._handle_blocked_callback)) + + logger.info("✅ Telegram Bot handlers registered via setup_bot()") diff --git a/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/core/config.py b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/core/config.py new file mode 100644 index 000000000..08103e595 --- /dev/null +++ b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/core/config.py @@ -0,0 +1,252 @@ +""" +OpenClaw v5.0 - Core Configuration +企業級 AIOps 閉環架構配置模組 + +使用 Pydantic Settings 進行嚴格環境變數驗證 +""" + +from functools import lru_cache +from typing import Literal +from pydantic import Field, field_validator +from pydantic_settings import BaseSettings + + +class Settings(BaseSettings): + """OpenClaw v5.0 配置 - 所有環境變數的單一真理來源""" + + # ============ 應用程式基本設定 ============ + APP_NAME: str = "OpenClaw" + APP_VERSION: str = "5.0.0" + APP_ENV: Literal["development", "staging", "production"] = "development" + DEBUG: bool = False + HOST: str = "0.0.0.0" + PORT: int = 8088 + + # ============ Telegram Bot 設定 ============ + TELEGRAM_BOT_TOKEN: str = Field(..., description="Telegram Bot API Token") + TELEGRAM_ADMIN_CHAT_ID: str = Field(..., description="管理員 Chat ID (用於告警)") + TELEGRAM_ALLOWED_USERS: str = Field( + default="", + description="允許操作的用戶 ID,逗號分隔" + ) + TELEGRAM_CALLBACK_FORWARD_ENABLED: bool = Field( + default=False, + description="明確啟用 AWOOOI 讀取型 callback HMAC forwarder", + ) + TELEGRAM_CALLBACK_FORWARD_URL: str = Field( + default="https://awoooi.wooo.work/api/v1/telegram/callback-forward", + description="AWOOOI callback ingress(只允許 HTTPS allowlisted endpoint)", + ) + TELEGRAM_CALLBACK_FORWARD_ALLOWED_HOSTS: str = Field( + default="awoooi.wooo.work", + description="允許的 callback forward host,逗號分隔", + ) + TELEGRAM_CALLBACK_FORWARD_TIMEOUT_SECONDS: float = Field( + default=3.0, + ge=0.25, + le=10.0, + description="callback forward 總 HTTP timeout 秒數", + ) + TELEGRAM_CALLBACK_FORWARD_MAX_BODY_BYTES: int = Field( + default=16384, + ge=1024, + le=32768, + description="canonical callback body 大小上限", + ) + + @field_validator("TELEGRAM_ALLOWED_USERS") + @classmethod + def parse_allowed_users(cls, v: str) -> str: + return v + + def get_allowed_user_ids(self) -> list[int]: + """解析允許的用戶 ID 列表""" + if not self.TELEGRAM_ALLOWED_USERS: + return [] + return [int(uid.strip()) for uid in self.TELEGRAM_ALLOWED_USERS.split(",") if uid.strip()] + + def get_callback_forward_allowed_hosts(self) -> frozenset[str]: + """解析 callback destination allowlist;空清單會讓 forwarder fail closed。""" + return frozenset( + host.strip().lower() + for host in self.TELEGRAM_CALLBACK_FORWARD_ALLOWED_HOSTS.split(",") + if host.strip() + ) + + # ============ SigNoz / OpenTelemetry 設定 ============ + OTEL_EXPORTER_OTLP_ENDPOINT: str = Field( + default="http://localhost:4317", + description="OpenTelemetry Collector gRPC 端點" + ) + OTEL_SERVICE_NAME: str = "clawbot" + SIGNOZ_API_URL: str = Field( + default="http://localhost:3301", + description="SigNoz API 端點" + ) + + # ============ n8n 整合設定 ============ + N8N_WEBHOOK_URL: str = Field( + default="http://localhost:5678", + description="n8n Webhook 基礎 URL" + ) + N8N_API_KEY: str = Field(default="", description="n8n API Key (可選)") + + # ============ Grafana OnCall 設定 ============ + GRAFANA_ONCALL_URL: str = Field( + default="", + description="Grafana OnCall Webhook URL" + ) + GRAFANA_ONCALL_API_KEY: str = Field(default="", description="Grafana OnCall API Key") + + # ============ LLM 設定 ============ + # 2026-04-01 ogt: 修復生產 AI 無法使用 — 預設改 ollama,OLLAMA_BASE_URL 指向 AI 主機 + # 原因: ANTHROPIC_API_KEY 未注入容器環境導致 "AI 暫時無法使用" 錯誤 + # fallback 順序: Ollama → Gemini → Claude (feedback_ai_fallback_order.md) + LLM_PROVIDER: Literal["anthropic", "openai", "ollama"] = "ollama" + ANTHROPIC_API_KEY: str = Field(default="", description="Anthropic Claude API Key") + OPENAI_API_KEY: str = Field(default="", description="OpenAI API Key") + # 2026-03-31 ogt: NVIDIA NIM / 自訂 OpenAI 相容端點 (空字串=使用預設 api.openai.com) + OPENAI_BASE_URL: str = Field( + default="", + description="OpenAI 相容 API Base URL (e.g. https://integrate.api.nvidia.com/v1)" + ) + OLLAMA_BASE_URL: str = Field( + default="http://192.168.0.188:11434", + description="Ollama LLM 端點 (AI 主機 192.168.0.188)" + ) + LLM_MODEL: str = Field( + default="qwen2.5:7b-instruct", + description="預設 LLM 模型 (Ollama: qwen2.5:7b-instruct)" + ) + + # ============ Git 設定 ============ + GITLAB_URL: str = Field( + default="http://192.168.0.110", + description="GitLab 伺服器 URL" + ) + GITLAB_TOKEN: str = Field(default="", description="GitLab Personal Access Token") + GITHUB_TOKEN: str = Field(default="", description="GitHub Personal Access Token") + GIT_DEFAULT_BRANCH: str = "main" + + # ============ Gitea 設定 (知識沉澱) ============ + GITEA_URL: str = Field( + default="http://192.168.0.110:3001", + description="Gitea 伺服器 URL" + ) + GITEA_TOKEN: str = Field(default="", description="Gitea Personal Access Token") + GITEA_OWNER: str = Field(default="wooo", description="Gitea 組織/用戶名") + GITEA_REPO: str = Field(default="wooo-aiops", description="預設 Issue 存放的 Repository") + + # ============ Ansible 設定 ============ + ANSIBLE_PLAYBOOKS_PATH: str = Field( + default="/opt/clawbot/ansible/playbooks", + description="Ansible Playbooks 目錄" + ) + ANSIBLE_INVENTORY_PATH: str = Field( + default="/opt/clawbot/ansible/inventory", + description="Ansible Inventory 檔案路徑" + ) + ANSIBLE_SSH_KEY_PATH: str = Field( + default="~/.ssh/id_rsa", + description="SSH 私鑰路徑" + ) + + # ============ Redis 設定 (用於狀態/快取) ============ + REDIS_URL: str = Field( + default="redis://localhost:6379/0", + description="Redis 連線 URL" + ) + + # ============ OPS.86 Enterprise AI Stack Feature Flags ============ + USE_INSTRUCTOR_PARSER: bool = Field( + default=True, + description="Feature Flag: 啟用 Instructor 結構化輸出 (設 false 回退到 legacy JSON 解析)" + ) + SHADOW_MODE_ENABLED: bool = Field( + default=False, + description="Shadow Mode: 同時執行新舊邏輯並比對結果 (用於驗證 Instructor 準確率)" + ) + + # ============ OPS.87 Semantic Cache 設定 ============ + SEMANTIC_CACHE_ENABLED: bool = Field( + default=True, + description="Feature Flag: 啟用語意快取層" + ) + SEMANTIC_CACHE_TTL_SECONDS: int = Field( + default=3600, + description="快取 TTL (秒) - 預設 1 小時" + ) + SEMANTIC_CACHE_SIMILARITY_THRESHOLD: float = Field( + default=0.92, + description="語意相似度閾值 (0.0-1.0) - 超過此值視為快取命中" + ) + SEMANTIC_CACHE_USE_VECTOR_SEARCH: bool = Field( + default=False, + description="使用 Redis Stack 向量搜索 (需升級 Redis)" + ) + EMBEDDING_MODEL: str = Field( + default="nomic-embed-text", + description="Ollama Embedding 模型名稱" + ) + + # ============ 安全設定 ============ + API_SECRET_KEY: str = Field( + default="change-me-in-production", + description="API 簽名密鑰" + ) + APPROVAL_TIMEOUT_SECONDS: int = Field( + default=300, + description="審批超時時間 (秒)" + ) + + # ============ HA 高可用設定 ============ + STANDBY_MODE: bool = Field( + default=False, + description="Standby 模式 - 不啟動 Telegram polling,僅運行 API" + ) + INSTANCE_ID: str = Field( + default="primary", + description="實例識別碼 (primary / standby)" + ) + PRIMARY_HEALTH_URL: str = Field( + default="http://192.168.0.188:8088/health", + description="Primary 實例健康檢查 URL" + ) + + # ============ L3 Dead Man's Switch (Healthchecks.io) ============ + HEALTHCHECKS_PING_URL: str = Field( + default="", + description="Healthchecks.io Ping URL (格式: https://hc-ping.com/your-uuid)" + ) + + # ============ 監控目標 ============ + MONITORED_SITES: str = Field( + default="https://aiops.wooo.work,https://mo.wooo.work", + description="監控的網站列表,逗號分隔" + ) + + def get_monitored_sites(self) -> list[str]: + """解析監控網站列表""" + return [s.strip() for s in self.MONITORED_SITES.split(",") if s.strip()] + + # ============ Playwright 設定 ============ + PLAYWRIGHT_HEADLESS: bool = True + PLAYWRIGHT_TIMEOUT_MS: int = 30000 + SCREENSHOT_PATH: str = "/tmp/clawbot/screenshots" + + model_config = { + "env_file": ".env", + "env_file_encoding": "utf-8", + "case_sensitive": True, + "extra": "ignore" + } + + +@lru_cache +def get_settings() -> Settings: + """取得快取的設定實例""" + return Settings() + + +# 快捷存取 +settings = get_settings() diff --git a/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/services/telegram_callback_forwarder.py b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/services/telegram_callback_forwarder.py new file mode 100644 index 000000000..8f367f8e7 --- /dev/null +++ b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/app/services/telegram_callback_forwarder.py @@ -0,0 +1,408 @@ +"""Fail-closed forwarding for AWOOOI-owned Telegram callback queries. + +ClawBot owns the Telegram ``getUpdates`` loop. AWOOOI may render a small set +of read-only buttons, so those callback updates have to cross the service +boundary without transferring the bot token or giving unknown/write actions a +new execution path. + +The forwarder deliberately serializes a minimal Telegram Update projection, +signs the exact compact body, follows no redirects, and treats only a 2xx +response as an acknowledgement. It never logs callback payloads, actor ids, +the signing key, or upstream response bodies. +""" + +from __future__ import annotations + +import hashlib +import hmac +import json +import os +import re +import stat +import time +from collections.abc import Mapping +from dataclasses import dataclass +from typing import Any +from urllib.parse import urlsplit + +import httpx + +AWOOOI_READ_CALLBACK_ACTIONS = frozenset( + { + "detail", + "history", + "reanalyze", + "drift_view_page", + "check_process", + "check_log_nginx", + "check_log_container", + "check_log_minio", + "check_port", + "check_health", + "check_pod_logs", + "describe_pod", + "open_signoz", + "open_flywheel", + "backup_check_host_disk", + "backup_check_jobs", + "backup_check_velero", + } +) + +CLAWBOT_APPROVAL_CALLBACK_PATTERN = r"^(?:repair|rag|approve|reject):" +AWOOOI_READ_CALLBACK_PATTERN = ( + r"^(?:" + + "|".join(sorted(map(re.escape, AWOOOI_READ_CALLBACK_ACTIONS))) + + r"):[A-Za-z0-9][A-Za-z0-9_.@/-]*$" +) + +_EXPECTED_PATH = "/api/v1/telegram/callback-forward" +_CONTROLLED_IDENTITY_PATH = ( + "/run/awoooi/openclaw-callback-forwarder/current-deployment.json" +) +_CONTROLLED_IDENTITY_MAX_BYTES = 4_096 +_CONTROLLED_IDENTITY_KEYS = frozenset( + {"schema_version", "catalog_id", "trace_id", "run_id", "work_item_id"} +) +_CONTROLLED_IDENTITY_SCHEMA = "openclaw_callback_deployment_identity_v1" +_CONTROLLED_CORRELATION_SCHEMA = "awoooi_controlled_callback_correlation_v1" +_CONTROLLED_CATALOG_ID = "ansible:188-openclaw-callback-forwarder" +_SAFE_CORRELATION_ID_RE = re.compile( + r"^[A-Za-z0-9](?:[A-Za-z0-9_.:-]{0,158}[A-Za-z0-9])?$" +) +_SAFE_VALUE_RE = re.compile(r"^[\x20-\x7e]+$") +_SAFE_ACTION_RE = re.compile(r"^[a-z0-9_]{1,64}$") +_SAFE_CALLBACK_VALUE_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.@/-]*$") +_MAX_CALLBACK_QUERY_ID_BYTES = 256 +_MAX_CHAT_INSTANCE_BYTES = 256 +_MAX_TEXT_BYTES = 8192 +_MAX_NAME_BYTES = 256 + + +class CallbackPayloadRejected(ValueError): + """Raised when an update cannot be safely projected for forwarding.""" + + +def load_controlled_correlation( + path: str = _CONTROLLED_IDENTITY_PATH, +) -> dict[str, str] | None: + """Read one bounded non-secret deployment identity without following links.""" + flags = os.O_RDONLY | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_NOFOLLOW", 0) + try: + descriptor = os.open(path, flags) + except OSError: + return None + try: + metadata = os.fstat(descriptor) + if ( + not stat.S_ISREG(metadata.st_mode) + or metadata.st_size <= 0 + or metadata.st_size > _CONTROLLED_IDENTITY_MAX_BYTES + ): + return None + with os.fdopen(descriptor, "r", encoding="utf-8", closefd=False) as stream: + raw = stream.read(_CONTROLLED_IDENTITY_MAX_BYTES + 1) + except (OSError, UnicodeError): + return None + finally: + os.close(descriptor) + if len(raw.encode("utf-8")) > _CONTROLLED_IDENTITY_MAX_BYTES: + return None + try: + identity = json.loads(raw) + except (TypeError, ValueError, json.JSONDecodeError): + return None + if not isinstance(identity, Mapping) or set(identity) != _CONTROLLED_IDENTITY_KEYS: + return None + if ( + identity.get("schema_version") != _CONTROLLED_IDENTITY_SCHEMA + or identity.get("catalog_id") != _CONTROLLED_CATALOG_ID + ): + return None + correlation = {"schema_version": _CONTROLLED_CORRELATION_SCHEMA} + for key in ("trace_id", "run_id", "work_item_id"): + value = identity.get(key) + if not isinstance(value, str) or not _SAFE_CORRELATION_ID_RE.fullmatch(value): + return None + correlation[key] = value + return correlation + + +@dataclass(frozen=True) +class CallbackForwardResult: + acknowledged: bool + reason: str + status_code: int | None = None + + +def callback_action(callback_data: str) -> str: + """Return a redaction-safe action prefix, or ``unknown``.""" + action = str(callback_data or "").split(":", 1)[0].strip() + return action if _SAFE_ACTION_RE.fullmatch(action) else "unknown" + + +def is_awoooi_read_callback(callback_data: str, *, max_data_bytes: int = 64) -> bool: + """Accept only the explicit read-only action allowlist and Telegram size.""" + raw = str(callback_data or "") + try: + encoded = raw.encode("utf-8") + except UnicodeError: + return False + if not raw or len(encoded) > max_data_bytes or not _SAFE_VALUE_RE.fullmatch(raw): + return False + action, separator, value = raw.partition(":") + return bool( + separator + and action in AWOOOI_READ_CALLBACK_ACTIONS + and _SAFE_CALLBACK_VALUE_RE.fullmatch(value) + ) + + +def _bounded_string(value: Any, *, max_bytes: int, required: bool = False) -> str | None: + if value is None: + if required: + raise CallbackPayloadRejected("required_field_missing") + return None + text = str(value) + if required and not text: + raise CallbackPayloadRejected("required_field_empty") + if len(text.encode("utf-8")) > max_bytes: + raise CallbackPayloadRejected("field_too_large") + return text + + +def _copy_optional_strings( + source: Mapping[str, Any], + target: dict[str, Any], + fields: tuple[str, ...], + *, + max_bytes: int, +) -> None: + for field in fields: + value = _bounded_string(source.get(field), max_bytes=max_bytes) + if value is not None: + target[field] = value + + +def canonical_callback_update( + update: Mapping[str, Any], + *, + max_data_bytes: int = 64, +) -> dict[str, Any]: + """Build the minimal callback Update accepted by the AWOOOI receiver.""" + if not isinstance(update, Mapping): + raise CallbackPayloadRejected("update_not_mapping") + update_id = update.get("update_id") + if not isinstance(update_id, int) or isinstance(update_id, bool) or update_id < 0: + raise CallbackPayloadRejected("invalid_update_id") + + query = update.get("callback_query") + if not isinstance(query, Mapping): + raise CallbackPayloadRejected("callback_query_missing") + data = _bounded_string(query.get("data"), max_bytes=max_data_bytes, required=True) + if data is None or not is_awoooi_read_callback(data, max_data_bytes=max_data_bytes): + raise CallbackPayloadRejected("action_not_allowlisted") + + actor = query.get("from") + if not isinstance(actor, Mapping): + raise CallbackPayloadRejected("actor_missing") + actor_id = actor.get("id") + if not isinstance(actor_id, int) or isinstance(actor_id, bool): + raise CallbackPayloadRejected("invalid_actor_id") + actor_projection: dict[str, Any] = {"id": actor_id} + if isinstance(actor.get("is_bot"), bool): + actor_projection["is_bot"] = actor["is_bot"] + _copy_optional_strings( + actor, + actor_projection, + ("first_name", "last_name", "username", "language_code"), + max_bytes=_MAX_NAME_BYTES, + ) + + query_projection: dict[str, Any] = { + "data": data, + "from": actor_projection, + "id": _bounded_string( + query.get("id"), max_bytes=_MAX_CALLBACK_QUERY_ID_BYTES, required=True + ), + } + chat_instance = _bounded_string( + query.get("chat_instance"), max_bytes=_MAX_CHAT_INSTANCE_BYTES + ) + if chat_instance is not None: + query_projection["chat_instance"] = chat_instance + + message = query.get("message") + if isinstance(message, Mapping): + message_id = message.get("message_id") + if not isinstance(message_id, int) or isinstance(message_id, bool): + raise CallbackPayloadRejected("invalid_message_id") + message_projection: dict[str, Any] = {"message_id": message_id} + date = message.get("date") + if isinstance(date, int) and not isinstance(date, bool): + message_projection["date"] = date + text = _bounded_string(message.get("text"), max_bytes=_MAX_TEXT_BYTES) + if text is not None: + message_projection["text"] = text + + chat = message.get("chat") + if not isinstance(chat, Mapping): + raise CallbackPayloadRejected("message_chat_missing") + chat_id = chat.get("id") + if not isinstance(chat_id, int) or isinstance(chat_id, bool): + raise CallbackPayloadRejected("invalid_chat_id") + chat_projection: dict[str, Any] = {"id": chat_id} + _copy_optional_strings( + chat, + chat_projection, + ("type", "title", "username", "first_name", "last_name"), + max_bytes=_MAX_NAME_BYTES, + ) + message_projection["chat"] = chat_projection + query_projection["message"] = message_projection + else: + inline_message_id = _bounded_string( + query.get("inline_message_id"), + max_bytes=_MAX_CALLBACK_QUERY_ID_BYTES, + ) + if inline_message_id is None: + raise CallbackPayloadRejected("message_identity_missing") + query_projection["inline_message_id"] = inline_message_id + + return {"callback_query": query_projection, "update_id": update_id} + + +def compact_json_body(payload: Mapping[str, Any], *, max_body_bytes: int) -> bytes: + body = json.dumps( + payload, + ensure_ascii=False, + separators=(",", ":"), + sort_keys=True, + ).encode("utf-8") + if len(body) > max_body_bytes: + raise CallbackPayloadRejected("body_too_large") + return body + + +def callback_signature(signing_key: str, timestamp: str, body: bytes) -> str: + digest = hmac.new( + signing_key.encode("utf-8"), + timestamp.encode("ascii") + b"." + body, + hashlib.sha256, + ).hexdigest() + return f"sha256={digest}" + + +def validate_forward_url(url: str, allowed_hosts: frozenset[str]) -> str: + parsed = urlsplit(str(url or "").strip()) + hostname = (parsed.hostname or "").lower() + if ( + parsed.scheme != "https" + or hostname not in allowed_hosts + or parsed.path != _EXPECTED_PATH + or parsed.username is not None + or parsed.password is not None + or parsed.query + or parsed.fragment + or (parsed.port not in {None, 443}) + ): + raise ValueError("forward_url_not_allowlisted") + return parsed.geturl() + + +class TelegramCallbackForwarder: + """Bounded HMAC client for the AWOOOI callback ingress.""" + + def __init__( + self, + *, + enabled: bool, + url: str, + signing_key: str, + allowed_hosts: frozenset[str] = frozenset({"awoooi.wooo.work"}), + timeout_seconds: float = 3.0, + max_body_bytes: int = 16_384, + max_data_bytes: int = 64, + controlled_identity_path: str = _CONTROLLED_IDENTITY_PATH, + ) -> None: + self.enabled = bool(enabled) + self.url = str(url or "") + self._signing_key = str(signing_key or "") + self.allowed_hosts = frozenset(host.lower() for host in allowed_hosts if host) + self.timeout_seconds = min(max(float(timeout_seconds), 0.25), 10.0) + self.max_body_bytes = min(max(int(max_body_bytes), 1_024), 32_768) + self.max_data_bytes = min(max(int(max_data_bytes), 1), 64) + self.controlled_identity_path = str(controlled_identity_path) + + async def forward( + self, + update: Mapping[str, Any], + *, + now: int | None = None, + client: httpx.AsyncClient | None = None, + ) -> CallbackForwardResult: + if not self.enabled: + return CallbackForwardResult(False, "forwarding_disabled") + if not self._signing_key: + return CallbackForwardResult(False, "signing_key_unavailable") + + try: + url = validate_forward_url(self.url, self.allowed_hosts) + payload = canonical_callback_update( + update, + max_data_bytes=self.max_data_bytes, + ) + correlation = load_controlled_correlation(self.controlled_identity_path) + if correlation is not None: + payload["controlled_correlation"] = correlation + body = compact_json_body(payload, max_body_bytes=self.max_body_bytes) + except (CallbackPayloadRejected, ValueError, UnicodeError): + return CallbackForwardResult(False, "payload_or_destination_rejected") + + timestamp = str(int(time.time()) if now is None else int(now)) + headers = { + "Content-Type": "application/json", + "X-Telegram-Forward-Schema": "telegram_callback_forward_v1", + "X-Telegram-Forward-Timestamp": timestamp, + "X-Telegram-Forward-Signature": callback_signature( + self._signing_key, + timestamp, + body, + ), + } + + owns_client = client is None + active_client = client or httpx.AsyncClient( + timeout=httpx.Timeout(self.timeout_seconds), + follow_redirects=False, + ) + try: + response = await active_client.post(url, content=body, headers=headers) + except httpx.TimeoutException: + return CallbackForwardResult(False, "upstream_timeout") + except httpx.HTTPError: + return CallbackForwardResult(False, "upstream_transport_error") + finally: + if owns_client: + await active_client.aclose() + + if 200 <= response.status_code < 300: + return CallbackForwardResult(True, "forward_acknowledged", response.status_code) + return CallbackForwardResult(False, "upstream_non_2xx", response.status_code) + + +__all__ = [ + "AWOOOI_READ_CALLBACK_ACTIONS", + "AWOOOI_READ_CALLBACK_PATTERN", + "CLAWBOT_APPROVAL_CALLBACK_PATTERN", + "CallbackForwardResult", + "CallbackPayloadRejected", + "TelegramCallbackForwarder", + "callback_action", + "callback_signature", + "canonical_callback_update", + "compact_json_body", + "is_awoooi_read_callback", + "load_controlled_correlation", + "validate_forward_url", +] diff --git a/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/manifest.json b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/manifest.json new file mode 100644 index 000000000..ce6c36430 --- /dev/null +++ b/infra/ansible/files/openclaw-callback-forwarder/8f367f8e7d6104c44b29bd5fa2d894db09ddf365/manifest.json @@ -0,0 +1,32 @@ +{ + "schema_version": "openclaw_callback_overlay_manifest_v1", + "source_repository": "gitea:wooo/clawbot-v5", + "source_commit": "ffc45311b6403b30cb6ae6e1aa83dae1a047e9c9", + "overlay_revision_kind": "git_blob_of_callback_forwarder", + "overlay_revision": "8f367f8e7d6104c44b29bd5fa2d894db09ddf365", + "deployment_mode": "host_user_owned_read_only_compose_override_no_build_no_pull", + "runtime_source_mutation_allowed": false, + "runtime_compose_override_write_allowed": true, + "host_privilege_escalation_required": false, + "runtime_secret_read_allowed": false, + "payloads": [ + { + "path": "app/bot/telegram.py", + "git_blob": "f9426b4b0ee3b685bbb155d7045c75b94a6bf0eb", + "sha256": "af9600f32305ae73184651a8d3437063d7c8d35e3c3eb0910a2e27ac874c9302" + }, + { + "path": "app/core/config.py", + "git_blob": "08103e595dd82052ddbd55622f43ca77f8ca2fff", + "sha256": "e1b8b633f7b4780d3078d6c42a55e49d200255d3597a34e0cd2fd649b84eaea1" + }, + { + "path": "app/services/telegram_callback_forwarder.py", + "git_blob": "8f367f8e7d6104c44b29bd5fa2d894db09ddf365", + "sha256": "8e8e0faf6322abb21d670ad9094f1805d7a4b3643da8fe88d74a8c9591066bf3" + } + ], + "superseded_by": null, + "expiry": "after_clawbot_runtime_source_is_reconciled_to_an_immutable_gitea_image", + "exit_condition": "remove_overlay_only_after_clean_immutable_image_and_real_callback_replay_are_verified" +} diff --git a/infra/ansible/playbooks/188-openclaw-callback-forwarder.yml b/infra/ansible/playbooks/188-openclaw-callback-forwarder.yml index 7134c0cf5..07c32fb50 100644 --- a/infra/ansible/playbooks/188-openclaw-callback-forwarder.yml +++ b/infra/ansible/playbooks/188-openclaw-callback-forwarder.yml @@ -16,8 +16,8 @@ vars: catalog_id: ansible:188-openclaw-callback-forwarder canonical_asset_id: service:openclaw:host188 - expected_source_sha: ffc45311b6403b30cb6ae6e1aa83dae1a047e9c9 - expected_manifest_sha256: dafeb467b744ee7e80dc043d51f9e5d7e6905a58103f1710d519b553b56c6e48 + expected_source_sha: 8f367f8e7d6104c44b29bd5fa2d894db09ddf365 + expected_manifest_sha256: 0e5edaf700d264b44cc0eae3b9429af642be0b1563a0bffa633942d30d2c2e1c allowed_working_directories: - /home/ollama/clawbot-v5 compose_project: clawbot @@ -29,6 +29,11 @@ payload_root: "{{ managed_base }}/payloads/{{ expected_source_sha }}" compose_override_path: "{{ openclaw_workdir }}/docker-compose.override.yml" receipt_root: "{{ managed_base }}/receipts" + callback_identity_path: "{{ receipt_root }}/current-deployment.json" + callback_identity_container_path: >- + /run/awoooi/openclaw-callback-forwarder/current-deployment.json + superseded_compose_override_sha256s: + - 8f6b7e8a0d6746c1331e99910656e64deebf67c9582bab3e1d36f71bb7a253c0 rollback_root: "{{ managed_base }}/rollback" trace_id: "{{ controlled_trace_id | default('') }}" run_id: "{{ controlled_run_id | default('') }}" @@ -60,7 +65,7 @@ {{ payload_root }}/app/services/telegram_callback_forwarder.py container_path: /app/app/services/telegram_callback_forwarder.py baseline_sha256: "" - target_sha256: 0d7693e34f63db4491d974f9f76a3a1c793d32cda193cb55ed3d8c86297b4fe3 + target_sha256: 8e8e0faf6322abb21d670ad9094f1805d7a4b3643da8fe88d74a8c9591066bf3 managed_directories: - path: "{{ managed_base }}" mode: "0755" @@ -428,6 +433,41 @@ check_mode: false no_log: true + - name: Inspect existing controlled callback deployment identity + ansible.builtin.stat: + path: "{{ callback_identity_path }}" + follow: false + register: callback_identity_before + check_mode: false + changed_when: false + no_log: true + + - name: Read existing controlled callback deployment identity for rollback + ansible.builtin.slurp: + src: "{{ callback_identity_path }}" + register: callback_identity_content_before + when: callback_identity_before.stat.exists | default(false) + check_mode: false + changed_when: false + no_log: true + + - name: Fail closed on foreign callback deployment identity + ansible.builtin.assert: + that: + - >- + not (callback_identity_before.stat.exists | default(false)) + or ( + callback_identity_before.stat.isreg | default(false) | bool + and not (callback_identity_before.stat.islnk | default(false) | bool) + and (callback_identity_before.stat.size | default(4097) | int) <= 4096 + and (callback_identity_before.stat.mode | default('')) == '0444' + and (callback_identity_before.stat.pw_name | default('')) == 'ollama' + and (callback_identity_before.stat.gr_name | default('')) == 'ollama' + ) + fail_msg: openclaw_callback_deployment_identity_foreign_owner + changed_when: false + no_log: true + - name: Inspect existing user-owned overlay payloads ansible.builtin.stat: path: "{{ item.overlay_path }}" @@ -481,6 +521,14 @@ - name: Define exact non-secret managed overlay documents ansible.builtin.set_fact: + callback_identity_content: | + { + "schema_version": "openclaw_callback_deployment_identity_v1", + "catalog_id": "{{ catalog_id }}", + "trace_id": "{{ trace_id }}", + "run_id": "{{ run_id }}", + "work_item_id": "{{ work_item_id }}" + } compose_override_content: | services: clawbot: @@ -494,6 +542,7 @@ - "{{ payload_root }}/app/bot/telegram.py:/app/app/bot/telegram.py:ro" - "{{ payload_root }}/app/core/config.py:/app/app/core/config.py:ro" - "{{ payload_root }}/app/services/telegram_callback_forwarder.py:/app/app/services/telegram_callback_forwarder.py:ro" + - "{{ callback_identity_path }}:{{ callback_identity_container_path }}:ro" changed_when: false - name: Fail closed on foreign pre-existing managed artifacts @@ -502,14 +551,27 @@ - >- not (managed_artifacts_before.results[0].stat.exists | default(false)) or ( - managed_artifact_contents_before.results - | selectattr('item.item', 'equalto', compose_override_path) - | map(attribute='content') - | map('b64decode') - | list - | first - | default('') - ) == compose_override_content + ( + ( + managed_artifact_contents_before.results + | selectattr('item.item', 'equalto', compose_override_path) + | map(attribute='content') + | map('b64decode') + | list + | first + | default('') + ) == compose_override_content + or ( + managed_artifact_contents_before.results + | selectattr('item.item', 'equalto', compose_override_path) + | map(attribute='content') + | map('b64decode') + | list + | first + | default('') + | hash('sha256') + ) in superseded_compose_override_sha256s + ) and (managed_artifacts_before.results[0].stat.isreg | default(false) | bool) and not (managed_artifacts_before.results[0].stat.islnk @@ -768,6 +830,14 @@ mode: "0444" content: "{{ compose_override_content }}" + - name: Install same-run non-secret callback deployment identity + ansible.builtin.copy: + dest: "{{ callback_identity_path }}" + owner: ollama + group: ollama + mode: "0444" + content: "{{ callback_identity_content }}" + - name: Mark the bounded service activation attempt ansible.builtin.set_fact: activation_attempted: true @@ -1159,6 +1229,22 @@ changed_when: false rescue: + - name: Remove callback identity created by the failed activation + ansible.builtin.file: + path: "{{ callback_identity_path }}" + state: absent + when: not (callback_identity_before.stat.exists | default(false)) + + - name: Restore exact callback identity after failed activation + ansible.builtin.copy: + dest: "{{ callback_identity_path }}" + owner: ollama + group: ollama + mode: "0444" + content: "{{ callback_identity_content_before.content | b64decode }}" + when: callback_identity_before.stat.exists | default(false) + no_log: true + - name: Remove managed Compose override after failed activation ansible.builtin.file: path: "{{ compose_override_path }}"