feat(avatar): honor square interaction authorization

This commit is contained in:
stefanfeng
2026-09-08 13:19:50 +08:00
parent 2a01a9946a
commit 434caac056
4 changed files with 432 additions and 18 deletions
+117 -16
View File
@@ -23,6 +23,7 @@ class SchedulerService:
from app.core.database import AsyncSessionLocal
logger.info("⚡ 立即触发互动任务")
async with AsyncSessionLocal() as session:
await self._sync_delegated_avatar_users(session)
try:
max_concurrent = int(await self._get_config(session, "max_concurrent_users", "5"))
except (TypeError, ValueError):
@@ -146,7 +147,9 @@ class SchedulerService:
async def _check_sessions(self):
"""定时校验登录状态"""
from app.services.news_service import news_service
from app.services.avatar_service import is_delegated_avatar_user
async with AsyncSessionLocal() as db:
await self._sync_delegated_avatar_users(db)
result = await db.execute(
select(VirtualUser).where(VirtualUser.status == 2, VirtualUser.is_enabled == 1)
)
@@ -154,7 +157,7 @@ class SchedulerService:
for user in users:
try:
valid = await news_service.check_session(db, user)
if not valid:
if not valid and not is_delegated_avatar_user(user):
logger.warning(f"用户 {user.account} 会话失效,尝试重登")
await news_service.login(db, user)
except Exception as e:
@@ -163,6 +166,7 @@ class SchedulerService:
async def _run_interactions(self):
"""执行互动任务"""
async with AsyncSessionLocal() as db:
await self._sync_delegated_avatar_users(db)
# 检查调度器开关
enabled = await self._get_config(db, "scheduler_enabled", "true")
if enabled != "true":
@@ -184,8 +188,11 @@ class SchedulerService:
logger.debug(f"[调度] 当前北京时间 {now_time} 不在互动时段 {start_str}-{end_str}")
return
# 获取最小互动间隔(秒)
min_interval = int(await self._get_config(db, "interact_min_interval", "300"))
# 获取互动间隔范围(秒),与调度设置页面字段保持一致
min_interval = await self._get_int_config(db, "interact_interval_min", 300)
max_interval = await self._get_int_config(db, "interact_interval_max", min_interval)
min_interval = max(0, min_interval)
max_interval = max(min_interval, max_interval)
# 获取最大并发
max_concurrent = int(await self._get_config(db, "max_concurrent_users", "5"))
@@ -204,7 +211,7 @@ class SchedulerService:
await self._try_login_users(db)
return
# 检查互动间隔:过滤掉最近 min_interval 秒内已互动的用户
# 每个用户在其最小/最大间隔内取得稳定随机值,直到下次互动后再变化
now_dt = datetime.now()
eligible = []
for u in all_users:
@@ -212,11 +219,17 @@ class SchedulerService:
eligible.append(u)
else:
elapsed = (now_dt - u.last_interact_at).total_seconds()
if elapsed >= min_interval:
interval = random.Random(
f"{u.id}:{u.last_interact_at.isoformat()}"
).randint(min_interval, max_interval)
if elapsed >= interval:
eligible.append(u)
if not eligible:
logger.debug(f"[调度] 所有 {len(all_users)} 个用户在 {min_interval}s 内已互动,跳过本次")
logger.debug(
f"[调度] 所有 {len(all_users)} 个用户尚未达到 "
f"{min_interval}-{max_interval}s 随机互动间隔,跳过本次"
)
return
# 按最后互动时间升序排序:最久没互动的用户优先
@@ -257,10 +270,12 @@ class SchedulerService:
async def _try_login_users(self, db):
"""尝试登录未登录的用户"""
from app.services.news_service import news_service
from app.services.avatar_service import AVATAR_ACCOUNT_PREFIX
result = await db.execute(
select(VirtualUser).where(
VirtualUser.status.in_([0, 3]),
VirtualUser.is_enabled == 1
VirtualUser.is_enabled == 1,
~VirtualUser.account.like(f"{AVATAR_ACCOUNT_PREFIX}%"),
).limit(3)
)
users = result.scalars().all()
@@ -275,6 +290,11 @@ class SchedulerService:
"""执行单用户互动 - 基于真实接口"""
from app.services.news_service import news_service
from app.services.ai_service import ai_service
from app.services.avatar_service import (
delegated_avatar_id,
get_square_interaction_permissions,
is_delegated_avatar_user,
)
async with AsyncSessionLocal() as db:
try:
@@ -289,6 +309,23 @@ class SchedulerService:
"interactions": [],
}
allowed_actions = {"like", "collect", "comment", "reply", "forward"}
if is_delegated_avatar_user(user):
allowed_actions = set(
get_square_interaction_permissions(delegated_avatar_id(user))
)
if not allowed_actions:
user.status = 0
user.is_enabled = 0
await db.commit()
return {
"user_id": user.id,
"account": user.account,
"status": "skipped",
"reason": "avatar_interaction_not_authorized",
"interactions": [],
}
# 检查今日评论限额
can_comment = True
if user.today_comment_count >= user.daily_comment_limit:
@@ -398,14 +435,53 @@ class SchedulerService:
interactions_done = []
action_failures = []
# ① 先记录阅读(每次必做,模拟真实用户打开文章)
done_on_this = today_done.get(news_id, set())
wants = {
"like": (
"like" in allowed_actions
and "like" not in done_on_this
and random.random() < like_prob
),
"collect": (
"collect" in allowed_actions
and "collect" not in done_on_this
and random.random() < collect_prob
),
"forward": (
"forward" in allowed_actions
and "forward" not in done_on_this
and random.random() < forward_prob
),
"reply": (
"reply" in allowed_actions
and can_comment
and personality is not None
and random.random() < reply_prob
),
"comment": (
"comment" in allowed_actions
and can_comment
and personality is not None
and not already_commented_this
and random.random() < comment_prob
),
}
if not any(wants.values()):
return {
"user_id": user.id,
"account": user.account,
"status": "skipped",
"reason": "no_actions_triggered",
"interactions": [],
"article_id": news_id,
"article_title": news_title,
}
# 只有动作命中调度概率后才打开文章
await news_service.read_news(db, user, news_id)
# 今日已对此文章做过的互动类型
done_on_this = today_done.get(news_id, set())
# ② 点赞(每篇文章每用户每天只点赞一次)
if "like" not in done_on_this and random.random() < like_prob:
if wants["like"]:
success, err = await news_service.like_news(db, user, news_id, org_id=article_org_id, to_user_id=news_author, title=news_title)
await self._save_record(db, user, news_id, news_title, "like", None, 0, success, err)
if success:
@@ -415,16 +491,17 @@ class SchedulerService:
action_failures.append({"type": "like", "error": err})
# ③ 收藏(每篇文章每用户每天只收藏一次)
if "collect" not in done_on_this and random.random() < collect_prob:
if wants["collect"]:
success, err = await news_service.collect_news(db, user, news_id, org_id=article_org_id, to_user_id=news_author, title=news_title)
await self._save_record(db, user, news_id, news_title, "collect", None, 0, success, err)
if success:
interactions_done.append("collect")
await self._incr_total(db, user_id)
else:
action_failures.append({"type": "collect", "error": err})
# ④ 转发(每篇文章每用户每天只转发一次)
if "forward" not in done_on_this and random.random() < forward_prob:
if wants["forward"]:
success, err = await news_service.forward_news(db, user, news_id)
await self._save_record(db, user, news_id, news_title, "forward", None, 0, success, err)
if success:
@@ -438,7 +515,7 @@ class SchedulerService:
style_prompt = personality.comment_style_prompt or ""
safe_word_max = min(personality.word_count_max, 80)
if random.random() < reply_prob:
if wants["reply"]:
reply_actions, reply_failures = await self._run_reply_interaction_chain(
db=db,
starter=user,
@@ -455,7 +532,7 @@ class SchedulerService:
action_failures.extend(reply_failures)
# 每篇文章每个用户每天只发一条顶层评论;回复不再要求先评论
if not already_commented_this and random.random() < comment_prob:
if wants["comment"]:
comment_text, tokens = await ai_service.generate_comment(
db, news_title, news_content,
style_prompt, personality.word_count_min, safe_word_max
@@ -679,6 +756,7 @@ class SchedulerService:
async with AsyncSessionLocal() as db:
try:
await self._sync_delegated_avatar_users(db)
now = datetime.now()
await db.execute(
update(PendingReplyTask)
@@ -706,6 +784,12 @@ class SchedulerService:
logger.error(f"待发送回复队列处理异常: {e}")
async def _process_pending_reply_task(self, db, task: PendingReplyTask, news_service, ai_service):
from app.services.avatar_service import (
delegated_avatar_id,
get_square_interaction_permissions,
is_delegated_avatar_user,
)
task.status = 1
task.locked_at = datetime.now()
task.attempts = (task.attempts or 0) + 1
@@ -716,6 +800,13 @@ class SchedulerService:
task.status = 3
task.last_error = "用户未登录或已禁用"
return
if (
is_delegated_avatar_user(actor)
and "reply" not in get_square_interaction_permissions(delegated_avatar_id(actor))
):
task.status = 3
task.last_error = "数字分身广场互动授权已撤销"
return
reply_result = await self._post_contextual_reply(
db=db,
@@ -858,6 +949,16 @@ class SchedulerService:
except (TypeError, ValueError):
return default
async def _sync_delegated_avatar_users(self, db):
from app.services.avatar_service import sync_square_interaction_users
try:
return await sync_square_interaction_users(db)
except Exception as exc:
await db.rollback()
logger.error(f"数字分身广场互动身份同步异常: {exc}")
return set()
async def _incr_total(self, db, user_id: int):
await db.execute(
update(VirtualUser).where(VirtualUser.id == user_id).values(