Files
ewoooc/services/nemotron_decision_canary_service.py
ogt 95682ee3bc
Some checks failed
CD Pipeline / deploy (push) Has been cancelled
feat(ai): verify RAG and Nemotron runtime canaries
2026-07-17 11:11:08 +08:00

460 lines
17 KiB
Python

"""Decision-only NemoTron runtime canary.
The canary exercises the production qwen3 tool-calling payload against an
approved GCP Ollama host, validates the returned decision, and stops before any
tool, Telegram, database, price, or insight write is allowed.
"""
from __future__ import annotations
import json
import os
import time
import uuid
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Mapping
import requests
from services.ollama_service import OLLAMA_HOST_PRIMARY, OLLAMA_HOST_SECONDARY
POLICY = "controlled_nemotron_decision_only_canary_v1"
CANARY_VERSION = "nemotron_decision_only_canary_v1"
NEMOTRON_OLLAMA_MODEL = os.getenv("NEMOTRON_OLLAMA_MODEL", "qwen3:14b")
NEMOTRON_OLLAMA_EXPECTED_DIGEST = os.getenv(
"NEMOTRON_OLLAMA_EXPECTED_DIGEST",
"bdbd181c33f2ed1b31c972991882db3cf4d192569092138a7d29e973cd9debe8",
).strip().lower()
DEFAULT_TIMEOUT_SEC = max(
30,
min(int(os.getenv("NEMOTRON_DECISION_CANARY_TIMEOUT_SEC", "300")), 600),
)
DEFAULT_MODEL_IDENTITY_TIMEOUT_SEC = max(
2,
min(int(os.getenv("NEMOTRON_MODEL_IDENTITY_TIMEOUT_SEC", "10")), 30),
)
DEFAULT_MAX_AGE_HOURS = max(
1,
min(int(os.getenv("NEMOTRON_DECISION_CANARY_MAX_AGE_HOURS", "26")), 168),
)
DEFAULT_OUTPUT_ROOT = os.getenv(
"NEMOTRON_DECISION_CANARY_RECEIPT_ROOT",
"/app/data/ai_automation/nemotron_decision_canary_receipts"
if Path("/app/data").exists()
else "runtime_artifacts/nemotron_decision_canary_receipts",
)
APPROVED_TOOLS = {
"trigger_price_alert",
"add_to_recommendation",
"flag_for_human_review",
}
@dataclass(frozen=True)
class _SyntheticThreat:
sku: str = "CANARY-NEMO-001"
name: str = "NemoTron 受控決策驗證品"
momo_price: float = 1200.0
pchome_price: float = 980.0
gap_pct: float = 22.4
sales_7d_delta_pct: float = -35.0
risk: str = "HIGH"
recommended_action: str = "依競價證據建立受控價格告警候選"
confidence: float = 0.91
match_type: str = "exact"
price_basis: str = "total_price"
alert_tier: str = "price_alert_exact"
match_score: float = 0.99
competitor_product_id: str = "CANARY-COMP-001"
competitor_product_name: str = "NemoTron 受控決策驗證品"
def _parse_datetime(value: Any) -> datetime | None:
if not value:
return None
try:
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
def _model_identity(host: str) -> dict[str, Any]:
started = time.monotonic()
try:
response = requests.get(
f"{host.rstrip('/')}/api/tags",
timeout=DEFAULT_MODEL_IDENTITY_TIMEOUT_SEC,
)
response.raise_for_status()
models = response.json().get("models") or []
except Exception as exc:
return {
"ok": False,
"host": host,
"model": NEMOTRON_OLLAMA_MODEL,
"digest": None,
"digest_matches": False,
"elapsed_ms": round((time.monotonic() - started) * 1000),
"error": f"{type(exc).__name__}: {str(exc)[:180]}",
}
target = NEMOTRON_OLLAMA_MODEL
target_without_latest = target.removesuffix(":latest")
matched: Mapping[str, Any] = {}
for item in models:
names = {
str(item.get("name") or "").strip(),
str(item.get("model") or "").strip(),
}
if target in names or target_without_latest in names:
matched = item
break
digest = str(matched.get("digest") or "").strip().lower()
digest_matches = bool(digest) and digest == NEMOTRON_OLLAMA_EXPECTED_DIGEST
return {
"ok": bool(matched) and digest_matches,
"host": host,
"model": target,
"digest": digest or None,
"expected_digest": NEMOTRON_OLLAMA_EXPECTED_DIGEST,
"digest_matches": digest_matches,
"parameter_size": str((matched.get("details") or {}).get("parameter_size") or ""),
"quantization_level": str(
(matched.get("details") or {}).get("quantization_level") or ""
),
"elapsed_ms": round((time.monotonic() - started) * 1000),
"error": None if matched else f"model_not_found:{target}",
}
def _latest_execution(root: Path, *, now: datetime | None = None) -> dict[str, Any]:
if not root.exists():
return {}
candidates = sorted(
root.glob("*/nemotron_decision_canary_receipt.json"),
key=lambda path: path.stat().st_mtime,
reverse=True,
)
if not candidates:
return {}
path = candidates[0]
try:
payload = json.loads(path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError):
return {}
generated_at = _parse_datetime(payload.get("generated_at"))
current = now or datetime.now(timezone.utc)
age_hours = (
(current - generated_at).total_seconds() / 3600
if generated_at is not None
else None
)
return {
"receipt_path": str(path),
"generated_at": payload.get("generated_at"),
"age_hours": round(age_hours, 3) if age_hours is not None else None,
"fresh": age_hours is not None and age_hours <= DEFAULT_MAX_AGE_HOURS,
"status": payload.get("status"),
"canary_passed": payload.get("canary_passed") is True,
"selected_host": payload.get("selected_host"),
"model": payload.get("model"),
"observed_digest": payload.get("observed_digest"),
"decision_tool": payload.get("decision_tool"),
"decision_sku": payload.get("decision_sku"),
"model_elapsed_ms": payload.get("model_elapsed_ms"),
"tool_execution_count": int(payload.get("tool_execution_count") or 0),
"database_call_performed": payload.get("database_call_performed") is True,
"writes_database": payload.get("writes_database") is True,
"writes_price_tables": payload.get("writes_price_tables") is True,
"writes_ai_insights": payload.get("writes_ai_insights") is True,
"telegram_sent": payload.get("telegram_sent") is True,
"rollback_terminal": payload.get("rollback_terminal"),
"terminal_status": payload.get("terminal_status"),
"acknowledgements": payload.get("acknowledgements") or {},
"run_identity": payload.get("run_identity") or {},
"error": payload.get("error"),
}
def _write_receipt(root: Path, run_id: str, payload: Mapping[str, Any]) -> str:
target = root / run_id / "nemotron_decision_canary_receipt.json"
target.parent.mkdir(parents=True, exist_ok=True)
target.write_text(
json.dumps(dict(payload), ensure_ascii=False, indent=2, sort_keys=True),
encoding="utf-8",
)
return str(target)
def acknowledge_nemotron_decision_canary(
receipt_path: str | Path,
*,
telegram_status: str,
telegram_sent: int = 0,
telegram_failed: int = 0,
) -> dict[str, Any]:
"""Atomically append the scheduler Telegram acknowledgement to a receipt."""
target = Path(receipt_path)
payload = json.loads(target.read_text(encoding="utf-8"))
acknowledgement = {
"status": str(telegram_status),
"sent": int(telegram_sent),
"failed": int(telegram_failed),
}
payload["acknowledgements"] = {
"telegram": acknowledgement,
"km": "not_applicable_decision_only_canary",
"rag": "not_applicable_decision_only_canary",
"mcp": "not_applicable_decision_only_canary",
"playbook": "nemotron_decision_canary_receipt_written",
}
closure = payload.setdefault("closure_receipt", {})
closure["telegram_acknowledgement"] = str(telegram_status)
payload["terminal_status"] = (
"verified_decision_only_no_write"
if payload.get("canary_passed") is True and telegram_status == "acknowledged"
else "partial"
)
temporary = target.with_suffix(target.suffix + ".tmp")
temporary.write_text(
json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True),
encoding="utf-8",
)
temporary.replace(target)
return payload
def run_nemotron_decision_canary(
*,
output_root: str | Path | None = None,
execute: bool = False,
write_receipt: bool = False,
timeout_sec: int | None = None,
trace_id: str | None = None,
run_id: str | None = None,
work_item_id: str = "AI-RUNTIME-NEMOTRON-CANARY-001",
) -> dict[str, Any]:
"""Read the latest receipt or execute one bounded decision-only model call."""
now = datetime.now(timezone.utc)
receipt_root = Path(output_root or DEFAULT_OUTPUT_ROOT)
resolved_run_id = str(run_id or f"nemotron-canary-{uuid.uuid4()}")
run_identity = {
"trace_id": str(trace_id or f"trace-{resolved_run_id}"),
"run_id": resolved_run_id,
"work_item_id": str(work_item_id or "AI-RUNTIME-NEMOTRON-CANARY-001"),
}
latest_before = _latest_execution(receipt_root, now=now)
if not execute:
verified = (
latest_before.get("canary_passed") is True
and latest_before.get("fresh") is True
)
return {
"success": True,
"policy": POLICY,
"canary_version": CANARY_VERSION,
"generated_at": now.isoformat(),
"status": "verified" if verified else "warning",
"execute": False,
"model": NEMOTRON_OLLAMA_MODEL,
"expected_digest": NEMOTRON_OLLAMA_EXPECTED_DIGEST,
"latest_execution": latest_before,
"controlled_apply": {
"network_call": False,
"model_call": False,
"tool_execution_count": 0,
"database_call_performed": False,
"database_write": False,
"price_table_write": False,
"ai_insights_write": False,
"telegram_send": False,
},
"next_machine_action": (
"continue_scheduled_nemotron_decision_canary"
if verified
else "execute_nemotron_decision_only_canary"
),
}
host_preflights = [
_model_identity(OLLAMA_HOST_PRIMARY),
_model_identity(OLLAMA_HOST_SECONDARY),
]
selected = next((item for item in host_preflights if item.get("ok")), None)
model_call_performed = False
model_elapsed_ms = 0
decisions: list[dict[str, Any]] = []
decision_format = "none"
error: str | None = None
if selected is None:
error = "no_approved_gcp_host_with_expected_nemotron_digest"
else:
from services.nemoton_dispatcher_service import (
_parse_content_fallback,
_parse_tool_calls_struct,
build_qwen3_dispatch_payload,
)
threat = _SyntheticThreat()
payload = build_qwen3_dispatch_payload(
[threat],
mcp_context="controlled decision-only canary; no live market data",
num_predict=160,
keep_alive=0,
)
payload["think"] = False
started = time.monotonic()
model_call_performed = True
try:
response = requests.post(
f"{str(selected['host']).rstrip('/')}/api/chat",
json=payload,
timeout=max(30, min(int(timeout_sec or DEFAULT_TIMEOUT_SEC), 600)),
)
response.raise_for_status()
body = response.json()
message = body.get("message") or {}
decisions = _parse_tool_calls_struct(message.get("tool_calls") or [])
if decisions:
decision_format = "tool_calls"
else:
decisions = _parse_content_fallback(str(message.get("content") or ""))
if decisions:
decision_format = "content_fallback"
except Exception as exc:
error = f"{type(exc).__name__}: {str(exc)[:240]}"
finally:
model_elapsed_ms = round((time.monotonic() - started) * 1000)
decision = decisions[0] if decisions else {}
decision_tool = str(decision.get("tool") or "")
decision_args = decision.get("args") if isinstance(decision.get("args"), dict) else {}
decision_sku = str((decision_args or {}).get("sku") or "")
expected_tool = "trigger_price_alert"
post_identity = _model_identity(str(selected["host"])) if selected else {}
checks = {
"approved_host_selected": selected is not None,
"expected_digest_verified_before": bool(selected and selected.get("digest_matches")),
"model_call_completed": model_call_performed and error is None,
"decision_present": bool(decisions),
"decision_tool_allowed": decision_tool in APPROVED_TOOLS,
"expected_decision_tool_selected": decision_tool == expected_tool,
"synthetic_sku_preserved": decision_sku == _SyntheticThreat.sku,
"expected_digest_stable_after": bool(post_identity.get("digest_matches")),
"tool_execution_absent": True,
"database_call_absent": True,
"telegram_send_absent": True,
}
canary_passed = all(checks.values())
status = "canary_passed" if canary_passed else "canary_failed"
execution_receipt: dict[str, Any] = {
"policy": POLICY,
"canary_version": CANARY_VERSION,
"generated_at": now.isoformat(),
"run_identity": run_identity,
"status": status,
"canary_passed": canary_passed,
"model": NEMOTRON_OLLAMA_MODEL,
"expected_digest": NEMOTRON_OLLAMA_EXPECTED_DIGEST,
"observed_digest": selected.get("digest") if selected else None,
"selected_host": selected.get("host") if selected else None,
"host_preflights": host_preflights,
"post_model_identity": post_identity,
"model_call_performed": model_call_performed,
"model_elapsed_ms": model_elapsed_ms,
"decision_format": decision_format,
"decision_count": len(decisions),
"decision_tool": decision_tool or None,
"decision_sku": decision_sku or None,
"expected_decision_tool": expected_tool,
"decision_argument_keys": sorted((decision_args or {}).keys()),
"checks": checks,
"check_count": len(checks),
"check_pass_count": sum(checks.values()),
"decision_only": True,
"tool_execution_count": 0,
"database_call_performed": False,
"writes_database": False,
"writes_price_tables": False,
"writes_ai_insights": False,
"telegram_sent": False,
"error": error,
"rollback_terminal": "decision_only_no_tool_or_data_write",
"source_of_truth_diff": {
"expected_model": NEMOTRON_OLLAMA_MODEL,
"expected_digest": NEMOTRON_OLLAMA_EXPECTED_DIGEST,
"observed_digest": selected.get("digest") if selected else None,
"expected_decision_tool": expected_tool,
"observed_decision_tool": decision_tool or None,
},
"closure_receipt": {
"sensor_source_receipt": bool(host_preflights),
"normalized_asset_identity": selected.get("host") if selected else None,
"source_of_truth_diff_recorded": True,
"ai_candidate_decision_recorded": bool(decisions),
"risk_policy_decision": "medium_decision_only_no_tool_execution",
"check_mode_passed": bool(selected),
"bounded_execution_performed": model_call_performed,
"independent_verifier_passed": canary_passed,
"rollback_or_no_write_terminal": "decision_only_no_tool_or_data_write",
"telegram_acknowledgement": "pending_scheduler_dispatch",
"learning_write_acknowledgement": "pending_receipt_write",
},
}
if write_receipt:
receipt_path = str(
receipt_root
/ resolved_run_id
/ "nemotron_decision_canary_receipt.json"
)
execution_receipt["receipt_path"] = receipt_path
execution_receipt["closure_receipt"][
"learning_write_acknowledgement"
] = "nemotron_decision_canary_receipt_written"
_write_receipt(receipt_root, resolved_run_id, execution_receipt)
latest_after = _latest_execution(receipt_root)
return {
"success": canary_passed,
"policy": POLICY,
"canary_version": CANARY_VERSION,
"generated_at": now.isoformat(),
"run_identity": run_identity,
"status": status,
"execute": True,
"execution": execution_receipt,
"latest_execution": latest_after,
"controlled_apply": {
"risk": "medium",
"network_call": True,
"model_call": model_call_performed,
"tool_execution_count": 0,
"database_call_performed": False,
"database_write": False,
"price_table_write": False,
"ai_insights_write": False,
"telegram_send": False,
"rollback_terminal": "decision_only_no_tool_or_data_write",
},
"next_machine_action": (
"continue_scheduled_nemotron_decision_canary"
if canary_passed
else "repair_nemotron_model_runtime_then_retry_decision_only_canary"
),
}
__all__ = [
"CANARY_VERSION",
"POLICY",
"acknowledge_nemotron_decision_canary",
"run_nemotron_decision_canary",
]