292 lines
10 KiB
Python
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()
|