ForcePilot/backend/package/yuxi/storage/postgres/manager.py
Wenjie Zhang 1efbf86bf0 refactor: 评估系统重构 - 后端统一 dataset/run 语义
- 评估数据集/题目/运行/逐题结果全量入库,JSONL 仅作为交换格式
- benchmark -> dataset, task -> run 术语统一
- 评估生成支持配置并发数
- 移除 base.py 中旧的 benchmarks_meta 元数据逻辑
- manager.py 统一表创建方法名
2026-05-26 17:40:17 +08:00

591 lines
29 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""PostgreSQL 数据库管理器 - 支持知识库和业务数据"""
import json
import os
from contextlib import asynccontextmanager
from psycopg_pool import AsyncConnectionPool
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker, create_async_engine
from sqlalchemy.orm import declarative_base
from yuxi.storage.postgres.models_business import Base as BusinessBase
from yuxi.storage.postgres.models_knowledge import Base as KnowledgeBase
from yuxi.utils import logger
from server.utils.singleton import SingletonMeta
# 合并两个 Base
CombinedBase = declarative_base()
# 继承所有表
for module in [KnowledgeBase, BusinessBase]:
for table_name in dir(module):
table = getattr(module, table_name)
if isinstance(table, type) and hasattr(table, "__tablename__"):
setattr(CombinedBase, table_name, table)
class PostgresManager(metaclass=SingletonMeta):
"""PostgreSQL 数据库管理器 - 支持知识库和业务数据"""
# 知识库 PostgreSQL URL 环境变量名
KB_DATABASE_URL_ENV = "POSTGRES_URL"
def __init__(self):
self.async_engine = None
self.AsyncSession = None
self.langgraph_pool = None
self._initialized = False
def initialize(self):
"""初始化数据库连接"""
if self._initialized:
return
db_url = os.getenv(self.KB_DATABASE_URL_ENV)
if not db_url:
logger.error(
f"环境变量 {self.KB_DATABASE_URL_ENV} 未设置,"
"请在 docker-compose.yml 或 .env 中配置 PostgreSQL 连接字符串"
)
return
try:
# 创建异步 SQLAlchemy 引擎
self.async_engine = create_async_engine(
db_url,
json_serializer=lambda obj: json.dumps(obj, ensure_ascii=False),
json_deserializer=json.loads,
pool_pre_ping=True,
pool_recycle=1800,
pool_size=10,
max_overflow=20,
)
# 创建异步会话工厂
self.AsyncSession = async_sessionmaker(
bind=self.async_engine,
class_=AsyncSession,
expire_on_commit=False,
)
# ==========================================
# 2. 为 LangGraph 专门初始化一个原生 psycopg_pool
# ==========================================
# ⚠️ 注意psycopg 不认识 "+asyncpg" 这样的 SQLAlchemy 方言标识。
# 如果你的 db_url 是 "postgresql+asyncpg://user:pwd@host/db"
# 需要把它清洗成标准的 "postgresql://user:pwd@host/db"
langgraph_db_url = db_url.replace("+asyncpg", "").replace("+psycopg", "")
# 创建 LangGraph 专属连接池
self.langgraph_pool = AsyncConnectionPool(
conninfo=langgraph_db_url,
max_size=10, # 根据你的 Agent 并发情况设置,通常 5-10 足够了
kwargs={"autocommit": True}, # LangGraph Checkpoint 强依赖 autocommit
)
self._initialized = True
logger.info(f"PostgreSQL manager initialized for knowledge base: {db_url.split('@')[0]}://***")
except Exception as e:
logger.error(f"Failed to initialize PostgreSQL manager: {e}")
# 不抛出异常,允许应用启动,但在使用时会报错
def _check_initialized(self):
"""检查是否已初始化"""
if not self._initialized:
raise RuntimeError("PostgreSQL manager not initialized. Please check configuration.")
async def create_tables(self):
"""创建所有表(知识库和业务表)"""
self._check_initialized()
async with self.async_engine.begin() as conn:
await conn.run_sync(KnowledgeBase.metadata.create_all)
await conn.run_sync(BusinessBase.metadata.create_all)
logger.info("PostgreSQL tables created/checked (knowledge + business)")
async def create_business_tables(self):
"""创建所有业务数据表"""
self._check_initialized()
async with self.async_engine.begin() as conn:
await conn.run_sync(BusinessBase.metadata.create_all)
logger.info("PostgreSQL business tables created/checked")
async def drop_tables(self):
"""删除所有表(慎用!)"""
self._check_initialized()
async with self.async_engine.begin() as conn:
await conn.run_sync(BusinessBase.metadata.drop_all)
await conn.run_sync(KnowledgeBase.metadata.drop_all)
logger.info("PostgreSQL tables dropped")
async def ensure_knowledge_schema(self):
"""确保知识库 schema 包含所有必要字段"""
self._check_initialized()
stmts = [
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS embedding_model_spec VARCHAR(512)",
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS llm_model_spec VARCHAR(512)",
"ALTER TABLE IF EXISTS knowledge_bases DROP COLUMN IF EXISTS embed_info",
"ALTER TABLE IF EXISTS knowledge_bases DROP COLUMN IF EXISTS llm_info",
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS query_params JSONB",
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS additional_params JSONB",
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS share_config JSONB",
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS mindmap JSONB",
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS sample_questions JSONB",
"ALTER TABLE IF EXISTS knowledge_bases ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS parent_id VARCHAR(64)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS original_filename VARCHAR(512)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS file_type VARCHAR(64)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS path VARCHAR(1024)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS minio_url VARCHAR(1024)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS markdown_file VARCHAR(1024)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS status VARCHAR(32)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS content_hash VARCHAR(128)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS file_size BIGINT",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS content_type VARCHAR(64)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS processing_params JSONB",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS is_folder BOOLEAN",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS error_message TEXT",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS created_by VARCHAR(64)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS updated_by VARCHAR(64)",
"ALTER TABLE IF EXISTS knowledge_files ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ",
"ALTER TABLE IF EXISTS evaluation_datasets ADD COLUMN IF NOT EXISTS created_by VARCHAR(64)",
"ALTER TABLE IF EXISTS evaluation_datasets ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ",
"ALTER TABLE IF EXISTS evaluation_datasets ADD COLUMN IF NOT EXISTS build_metadata JSONB",
"ALTER TABLE IF EXISTS evaluation_runs ADD COLUMN IF NOT EXISTS metrics JSONB",
"ALTER TABLE IF EXISTS evaluation_runs ADD COLUMN IF NOT EXISTS overall_score DOUBLE PRECISION",
"ALTER TABLE IF EXISTS evaluation_runs ADD COLUMN IF NOT EXISTS total_items INTEGER",
"ALTER TABLE IF EXISTS evaluation_runs ADD COLUMN IF NOT EXISTS completed_items INTEGER",
"ALTER TABLE IF EXISTS evaluation_runs ADD COLUMN IF NOT EXISTS started_at TIMESTAMPTZ",
"ALTER TABLE IF EXISTS evaluation_runs ADD COLUMN IF NOT EXISTS completed_at TIMESTAMPTZ",
"ALTER TABLE IF EXISTS evaluation_runs ADD COLUMN IF NOT EXISTS created_by VARCHAR(64)",
"ALTER TABLE IF EXISTS evaluation_run_items ADD COLUMN IF NOT EXISTS gold_chunk_ids JSONB",
"ALTER TABLE IF EXISTS evaluation_run_items ADD COLUMN IF NOT EXISTS gold_answer TEXT",
"ALTER TABLE IF EXISTS evaluation_run_items ADD COLUMN IF NOT EXISTS generated_answer TEXT",
"ALTER TABLE IF EXISTS evaluation_run_items ADD COLUMN IF NOT EXISTS retrieved_chunks JSONB",
"ALTER TABLE IF EXISTS evaluation_run_items ADD COLUMN IF NOT EXISTS metrics JSONB",
"ALTER TABLE IF EXISTS evaluation_run_items ADD COLUMN IF NOT EXISTS created_at TIMESTAMPTZ",
"""
CREATE TABLE IF NOT EXISTS evaluation_datasets (
id SERIAL PRIMARY KEY,
dataset_id VARCHAR(64) NOT NULL UNIQUE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
name VARCHAR(255) NOT NULL,
description TEXT,
item_count INTEGER DEFAULT 0,
has_gold_chunks BOOLEAN DEFAULT FALSE,
has_gold_answers BOOLEAN DEFAULT FALSE,
build_metadata JSONB,
created_by VARCHAR(64),
created_at TIMESTAMPTZ DEFAULT NOW(),
updated_at TIMESTAMPTZ DEFAULT NOW()
)
""",
"""
CREATE TABLE IF NOT EXISTS evaluation_dataset_items (
id SERIAL PRIMARY KEY,
item_id VARCHAR(64) NOT NULL UNIQUE,
dataset_id VARCHAR(64) NOT NULL REFERENCES evaluation_datasets(dataset_id) ON DELETE CASCADE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
item_index INTEGER NOT NULL,
query_text TEXT NOT NULL,
gold_chunk_ids JSONB,
gold_answer TEXT,
created_at TIMESTAMPTZ DEFAULT NOW(),
CONSTRAINT uq_evaluation_dataset_items_dataset_index UNIQUE (dataset_id, item_index)
)
""",
"""
CREATE TABLE IF NOT EXISTS evaluation_runs (
id SERIAL PRIMARY KEY,
run_id VARCHAR(64) NOT NULL UNIQUE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
dataset_id VARCHAR(64) REFERENCES evaluation_datasets(dataset_id) ON DELETE SET NULL,
status VARCHAR(32) DEFAULT 'running',
retrieval_config JSONB,
metrics JSONB,
overall_score DOUBLE PRECISION,
total_items INTEGER DEFAULT 0,
completed_items INTEGER DEFAULT 0,
started_at TIMESTAMPTZ DEFAULT NOW(),
completed_at TIMESTAMPTZ,
created_by VARCHAR(64)
)
""",
"""
CREATE TABLE IF NOT EXISTS evaluation_run_items (
id SERIAL PRIMARY KEY,
run_id VARCHAR(64) NOT NULL REFERENCES evaluation_runs(run_id) ON DELETE CASCADE,
dataset_item_id VARCHAR(64) REFERENCES evaluation_dataset_items(item_id) ON DELETE SET NULL,
item_index INTEGER NOT NULL,
query_text TEXT NOT NULL,
gold_chunk_ids JSONB,
gold_answer TEXT,
generated_answer TEXT,
retrieved_chunks JSONB,
metrics JSONB,
created_at TIMESTAMPTZ DEFAULT NOW(),
CONSTRAINT uq_evaluation_run_items_run_index UNIQUE (run_id, item_index)
)
""",
"""
CREATE TABLE IF NOT EXISTS knowledge_chunks (
id SERIAL PRIMARY KEY,
chunk_id VARCHAR(128) NOT NULL UNIQUE,
file_id VARCHAR(64) NOT NULL REFERENCES knowledge_files(file_id) ON DELETE CASCADE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
chunk_index INTEGER NOT NULL,
content TEXT NOT NULL,
start_char_pos INTEGER,
end_char_pos INTEGER,
start_token_pos INTEGER,
end_token_pos INTEGER,
graph_indexed BOOLEAN DEFAULT FALSE,
ent_ids JSONB,
tags JSONB,
extraction_result JSONB,
created_at TIMESTAMPTZ DEFAULT NOW(),
updated_at TIMESTAMPTZ DEFAULT NOW()
)
""",
"ALTER TABLE IF EXISTS knowledge_chunks ADD COLUMN IF NOT EXISTS extraction_result JSONB",
"""
CREATE TABLE IF NOT EXISTS knowledge_graph_entities (
id SERIAL PRIMARY KEY,
entity_id VARCHAR(64) NOT NULL UNIQUE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
normalized_name VARCHAR(512) NOT NULL,
label VARCHAR(128) NOT NULL,
name VARCHAR(512) NOT NULL,
attributes JSONB,
created_at TIMESTAMPTZ DEFAULT NOW(),
updated_at TIMESTAMPTZ DEFAULT NOW(),
CONSTRAINT uq_knowledge_graph_entities_identity UNIQUE (db_id, normalized_name, label)
)
""",
"""
CREATE TABLE IF NOT EXISTS knowledge_graph_entity_mentions (
id SERIAL PRIMARY KEY,
entity_id VARCHAR(64) NOT NULL REFERENCES knowledge_graph_entities(entity_id) ON DELETE CASCADE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
file_id VARCHAR(64) NOT NULL REFERENCES knowledge_files(file_id) ON DELETE CASCADE,
chunk_id VARCHAR(128) NOT NULL REFERENCES knowledge_chunks(chunk_id) ON DELETE CASCADE,
created_at TIMESTAMPTZ DEFAULT NOW(),
CONSTRAINT uq_knowledge_graph_entity_mentions_entity_chunk UNIQUE (entity_id, chunk_id)
)
""",
"""
CREATE TABLE IF NOT EXISTS knowledge_graph_triples (
id SERIAL PRIMARY KEY,
triple_id VARCHAR(64) NOT NULL UNIQUE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
source_entity_id VARCHAR(64) NOT NULL REFERENCES knowledge_graph_entities(entity_id) ON DELETE CASCADE,
target_entity_id VARCHAR(64) NOT NULL REFERENCES knowledge_graph_entities(entity_id) ON DELETE CASCADE,
relation_type VARCHAR(256) NOT NULL,
content TEXT NOT NULL,
created_at TIMESTAMPTZ DEFAULT NOW(),
updated_at TIMESTAMPTZ DEFAULT NOW()
)
""",
"""
CREATE TABLE IF NOT EXISTS knowledge_graph_triple_mentions (
id SERIAL PRIMARY KEY,
triple_id VARCHAR(64) NOT NULL REFERENCES knowledge_graph_triples(triple_id) ON DELETE CASCADE,
db_id VARCHAR(80) NOT NULL REFERENCES knowledge_bases(db_id) ON DELETE CASCADE,
file_id VARCHAR(64) NOT NULL REFERENCES knowledge_files(file_id) ON DELETE CASCADE,
chunk_id VARCHAR(128) NOT NULL REFERENCES knowledge_chunks(chunk_id) ON DELETE CASCADE,
text TEXT,
extractor_type VARCHAR(128),
created_at TIMESTAMPTZ DEFAULT NOW(),
CONSTRAINT uq_knowledge_graph_triple_mentions_triple_chunk UNIQUE (triple_id, chunk_id)
)
""",
# 扩展 db_id 字段长度以支持最长 75 字符的 IDkb_private_ + 64字符hash
"ALTER TABLE IF EXISTS knowledge_bases ALTER COLUMN db_id TYPE VARCHAR(80)",
"ALTER TABLE IF EXISTS knowledge_files ALTER COLUMN db_id TYPE VARCHAR(80)",
"ALTER TABLE IF EXISTS evaluation_datasets ALTER COLUMN db_id TYPE VARCHAR(80)",
"ALTER TABLE IF EXISTS evaluation_dataset_items ALTER COLUMN db_id TYPE VARCHAR(80)",
"ALTER TABLE IF EXISTS evaluation_runs ALTER COLUMN db_id TYPE VARCHAR(80)",
"CREATE INDEX IF NOT EXISTS idx_kb_type ON knowledge_bases(kb_type)",
"CREATE INDEX IF NOT EXISTS idx_kb_name ON knowledge_bases(name)",
"CREATE INDEX IF NOT EXISTS idx_kf_db_id ON knowledge_files(db_id)",
"CREATE INDEX IF NOT EXISTS idx_kf_parent ON knowledge_files(parent_id)",
"CREATE INDEX IF NOT EXISTS idx_kf_status ON knowledge_files(status)",
"CREATE INDEX IF NOT EXISTS idx_kf_hash ON knowledge_files(content_hash)",
"CREATE INDEX IF NOT EXISTS ix_evaluation_datasets_db_id ON evaluation_datasets(db_id)",
(
"CREATE INDEX IF NOT EXISTS ix_evaluation_dataset_items_dataset_index "
"ON evaluation_dataset_items(dataset_id, item_index)"
),
"CREATE INDEX IF NOT EXISTS ix_evaluation_dataset_items_db_id ON evaluation_dataset_items(db_id)",
"CREATE INDEX IF NOT EXISTS ix_evaluation_runs_db_id ON evaluation_runs(db_id)",
"CREATE INDEX IF NOT EXISTS ix_evaluation_runs_status ON evaluation_runs(status)",
"CREATE INDEX IF NOT EXISTS ix_evaluation_runs_started ON evaluation_runs(started_at DESC)",
"CREATE INDEX IF NOT EXISTS ix_evaluation_run_items_run_index ON evaluation_run_items(run_id, item_index)",
"CREATE UNIQUE INDEX IF NOT EXISTS uq_knowledge_chunks_chunk_id ON knowledge_chunks(chunk_id)",
"CREATE INDEX IF NOT EXISTS ix_knowledge_chunks_file_id ON knowledge_chunks(file_id)",
"CREATE INDEX IF NOT EXISTS ix_knowledge_chunks_db_id ON knowledge_chunks(db_id)",
"CREATE INDEX IF NOT EXISTS ix_knowledge_chunks_graph_indexed ON knowledge_chunks(graph_indexed)",
(
"CREATE UNIQUE INDEX IF NOT EXISTS uq_knowledge_graph_entities_entity_id "
"ON knowledge_graph_entities(entity_id)"
),
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_entities_db_id ON knowledge_graph_entities(db_id)",
(
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_entity_mentions_db_id "
"ON knowledge_graph_entity_mentions(db_id)"
),
(
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_entity_mentions_file_id "
"ON knowledge_graph_entity_mentions(file_id)"
),
(
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_entity_mentions_chunk_id "
"ON knowledge_graph_entity_mentions(chunk_id)"
),
(
"CREATE UNIQUE INDEX IF NOT EXISTS uq_knowledge_graph_triples_triple_id "
"ON knowledge_graph_triples(triple_id)"
),
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_triples_db_id ON knowledge_graph_triples(db_id)",
(
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_triple_mentions_db_id "
"ON knowledge_graph_triple_mentions(db_id)"
),
(
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_triple_mentions_file_id "
"ON knowledge_graph_triple_mentions(file_id)"
),
(
"CREATE INDEX IF NOT EXISTS ix_knowledge_graph_triple_mentions_chunk_id "
"ON knowledge_graph_triple_mentions(chunk_id)"
),
]
async with self.async_engine.begin() as conn:
for stmt in stmts:
await conn.execute(text(stmt))
async def ensure_business_schema(self):
"""确保业务 schema 包含后续新增字段(运行时 schema 演进)。"""
self._check_initialized()
stmts = [
"ALTER TABLE IF EXISTS skills ADD COLUMN IF NOT EXISTS tool_dependencies JSONB DEFAULT '[]'::jsonb",
"ALTER TABLE IF EXISTS skills ADD COLUMN IF NOT EXISTS mcp_dependencies JSONB DEFAULT '[]'::jsonb",
"ALTER TABLE IF EXISTS skills ADD COLUMN IF NOT EXISTS skill_dependencies JSONB DEFAULT '[]'::jsonb",
"ALTER TABLE IF EXISTS skills ADD COLUMN IF NOT EXISTS version VARCHAR(64)",
"ALTER TABLE IF EXISTS skills ADD COLUMN IF NOT EXISTS is_builtin BOOLEAN NOT NULL DEFAULT FALSE",
"ALTER TABLE IF EXISTS skills ADD COLUMN IF NOT EXISTS content_hash VARCHAR(128)",
"ALTER TABLE IF EXISTS subagents ADD COLUMN IF NOT EXISTS enabled BOOLEAN NOT NULL DEFAULT TRUE",
"ALTER TABLE IF EXISTS conversations ADD COLUMN IF NOT EXISTS is_pinned BOOLEAN NOT NULL DEFAULT FALSE",
"ALTER TABLE IF EXISTS mcp_servers ADD COLUMN IF NOT EXISTS env JSONB",
"""
ALTER TABLE IF EXISTS agent_configs
ADD COLUMN IF NOT EXISTS uid VARCHAR
""",
"""
UPDATE agent_configs ac
SET uid = u.uid
FROM users u
WHERE ac.uid IS NULL
AND ac.created_by ~ '^[0-9]+$'
AND u.id = ac.created_by::integer
""",
"""
UPDATE agent_configs ac
SET uid = u.uid
FROM users u
WHERE ac.uid IS NULL
AND ac.created_by = u.uid
""",
"""
DO $$
BEGIN
IF EXISTS (
SELECT 1 FROM information_schema.columns
WHERE table_name = 'agent_configs' AND column_name = 'department_id'
) THEN
EXECUTE $sql$
UPDATE agent_configs ac
SET uid = (
SELECT u.uid
FROM users u
WHERE u.department_id = ac.department_id
AND u.is_deleted = 0
ORDER BY
CASE WHEN u.role = 'superadmin' THEN 0 WHEN u.role = 'admin' THEN 1 ELSE 2 END,
u.id ASC
LIMIT 1
)
WHERE ac.uid IS NULL
$sql$;
END IF;
END $$;
""",
"DELETE FROM agent_configs WHERE uid IS NULL",
"DROP INDEX IF EXISTS uq_agent_configs_department_agent_default",
"DROP INDEX IF EXISTS ix_agent_configs_department_id",
"ALTER TABLE IF EXISTS agent_configs DROP CONSTRAINT IF EXISTS uq_agent_configs_department_agent_name",
"ALTER TABLE IF EXISTS agent_configs DROP CONSTRAINT IF EXISTS agent_configs_department_id_fkey",
"ALTER TABLE IF EXISTS agent_configs DROP COLUMN IF EXISTS department_id",
"ALTER TABLE IF EXISTS agent_configs ALTER COLUMN uid SET NOT NULL",
"""
WITH ranked AS (
SELECT id, ROW_NUMBER() OVER (PARTITION BY uid, agent_id, name ORDER BY id) AS rn
FROM agent_configs
)
UPDATE agent_configs ac
SET name = LEFT(ac.name, 90) || '-' || ac.id::text
FROM ranked
WHERE ac.id = ranked.id AND ranked.rn > 1
""",
"""
WITH ranked AS (
SELECT id, ROW_NUMBER() OVER (PARTITION BY uid, agent_id ORDER BY id) AS rn
FROM agent_configs
WHERE is_default IS TRUE
)
UPDATE agent_configs ac
SET is_default = FALSE
FROM ranked
WHERE ac.id = ranked.id AND ranked.rn > 1
""",
"""
DO $$
BEGIN
IF NOT EXISTS (
SELECT 1 FROM pg_constraint WHERE conname = 'agent_configs_uid_fkey'
) THEN
ALTER TABLE agent_configs
ADD CONSTRAINT agent_configs_uid_fkey FOREIGN KEY (uid) REFERENCES users(uid);
END IF;
END $$;
""",
"""
DO $$
BEGIN
IF NOT EXISTS (
SELECT 1 FROM pg_constraint WHERE conname = 'uq_agent_configs_uid_agent_name'
) THEN
ALTER TABLE agent_configs
ADD CONSTRAINT uq_agent_configs_uid_agent_name UNIQUE (uid, agent_id, name);
END IF;
END $$;
""",
"CREATE INDEX IF NOT EXISTS ix_agent_configs_uid ON agent_configs(uid)",
"""
CREATE UNIQUE INDEX IF NOT EXISTS uq_agent_configs_uid_agent_default
ON agent_configs(uid, agent_id)
WHERE is_default IS TRUE
""",
"""
CREATE TABLE IF NOT EXISTS model_providers (
id SERIAL PRIMARY KEY,
provider_id VARCHAR(100) NOT NULL UNIQUE,
display_name VARCHAR(100) NOT NULL,
provider_type VARCHAR(32) NOT NULL DEFAULT 'openai',
default_protocol VARCHAR(64),
base_url VARCHAR(500) NOT NULL,
embedding_base_url VARCHAR(500),
rerank_base_url VARCHAR(500),
models_endpoint VARCHAR(200),
embedding_models_endpoint VARCHAR(200),
rerank_models_endpoint VARCHAR(200),
api_key_env VARCHAR(128),
api_key VARCHAR(500),
capabilities JSONB NOT NULL DEFAULT '[]'::jsonb,
enabled_models JSONB NOT NULL DEFAULT '[]'::jsonb,
headers_json JSONB,
extra_json JSONB,
is_enabled BOOLEAN NOT NULL DEFAULT TRUE,
is_builtin BOOLEAN NOT NULL DEFAULT FALSE,
created_by VARCHAR(100),
updated_by VARCHAR(100),
created_at TIMESTAMPTZ DEFAULT NOW(),
updated_at TIMESTAMPTZ DEFAULT NOW()
)
""",
"""
CREATE TABLE IF NOT EXISTS agent_runs (
id VARCHAR(64) PRIMARY KEY,
thread_id VARCHAR(64) NOT NULL,
agent_id VARCHAR(64) NOT NULL,
uid VARCHAR(64) NOT NULL,
status VARCHAR(32) NOT NULL DEFAULT 'pending',
request_id VARCHAR(64) NOT NULL UNIQUE,
input_payload JSONB NOT NULL DEFAULT '{}'::jsonb,
error_type VARCHAR(64),
error_message TEXT,
started_at TIMESTAMPTZ,
finished_at TIMESTAMPTZ,
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
)
""",
"CREATE INDEX IF NOT EXISTS idx_agent_runs_uid_created ON agent_runs(uid, created_at DESC)",
"CREATE INDEX IF NOT EXISTS idx_agent_runs_thread_created ON agent_runs(thread_id, created_at DESC)",
"CREATE INDEX IF NOT EXISTS idx_agent_runs_status_updated ON agent_runs(status, updated_at)",
"CREATE INDEX IF NOT EXISTS ix_conversations_is_pinned ON conversations(is_pinned)",
"CREATE UNIQUE INDEX IF NOT EXISTS ix_model_providers_provider_id ON model_providers(provider_id)",
"CREATE INDEX IF NOT EXISTS ix_model_providers_is_enabled ON model_providers(is_enabled)",
]
async with self.async_engine.begin() as conn:
for stmt in stmts:
await conn.execute(text(stmt))
@property
def is_postgresql(self) -> bool:
"""检查是否是 PostgreSQL 数据库"""
if not self._initialized:
return False
return self.async_engine.dialect.name == "postgresql"
async def get_async_session(self) -> AsyncSession:
"""获取异步数据库会话"""
self.initialize() # 确保已初始化
return self.AsyncSession()
@asynccontextmanager
async def get_async_session_context(self):
"""获取异步数据库会话的上下文管理器"""
self.initialize() # 确保已初始化
session = self.AsyncSession()
try:
yield session
await session.commit()
except Exception as e:
await session.rollback()
logger.error(f"PostgreSQL async operation failed: {e}")
raise
finally:
await session.close()
async def close(self):
"""关闭引擎"""
if self.async_engine:
await self.async_engine.dispose()
if self.langgraph_pool:
await self.langgraph_pool.close()
async def async_check_first_run(self):
"""检查是否首次运行(异步版本)- 检查用户表是否有数据"""
from sqlalchemy import func, select
self._check_initialized()
async with self.get_async_session_context() as session:
from yuxi.storage.postgres.models_business import User
result = await session.execute(select(func.count(User.id)))
count = result.scalar()
return count == 0
async def commit(self):
"""提交当前会话"""
self._check_initialized()
async with self.get_async_session_context():
pass # commit is automatic in context manager
# 创建全局 PostgreSQL 管理器实例
pg_manager = PostgresManager()