Files
llm/plugins/1c/agent/agent_server.py
T
2026-08-14 09:40:51 +03:00

2178 lines
85 KiB
Python

#!/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())