diff --git a/apps/api/src/services/telegram_controlled_action_work_item.py b/apps/api/src/services/telegram_controlled_action_work_item.py new file mode 100644 index 000000000..47a364346 --- /dev/null +++ b/apps/api/src/services/telegram_controlled_action_work_item.py @@ -0,0 +1,483 @@ +"""Durable, fail-closed work items for SRE war-room action requests. + +This module turns a verified Telegram conversation turn into either a Codex +development candidate or an asset-identity drift item. It deliberately does +not call an AI provider, executor, Agent99, or any runtime mutation path. +""" + +from __future__ import annotations + +import base64 +import hashlib +import json +import re +import unicodedata +from collections.abc import Awaitable, Callable, Mapping +from dataclasses import dataclass +from html import escape +from typing import Any +from uuid import UUID + +import structlog +from sqlalchemy import text + +from src.services.audit_sink import sanitize +from src.services.service_registry import get_service_registry + +logger = structlog.get_logger(__name__) + +_RECEIPT_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:@+-]{0,191}$") +_TARGET_TOKEN = re.compile(r"[A-Za-z0-9][A-Za-z0-9._:-]{1,127}") +_NON_TARGET_TOKENS = { + "apply", + "clear", + "execute", + "fix", + "now", + "please", + "restart", + "rollback", + "scale", + "service", + "up", + "down", +} +_ACTION_MARKERS: tuple[tuple[str, tuple[str, ...]], ...] = ( + ("restart", ("重啟", "重新啟動", "restart")), + ("rollback", ("回滾", "rollback")), + ("scale_up", ("擴容", "scale up", "scale-up")), + ("scale_down", ("縮容", "scale down", "scale-down")), + ("clear_cache", ("清快取", "清除快取", "clear cache")), + ("controlled_investigation", ("調查", "診斷", "investigate", "diagnose")), + ("controlled_action", ("修復", "執行", "fix", "execute")), +) + +PersistWorkItem = Callable[[dict[str, Any]], Awaitable[dict[str, Any]]] + + +def _normalized_request(message: str) -> str: + return " ".join( + unicodedata.normalize("NFKC", str(message or "")).casefold().split() + ) + + +def _stable_digest(*values: str) -> str: + canonical = json.dumps(values, ensure_ascii=False, separators=(",", ":")) + digest = hashlib.sha256(canonical.encode("utf-8")).digest() + return base64.urlsafe_b64encode(digest).decode("ascii").rstrip("=") + + +def _request_digest(value: str) -> str: + digest = hashlib.sha256(value.encode("utf-8")).digest() + return base64.urlsafe_b64encode(digest).decode("ascii").rstrip("=") + + +def _unresolved_identity(value: str, *, reason: str) -> dict[str, Any]: + normalized = str(value or "missing-target").strip().casefold() or "missing-target" + digest = hashlib.sha256(normalized.encode("utf-8")).hexdigest()[:12].upper() + return { + "schema_version": "canonical_asset_identity_v2", + "resolution_status": "asset_identity_unresolved", + "resolution_reason": reason, + "normalized_identity": normalized, + "canonical_id": None, + "asset_domain": "unknown", + "executor": None, + "verifier": None, + "controlled_apply_allowed": False, + "drift_work_item_id": f"AIA-ASSET-DRIFT-{digest}", + } + + +@dataclass(frozen=True) +class ControlledActionWorkItemResult: + """Public-safe receipt returned to the Telegram gateway.""" + + accepted: bool + created: bool + deduplicated: bool + candidate_created: bool + status: str + trace_id: str + work_item_id: str + receipt_id: str + fingerprint: str + canonical_asset_id: str + asset_domain: str + next_safe_action: str + + def render_html(self) -> str: + if not self.accepted: + return ( + "⛔ 受控 work item 未建立\n" + f"trace={escape(self.trace_id)}\n" + f"status={escape(self.status)}\n" + f"下一安全動作:{escape(self.next_safe_action)}\n\n" + "provider=none | paid_provider=false | executor_invoked=false | " + "Agent99=false | runtime_mutation=false" + ) + + title = ( + "資產 identity drift item 已建立" + if self.status == "asset_identity_unresolved" + else "受控 Codex work item 已存在,未重建" + if self.deduplicated + else "受控 Codex work item 已建立" + ) + return ( + f"🧭 {title}\n" + f"trace={escape(self.trace_id)}\n" + f"work_item={escape(self.work_item_id)}\n" + f"status={escape(self.status)}\n" + f"asset={escape(self.canonical_asset_id or 'unresolved')}\n" + f"domain={escape(self.asset_domain)}\n" + f"下一安全動作:{escape(self.next_safe_action)}\n\n" + f"receipt={escape(self.receipt_id)} | fingerprint=" + f"{escape(self.fingerprint[:16])} | provider=none | paid_provider=false | " + "executor_invoked=false | Agent99=false | runtime_mutation=false" + ) + + +async def persist_controlled_action_work_item(record: dict[str, Any]) -> dict[str, Any]: + """Insert one deterministic internal event, returning the existing row on conflict.""" + + from src.db.base import get_db_context + + fingerprint = str(record["fingerprint"]) + content_hash = hashlib.sha256(fingerprint.encode("utf-8")).hexdigest() + work_item_id = str(record["work_item_id"]) + provider_event_id = f"sre-controlled-action:{fingerprint}" + run_id = UUID(str(record["run_id"])) + envelope = sanitize(record) + source_envelope = json.dumps(envelope, ensure_ascii=False, default=str) + preview = f"{record['kind']}:{record['status']}:{work_item_id}"[:256] + + async with get_db_context("awoooi") as db: + inserted = await db.execute( + text( + """ + INSERT INTO awooop_conversation_event ( + project_id, channel_type, provider_event_id, + run_id, content_type, content_hash, content_preview, + content_redacted, redaction_version, source_envelope, + is_duplicate, received_at + ) VALUES ( + 'awoooi', 'internal', :provider_event_id, + :run_id, 'command', :content_hash, :preview, + :preview, 'audit_sink_v1', CAST(:source_envelope AS jsonb), + FALSE, NOW() + ) + ON CONFLICT (project_id, channel_type, provider_event_id) + DO NOTHING + RETURNING event_id + """ + ), + { + "provider_event_id": provider_event_id, + "run_id": run_id, + "content_hash": content_hash, + "preview": preview, + "source_envelope": source_envelope, + }, + ) + row = inserted.fetchone() + if row is not None: + return {"receipt_id": str(row[0]), "created": True} + + existing = await db.execute( + text( + """ + SELECT event_id, source_envelope + FROM awooop_conversation_event + WHERE project_id = 'awoooi' + AND channel_type = 'internal' + AND provider_event_id = :provider_event_id + LIMIT 1 + """ + ), + {"provider_event_id": provider_event_id}, + ) + existing_row = existing.fetchone() + if existing_row is None: + raise RuntimeError("controlled_action_dedupe_receipt_missing") + stored_envelope = existing_row[1] + if isinstance(stored_envelope, str): + stored_envelope = json.loads(stored_envelope) + if ( + not isinstance(stored_envelope, Mapping) + or str(stored_envelope.get("work_item_id") or "") != work_item_id + ): + raise RuntimeError("controlled_action_dedupe_receipt_mismatch") + return {"receipt_id": str(existing_row[0]), "created": False} + + +class TelegramControlledActionWorkItemService: + """Create only a candidate/drift receipt from a verified conversation turn.""" + + def __init__( + self, + *, + registry: Any | None = None, + persist_work_item: PersistWorkItem | None = None, + ) -> None: + self._registry = registry or get_service_registry() + self._persist_work_item = ( + persist_work_item or persist_controlled_action_work_item + ) + + @staticmethod + def _action_kind(normalized_message: str) -> str: + for action, markers in _ACTION_MARKERS: + if any(marker in normalized_message for marker in markers): + return action + return "controlled_action" + + def _resolve_message_identity(self, message: str) -> dict[str, Any]: + normalized = _normalized_request(message) + tokens = [ + token.casefold() + for token in _TARGET_TOKEN.findall(normalized) + if token.casefold() not in _NON_TARGET_TOKENS + and not token.casefold().startswith("@") + ] + resolved: dict[str, dict[str, Any]] = {} + for token in dict.fromkeys(tokens): + try: + identity = self._registry.resolve_identity(token) + except Exception as exc: # noqa: BLE001 - identity must fail closed + logger.warning( + "telegram_controlled_action_identity_lookup_failed", + error=type(exc).__name__, + ) + return _unresolved_identity( + f"lookup-error:{_stable_digest(normalized)[:16]}", + reason="registry_lookup_failed", + ) + if identity.get("resolution_status") == "resolved": + canonical_id = str(identity.get("canonical_id") or "") + domain = str(identity.get("asset_domain") or "") + if canonical_id and domain and domain != "unknown": + resolved[canonical_id] = dict(identity) + + if len(resolved) == 1: + return next(iter(resolved.values())) + if len(resolved) > 1: + return _unresolved_identity( + f"ambiguous:{_stable_digest(*sorted(resolved))[:16]}", + reason="multiple_canonical_assets_matched", + ) + + target_hint = ( + tokens[0] + if len(tokens) == 1 + else (f"unparsed:{_stable_digest(normalized)[:16]}") + ) + try: + identity = self._registry.resolve_identity(target_hint) + except Exception as exc: # noqa: BLE001 - identity must fail closed + logger.warning( + "telegram_controlled_action_identity_unresolved", + error=type(exc).__name__, + ) + return _unresolved_identity(target_hint, reason="registry_lookup_failed") + if identity.get("resolution_status") != "resolved": + unresolved = dict(identity) + unresolved.setdefault("resolution_reason", "exact_alias_not_found") + unresolved["controlled_apply_allowed"] = False + unresolved["executor"] = None + unresolved["verifier"] = None + return unresolved + return _unresolved_identity(target_hint, reason="identity_contract_incomplete") + + async def create_from_verified_message( + self, + *, + message: str, + intent: str, + inbound_receipt: Mapping[str, Any], + canonical_chat_identity: Mapping[str, Any], + ) -> ControlledActionWorkItemResult: + """Persist a candidate only after exact chat and inbound receipt validation.""" + + trace_id = str(inbound_receipt.get("trace_id") or "none") + receipt_fields = ("event_id", "trace_id", "run_id") + receipt_verified = all( + _RECEIPT_ID.fullmatch(str(inbound_receipt.get(field) or "")) + for field in receipt_fields + ) + identity_verified = bool( + canonical_chat_identity.get("status") == "verified" + and canonical_chat_identity.get("room") == "awoooi_sre_war_room" + ) + if ( + intent != "controlled_action_request" + or not receipt_verified + or not identity_verified + ): + return ControlledActionWorkItemResult( + accepted=False, + created=False, + deduplicated=False, + candidate_created=False, + status="identity_or_ingress_unverified", + trace_id=trace_id, + work_item_id="none", + receipt_id="none", + fingerprint="none", + canonical_asset_id="", + asset_domain="unknown", + next_safe_action=( + "先完成 canonical chat identity 與 durable ingress receipt 驗證。" + ), + ) + + normalized_message = _normalized_request(message) + if not normalized_message: + return ControlledActionWorkItemResult( + accepted=False, + created=False, + deduplicated=False, + candidate_created=False, + status="empty_controlled_action_request", + trace_id=trace_id, + work_item_id="none", + receipt_id="none", + fingerprint="none", + canonical_asset_id="", + asset_domain="unknown", + next_safe_action="提供一個具名 canonical asset 與受控動作。", + ) + + asset_identity = self._resolve_message_identity(normalized_message) + resolved = asset_identity.get("resolution_status") == "resolved" + canonical_asset_id = str(asset_identity.get("canonical_id") or "") + asset_domain = str(asset_identity.get("asset_domain") or "unknown") + action_kind = self._action_kind(normalized_message) + request_digest = _request_digest(normalized_message) + fingerprint = _stable_digest( + "telegram_controlled_action_work_item_v1", + action_kind, + canonical_asset_id + if resolved + else str(asset_identity.get("normalized_identity") or "unresolved"), + request_digest, + ) + + if resolved: + work_item_id = f"CODEX-SRE-{fingerprint[:16].upper()}" + status = "candidate_recorded" + kind = "codex_development_work_item" + candidate_created = True + next_safe_action = ( + "由 Codex 準備同 domain 的 no-write investigation/candidate;" + "apply 仍須通過 typed policy、單一 executor 與 independent verifier。" + ) + else: + work_item_id = str(asset_identity.get("drift_work_item_id") or "") + if not work_item_id: + work_item_id = _unresolved_identity( + str(asset_identity.get("normalized_identity") or "missing-target"), + reason="drift_work_item_id_missing", + )["drift_work_item_id"] + status = "asset_identity_unresolved" + kind = "asset_identity_drift_work_item" + candidate_created = False + canonical_asset_id = "" + asset_domain = "unknown" + next_safe_action = ( + "先補齊 canonical asset/domain mapping;完成前維持 fail-closed," + "禁止 fallback AUTO。" + ) + + record = { + "schema_version": "telegram_controlled_action_work_item_v1", + "kind": kind, + "status": status, + "work_item_id": work_item_id, + "fingerprint": fingerprint, + "request_digest": request_digest, + "trace_id": trace_id, + "run_id": str(inbound_receipt["run_id"]), + "inbound_event_id": str(inbound_receipt["event_id"]), + "intent": intent, + "action_kind": action_kind, + "canonical_asset_identity": { + "resolution_status": asset_identity.get("resolution_status"), + "resolution_reason": asset_identity.get("resolution_reason"), + "canonical_id": canonical_asset_id or None, + "asset_domain": asset_domain, + "executor": asset_identity.get("executor") if resolved else None, + "verifier": asset_identity.get("verifier") if resolved else None, + "drift_work_item_id": ( + asset_identity.get("drift_work_item_id") if not resolved else None + ), + }, + "next_safe_action": next_safe_action, + "policy": { + "provider_call_allowed": False, + "paid_provider_call_allowed": False, + "runtime_mutation_allowed": False, + "executor_invocation_allowed": False, + "agent99_dispatch_allowed": False, + "cross_domain_execution_allowed": False, + "fallback_auto_allowed": False, + }, + "raw_message_recorded": False, + "source_commitment": "AIA-CONV-030", + } + try: + durable_receipt = await self._persist_work_item(record) + receipt_id = str(durable_receipt.get("receipt_id") or "") + created = durable_receipt.get("created") is True + if not _RECEIPT_ID.fullmatch(receipt_id): + raise RuntimeError("controlled_action_receipt_invalid") + except Exception as exc: # noqa: BLE001 - no receipt means no work item + logger.warning( + "telegram_controlled_action_work_item_persist_failed", + trace_id=trace_id, + error=type(exc).__name__, + ) + return ControlledActionWorkItemResult( + accepted=False, + created=False, + deduplicated=False, + candidate_created=False, + status="durable_work_item_receipt_failed", + trace_id=trace_id, + work_item_id="none", + receipt_id="none", + fingerprint=fingerprint, + canonical_asset_id=canonical_asset_id, + asset_domain=asset_domain, + next_safe_action=( + "先恢復 durable work-item receipt;不得呼叫 provider 或 " + "runtime executor。" + ), + ) + + return ControlledActionWorkItemResult( + accepted=True, + created=created, + deduplicated=not created, + candidate_created=candidate_created, + status=status, + trace_id=trace_id, + work_item_id=work_item_id, + receipt_id=receipt_id, + fingerprint=fingerprint, + canonical_asset_id=canonical_asset_id, + asset_domain=asset_domain, + next_safe_action=next_safe_action, + ) + + +_service: TelegramControlledActionWorkItemService | None = None + + +def get_telegram_controlled_action_work_item_service() -> ( + TelegramControlledActionWorkItemService +): + global _service + if _service is None: + _service = TelegramControlledActionWorkItemService() + return _service diff --git a/apps/api/src/services/telegram_gateway.py b/apps/api/src/services/telegram_gateway.py index 43bda9677..923c6d907 100644 --- a/apps/api/src/services/telegram_gateway.py +++ b/apps/api/src/services/telegram_gateway.py @@ -6776,7 +6776,10 @@ class TelegramGateway: provider_event_id=provider_event_id, channel_user_id=f"sha256:{user_hash[:16]}", channel_chat_id=str(chat_id), - content_type="sre_group_text", + # ``awooop_conversation_event`` accepts the canonical + # channel-hub content types. The SRE-room subtype lives in + # the source envelope instead of bypassing that DB contract. + content_type="text", raw_content=str(text), source_envelope=source_envelope, run_id=run_id, @@ -14808,6 +14811,49 @@ class TelegramGateway: ) return + canonical_chat_identity = { + "status": "verified", + "room": "awoooi_sre_war_room", + } + if intent == "controlled_action_request": + # A controlled action turn is deterministic and candidate-only. + # It must stop before RAG/providers and may never dispatch an + # executor, Agent99, or another domain from the conversation path. + from src.services.telegram_controlled_action_work_item import ( + get_telegram_controlled_action_work_item_service, + ) + + work_item_result = await ( + get_telegram_controlled_action_work_item_service() + .create_from_verified_message( + message=clean_text, + intent=intent, + inbound_receipt=inbound_receipt, + canonical_chat_identity=canonical_chat_identity, + ) + ) + sender = ( + self.send_as_nemotron + if mention_nemo and not mention_openclaw + else self.send_as_openclaw + ) + await sender( + text=work_item_result.render_html(), + reply_to_message_id=message_id, + inbound_chat_id=chat_id, + ) + logger.info( + "telegram_controlled_action_work_item_handled", + trace_id=inbound_receipt["trace_id"], + work_item_id=work_item_result.work_item_id, + status=work_item_result.status, + created=work_item_result.created, + deduplicated=work_item_result.deduplicated, + provider_call_performed=False, + infrastructure_remediation_write_allowed=False, + ) + return + rag_receipt = await chat_mgr.get_rag_context(clean_text) if not isinstance(rag_receipt, dict): rag_receipt = { @@ -14825,10 +14871,7 @@ class TelegramGateway: route_context = { **inbound_receipt, "source_channel": "telegram_sre_war_room", - "canonical_chat_identity": { - "status": "verified", - "room": "awoooi_sre_war_room", - }, + "canonical_chat_identity": canonical_chat_identity, "inbound_event_id": inbound_receipt["event_id"], "intent": intent, "rag_retrieval_receipt": { diff --git a/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py b/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py index d68daab18..93fd48de8 100644 --- a/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py +++ b/apps/api/tests/test_sre_k3s_controlled_automation_work_items_api.py @@ -155,9 +155,9 @@ def test_loader_returns_fixed_architecture_provider_order_and_agent99_bridge() - "active_or_completed_commitments": 68, "by_status": { "analysis_or_governance_complete": 6, - "source_implemented_runtime_pending": 27, + "source_implemented_runtime_pending": 28, "in_progress": 32, - "not_started_or_no_current_evidence": 3, + "not_started_or_no_current_evidence": 2, "superseded": 2, }, "product_runtime_closed_commitments": 0, @@ -170,6 +170,12 @@ def test_loader_returns_fixed_architecture_provider_order_and_agent99_bridge() - assert commitments["AIA-CONV-029"]["status"] == ( "source_implemented_runtime_pending" ) + assert commitments["AIA-CONV-030"]["status"] == ( + "source_implemented_runtime_pending" + ) + assert "fingerprint-deduplicated" in " ".join( + commitments["AIA-CONV-030"]["source_evidence"] + ) assert "Host112" in commitments["AIA-CONV-049"]["title"] assert commitments["AIA-CONV-049"]["status"] == ( "source_implemented_runtime_pending" diff --git a/apps/api/tests/test_telegram_controlled_action_work_item.py b/apps/api/tests/test_telegram_controlled_action_work_item.py new file mode 100644 index 000000000..641b2c05d --- /dev/null +++ b/apps/api/tests/test_telegram_controlled_action_work_item.py @@ -0,0 +1,309 @@ +from __future__ import annotations + +from contextlib import asynccontextmanager +from types import SimpleNamespace +from unittest.mock import AsyncMock +from uuid import UUID + +import pytest + +from src.services import chat_manager as chat_module +from src.services import telegram_controlled_action_work_item as work_item_module +from src.services import telegram_gateway as gateway_module +from src.services.chat_manager import ChatManager +from src.services.telegram_controlled_action_work_item import ( + TelegramControlledActionWorkItemService, + persist_controlled_action_work_item, +) +from src.services.telegram_gateway import TelegramGateway + + +class _Registry: + def resolve_identity(self, name: str) -> dict: + if name == "alertmanager": + return { + "resolution_status": "resolved", + "canonical_id": "service:alertmanager", + "asset_domain": "docker_container", + "executor": "host_ansible_executor", + "verifier": "container_runtime_independent_verifier", + "controlled_apply_allowed": False, + "drift_work_item_id": None, + } + return { + "resolution_status": "asset_identity_unresolved", + "normalized_identity": name, + "canonical_id": None, + "asset_domain": "unknown", + "executor": None, + "verifier": None, + "controlled_apply_allowed": False, + "drift_work_item_id": "AIA-ASSET-DRIFT-UNKNOWN001", + } + + +class _MemoryStore: + def __init__(self) -> None: + self.records: dict[str, tuple[str, dict]] = {} + + async def __call__(self, record: dict) -> dict: + fingerprint = record["fingerprint"] + existing = self.records.get(fingerprint) + if existing is not None: + return {"receipt_id": existing[0], "created": False} + receipt_id = f"work-item-receipt-{len(self.records) + 1}" + self.records[fingerprint] = (receipt_id, record) + return {"receipt_id": receipt_id, "created": True} + + +def _inbound_receipt(*, suffix: str = "1") -> dict[str, str]: + return { + "event_id": f"event-{suffix}", + "trace_id": f"telegram-message:event-{suffix}", + "run_id": "9e39daf8-83e4-5fe8-a3d2-aebaa969a0d1", + } + + +def _chat_identity(*, verified: bool = True) -> dict[str, str]: + return { + "status": "verified" if verified else "unverified", + "room": "awoooi_sre_war_room", + } + + +@pytest.mark.asyncio +async def test_duplicate_controlled_action_reuses_one_durable_work_item() -> None: + store = _MemoryStore() + service = TelegramControlledActionWorkItemService( + registry=_Registry(), + persist_work_item=store, + ) + + first = await service.create_from_verified_message( + message="請重啟 alertmanager", + intent="controlled_action_request", + inbound_receipt=_inbound_receipt(suffix="1"), + canonical_chat_identity=_chat_identity(), + ) + duplicate = await service.create_from_verified_message( + message="請重啟 alertmanager", + intent="controlled_action_request", + inbound_receipt=_inbound_receipt(suffix="2"), + canonical_chat_identity=_chat_identity(), + ) + + assert first.accepted is True + assert first.created is True + assert duplicate.accepted is True + assert duplicate.created is False + assert duplicate.deduplicated is True + assert duplicate.work_item_id == first.work_item_id + assert duplicate.receipt_id == first.receipt_id + assert len(store.records) == 1 + record = next(iter(store.records.values()))[1] + assert record["kind"] == "codex_development_work_item" + assert record["canonical_asset_identity"]["canonical_id"] == ( + "service:alertmanager" + ) + assert record["raw_message_recorded"] is False + assert set(record["policy"].values()) == {False} + + +@pytest.mark.asyncio +async def test_unverified_chat_identity_creates_no_work_item() -> None: + store = _MemoryStore() + service = TelegramControlledActionWorkItemService( + registry=_Registry(), + persist_work_item=store, + ) + + result = await service.create_from_verified_message( + message="請重啟 alertmanager", + intent="controlled_action_request", + inbound_receipt=_inbound_receipt(), + canonical_chat_identity=_chat_identity(verified=False), + ) + + assert result.accepted is False + assert result.status == "identity_or_ingress_unverified" + assert result.work_item_id == "none" + assert store.records == {} + + +@pytest.mark.asyncio +async def test_unknown_asset_creates_only_fail_closed_drift_item() -> None: + store = _MemoryStore() + service = TelegramControlledActionWorkItemService( + registry=_Registry(), + persist_work_item=store, + ) + + result = await service.create_from_verified_message( + message="請重啟 mystery-service", + intent="controlled_action_request", + inbound_receipt=_inbound_receipt(), + canonical_chat_identity=_chat_identity(), + ) + + assert result.accepted is True + assert result.status == "asset_identity_unresolved" + assert result.candidate_created is False + assert result.work_item_id == "AIA-ASSET-DRIFT-UNKNOWN001" + record = next(iter(store.records.values()))[1] + assert record["kind"] == "asset_identity_drift_work_item" + assert record["canonical_asset_identity"]["executor"] is None + assert record["canonical_asset_identity"]["verifier"] is None + assert record["policy"]["fallback_auto_allowed"] is False + assert record["policy"]["runtime_mutation_allowed"] is False + + +@pytest.mark.asyncio +async def test_default_store_uses_idempotent_internal_command_event( + monkeypatch: pytest.MonkeyPatch, +) -> None: + event_id = UUID("bfa2fc3e-3fcb-4a2f-adf9-3c8f0ec7844c") + inserted = SimpleNamespace(fetchone=lambda: (event_id,)) + db = SimpleNamespace(execute=AsyncMock(return_value=inserted)) + + @asynccontextmanager + async def _db_context(project_id: str): + assert project_id == "awoooi" + yield db + + from src.db import base as base_module + + monkeypatch.setattr(base_module, "get_db_context", _db_context) + record = { + "kind": "codex_development_work_item", + "status": "candidate_recorded", + "work_item_id": "CODEX-SRE-1234", + "fingerprint": "a" * 43, + "run_id": "9e39daf8-83e4-5fe8-a3d2-aebaa969a0d1", + "policy": {"runtime_mutation_allowed": False}, + "raw_message_recorded": False, + } + + receipt = await persist_controlled_action_work_item(record) + + assert receipt == {"receipt_id": str(event_id), "created": True} + sql = str(db.execute.await_args.args[0]) + assert "'internal'" in sql + assert "'command'" in sql + assert "ON CONFLICT" in sql + params = db.execute.await_args.args[1] + assert len(params["content_hash"]) == 64 + assert "請重啟" not in params["source_envelope"] + + +@pytest.mark.asyncio +async def test_gateway_controlled_action_stops_before_rag_provider_or_runtime( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setattr(gateway_module.settings, "HERMES_NL_ENABLED", False) + manager = ChatManager() + monkeypatch.setattr(chat_module, "get_chat_manager", lambda: manager) + monkeypatch.setattr( + manager, + "get_rag_context", + AsyncMock(side_effect=AssertionError("controlled action must not read RAG")), + ) + monkeypatch.setattr( + manager, + "get_system_context", + AsyncMock( + side_effect=AssertionError("controlled action must not read runtime") + ), + ) + call_openclaw = AsyncMock( + side_effect=AssertionError("controlled action must not call provider") + ) + call_nemotron = AsyncMock( + side_effect=AssertionError("controlled action must not call provider") + ) + monkeypatch.setattr(manager, "_call_openclaw", call_openclaw) + monkeypatch.setattr(manager, "_call_nemotron", call_nemotron) + + store = _MemoryStore() + service = TelegramControlledActionWorkItemService( + registry=_Registry(), + persist_work_item=store, + ) + monkeypatch.setattr( + work_item_module, + "get_telegram_controlled_action_work_item_service", + lambda: service, + ) + + gateway = object.__new__(TelegramGateway) + monkeypatch.setattr( + gateway, + "mirror_sre_group_message_received", + AsyncMock(return_value=_inbound_receipt()), + ) + send_as_openclaw = AsyncMock(return_value={"ok": True}) + monkeypatch.setattr(gateway, "send_as_openclaw", send_as_openclaw) + + await gateway._handle_group_message( + "@OpenClawAwoooI_Bot 請重啟 alertmanager", + user_id=99, + username="operator", + chat_id=-100123, + message_id=459, + update_id=9003, + ) + + manager.get_rag_context.assert_not_awaited() + manager.get_system_context.assert_not_awaited() + call_openclaw.assert_not_awaited() + call_nemotron.assert_not_awaited() + send_as_openclaw.assert_awaited_once() + outbound = send_as_openclaw.await_args.kwargs["text"] + assert "work_item=CODEX-SRE-" in outbound + assert "status=candidate_recorded" in outbound + assert "provider=none" in outbound + assert "paid_provider=false" in outbound + assert "executor_invoked=false" in outbound + assert "Agent99=false" in outbound + assert "runtime_mutation=false" in outbound + assert len(store.records) == 1 + + +@pytest.mark.asyncio +async def test_sre_group_ingress_uses_channel_hub_text_contract( + monkeypatch: pytest.MonkeyPatch, +) -> None: + @asynccontextmanager + async def _db_context(project_id: str): + assert project_id == "awoooi" + yield object() + + from src.db import base as base_module + from src.services import channel_hub + + event_id = UUID("5ca599f7-9b18-47ea-a3e6-33ee71d067a1") + mirror = AsyncMock(return_value=event_id) + monkeypatch.setattr(base_module, "get_db_context", _db_context) + monkeypatch.setattr(channel_hub, "mirror_inbound_event", mirror) + monkeypatch.setattr(gateway_module.settings, "SRE_GROUP_CHAT_ID", "-100123") + + gateway = object.__new__(TelegramGateway) + receipt = await gateway.mirror_sre_group_message_received( + update_id=9004, + text="請重啟 alertmanager", + user_id=99, + username="operator", + chat_id=-100123, + message_id=460, + intent="controlled_action_request", + ) + + assert receipt is not None + assert receipt["event_id"] == str(event_id) + assert mirror.await_args.kwargs["content_type"] == "text" + envelope = mirror.await_args.kwargs["source_envelope"] + assert envelope["telegram_sre_group_message"]["intent"] == ( + "controlled_action_request" + ) + assert ( + envelope["telegram_sre_group_message"]["executor_invocation_allowed"] is False + ) diff --git a/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json b/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json index 3a77c6caa..cb8d13f92 100644 --- a/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json +++ b/docs/operations/sre-ai-agent-conversation-commitments.snapshot.json @@ -46,7 +46,7 @@ {"id":"AIA-CONV-027","category":"telegram","status":"source_implemented_runtime_pending","title":"修復重啟、擴縮容、回滾、清快取等無效按鈕","linked_work_items":["AIA-SRE-017"],"terminal_condition":"每個顯示按鈕都有 registry action、授權、callback、executor/verifier 或 fail-closed receipt。"}, {"id":"AIA-CONV-028","category":"telegram","status":"source_implemented_runtime_pending","title":"串通 callback ingress、dispatch、verifier 與 durable action receipt","linked_work_items":["AIA-SRE-009","AIA-SRE-015","AIA-SRE-017"],"terminal_condition":"recipient callback 與 same-run controlled action receipt 可由 production readback 證明。"}, {"id":"AIA-CONV-029","category":"telegram","status":"source_implemented_runtime_pending","title":"AwoooI SRE 戰情室 Bot 自動理解並正確回覆使用者訊息","linked_work_items":["AIA-SRE-012","AIA-SRE-013","AIA-SRE-017"],"terminal_condition":"群組訊息經身份、意圖、RAG 與 policy 後產生可追溯回覆,不接受未授權 runtime mutation。","source_evidence":["telegram group text durable ingress receipt","deterministic intent classification","bounded sanitized KM/RAG retrieval","traceable no-runtime-mutation provider receipt"]}, - {"id":"AIA-CONV-030","category":"telegram","status":"not_started_or_no_current_evidence","title":"對話可啟動受控調查/修復並建立 Codex 開發 work item","linked_work_items":["AIA-SRE-001","AIA-SRE-015","AIA-SRE-017"],"terminal_condition":"討論結果轉為去重 work item;production action 仍走 typed executor/verifier。"}, + {"id":"AIA-CONV-030","category":"telegram","status":"source_implemented_runtime_pending","title":"對話可啟動受控調查/修復並建立 Codex 開發 work item","linked_work_items":["AIA-SRE-001","AIA-SRE-015","AIA-SRE-017"],"terminal_condition":"討論結果轉為去重 work item;production action 仍走 typed executor/verifier。","source_evidence":["verified controlled_action_request ingress gate","exact canonical asset/domain resolution or asset_identity_unresolved drift item","fingerprint-deduplicated durable internal work-item receipt","deterministic no-provider no-runtime-mutation Telegram response"]}, {"id":"AIA-CONV-031","category":"telegram","status":"in_progress","title":"同 fingerprint 去重、聚合、抑噪及 resolved/recovered 更新","linked_work_items":["AIA-SRE-017"],"terminal_condition":"同事件只維護 canonical lifecycle,恢復卡有 destination-bound provider acknowledgement。"}, {"id":"AIA-CONV-032","category":"named_incident","status":"in_progress","title":"INC-20260711-11C751 cold-start-gate 補 PlayBook、transport、rollback 與 verifier","linked_work_items":["AIA-SRE-004","AIA-SRE-006","AIA-SRE-015"],"terminal_condition":"namespace identity 修正後完成一次 bounded same-run repair 或明確 no-write terminal。"}, @@ -98,9 +98,9 @@ "active_or_completed_commitments": 68, "by_status": { "analysis_or_governance_complete": 6, - "source_implemented_runtime_pending": 27, + "source_implemented_runtime_pending": 28, "in_progress": 32, - "not_started_or_no_current_evidence": 3, + "not_started_or_no_current_evidence": 2, "superseded": 2 }, "product_runtime_closed_commitments": 0