Files
video-gen/video-gen-api/app/tasks/pre_test_result_task.py
2026-06-26 11:00:56 +08:00

124 lines
4.4 KiB
Python

from datetime import datetime, timezone
import asyncio
import json
from sqlalchemy import select, update
from sqlalchemy.ext.asyncio import AsyncSession
from app.models.base import async_session
from app.models.resources_material import ResourcesMaterial
from app.models.user_oauth import UserOAuth
from app.utils.douyinApi import DouyinApi
from app.utils.logger import get_logger
logger = get_logger("pre_test_result_task", "pre_test_result_task")
douyin_api = DouyinApi()
#获取前测结果并更新数据库,计划任务,每2分钟执行一次
async def poll_pre_test_results():
"""每2分钟轮询前测结果并更新数据库"""
logger.info("Pre-test result polling task started")
while True:
try:
await process_pending_pre_tests()
await asyncio.sleep(120)
except asyncio.CancelledError:
logger.info("Pre-test result polling task cancelled")
break
except Exception as e:
logger.error(f"Error in poll_pre_test_results: {e}")
await asyncio.sleep(120)
async def process_pending_pre_tests():
"""处理所有待查询的前测任务"""
async with async_session() as db:
result = await db.execute(
select(ResourcesMaterial).where(
ResourcesMaterial.status == "PENDING",
ResourcesMaterial.task_id.is_not(None),
ResourcesMaterial.deleted_at.is_(None),
)
)
pending_materials = result.scalars().all()
if not pending_materials:
return
logger.info(f"Found {len(pending_materials)} pending pre-test tasks to process")
for material in pending_materials:
try:
await update_single_pre_test_result(db, material)
except Exception as e:
logger.error(f"Error processing pre-test for material {material.id}: {e}")
async def update_single_pre_test_result(db: AsyncSession, material: ResourcesMaterial):
"""更新单个素材的前测结果"""
oauth = await db.execute(
select(UserOAuth).where(
UserOAuth.id == material.oauth_id,
UserOAuth.deleted_at.is_(None),
)
)
oauth = oauth.scalar_one_or_none()
if not oauth:
logger.error(f"OAuth record not found for material {material.id}")
await db.execute(
update(ResourcesMaterial).where(ResourcesMaterial.id == material.id).values(
status="FAILED",
note="授权记录不存在",
)
)
await db.commit()
return
params = {
"advertiser_id": int(material.advertiser_id),
"task_ids": json.dumps([int(material.task_id)]),
}
try:
response = await douyin_api.get_material_pre_test_result(oauth.id, params)
except Exception as e:
logger.error(f"Failed to get pre-test result for task {material.task_id}: {e}")
return
code = response.get("code", -1)
if code != 0:
logger.error(f"API error for task {material.task_id}: {response.get('message', 'Unknown error')}")
return
data = response.get("data", {})
task_details = data.get("task_list", [])
if not task_details:
return
task_detail = task_details[0]
status = task_detail.get("status")
pre_result = {
"video_id": task_detail.get("video_id") or None,
"advertiser_id": task_detail.get("advertiser_id") or None,
"material_id": task_detail.get("material_id") or None,
"is_ad_high_quality_material": task_detail.get("is_ad_high_quality_material") or None,
"is_ecp_high_quality_material": task_detail.get("is_ecp_high_quality_material") or None,
"is_inefficient_material": task_detail.get("is_inefficient_material") or None,
"is_first_publish_material": task_detail.get("is_first_publish_material") or None,
"not_ad_high_quality_reason": task_detail.get("not_ad_high_quality_reason") or None,
"not_ecp_high_quality_reason": task_detail.get("not_ecp_high_quality_reason") or None,
"is_local_high_quality_material": task_detail.get("is_local_high_quality_material") or None,
}
await db.execute(
update(ResourcesMaterial).where(ResourcesMaterial.id == material.id).values(
status=status,
pre_result=json.dumps(pre_result, ensure_ascii=False),
)
)
await db.commit()