#!/usr/bin/env python3 from __future__ import annotations import argparse import json import os import socket import mimetypes import sqlite3 import sys import threading import time import uuid from datetime import datetime, timezone from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from typing import Any from urllib.parse import parse_qs, urlparse from urllib.request import Request, urlopen import urllib.error from urllib.error import HTTPError, URLError ROOT = Path(__file__).resolve().parents[3] SCRIPTS_DIR = ROOT / "scripts" if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) if str(SCRIPTS_DIR) not in sys.path: sys.path.insert(0, str(SCRIPTS_DIR)) from core.observability import JsonlAuditStore, resolve_audit_root, resolve_request_id, resolve_trace_id, sanitize_for_logging # noqa: E402 from ask_1c_rag import DEFAULT_INDEX as DEFAULT_1C_RAG_INDEX # noqa: E402 from ask_1c_rag import DEFAULT_RAG_PROMPT as DEFAULT_1C_RAG_PROMPT # noqa: E402 from ask_1c_rag import format_context, render_prompt # noqa: E402 from common import read_json, search_lexical_index # noqa: E402 from rag_profiles import resolve_rag_profile # noqa: E402 DEFAULT_HOST = "127.0.0.1" DEFAULT_PORT = 8090 DEFAULT_PROVIDER = "default" DEFAULT_PLUGIN = "1c" WEB_UI_DIR = ROOT / "plugins" / "1c" / "agent" / "web" REGISTRY_INDEX = ROOT / "registry" / "index.json" RUNTIME_PROFILES_PATH = ROOT / "config" / "runtime_profiles.json" DEFAULT_SYSTEM_PROMPT_PATH = ROOT / "plugins" / "1c" / "prompts" / "system.md" SERVICE_STARTED_AT = datetime.now(timezone.utc) TASK_PLUGIN_MAP = { "text": "text", "chat": "text", "summarization": "text", "code": "text", "tool-use": "text", "translation": "translation", "speech-to-text": "audio", "speech-translation": "audio", "video": "video", "image-understanding": "video", "document-understanding": "video", "visual-question-answering": "video", "image-generation": "image", "image-editing": "image", "inpainting": "image", "1c": "1c", "1c-rag": "1c", "bsl-code": "1c", "metadata-safety": "1c", "1c-query": "1c", } STATUS_RANK = { "production": 0, "staging": 1, "candidate": 2, "draft": 3, "archived": 9, } def now_iso() -> str: return datetime.now(timezone.utc).isoformat() def safe_request_json(handler: BaseHTTPRequestHandler) -> dict: length = int(handler.headers.get("Content-Length") or 0) raw = handler.rfile.read(length).decode("utf-8") if length else "{}" if not raw.strip(): return {} payload = json.loads(raw) if not isinstance(payload, dict): raise ValueError("request body must be JSON object") return payload def json_response(handler: BaseHTTPRequestHandler, status: int, payload: dict[str, Any]) -> None: body = json.dumps(payload, ensure_ascii=False).encode("utf-8") handler.send_response(status) handler.send_header("Content-Type", "application/json; charset=utf-8") handler.send_header("Content-Length", str(len(body))) handler.send_header("Access-Control-Allow-Origin", "*") handler.send_header("Access-Control-Allow-Methods", "GET, POST, PATCH, DELETE, OPTIONS") handler.send_header("Access-Control-Allow-Headers", "Content-Type, Authorization, X-Request-Id") handler.end_headers() handler.wfile.write(body) def next_trace_id() -> str: return uuid.uuid4().hex def normalize_limit(value: Any, *, default: int = 200, max_value: int = 500) -> int: if value is None: return default try: limit = int(value) except (TypeError, ValueError): raise ValueError("limit must be integer") if limit < 1: raise ValueError("limit must be >= 1") return min(limit, max_value) def normalize_float(value: Any, *, field: str, default: float | None = None) -> float: if value is None: if default is None: raise ValueError(f"{field} is required") return float(default) try: return float(value) except (TypeError, ValueError): raise ValueError(f"{field} must be number") def normalize_int(value: Any, *, field: str, default: int | None = None) -> int: if value is None: if default is None: raise ValueError(f"{field} is required") return int(default) try: return int(value) except (TypeError, ValueError): raise ValueError(f"{field} must be integer") def request_trace_id(handler: BaseHTTPRequestHandler) -> str: return resolve_trace_id(handler.headers) def request_request_id(handler: BaseHTTPRequestHandler, trace_id: str) -> str: return resolve_request_id(handler.headers, fallback_trace_id=trace_id) def usage_from_provider_response(raw: dict[str, Any] | None) -> dict[str, int]: usage = (raw or {}).get("usage") if isinstance(raw, dict) else None if not isinstance(usage, dict): return {} normalized: dict[str, int] = {} for key in ("prompt_tokens", "completion_tokens", "total_tokens"): value = usage.get(key) try: if value is not None: normalized[key] = int(value) except (TypeError, ValueError): continue return normalized def json_error(handler: BaseHTTPRequestHandler, status: int, message: str, *, code: str | None = None, trace_id: str) -> None: payload: dict[str, Any] = { "status": "error", "error": {"message": message}, "trace_id": trace_id, } if code: payload["error"]["code"] = code json_response(handler, status, payload) def normalize_base_url(value: str) -> str: parsed = urlparse((value or "").strip()) if parsed.scheme not in {"http", "https"} or not parsed.hostname: raise ValueError("base_url must be http/https URL") return f"{parsed.scheme}://{parsed.hostname}:{parsed.port}" if parsed.port else f"{parsed.scheme}://{parsed.hostname}" def read_path_text(path: Path, fallback: str) -> str: if not path.exists(): return fallback return path.read_text(encoding="utf-8") def _resolve_runtime_snapshot( *, project: dict[str, Any] | None, chat: dict[str, Any] | None, providers: dict[str, ProviderConfig], ) -> dict[str, Any]: if not chat: raise ValueError("chat not found") provider_id = str(chat.get("provider_id") or DEFAULT_PROVIDER) if provider_id not in providers: raise ValueError(f"unknown provider_id={provider_id}") provider = providers[provider_id] selected = route_for_model(chat.get("model_id") or chat.get("served_model_name")) served_model = str( chat.get("served_model_name") or selected.get("served_model_name") or selected.get("model", {}).get("served_model_name") or provider.model ) model_base_url = str( chat.get("base_url") or selected.get("base_url") or provider.base_url ) return { "project": { "id": chat.get("project_id"), "name": str((project or {}).get("name") or ""), }, "chat": { "id": chat.get("id"), "title": str(chat.get("title") or ""), }, "provider": { "id": provider_id, "type": provider.type, "base_url": provider.base_url, "model": provider.model, }, "model": { "chat_model_id": chat.get("model_id"), "selected_model_id": str(selected.get("model", {}).get("id") or ""), "served_model": served_model, "base_url": model_base_url, "route": selected, }, "runtime_settings": { "temperature": chat.get("temperature"), "max_tokens": chat.get("max_tokens"), "rag_profile": chat.get("rag_profile"), "rag_limit": chat.get("rag_limit"), "system_prompt": chat.get("system_prompt"), }, "adapter": { "url": os.environ.get("ONEC_ADAPTER_URL"), }, } def _build_turn_payload( *, route: dict[str, Any], provider_id: str, provider: ProviderConfig, served_model: str, model_base_url: str, temperature: float, max_tokens: int, history_limit: int, model_messages: list[dict[str, Any]], ) -> dict[str, Any]: return { "provider_id": provider_id, "provider": { "type": provider.type, "base_url": provider.base_url, "model": provider.model, }, "model": { "id": str(route.get("model", {}).get("id") or ""), "served_model": served_model, "base_url": model_base_url, }, "routing": route, "runtime": { "temperature": temperature, "max_tokens": max_tokens, "history_limit": history_limit, }, "messages": model_messages, } def _task_list(model: dict[str, Any]) -> list[str]: tasks = model.get("task") or model.get("tasks") or [] if isinstance(tasks, str): return [tasks] if isinstance(tasks, tuple): tasks = list(tasks) if isinstance(tasks, list): return [str(task) for task in tasks if task is not None] return [] def model_plugins(model: dict[str, Any]) -> list[str]: plugins: list[str] = [] for task in _task_list(model): plugin = TASK_PLUGIN_MAP.get(str(task)) if plugin and plugin not in plugins: plugins.append(plugin) return plugins or ["text"] def plugin_model_score(model: dict[str, Any], plugin: str) -> tuple[int, int, str]: tasks = set(model.get("task") or model.get("tasks") or []) status = str(model.get("status") or "") model_id = str(model.get("id") or "") runtime = str(model.get("runtime") or "") model_type = str(model.get("type") or "") quantization = str(model.get("quantization") or "") score = 0 if plugin == "text": if tasks & {"text", "chat"}: score -= 5 if tasks & {"1c", "1c-rag", "bsl-code"}: score += 6 elif plugin == "1c": if tasks & {"1c", "1c-rag", "bsl-code", "1c-query"}: score -= 6 if "qwen3-coder" in model_id: score -= 8 if quantization == "Q6_K" or "q6" in model_id: score -= 3 if runtime == "llama.cpp": score -= 2 if model_type == "lora-adapter" or model_id.endswith("-lora-v1"): score += 8 return (score, STATUS_RANK.get(status, 8), model_id) def model_catalog() -> dict[str, Any]: if not REGISTRY_INDEX.exists(): return {"models": []} data = read_json(REGISTRY_INDEX) return data if isinstance(data, dict) else {"models": []} def load_runtime_profiles() -> dict[str, Any]: if not RUNTIME_PROFILES_PATH.exists(): return {"default": "", "profiles": {}} data = read_json(RUNTIME_PROFILES_PATH) profiles = data.get("profiles") or {} if not isinstance(profiles, dict): profiles = {} normalized: dict[str, dict[str, Any]] = {} for profile_id, profile in profiles.items(): if not isinstance(profile, dict): continue item = dict(profile) item["id"] = str(item.get("id") or profile_id) normalized[str(profile_id)] = item return {"default": str(data.get("default") or next(iter(normalized), "")), "profiles": normalized} def runtime_profile_target(profile: dict[str, Any], *, model: dict[str, Any], model_id: str, plugin: str) -> dict[str, Any]: override = (profile.get("model_overrides") or {}).get(model_id) or {} base_url = override.get("base_url") or (profile.get("endpoints") or {}).get(plugin) if not base_url: raise ValueError(f"runtime profile `{profile.get('id')}` has no endpoint for plugin `{plugin}`") served_model_name = override.get("served_model_name") or model.get("served_model_name") or model_id return { "base_url": normalize_base_url(str(base_url)), "served_model_name": str(served_model_name), "container_name": override.get("container_name"), "host": profile.get("host"), "role": profile.get("role"), } def route_for_model(model_id: str | None = None, *, plugin: str = DEFAULT_PLUGIN) -> dict[str, Any]: catalog = model_catalog() all_models = catalog.get("models") or [] plugin_candidates = [model for model in all_models if plugin in model_plugins(model)] selected = None if model_id: wanted = str(model_id) selected = next((m for m in all_models if str(m.get("id")) == wanted), None) if not selected: raise ValueError(f"model `{wanted}` not found") if plugin not in model_plugins(selected): raise ValueError(f"model `{wanted}` is not configured for plugin={plugin}") if selected is None: if not plugin_candidates: raise ValueError(f"no model route for plugin={plugin}") plugin_candidates.sort(key=lambda m: plugin_model_score(m, plugin)) selected = plugin_candidates[0] if not selected: raise ValueError(f"no model route for plugin={plugin}") route = { "base_url": "", "served_model_name": str(selected.get("served_model_name") or selected.get("id") or ""), "container_name": None, "host": None, "role": None, } runtime_data = load_runtime_profiles() profile = runtime_data["profiles"].get(runtime_data.get("default", "")) if profile: try: route = { **route, **runtime_profile_target(profile, model=selected, model_id=str(selected.get("id") or ""), plugin=plugin), } except Exception: pass return { "plugin": plugin, "model": selected, "base_url": route["base_url"], "served_model_name": route["served_model_name"], "container_name": route.get("container_name"), "host": route.get("host"), "role": route.get("role"), } class ProviderConfig: def __init__(self, name: str, data: dict[str, Any]) -> None: self.name = name self.type = str(data.get("type") or "openai-compatible") self.base_url = normalize_base_url(str(data.get("base_url") or "")) self.model = str(data.get("model") or data.get("served_model_name") or "") self.api_key_env = str(data.get("api_key_env") or "") self.timeout = int(data.get("timeout") or 30) def load_provider_configs() -> dict[str, ProviderConfig]: raw = os.environ.get("ONEC_AGENT_PROVIDERS", "").strip() parsed: dict[str, dict[str, Any]] = {} if raw: payload = json.loads(raw) if not isinstance(payload, dict): raise ValueError("ONEC_AGENT_PROVIDERS must be JSON object") parsed = {key: dict(value or {}) for key, value in payload.items() if isinstance(value, dict)} fallback = { "default": { "type": "openai-compatible", "base_url": os.environ.get("ONEC_AGENT_DEFAULT_BASE_URL", "http://docker-gpu.cin.su:8000"), "model": os.environ.get("ONEC_AGENT_DEFAULT_MODEL", "qwen3-4b-instruct-2507"), } } fallback.update(parsed) return {name: ProviderConfig(name=name, data=cfg) for name, cfg in fallback.items()} class Store: def __init__(self, db_path: Path) -> None: self.path = db_path self.path.parent.mkdir(parents=True, exist_ok=True) self._lock = threading.Lock() self._conn = sqlite3.connect(self.path, check_same_thread=False) self._conn.row_factory = sqlite3.Row self._conn.execute("PRAGMA foreign_keys=ON") self._conn.execute("PRAGMA journal_mode=WAL") self._initialize() def _initialize(self) -> None: with self._lock: self._conn.executescript( """ CREATE TABLE IF NOT EXISTS projects ( id TEXT PRIMARY KEY, name TEXT NOT NULL, description TEXT NOT NULL, metadata_json TEXT NOT NULL DEFAULT "{}", created_at TEXT NOT NULL, updated_at TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS chats ( id TEXT PRIMARY KEY, project_id TEXT NOT NULL, title TEXT NOT NULL, model_id TEXT, provider_id TEXT NOT NULL, base_url TEXT, served_model_name TEXT, temperature REAL NOT NULL DEFAULT 0.2, max_tokens INTEGER NOT NULL DEFAULT 1200, rag_profile TEXT NOT NULL DEFAULT "auto", rag_limit INTEGER, system_prompt TEXT NOT NULL, metadata_json TEXT NOT NULL DEFAULT "{}", created_at TEXT NOT NULL, updated_at TEXT NOT NULL, FOREIGN KEY(project_id) REFERENCES projects(id) ON DELETE CASCADE ); CREATE TABLE IF NOT EXISTS messages ( id TEXT PRIMARY KEY, project_id TEXT NOT NULL, chat_id TEXT NOT NULL, role TEXT NOT NULL, content TEXT NOT NULL, payload_json TEXT NOT NULL DEFAULT "{}", created_at TEXT NOT NULL, FOREIGN KEY(chat_id) REFERENCES chats(id) ON DELETE CASCADE ); CREATE INDEX IF NOT EXISTS idx_messages_chat ON messages(chat_id, created_at); """ ) self._conn.commit() def create_project(self, *, name: str, description: str, metadata: dict[str, Any] | None = None) -> dict[str, Any]: project_id = uuid.uuid4().hex now = now_iso() with self._lock: self._conn.execute( "INSERT INTO projects (id, name, description, metadata_json, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?)", (project_id, name, description, json.dumps(metadata or {}, ensure_ascii=False), now, now), ) self._conn.commit() return { "id": project_id, "name": name, "description": description, "metadata": metadata or {}, "created_at": now, "updated_at": now, } def update_project(self, project_id: str, *, name: str | None = None, description: str | None = None, metadata: dict[str, Any] | None = None) -> dict[str, Any] | None: current = self.get_project(project_id) if not current: return None fields: list[str] = [] values: list[Any] = [] if name is not None: fields.append("name = ?") values.append(name) if description is not None: fields.append("description = ?") values.append(description) if metadata is not None: fields.append("metadata_json = ?") values.append(json.dumps(metadata, ensure_ascii=False)) if not fields: return current fields.append("updated_at = ?") values.extend([now_iso(), project_id]) with self._lock: self._conn.execute(f"UPDATE projects SET {', '.join(fields)} WHERE id = ?", tuple(values)) self._conn.commit() return self.get_project(project_id) def delete_project(self, project_id: str) -> bool: with self._lock: cursor = self._conn.execute("DELETE FROM projects WHERE id = ?", (project_id,)) self._conn.commit() return cursor.rowcount > 0 def list_projects(self) -> list[dict[str, Any]]: with self._lock: rows = self._conn.execute("SELECT * FROM projects ORDER BY updated_at DESC").fetchall() return [ { "id": row["id"], "name": row["name"], "description": row["description"], "metadata": json.loads(row["metadata_json"] or "{}"), "created_at": row["created_at"], "updated_at": row["updated_at"], } for row in rows ] def get_project(self, project_id: str) -> dict[str, Any] | None: with self._lock: row = self._conn.execute("SELECT * FROM projects WHERE id = ?", (project_id,)).fetchone() if not row: return None return { "id": row["id"], "name": row["name"], "description": row["description"], "metadata": json.loads(row["metadata_json"] or "{}"), "created_at": row["created_at"], "updated_at": row["updated_at"], } def create_chat( self, *, project_id: str, title: str, model_id: str | None, provider_id: str, base_url: str | None, served_model_name: str | None, temperature: float, max_tokens: int, rag_profile: str, rag_limit: int | None, system_prompt: str, metadata: dict[str, Any] | None = None, ) -> dict[str, Any]: if not self.get_project(project_id): raise FileNotFoundError("project not found") chat_id = uuid.uuid4().hex now = now_iso() with self._lock: self._conn.execute( "INSERT INTO chats (id, project_id, title, model_id, provider_id, base_url, served_model_name, temperature, max_tokens, rag_profile, rag_limit, system_prompt, metadata_json, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", (chat_id, project_id, title, model_id, provider_id, base_url, served_model_name, temperature, max_tokens, rag_profile, rag_limit, system_prompt, json.dumps(metadata or {}, ensure_ascii=False), now, now), ) self._conn.commit() return self.get_chat(project_id, chat_id) # type: ignore[misc] def update_chat_config( self, project_id: str, chat_id: str, *, title: str | None = None, model_id: str | None = None, provider_id: str | None = None, base_url: str | None = None, served_model_name: str | None = None, temperature: float | None = None, max_tokens: int | None = None, rag_profile: str | None = None, rag_limit: int | None = None, system_prompt: str | None = None, metadata: dict[str, Any] | None = None, ) -> dict[str, Any] | None: current = self.get_chat(project_id, chat_id) if not current: return None payload = { "title": title, "model_id": model_id, "provider_id": provider_id, "base_url": base_url, "served_model_name": served_model_name, "temperature": temperature, "max_tokens": max_tokens, "rag_profile": rag_profile, "rag_limit": rag_limit, "system_prompt": system_prompt, "metadata": metadata, } update_payload: dict[str, Any] = {name: value for name, value in payload.items() if value is not None} if not update_payload: return current self.update_chat(project_id, chat_id, **update_payload) return self.get_chat(project_id, chat_id) def delete_chat(self, project_id: str, chat_id: str) -> bool: with self._lock: cursor = self._conn.execute("DELETE FROM chats WHERE id = ? AND project_id = ?", (chat_id, project_id)) self._conn.commit() return cursor.rowcount > 0 def get_chat(self, project_id: str, chat_id: str) -> dict[str, Any] | None: with self._lock: row = self._conn.execute( "SELECT * FROM chats WHERE id = ? AND project_id = ?", (chat_id, project_id), ).fetchone() if not row: return None return { "id": row["id"], "project_id": row["project_id"], "title": row["title"], "model_id": row["model_id"], "provider_id": row["provider_id"], "base_url": row["base_url"], "served_model_name": row["served_model_name"], "temperature": row["temperature"], "max_tokens": row["max_tokens"], "rag_profile": row["rag_profile"], "rag_limit": row["rag_limit"], "system_prompt": row["system_prompt"], "metadata": json.loads(row["metadata_json"] or "{}"), "created_at": row["created_at"], "updated_at": row["updated_at"], } def list_chats(self, project_id: str) -> list[dict[str, Any]]: with self._lock: rows = self._conn.execute("SELECT * FROM chats WHERE project_id = ? ORDER BY updated_at DESC", (project_id,)).fetchall() return [ { "id": row["id"], "project_id": row["project_id"], "title": row["title"], "model_id": row["model_id"], "provider_id": row["provider_id"], "base_url": row["base_url"], "served_model_name": row["served_model_name"], "temperature": row["temperature"], "max_tokens": row["max_tokens"], "rag_profile": row["rag_profile"], "rag_limit": row["rag_limit"], "created_at": row["created_at"], "updated_at": row["updated_at"], } for row in rows ] def update_chat(self, project_id: str, chat_id: str, **payload: Any) -> None: current = self.get_chat(project_id, chat_id) if not current: return fields = [] values: list[Any] = [] for key, value in payload.items(): if key not in {"model_id", "provider_id", "base_url", "served_model_name", "temperature", "max_tokens", "rag_profile", "rag_limit", "system_prompt", "title", "metadata"}: continue if key == "metadata": fields.append("metadata_json = ?") values.append(json.dumps(value or {}, ensure_ascii=False)) else: fields.append(f"{key} = ?") values.append(value) if not fields: return fields.append("updated_at = ?") values.extend([now_iso(), chat_id, project_id]) with self._lock: self._conn.execute(f"UPDATE chats SET {', '.join(fields)} WHERE id = ? AND project_id = ?", tuple(values)) self._conn.commit() def add_message(self, project_id: str, chat_id: str, role: str, content: str, payload: dict[str, Any] | None = None) -> dict[str, Any]: message_id = uuid.uuid4().hex now = now_iso() payload_json = json.dumps(payload or {}, ensure_ascii=False) with self._lock: self._conn.execute( "INSERT INTO messages (id, project_id, chat_id, role, content, payload_json, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)", (message_id, project_id, chat_id, role, content, payload_json, now), ) self._conn.execute( "UPDATE chats SET updated_at = ? WHERE id = ? AND project_id = ?", (now, chat_id, project_id), ) self._conn.commit() return { "id": message_id, "project_id": project_id, "chat_id": chat_id, "role": role, "content": content, "payload": payload or {}, "created_at": now, } def list_messages(self, project_id: str, chat_id: str, *, limit: int = 200) -> list[dict[str, Any]]: chat = self.get_chat(project_id, chat_id) if not chat: raise FileNotFoundError("chat not found") with self._lock: rows = self._conn.execute( "SELECT * FROM messages WHERE project_id = ? AND chat_id = ? ORDER BY created_at DESC LIMIT ?", (project_id, chat_id, limit), ).fetchall() rows = list(reversed(rows)) return [ { "id": row["id"], "project_id": row["project_id"], "chat_id": row["chat_id"], "role": row["role"], "content": row["content"], "payload": json.loads(row["payload_json"] or "{}"), "created_at": row["created_at"], } for row in rows ] def delete_messages(self, project_id: str, chat_id: str) -> bool: chat = self.get_chat(project_id, chat_id) if not chat: return False with self._lock: self._conn.execute( "DELETE FROM messages WHERE project_id = ? AND chat_id = ?", (project_id, chat_id), ) self._conn.commit() return True def stats(self) -> dict[str, int]: with self._lock: projects = self._conn.execute("SELECT COUNT(1) AS total FROM projects").fetchone() chats = self._conn.execute("SELECT COUNT(1) AS total FROM chats").fetchone() messages = self._conn.execute("SELECT COUNT(1) AS total FROM messages").fetchone() return { "projects": int(projects[0]) if projects else 0, "chats": int(chats[0]) if chats else 0, "messages": int(messages[0]) if messages else 0, } def build_rag_prompt(question: str, *, profile_name: str = "auto", limit: int | None = None) -> dict[str, Any]: if not profile_name: profile_name = "auto" if profile_name == "default": profile_name = "auto" if not DEFAULT_1C_RAG_INDEX.exists(): raise FileNotFoundError(f"1C RAG index is missing: {DEFAULT_1C_RAG_INDEX}") index = read_json(DEFAULT_1C_RAG_INDEX) profile = resolve_rag_profile(profile_name, question) results = search_lexical_index( index, question, limit=int(limit or profile["limit"]), candidate_limit=int(profile["candidate_limit"]), dedupe_by_document=bool(profile["dedupe_by_document"]), min_score=float(profile["min_score"]), source_types=profile["source_types"], ) context = format_context(results, max_chars=int(profile["max_context_chars"])) return { "prompt": render_prompt(DEFAULT_1C_RAG_PROMPT, context=context, question=question), "profile": profile["id"], "context_count": len(results), "sources": [ { "source_path": result["document"].get("source_path"), "title": result["document"].get("title"), "chunk_index": result["document"].get("chunk_index"), "score": round(float(result.get("score") or 0), 4), } for result in results ], } def truthy(value: Any, *, default: bool = False) -> bool: if value is None: return default if isinstance(value, bool): return value if isinstance(value, (int, float)): return value != 0 if isinstance(value, str): normalized = value.strip().casefold() if normalized in {"1", "true", "yes", "y", "on", "да"}: return True if normalized in {"0", "false", "no", "n", "off", "нет"}: return False return default def looks_like_runtime_question(message: str) -> bool: text = message.casefold() runtime_terms = ( "runtime", "конфигурац", "какая модель", "какой модель", "модель ии", "подключен", "провайдер", "adapter", "адаптер", "rag", "раг", "base_url", ) return any(term in text for term in runtime_terms) and ( "модель" in text or "провайдер" in text or "адаптер" in text or "adapter" in text or "rag" in text or "раг" in text ) def render_runtime_answer(snapshot: dict[str, Any]) -> str: provider = snapshot["provider"] model = snapshot["model"] settings = snapshot["runtime_settings"] adapter = snapshot["adapter"] lines = [ "Текущая runtime-конфигурация чата:", f"- provider: {provider.get('id')} ({provider.get('type')})", f"- model_id: {model.get('chat_model_id') or model.get('selected_model_id')}", f"- served_model: {model.get('served_model')}", f"- base_url: {model.get('base_url')}", f"- rag_profile: {settings.get('rag_profile')}", f"- rag_limit: {settings.get('rag_limit')}", f"- temperature: {settings.get('temperature')}", f"- max_tokens: {settings.get('max_tokens')}", f"- adapter_url: {adapter.get('url') or ''}", ] return "\n".join(lines) def build_system_content(default_prompt: str, chat_prompt: str | None) -> str: parts = [part.strip() for part in (default_prompt, chat_prompt or "") if part and part.strip()] return "\n\n".join(parts) def adapter_result_prompt(tool_results: list[dict[str, Any]]) -> str: return ( "Результаты инструментов 1C adapter для текущего вопроса:\n" + json.dumps(tool_results, ensure_ascii=False, indent=2) ) def call_openai_compatible( provider: ProviderConfig, model: str, messages: list[dict], *, temperature: float, max_tokens: int, base_urls: list[str] | None = None, ) -> dict[str, Any]: candidates = list(base_urls or []) if provider.base_url and provider.base_url not in candidates: candidates.append(provider.base_url) if not candidates: raise ValueError("provider base url is not configured") headers: dict[str, str] = {"Content-Type": "application/json"} if provider.api_key_env: api_key = os.environ.get(provider.api_key_env, "") if api_key: headers["Authorization"] = f"Bearer {api_key}" payload = { "model": model, "messages": messages, "temperature": temperature, "max_tokens": max_tokens, } body = json.dumps(payload, ensure_ascii=False).encode("utf-8") last_error: BaseException | None = None started = time.perf_counter() for base_url in candidates: request = Request( f"{normalize_base_url(base_url)}/v1/chat/completions", data=body, headers=headers, method="POST", ) try: with urlopen(request, timeout=provider.timeout) as response: raw = json.loads(response.read().decode("utf-8")) choices = raw.get("choices") or [] if not choices: raise ValueError("provider returned no choices") msg = choices[0].get("message") or {} content = (msg.get("content") or "").strip() if not content: raise ValueError("provider returned empty content") return { "text": content, "raw": raw, "latency_ms": round((time.perf_counter() - started) * 1000), "provider_base_url": base_url, } except HTTPError as exc: raw = exc.read().decode("utf-8", errors="replace") if exc.fp else "" raise ValueError(f"provider HTTP {exc.code}: {raw}") from exc except (URLError, socket.timeout) as exc: last_error = exc continue except ValueError: raise raise ValueError(str(last_error) if last_error else "provider returned invalid response") def call_model( provider: ProviderConfig, model: str, messages: list[dict], *, temperature: float, max_tokens: int, base_url: str | None = None, ) -> dict[str, Any]: if provider.type == "openai-compatible": return call_openai_compatible( provider, model, messages, temperature=temperature, max_tokens=max_tokens, base_urls=[base_url] if base_url else None, ) raise ValueError(f"unsupported provider type: {provider.type}") ADAPTER_BASE_ID_EXEMPT_METHODS = {"health", "help.methods"} ADAPTER_BASE_ID_REQUIRED_PREFIXES = ( "access.", "code.", "metadata.", "modules.", "query.", "storage.", "templates.", ) AGENT_FORBIDDEN_TECHNICAL_SELECTOR_FIELDS = { "table", "file_name", "file_names", "module_ref", "module_id", "stream_index", "bsl_offset", "cas_key", "storage_key", "include_storage", "guid", "object_guid", "form_guid", "extension_guid", } AGENT_CONFIGURATION_METHOD_PREFIXES = ("metadata.", "modules.", "code.", "templates.", "extension.") def adapter_method_requires_base_id(method: str) -> bool: method = str(method or "").strip() if method in ADAPTER_BASE_ID_EXEMPT_METHODS: return False return method.startswith(ADAPTER_BASE_ID_REQUIRED_PREFIXES) def validate_adapter_call(method: str, params: dict[str, Any] | None) -> None: if not isinstance(params, dict): raise ValueError("adapter params must be object") if adapter_method_requires_base_id(method) and not str(params.get("base_id") or "").strip(): raise ValueError(f"adapter method {method} requires params.base_id") def agent_technical_selector_fields(value: Any) -> list[str]: """Reject storage coordinates even when a caller nests them in JSON.""" found: set[str] = set() if isinstance(value, dict): for key, nested in value.items(): if key in AGENT_FORBIDDEN_TECHNICAL_SELECTOR_FIELDS: found.add(key) found.update(agent_technical_selector_fields(nested)) elif isinstance(value, list): for nested in value: found.update(agent_technical_selector_fields(nested)) return sorted(found) def prepare_agent_adapter_call(method: str, params: dict[str, Any]) -> dict[str, Any]: """Keep the agent on public metadata selectors rather than SQL routes.""" prepared = dict(params) diagnostic_allowed = str(os.environ.get("ONEC_AGENT_ALLOW_DIAGNOSTIC") or "").strip().casefold() in {"1", "true", "yes", "on"} technical = agent_technical_selector_fields(prepared) if technical and not diagnostic_allowed: raise ValueError( "agent adapter calls require public names/selectors; forbidden technical fields: " + ", ".join(technical) ) if method.startswith(AGENT_CONFIGURATION_METHOD_PREFIXES): prepared.setdefault("configuration_view", "effective_working") prepared.setdefault("source_state", "working") return prepared def call_adapter(method: str, params: dict[str, Any] | None, *, base_url: str | None = None) -> dict[str, Any]: """Call the adapter through its public MCP boundary, never its SQL REST surface.""" params = prepare_agent_adapter_call(method, params or {}) validate_adapter_call(method, params) mcp_url = normalize_base_url(base_url or os.environ.get("ONEC_MCP_URL", "http://docker.cin.su:8021")) headers = {"Content-Type": "application/json", "Accept": "application/json, text/event-stream"} initialize = { "jsonrpc": "2.0", "id": f"onec-agent-init-{uuid.uuid4().hex}", "method": "initialize", "params": { "protocolVersion": "2025-06-18", "capabilities": {}, "clientInfo": {"name": "onec-agent", "version": "1"}, }, } init_request = Request( f"{mcp_url}/mcp", data=json.dumps(initialize, ensure_ascii=False).encode("utf-8"), headers=headers, method="POST", ) try: with urlopen(init_request, timeout=30) as response: init_raw = json.loads(response.read().decode("utf-8")) session_id = response.headers.get("Mcp-Session-Id") except urllib.error.HTTPError as exc: body = exc.read().decode("utf-8", errors="replace") raise ValueError(f"MCP initialize returned HTTP {exc.code}: {body}") from exc if not isinstance(init_raw, dict) or not isinstance(init_raw.get("result"), dict): raise ValueError("MCP initialize response is not JSON-RPC success") call_headers = dict(headers) if session_id: call_headers["Mcp-Session-Id"] = session_id call = { "jsonrpc": "2.0", "id": f"onec-agent-call-{uuid.uuid4().hex}", "method": "tools/call", "params": {"name": "onec_request", "arguments": {"method": method, "payload": params}}, } request = Request( f"{mcp_url}/mcp", data=json.dumps(call, ensure_ascii=False).encode("utf-8"), headers=call_headers, method="POST", ) try: with urlopen(request, timeout=120) as response: raw = json.loads(response.read().decode("utf-8")) except urllib.error.HTTPError as exc: body = exc.read().decode("utf-8", errors="replace") raise ValueError(f"MCP tool call returned HTTP {exc.code}: {body}") from exc result = raw.get("result") if isinstance(raw, dict) and isinstance(raw.get("result"), dict) else None content = result.get("content") if isinstance(result, dict) and isinstance(result.get("content"), list) else [] text = content[0].get("text") if content and isinstance(content[0], dict) else None if not isinstance(text, str): raise ValueError("MCP tool response has no JSON text content") try: decoded = json.loads(text) except json.JSONDecodeError as exc: raise ValueError("MCP tool response text is not JSON") from exc if not isinstance(decoded, dict): raise ValueError("MCP tool response payload is not an object") return decoded class AgentHandler(BaseHTTPRequestHandler): def __init__( self, *args: Any, store: Store, providers: dict[str, ProviderConfig], system_prompt: str, audit_store: JsonlAuditStore | None = None, **kwargs: Any, ) -> None: self.store = store self.providers = providers self.system_prompt = system_prompt self.audit_store = audit_store self._request_started_at: float | None = None self._request_trace_id: str | None = None self._request_id: str | None = None super().__init__(*args, **kwargs) def _parts(self) -> list[str]: return [segment for segment in self.path.split("?")[0].strip("/").split("/") if segment] def _ensure_request_context(self) -> None: if self._request_started_at is None: self._request_started_at = time.perf_counter() if self._request_trace_id is None: self._request_trace_id = request_trace_id(self) if self._request_id is None: self._request_id = request_request_id(self, self._request_trace_id) def _trace_id(self) -> str: self._ensure_request_context() return str(self._request_trace_id) def _request_id_value(self) -> str: self._ensure_request_context() return str(self._request_id) def _request_duration_ms(self) -> int: self._ensure_request_context() return int((time.perf_counter() - float(self._request_started_at or time.perf_counter())) * 1000) def _emit_access_event(self, status_code: int, *, error_code: str | None = None) -> None: if not self.audit_store: return self.audit_store.write_event( "access_events", sanitize_for_logging( { "timestamp": now_iso(), "service": "onec-agent", "trace_id": self._trace_id(), "request_id": self._request_id_value(), "method": self.command, "path": self.path, "path_template": self.path.split("?")[0], "status_code": status_code, "duration_ms": self._request_duration_ms(), "request_size_bytes": int(self.headers.get("Content-Length") or 0), "client_type": "http", "error_code": error_code or "", } ), ) def _emit_turn_audit( self, *, project_id: str, chat_id: str, user_saved: dict[str, Any], assistant: dict[str, Any], outcome: str, failure_type: str, route: dict[str, Any], guardrail: dict[str, Any] | None = None, error_code: str | None = None, error_message: str | None = None, ) -> str: turn_id = uuid.uuid4().hex if not self.audit_store: return turn_id self.audit_store.write_event( "turn_audit", sanitize_for_logging( { "turn_id": turn_id, "timestamp": now_iso(), "project_id": project_id, "chat_id": chat_id, "plugin": str(route.get("plugin") or DEFAULT_PLUGIN), "service": "onec-agent", "trace_id": self._trace_id(), "request_id": self._request_id_value(), "user_message_id": user_saved.get("id"), "assistant_message_id": assistant.get("id"), "user_text": user_saved.get("content"), "assistant_text": assistant.get("content"), "outcome": outcome, "failure_type": failure_type, "duration_ms": self._request_duration_ms(), "human_review_status": "pending", "guardrail": guardrail or {}, "error_code": error_code or "", "error_message": error_message or "", } ), ) return turn_id def _emit_failed_turn_audit( self, *, project_id: str, chat_id: str, user_text: str, route: dict[str, Any], failure_type: str, error_code: str, error_message: str, ) -> str: placeholder_user = { "id": "", "content": user_text, } placeholder_assistant = { "id": "", "content": "", } return self._emit_turn_audit( project_id=project_id, chat_id=chat_id, user_saved=placeholder_user, assistant=placeholder_assistant, outcome="failure", failure_type=failure_type, route=route, error_code=error_code, error_message=error_message, ) def _emit_model_call_event( self, *, turn_id: str, route: dict[str, Any], provider_id: str, provider: ProviderConfig, served_model: str, model_base_url: str, temperature: float, max_tokens: int, model_messages: list[dict[str, Any]], answer: dict[str, Any], ) -> None: if not self.audit_store: return raw = answer.get("raw") if isinstance(answer, dict) else {} usage = usage_from_provider_response(raw if isinstance(raw, dict) else {}) finish_reason = None if isinstance(raw, dict): choices = raw.get("choices") or [] if choices and isinstance(choices[0], dict): finish_reason = choices[0].get("finish_reason") self.audit_store.write_event( "model_calls", sanitize_for_logging( { "model_call_id": uuid.uuid4().hex, "turn_id": turn_id, "timestamp": now_iso(), "trace_id": self._trace_id(), "request_id": self._request_id_value(), "provider_id": provider_id, "provider_type": provider.type, "base_url": answer.get("provider_base_url", model_base_url), "model_registry_id": str(route.get("model", {}).get("id") or ""), "served_model_name": served_model, "route_name": str(route.get("plugin") or DEFAULT_PLUGIN), "temperature": temperature, "max_tokens": max_tokens, "prompt_messages_json": model_messages, "response_json": raw if isinstance(raw, dict) else {}, "prompt_tokens": usage.get("prompt_tokens"), "completion_tokens": usage.get("completion_tokens"), "total_tokens": usage.get("total_tokens"), "cost_estimate": None, "latency_ms": answer.get("latency_ms"), "finish_reason": finish_reason, "cache_hit": False, } ), ) def _emit_tool_call_event( self, *, turn_id: str | None, tool_family: str, tool_name: str, target_service: str, request_json: dict[str, Any], response_json: dict[str, Any], status: str, duration_ms: int, retry_count: int = 0, ) -> None: if not self.audit_store: return self.audit_store.write_event( "tool_calls", sanitize_for_logging( { "tool_call_id": uuid.uuid4().hex, "turn_id": turn_id or "", "timestamp": now_iso(), "trace_id": self._trace_id(), "request_id": self._request_id_value(), "tool_family": tool_family, "tool_name": tool_name, "target_service": target_service, "request_json": request_json, "response_json": response_json, "status": status, "duration_ms": duration_ms, "retry_count": retry_count, } ), ) def _emit_retrieval_event( self, *, turn_id: str | None, plugin: str, profile: str, query_text: str, top_k: int | None, rag_data: dict[str, Any], retrieval_latency_ms: int, ) -> None: if not self.audit_store: return sources = rag_data.get("sources") or [] prompt_text = str(rag_data.get("prompt") or "") self.audit_store.write_event( "retrieval_events", sanitize_for_logging( { "retrieval_id": uuid.uuid4().hex, "turn_id": turn_id or "", "timestamp": now_iso(), "trace_id": self._trace_id(), "request_id": self._request_id_value(), "plugin": plugin, "profile": profile, "index_version": str(DEFAULT_1C_RAG_INDEX), "query_text": query_text, "top_k": top_k, "sources_json": sources, "context_chars": len(prompt_text), "retrieval_latency_ms": retrieval_latency_ms, } ), ) def _bad_request(self, message: str, *, code: str | None = None, status: int = 400) -> None: json_error(self, status, message, code=code, trace_id=self._trace_id()) self._emit_access_event(status, error_code=code) def do_OPTIONS(self) -> None: self.send_response(204) if self.path.startswith("/ui"): self.send_header("Access-Control-Allow-Origin", "*") self.send_header("Access-Control-Allow-Methods", "GET, OPTIONS") self.send_header("Access-Control-Allow-Headers", "Content-Type, Authorization, X-Request-Id") self.send_header("Content-Length", "0") self.end_headers() return super_headers = { "Access-Control-Allow-Origin": "*", "Access-Control-Allow-Methods": "GET, POST, PATCH, DELETE, OPTIONS", "Access-Control-Allow-Headers": "Content-Type, Authorization, X-Request-Id", } for key, value in super_headers.items(): self.send_header(key, value) self.end_headers() def _send_not_found(self, message: str = "not found") -> None: json_error(self, 404, message, code="not_found", trace_id=self._trace_id()) self._emit_access_event(404, error_code="not_found") def _send_ok(self, status: int, payload: dict[str, Any]) -> None: payload["trace_id"] = self._trace_id() json_response(self, status, payload) self._emit_access_event(status) def _send_no_content(self, status: int = 204) -> None: self.send_response(status) self.send_header("Access-Control-Allow-Origin", "*") self.send_header("Access-Control-Allow-Methods", "GET, POST, PATCH, DELETE, OPTIONS") self.send_header("Access-Control-Allow-Headers", "Content-Type, Authorization, X-Request-Id") self.send_header("Content-Length", "0") self.end_headers() self._emit_access_event(status) def _send_created(self, payload: dict[str, Any]) -> None: self._send_ok(201, payload) def _send_health(self) -> None: self._send_ok(200, {"status": "ok", "service": "onec-agent", "time": now_iso()}) def _send_state(self) -> None: uptime = (datetime.now(timezone.utc) - SERVICE_STARTED_AT).total_seconds() provider_ids = sorted(self.providers.keys()) self._send_ok( 200, { "service": "onec-agent", "started_at": SERVICE_STARTED_AT.isoformat(), "uptime_seconds": round(uptime, 3), "state": self.store.stats(), "providers": { "count": len(provider_ids), "ids": provider_ids, }, "paths": { "adapter_url": os.environ.get("ONEC_ADAPTER_URL"), "database": str(self.store.path), }, }, ) def _send_static_file(self, request_path: str) -> bool: parsed = urlparse(request_path) rel = parsed.path.lstrip("/") if rel == "ui": rel = "ui/index.html" if not rel.startswith("ui/"): return False rel = rel.removeprefix("ui/") if not rel: rel = "index.html" target = (WEB_UI_DIR / rel).resolve() base = WEB_UI_DIR.resolve() if not (str(target).startswith(f"{base}{os.sep}") or str(target) == str(base)): self._send_not_found("invalid ui path") return True if not target.exists() or not target.is_file(): self._send_not_found("ui file not found") return True content = target.read_bytes() mime_type, _encoding = mimetypes.guess_type(str(target)) self.send_response(200) self.send_header("Content-Type", mime_type or "application/octet-stream") self.send_header("Content-Length", str(len(content))) self.send_header("Access-Control-Allow-Origin", "*") self.end_headers() self.wfile.write(content) return True def do_GET(self) -> None: if self.path in {"/health", "/v1/health"}: self._send_health() return if self.path == "/": self._send_static_file("/ui/") return if self._send_static_file(self.path): return parts = self._parts() if len(parts) == 2 and parts[0] == "v1" and parts[1] == "projects": self._send_ok(200, {"projects": self.store.list_projects()}) return if len(parts) == 3 and parts[0] == "v1" and parts[1] == "projects": project = self.store.get_project(parts[2]) if not project: self._send_not_found("project not found") return self._send_ok(200, {"project": project}) return if len(parts) == 4 and parts[0] == "v1" and parts[1] == "projects" and parts[3] == "chats": if not self.store.get_project(parts[2]): self._send_not_found("project not found") return self._send_ok(200, {"project_id": parts[2], "chats": self.store.list_chats(parts[2])}) return if len(parts) == 5 and parts[0] == "v1" and parts[1] == "projects" and parts[3] == "chats": chat = self.store.get_chat(parts[2], parts[4]) if not chat: self._send_not_found("chat not found") return self._send_ok(200, {"chat": chat}) return if len(parts) == 6 and parts[0] == "v1" and parts[1] == "projects" and parts[3] == "chats" and parts[5] == "messages": parsed = parse_qs(urlparse(self.path).query) try: limit = normalize_limit(parsed.get("limit", [None])[0], default=200) except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return try: messages = self.store.list_messages(parts[2], parts[4], limit=limit) except FileNotFoundError: self._send_not_found("chat not found") return self._send_ok(200, {"project_id": parts[2], "chat_id": parts[4], "messages": messages}) return if len(parts) == 6 and parts[0] == "v1" and parts[1] == "projects" and parts[3] == "chats" and parts[5] == "runtime": project = self.store.get_project(parts[2]) if not project: self._send_not_found("project not found") return chat = self.store.get_chat(parts[2], parts[4]) if not chat: self._send_not_found("chat not found") return try: snapshot = _resolve_runtime_snapshot(project=project, chat=chat, providers=self.providers) except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return self._send_ok(200, {"project_id": parts[2], "chat_id": parts[4], "runtime": snapshot}) return if len(parts) == 2 and parts[0] == "v1" and parts[1] == "models": self._send_ok(200, {"models": model_catalog().get("models") or []}) return if len(parts) == 2 and parts[0] == "v1" and parts[1] == "providers": payload = {name: {"type": p.type, "base_url": p.base_url, "model": p.model} for name, p in self.providers.items()} self._send_ok(200, {"providers": payload}) return if len(parts) == 2 and parts[0] == "v1" and parts[1] == "state": self._send_state() return self._send_not_found("not found") def do_POST(self) -> None: if self.path.startswith("/v1/health"): self._send_health() return parts = self._parts() if len(parts) < 2 or parts[0] != "v1": self._bad_request("unsupported path") return try: payload = safe_request_json(self) except (ValueError, json.JSONDecodeError) as exc: self._bad_request(str(exc), code="invalid_argument") return if len(parts) == 2 and parts[1] == "projects": name = str(payload.get("name") or "").strip() if not name: self._bad_request("name is required", code="invalid_argument") return if payload.get("metadata") is not None and not isinstance(payload.get("metadata"), dict): self._bad_request("metadata must be object", code="invalid_argument") return description = str(payload.get("description") or "") project = self.store.create_project(name=name, description=description, metadata=payload.get("metadata") if isinstance(payload.get("metadata"), dict) else None) self._send_created({"project": project}) return if len(parts) == 4 and parts[1] == "projects" and parts[3] == "chats": if not self.store.get_project(parts[2]): self._send_not_found("project not found") return title = str(payload.get("title") or "1c chat") model_id = payload.get("model_id") provider_id = str(payload.get("provider_id") or DEFAULT_PROVIDER) if provider_id not in self.providers: self._bad_request(f"unknown provider_id={provider_id}", code="invalid_argument") return try: route = route_for_model(model_id if model_id else None) except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return try: temperature = normalize_float(payload.get("temperature"), field="temperature", default=0.2) max_tokens = normalize_int(payload.get("max_tokens"), field="max_tokens", default=1200) except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return rag_profile = str(payload.get("rag_profile") or "auto") rag_limit = payload.get("rag_limit") try: rag_limit = normalize_int(rag_limit, field="rag_limit") if rag_limit is not None else None except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return if not payload.get("metadata") is None and not isinstance(payload.get("metadata"), dict): self._bad_request("metadata must be object", code="invalid_argument") return system_prompt = str(payload.get("system_prompt") or "") chat = self.store.create_chat( project_id=parts[2], title=title, model_id=str(model_id) if model_id else None, provider_id=provider_id, base_url=str(route.get("base_url") or ""), served_model_name=str(route.get("served_model_name") or ""), temperature=temperature, max_tokens=max_tokens, rag_profile=rag_profile, rag_limit=rag_limit, system_prompt=system_prompt, metadata=payload.get("metadata") if isinstance(payload.get("metadata"), dict) else None, ) self._send_created({"chat": chat, "routing": route}) return if len(parts) == 6 and parts[1] == "projects" and parts[3] == "chats" and parts[5] in {"messages", "turn", "onec-tool"}: project_id = parts[2] chat_id = parts[4] if not self.store.get_chat(project_id, chat_id): self._send_not_found("chat not found") return if parts[5] == "messages": role = str(payload.get("role") or "user").strip() if role not in {"user", "assistant", "system", "tool"}: self._bad_request("role must be user|assistant|system|tool", code="invalid_argument") return content = str(payload.get("content") or "").strip() if not content: self._bad_request("content is required", code="invalid_argument") return message = self.store.add_message(project_id, chat_id, role=role, content=content, payload=payload.get("payload") if isinstance(payload.get("payload"), dict) else None) self._send_created({"message": message}) return if parts[5] == "onec-tool": method = str(payload.get("method") or "").strip() if not method: self._bad_request("method is required", code="invalid_argument") return params = payload.get("params") if isinstance(payload.get("params"), dict) else {} try: result = call_adapter(method, params) except (ValueError, urllib.error.URLError) as exc: self._bad_request(str(exc), code="adapter_error") return msg = self.store.add_message(project_id, chat_id, role="assistant", content=f"tool:{method}", payload={"tool": method, "result": result, "status": "ok"}) self._send_ok(200, {"tool_result": result, "message": msg}) return user_message = str(payload.get("message") or payload.get("content") or "").strip() if not user_message: self._bad_request("message is required", code="invalid_argument") return chat = self.store.get_chat(project_id, chat_id) assert chat is not None provider_id = str(payload.get("provider_id") or chat["provider_id"] or DEFAULT_PROVIDER) if provider_id not in self.providers: self._bad_request(f"unknown provider_id={provider_id}", code="invalid_argument") return provider = self.providers[provider_id] try: selected = route_for_model(chat.get("model_id") or payload.get("model_id")) except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return served_model = str( payload.get("served_model_name") or chat.get("served_model_name") or selected.get("served_model_name") or provider.model ) model_base_url = str( payload.get("base_url") or selected.get("base_url") or chat.get("base_url") or provider.base_url ) project = self.store.get_project(project_id) or {} try: temperature = normalize_float( payload.get("temperature") if payload.get("temperature") is not None else chat["temperature"], field="temperature", ) max_tokens = normalize_int( payload.get("max_tokens") if payload.get("max_tokens") is not None else chat["max_tokens"], field="max_tokens", ) history_limit = normalize_int(payload.get("history_limit"), field="history_limit", default=60) except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return history_limit = min(history_limit, 200) history = self.store.list_messages(project_id, chat_id, limit=history_limit) model_messages: list[dict[str, Any]] = [] system_content = build_system_content(self.system_prompt, chat.get("system_prompt")) if system_content: model_messages.append({"role": "system", "content": system_content}) if looks_like_runtime_question(user_message): try: runtime_snapshot = _resolve_runtime_snapshot(project=project, chat=chat, providers=self.providers) except ValueError as exc: self._bad_request(str(exc), code="invalid_argument") return model_messages.append({"role": "user", "content": user_message}) user_saved = self.store.add_message( project_id, chat_id, role="user", content=user_message, payload={ "mode": "turn", "transport": { "outbound": _build_turn_payload( route=selected, provider_id=provider_id, provider=provider, served_model=served_model, model_base_url=model_base_url, temperature=temperature, max_tokens=max_tokens, history_limit=history_limit, model_messages=model_messages, ), }, }, ) final_text = render_runtime_answer(runtime_snapshot) guardrail = {"type": "runtime_direct_answer"} assistant = self.store.add_message( project_id, chat_id, role="assistant", content=final_text, payload={ "guardrail": guardrail, "raw": {"reason": "runtime_question_fastpath", "runtime": runtime_snapshot}, "transport": { "outbound": _build_turn_payload( route=selected, provider_id=provider_id, provider=provider, served_model=served_model, model_base_url=model_base_url, temperature=temperature, max_tokens=max_tokens, history_limit=history_limit, model_messages=model_messages, ), "inbound": { "provider_response": {"reason": "runtime_question_fastpath"}, }, }, }, ) turn_id = self._emit_turn_audit( project_id=project_id, chat_id=chat_id, user_saved=user_saved, assistant=assistant, outcome="success", failure_type="none", route=selected, guardrail=guardrail, ) self._send_ok( 200, { "project_id": project_id, "chat_id": chat_id, "user_message": user_saved, "assistant_message": assistant, "guardrail": guardrail, "turn_id": turn_id, }, ) return adapter_calls = payload.get("adapter_calls") if adapter_calls is None: adapter_calls = [] if not isinstance(adapter_calls, list): self._bad_request("adapter_calls must be array", code="invalid_argument") return tool_results: list[dict[str, Any]] = [] for index, call in enumerate(adapter_calls): if not isinstance(call, dict): self._bad_request(f"adapter_calls[{index}] must be object", code="invalid_argument") return method = str(call.get("method") or "").strip() if not method: self._bad_request(f"adapter_calls[{index}].method is required", code="invalid_argument") return params = call.get("params") if isinstance(call.get("params"), dict) else {} tool_started = time.perf_counter() try: result = call_adapter(method, params) except (ValueError, urllib.error.URLError) as exc: self._emit_failed_turn_audit( project_id=project_id, chat_id=chat_id, user_text=user_message, route=selected, failure_type="adapter_error", error_code="adapter_error", error_message=str(exc), ) self._bad_request(str(exc), code="adapter_error") return tool_results.append( { "method": method, "params": params, "result": result, "status": "ok", "duration_ms": int((time.perf_counter() - tool_started) * 1000), } ) if tool_results: model_messages.append({"role": "system", "content": adapter_result_prompt(tool_results)}) rag_data: dict[str, Any] | None = None use_rag = truthy(payload.get("use_rag"), default=True) if use_rag: rag_started = time.perf_counter() try: rag_limit_value = payload.get("rag_limit") if payload.get("rag_limit") is not None else chat.get("rag_limit") rag_limit = normalize_int(rag_limit_value, field="rag_limit") if rag_limit_value is not None else None rag_profile = str(payload.get("rag_profile") or chat.get("rag_profile") or "auto") rag_data = build_rag_prompt(user_message, profile_name=rag_profile, limit=rag_limit) except (ValueError, FileNotFoundError) as exc: self._emit_failed_turn_audit( project_id=project_id, chat_id=chat_id, user_text=user_message, route=selected, failure_type="rag_error", error_code="rag_error", error_message=str(exc), ) self._bad_request(str(exc), code="rag_error") return model_messages.append({"role": "system", "content": rag_data["prompt"]}) for row in history: if row["role"] in {"user", "assistant"}: model_messages.append({"role": row["role"], "content": row["content"]}) if chat["provider_id"] != provider_id: self.store.update_chat(project_id, chat_id, provider_id=provider_id) model_messages.append({"role": "user", "content": user_message}) user_saved = self.store.add_message( project_id, chat_id, role="user", content=user_message, payload={ "mode": "turn", "rag": rag_data, "tools": tool_results, "transport": { "outbound": _build_turn_payload( route=selected, provider_id=provider_id, provider=provider, served_model=served_model, model_base_url=model_base_url, temperature=temperature, max_tokens=max_tokens, history_limit=history_limit, model_messages=model_messages, ), }, }, ) try: answer = call_model( provider, served_model, model_messages, temperature=temperature, max_tokens=max_tokens, base_url=model_base_url, ) except (ValueError, urllib.error.URLError) as exc: self._emit_failed_turn_audit( project_id=project_id, chat_id=chat_id, user_text=user_message, route=selected, failure_type="model_error", error_code="model_error", error_message=str(exc), ) self._bad_request(str(exc), code="model_error") return final_text = answer["text"] or "" assistant = self.store.add_message( project_id, chat_id, role="assistant", content=final_text or "", payload={ "provider": provider_id, "base_url": answer.get("provider_base_url", provider.base_url), "model": served_model, "route": selected, "raw": answer["raw"], "latency_ms": answer["latency_ms"], "rag": rag_data, "tools": tool_results, "transport": { "outbound": _build_turn_payload( route=selected, provider_id=provider_id, provider=provider, served_model=served_model, model_base_url=model_base_url, temperature=temperature, max_tokens=max_tokens, history_limit=history_limit, model_messages=model_messages, ), "inbound": { "provider_response": answer["raw"], }, }, }, ) turn_id = self._emit_turn_audit( project_id=project_id, chat_id=chat_id, user_saved=user_saved, assistant=assistant, outcome="success", failure_type="none", route=selected, ) if tool_results: for tool in tool_results: self._emit_tool_call_event( turn_id=turn_id, tool_family="adapter", tool_name=str(tool.get("method") or ""), target_service="1c-adapter", request_json={"method": tool.get("method"), "params": tool.get("params")}, response_json=tool.get("result") if isinstance(tool.get("result"), dict) else {}, status=str(tool.get("status") or "ok"), duration_ms=int(tool.get("duration_ms") or 0), ) if rag_data: self._emit_retrieval_event( turn_id=turn_id, plugin=str(selected.get("plugin") or DEFAULT_PLUGIN), profile=str(rag_data.get("profile") or payload.get("rag_profile") or chat.get("rag_profile") or "auto"), query_text=user_message, top_k=rag_data.get("context_count"), rag_data=rag_data, retrieval_latency_ms=int((time.perf_counter() - rag_started) * 1000), ) self._emit_model_call_event( turn_id=turn_id, route=selected, provider_id=provider_id, provider=provider, served_model=served_model, model_base_url=model_base_url, temperature=temperature, max_tokens=max_tokens, model_messages=model_messages, answer=answer, ) self.store.update_chat( project_id, chat_id, model_id=str(chat.get("model_id") or payload.get("model_id") or selected.get("model", {}).get("id")), served_model_name=served_model, base_url=answer.get("provider_base_url", model_base_url), temperature=temperature, max_tokens=max_tokens, ) self._send_ok( 200, { "project_id": project_id, "chat_id": chat_id, "user_message": user_saved, "assistant_message": assistant, "rag": rag_data, "tools": tool_results, "turn_id": turn_id, }, ) return self._send_not_found("not found") def do_PATCH(self) -> None: parts = self._parts() if len(parts) < 2 or parts[0] != "v1": self._send_not_found("unsupported path") return try: payload = safe_request_json(self) except (ValueError, json.JSONDecodeError) as exc: self._bad_request(str(exc), code="invalid_argument") return if len(parts) == 3 and parts[1] == "projects": project = self.store.get_project(parts[2]) if not project: self._send_not_found("project not found") return metadata = payload.get("metadata") if metadata is not None and not isinstance(metadata, dict): self._bad_request("metadata must be object", code="invalid_argument") return updated = self.store.update_project( parts[2], name=(str(payload.get("name")).strip() if payload.get("name") is not None else None), description=(str(payload.get("description")) if payload.get("description") is not None else None), metadata=metadata if isinstance(metadata, dict) else None, ) if not updated: self._send_not_found("project not found") return self._send_ok(200, {"project": updated}) return if len(parts) == 5 and parts[1] == "projects" and parts[3] == "chats": if not self.store.get_chat(parts[2], parts[4]): self._send_not_found("chat not found") return metadata = payload.get("metadata") if metadata is not None and not isinstance(metadata, dict): self._bad_request("metadata must be object", code="invalid_argument") return updated = self.store.update_chat_config( parts[2], parts[4], title=(str(payload.get("title")) if payload.get("title") is not None else None), model_id=(str(payload.get("model_id")) if payload.get("model_id") is not None else None), provider_id=(str(payload.get("provider_id")) if payload.get("provider_id") is not None else None), base_url=(str(payload.get("base_url")) if payload.get("base_url") is not None else None), served_model_name=(str(payload.get("served_model_name")) if payload.get("served_model_name") is not None else None), temperature=(normalize_float(payload.get("temperature"), field="temperature") if payload.get("temperature") is not None else None), max_tokens=(normalize_int(payload.get("max_tokens"), field="max_tokens") if payload.get("max_tokens") is not None else None), rag_profile=(str(payload.get("rag_profile")) if payload.get("rag_profile") is not None else None), rag_limit=(normalize_int(payload.get("rag_limit"), field="rag_limit") if payload.get("rag_limit") is not None else None), system_prompt=(str(payload.get("system_prompt")) if payload.get("system_prompt") is not None else None), metadata=metadata, ) if not updated: self._send_not_found("chat not found") return self._send_ok(200, {"chat": updated}) return self._send_not_found("not found") def do_DELETE(self) -> None: parts = self._parts() if len(parts) == 3 and parts[0] == "v1" and parts[1] == "projects": if not self.store.delete_project(parts[2]): self._send_not_found("project not found") return self._send_no_content() return if len(parts) == 5 and parts[0] == "v1" and parts[1] == "projects" and parts[3] == "chats": if not self.store.delete_chat(parts[2], parts[4]): self._send_not_found("chat not found") return self._send_no_content() return if len(parts) == 6 and parts[0] == "v1" and parts[1] == "projects" and parts[3] == "chats" and parts[5] == "messages": if not self.store.delete_messages(parts[2], parts[4]): self._send_not_found("chat not found") return self._send_no_content() return self._send_not_found("not found") def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description="Run independent 1C agent.") parser.add_argument("--host", default=os.environ.get("ONEC_AGENT_HOST", DEFAULT_HOST)) parser.add_argument("--port", type=int, default=int(os.environ.get("ONEC_AGENT_PORT", str(DEFAULT_PORT)))) return parser.parse_args() def main() -> int: args = parse_args() store = Store(Path(os.environ.get("ONEC_AGENT_DB_PATH", "./data/onec-agent.db"))) providers = load_provider_configs() system_prompt = read_path_text(DEFAULT_SYSTEM_PROMPT_PATH, "") audit_store = JsonlAuditStore(resolve_audit_root(ROOT, service="onec-agent", env_var="ONEC_AGENT_AUDIT_DIR"), service="onec-agent") def _handler(*handler_args: Any, **handler_kwargs: Any) -> None: AgentHandler( *handler_args, store=store, providers=providers, system_prompt=system_prompt, audit_store=audit_store, **handler_kwargs, ) server = ThreadingHTTPServer((args.host, args.port), _handler) print(f"onec-agent listening on http://{args.host}:{args.port}") print(f"db={store.path}") try: server.serve_forever() except KeyboardInterrupt: return 0 finally: server.server_close() return 0 if __name__ == "__main__": raise SystemExit(main())