Files
ai-talk-callback/app/database.py
2025-12-04 08:49:34 +08:00

95 lines
3.8 KiB
Python

from datetime import datetime
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession, async_sessionmaker
from sqlalchemy.orm import DeclarativeBase
from sqlalchemy import Column, String, Integer, DateTime, Text, JSON, Index, text
from app.config import settings
from app.logger import get_logger
from sqlalchemy.sql.expression import func
logger = get_logger("database")
class Base(DeclarativeBase):
pass
class CallbackFailureLog(Base):
__tablename__ = "callback_failure_logs"
id = Column(Integer, primary_key=True, autoincrement=True, comment="日志ID")
site_id = Column(String(100), nullable=False, comment="站点ID") # 新增siteId字段
remote_address = Column(String(45), nullable=True, comment="客户端IP地址")
server_ip = Column(String(45), nullable=True, comment="服务器IP地址")
request_url = Column(String(500), nullable=False, comment="请求URL")
request_headers = Column(JSON, nullable=False, comment="请求头")
request_body = Column(JSON, nullable=False, comment="请求体")
created_at = Column(DateTime, default=datetime.now(), comment="创建时间")
class CallbackFailureData(Base):
__tablename__ = "callback_failure_data"
id = Column(Integer, primary_key=True, autoincrement=True, comment="数据ID")
callback_failure_log_id = Column(Integer, nullable=False, comment="回调失败日志ID")
phone_number = Column(String(20), nullable=False, comment="手机号")
task_id = Column(String(100), nullable=False, comment="任务ID")
status = Column(Integer, nullable=False, comment="状态")
status_description = Column(String(200), nullable=False, comment="状态描述")
raw_data = Column(Text, nullable=False, comment="原始回调数据JSON字符串")
calldate = Column(DateTime, nullable=True, comment="通话日期")
created_at = Column(DateTime, default=datetime.now(), comment="创建时间")
class ExternalApiLog(Base):
__tablename__ = "external_api_logs" # 重命名表,更通用
id = Column(Integer, primary_key=True, autoincrement=True, comment="日志ID")
callback_failure_log_id = Column(Integer, nullable=False, comment="回调失败日志ID")
request_url = Column(String(500), nullable=False, comment="外部接口请求URL")
request_headers = Column(JSON, nullable=False, comment="外部接口请求头")
request_body = Column(JSON, nullable=False, comment="外部接口请求体")
response_status = Column(Integer, nullable=False, comment="外部接口响应状态码")
response_headers = Column(JSON, nullable=False, comment="外部接口响应头")
response_body = Column(Text, nullable=False, comment="外部接口响应体")
retry_count = Column(Integer, default=0, comment="重试次数")
created_at = Column(DateTime, default=datetime.now(), comment="创建时间")
# 创建数据库引擎
engine = create_async_engine(settings.database_url, echo=settings.debug, future=True)
# 创建会话工厂
AsyncSessionLocal = async_sessionmaker(
engine, class_=AsyncSession, expire_on_commit=False
)
async def get_db():
async with AsyncSessionLocal() as session:
try:
yield session
finally:
await session.close()
async def init_db():
logger.info("📊 初始化数据库表结构...")
try:
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
# 检查并创建手机号索引
await conn.execute(
text(
"""
CREATE INDEX IF NOT EXISTS idx_phone_number
ON callback_failure_data (phone_number)
"""
)
)
logger.info("✅ 数据库表结构初始化完成")
except Exception as e:
logger.error(f"❌ 数据库初始化失败: {e}")
raise