Files
huihuiSquare/digital-avatar-app/backend/main.py
T

292 lines
10 KiB
Python

from fastapi import FastAPI
from fastapi.middleware.cors import CORSMiddleware
import importlib.util
import logging
import os
from apscheduler.schedulers.asyncio import AsyncIOScheduler
from apscheduler.triggers.interval import IntervalTrigger
from database import engine, init_db, SessionLocal
from models import Avatar, Authorization, Organization, TokenAccount, TokenPlan, User
from fastapi.staticfiles import StaticFiles
import routers.avatars
import routers.tokens
import routers.authorizations
import routers.organizations
import routers.knowledge
import routers.huihui_auth
import routers.chat
import routers.takeover
from responses import ok
from services.chat_attachment_service import purge_expired_chat_attachments
from services.knowledge_vectorizer import knowledge_vectorizer
from services.token_billing import DEFAULT_TOKEN_GRANT, release_stale_reservations
logger = logging.getLogger(__name__)
takeover_scheduler = None
maintenance_scheduler = None
app = FastAPI(title="会会数字分身 API", version="1.0.0")
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=False,
allow_methods=["*"],
allow_headers=["*"],
)
app.include_router(routers.avatars.router, prefix="/api")
app.include_router(routers.tokens.router, prefix="/api")
app.include_router(routers.authorizations.router, prefix="/api")
app.include_router(routers.organizations.router, prefix="/api")
app.include_router(routers.knowledge.router, prefix="/api")
app.include_router(routers.huihui_auth.router, prefix="/api")
app.include_router(routers.chat.router, prefix="/api")
app.include_router(routers.takeover.router, prefix="/api")
UPLOAD_DIR = routers.knowledge.UPLOAD_DIR
os.makedirs(UPLOAD_DIR, exist_ok=True)
app.mount("/api/files", StaticFiles(directory=UPLOAD_DIR), name="knowledge-files")
@app.get("/api/health")
def health():
checks = _runtime_checks()
return ok({
"status": "ok" if all(checks.values()) else "degraded",
"gitSha": os.getenv("APP_GIT_SHA", "unknown"),
"buildTime": os.getenv("APP_BUILD_TIME", "unknown"),
"checks": checks,
})
def _runtime_checks():
return {
"database": _database_is_ready(),
"uploads": os.path.isdir(UPLOAD_DIR) and os.access(UPLOAD_DIR, os.W_OK),
"pdfOcr": importlib.util.find_spec("pymupdf") is not None,
}
def _database_is_ready():
try:
with engine.connect() as connection:
connection.exec_driver_sql("SELECT 1")
return True
except Exception:
logger.exception("Database readiness check failed")
return False
def seed():
db = SessionLocal()
try:
plan_specs = [
{"id": "1", "name": "基础套餐", "amount": 2_000_000, "price": 10, "badge": "", "desc": "2M 积分"},
{"id": "2", "name": "标准套餐", "amount": 20_000_000, "price": 100, "badge": "常用", "desc": "20M 积分"},
{"id": "3", "name": "专业套餐", "amount": 250_000_000, "price": 1000, "badge": "加赠25%", "desc": "250M 积分"},
{"id": "4", "name": "企业套餐", "amount": 2_500_000_000, "price": 10000, "badge": "企业推荐", "desc": "2500M 积分"},
]
for spec in plan_specs:
plan = db.query(TokenPlan).filter(TokenPlan.id == spec["id"]).first()
if plan is None:
db.add(TokenPlan(**spec))
else:
for key, value in spec.items():
setattr(plan, key, value)
for user in db.query(User).all():
account = db.query(TokenAccount).filter(TokenAccount.user_id == user.id).first()
if account is None:
db.add(TokenAccount(
user_id=user.id,
balance=DEFAULT_TOKEN_GRANT,
total_granted=DEFAULT_TOKEN_GRANT,
total_consumed=0,
))
if db.query(Avatar).count() == 0:
avatar = Avatar(
name="我的数字分身",
display_name="会会助手",
description="我是您的AI数字分身,可以帮您管理日程、回复消息、处理任务。",
emoji="🤖",
status="active",
token_balance=0,
config={
"replyStyle": "professional",
"creativity": 50,
"rigor": 50,
"humor": 30,
"responseLength": "medium",
"systemPrompt": "",
},
)
db.add(avatar)
db.commit()
db.refresh(avatar)
if db.query(Authorization).count() == 0:
auths = [
Authorization(avatar_id=avatar.id, target_type="application", target_name="微信小程序", permissions=["read", "reply"], status="active"),
Authorization(avatar_id=avatar.id, target_type="user", target_name="张三", permissions=["read"], status="active"),
Authorization(avatar_id=avatar.id, target_type="organization", target_name="产品团队", permissions=["read", "edit"], status="inactive"),
]
db.add_all(auths)
if db.query(Organization).count() == 0:
orgs = [
Organization(name="会会增长团队", description="负责会会产品的增长与运营", emoji="🚀", org_type="team", member_count=12),
Organization(name="AI 实验室", description="探索前沿 AI 能力", emoji="💡", org_type="company", member_count=8),
]
db.add_all(orgs)
db.commit()
release_stale_reservations(db)
finally:
db.close()
@app.on_event("startup")
def on_startup():
global takeover_scheduler
init_db()
seed()
knowledge_vectorizer.start()
# Release stale resources when startup is invoked again by a reload/test.
stop_takeover_scheduler()
stop_maintenance_scheduler()
try:
start_maintenance_scheduler()
except Exception as exc:
stop_maintenance_scheduler()
logger.warning(
"Failed to initialize chat attachment cleanup, app will continue: %s",
exc,
)
# --- Takeover scheduler ---
try:
# BOXIM production endpoints are intentionally separate from the login API.
from services.boxim_client import BoxIMClient
boxim_config = {
"HUIHUI_PLATFORM_BASE_URL": os.getenv(
"HUIHUI_PLATFORM_BASE_URL", "https://open.99hui.com/api"
),
"BOXIM_API_BASE_URL": os.getenv(
"BOXIM_API_BASE_URL", "https://im.99hui.com/api"
),
"HUIHUI_APP_ID": os.getenv("HUIHUI_APP_ID", ""),
"HUIHUI_ACCESS_ID": os.getenv("HUIHUI_ACCESS_ID", ""),
"HUIHUI_ACCESS_SECRET": os.getenv("HUIHUI_ACCESS_SECRET", ""),
"BOXIM_TIMEOUT_SECONDS": os.getenv("BOXIM_TIMEOUT_SECONDS", "20"),
}
boxim_client = BoxIMClient(boxim_config)
from services.takeover_service import TakeoverService
takeover_service = TakeoverService(
SessionLocal,
boxim_client,
poll_concurrency=int(os.getenv("BOXIM_POLL_CONCURRENCY", "8")),
max_message_age_seconds=int(
os.getenv("BOXIM_MAX_MESSAGE_AGE_SECONDS", "600")
),
)
poll_interval = max(0.5, float(os.getenv("BOXIM_POLL_INTERVAL_SECONDS", "1")))
takeover_scheduler = AsyncIOScheduler()
takeover_scheduler.add_job(
takeover_service.poll_messages,
trigger=IntervalTrigger(seconds=poll_interval),
id="takeover_message_poll",
max_instances=1,
coalesce=True,
)
process_interval = max(
0.25, float(os.getenv("TAKEOVER_PROCESS_INTERVAL_SECONDS", "0.5"))
)
takeover_scheduler.add_job(
takeover_service.process_reply_tasks,
trigger=IntervalTrigger(seconds=process_interval),
id="takeover_reply_process",
max_instances=1,
coalesce=True,
)
takeover_scheduler.start()
logger.info(
"BOXIM takeover scheduler started (poll=%ss, process=%ss)",
poll_interval,
process_interval,
)
except Exception as e:
stop_takeover_scheduler()
logger.warning(f"Failed to initialize takeover scheduler, app will continue without it: {e}")
def stop_takeover_scheduler():
global takeover_scheduler
if takeover_scheduler is not None:
try:
if takeover_scheduler.running:
takeover_scheduler.shutdown(wait=False)
except Exception as e:
logger.warning(f"Failed to stop takeover scheduler cleanly: {e}")
finally:
takeover_scheduler = None
def purge_expired_chat_attachments_job():
db = SessionLocal()
try:
count = purge_expired_chat_attachments(db)
if count:
logger.info("Purged %s expired chat image attachment(s)", count)
except Exception as exc:
db.rollback()
logger.warning("Failed to purge expired chat image attachments: %s", exc)
finally:
db.close()
def start_maintenance_scheduler():
global maintenance_scheduler
purge_expired_chat_attachments_job()
interval_minutes = max(
5, min(1440, int(os.getenv("CHAT_ATTACHMENT_CLEANUP_MINUTES", "60")))
)
maintenance_scheduler = AsyncIOScheduler()
maintenance_scheduler.add_job(
purge_expired_chat_attachments_job,
trigger=IntervalTrigger(minutes=interval_minutes),
id="chat_attachment_cleanup",
max_instances=1,
coalesce=True,
)
maintenance_scheduler.start()
def stop_maintenance_scheduler():
global maintenance_scheduler
if maintenance_scheduler is not None:
try:
if maintenance_scheduler.running:
maintenance_scheduler.shutdown(wait=False)
except Exception as exc:
logger.warning("Failed to stop maintenance scheduler cleanly: %s", exc)
finally:
maintenance_scheduler = None
@app.on_event("shutdown")
def on_shutdown():
stop_takeover_scheduler()
stop_maintenance_scheduler()