diff --git a/src/cluster/video.py b/src/cluster/video.py index dc57a95..1468fe3 100644 --- a/src/cluster/video.py +++ b/src/cluster/video.py @@ -22,24 +22,15 @@ app = modal.App( ]) with downloader_image.imports(): - import os, httpx, requests + import os, httpx import sentry_sdk from sentry_sdk.integrations.loguru import LoguruIntegration - import math - import uuid - from io import BytesIO - from PIL import ImageDraw, Image, ImageFont - from typing import List, Dict + from typing import List from loguru import logger - from datetime import datetime, timedelta, timezone - from modal import current_function_call_id - from httpx import Timeout from BowongModalFunctions.config import WorkerConfig - from BowongModalFunctions.utils.SentryUtils import SentryUtils from BowongModalFunctions.utils.KVCache import MediaSourceKVCache, LiveProductKVCache from BowongModalFunctions.models.media_model import MediaSource - from BowongModalFunctions.models.web_model import SentryTransactionInfo, LiveProduct, LiveProductCaches config = WorkerConfig() @@ -117,357 +108,7 @@ with downloader_image.imports(): from .video_apps import * - @app.function(cpu=(0.5, 16), max_containers=config.video_downloader_concurrency, timeout=240) - @modal.concurrent(max_inputs=200) - async def make_image_grid_upload(pic_info_list: List[Dict[str, str]], - image_size: int, - text_height: int, - font_size: int, - padding: int, - separator: int, - google_api_key: str, - sentry_trace: SentryTransactionInfo) -> str: - def create_image_grid(image_info_list: List[Dict[str, str]], - output_path: str, - image_size: int = 450, - text_height: int = 40, - font_size: int = 18, - padding: int = 5, - separator: int = 5) -> str | None: - """ - 创建一个包含图片和文字说明的马赛克拼图,最大36张图 - - :param image_info_list: 包含图片信息的字典列表,每个字典包含 "title" 和 "cover" 键 - :param output_path: 输出图片的保存路径 - :param image_size: 单个图片网格的尺寸/像素 - :param text_height: 文本框的高度/像素 - :param font_size: 文本尺寸/像素 - :param padding: 文本距离文本框边缘距离/像素 - :param separator: 分割线宽度/像素 - :return: 图片路径 - """ - try: - cell_height = image_size + text_height - cell_width = image_size - # 提取图片路径和文字说明 - image_paths = [] - captions = [] - for info in image_info_list: - captions.append(info["title"]) - if "resize:200" in info["cover"]: - info["cover"] = info["cover"].replace("resize:200", f"resize:{image_size}") - image_paths.append(info["cover"]) - - # 检查输入 - if len(image_paths) != len(captions): - raise ValueError("图片数量与文字说明数量不匹配") - - # 处理图片(包括从网络获取) - loaded_images = [] - for path in image_paths: - if path.startswith(('http://', 'https://')): - try: - response = requests.get(path) - response.raise_for_status() - img = Image.open(BytesIO(response.content)) - loaded_images.append(img) - except requests.RequestException as e: - raise FileNotFoundError(f"无法从网络获取图片: {path}, 错误: {e}") - elif os.path.exists(path): - img = Image.open(path) - loaded_images.append(img) - else: - raise FileNotFoundError(f"找不到图片文件: {path}") - - # 计算网格大小 - n = len(image_paths) - max_cols = min(6, n) # 最大列数为6或图片数量 - min_rows = math.ceil(n / max_cols) # 最小行数 - - # 找到最接近正方形的网格布局 - best_ratio = float('inf') - best_rows, best_cols = min_rows, max_cols - - for rows in range(min_rows, 7): # 最多6行 - cols = math.ceil(n / rows) - if cols > 6: # 列数不能超过6 - continue - - ratio = abs(rows / cols - 1) # 越接近1越好 - if ratio < best_ratio: - best_ratio = ratio - best_rows, best_cols = rows, cols - - # 检查是否超出最大网格限制 - if best_rows > 6 or best_cols > 6: - raise ValueError(f"图片数量({n})过多,无法在6x6网格内合理展示") - - # 创建画布 - canvas_width = best_cols * (cell_width + separator) - separator - canvas_height = best_rows * (cell_height + separator) - separator - canvas = Image.new('RGB', (canvas_width, canvas_height), color='white') - - # 优化分割线绘制 - 先绘制所有分割线 - draw = ImageDraw.Draw(canvas) - - # 绘制垂直分割线 - for col in range(best_cols - 1): - x = (col + 1) * (cell_width + separator) - separator // 2 - draw.line([(x, 0), (x, canvas_height)], fill=(0, 0, 0), width=separator) - - # 绘制水平分割线 - for row in range(best_rows - 1): - y = (row + 1) * (cell_height + separator) - separator // 2 - draw.line([(0, y), (canvas_width, y)], fill=(0, 0, 0), width=separator) - - # 尝试加载支持中文的字体 - font = None - # 尝试常用的中文字体 - chinese_fonts = [ - "simhei.ttf", # 黑体 - "simsun.ttc", # 宋体 - "microsoftyahei.ttf", # 微软雅黑 - "arial.ttf" # 最后尝试Arial - ] - - for font_name in chinese_fonts: - try: - font = ImageFont.truetype(font_name, font_size) - break - except IOError: - continue - - # 如果所有字体都失败,使用默认字体 - if font is None: - font = ImageFont.load_default() - print("警告: 无法加载任何指定字体,使用默认字体") - - # 绘制图片和文字 - for i, (img, caption) in enumerate(zip(loaded_images, captions)): - row = i // best_cols - col = i % best_cols - - # 计算位置 - x = col * (cell_width + separator) - y = row * (cell_height + separator) - - # 保持原始长宽比 - img_width, img_height = img.size - - # 计算等比例缩放后的尺寸 - ratio = min(image_size / img_width, image_size / img_height) - new_width = int(img_width * ratio) - new_height = int(img_height * ratio) - - # 调整图片大小 - img = img.resize((new_width, new_height), Image.LANCZOS) - - # 计算图片在单元格内的居中位置(只考虑图片区域,不包括文字区域) - img_x = x + (cell_width - new_width) // 2 - - # 计算图片区域的起始位置(文字区域下方) - image_area_start_y = y + min(text_height, text_height) - - # 计算图片区域的高度(单元格高度减去文字区域高度) - image_area_height = cell_height - (image_area_start_y - y) - - # 计算图片在图片区域内的垂直居中位置 - img_y = image_area_start_y + (image_area_height - new_height) // 2 - - # 确保图片不会超出网格底部 - if img_y + new_height > y + cell_height: - img_y = y + cell_height - new_height - - # 将图片粘贴到画布上 - canvas.paste(img, (img_x, img_y)) - - # 添加文字区域 (透明背景) - draw = ImageDraw.Draw(canvas) - - # 改进的文字自动换行处理(支持中文) - max_width = cell_width - 2 * padding - lines = [] - current_line = "" - - for char in caption: - test_line = current_line + char - left, top, right, bottom = font.getbbox(test_line) - test_width = right - left - - if test_width <= max_width: - current_line = test_line - else: - lines.append(current_line) - current_line = char - - if current_line: - lines.append(current_line) - - # 计算文字总行高 - line_height = font_size # 使用固定行高 - total_text_height = len(lines) * line_height - - # 绘制多行文字 - 不绘制背景,直接绘制文字 - for j, line in enumerate(lines): - left, top, right, bottom = font.getbbox(line) - text_width = right - left - - # 水平居中 - text_x = x + padding + (cell_width - 2 * padding - text_width) // 2 - - # 垂直位置 - 从单元格顶部开始,考虑padding - text_y_offset = y + padding + j * line_height - - # 绘制文字边框 - offsets = [(-1, -1), (0, -1), (1, -1), (-1, 0), (1, 0), (-1, 1), (0, 1), (1, 1)] - for offset in offsets: - draw.text((text_x + offset[0], text_y_offset + offset[1]), line, fill=(255, 255, 255), - font=font) - - # 绘制文字 - draw.text((text_x, text_y_offset), line, fill=(0, 0, 0), font=font) - - # 保存结果 - canvas.save(output_path) - return output_path - except Exception as e: - logger.exception(f"拼图错误 {e}") - return None - - def upload(file_path, google_api_key): - with open(file_path, "rb") as image: - image = image.read() - content_length = len(image) - if file_path.split('.')[-1] == "jpg": - content_type = f"image/jpeg" - elif file_path.split('.')[-1] == "png": - content_type = f"image/png" - elif file_path.split('.')[-1] == "gif": - content_type = f"image/gif" - elif file_path.split('.')[-1] == "webp": - content_type = f"image/webp" - else: - raise Exception(f"不支持的文件格式{file_path.split('.')[-1]}") - filename = file_path.split("\\")[-1] - logger.info(f"Uploading name = {filename}, size = {content_length}, type = {content_type} to google file") - with httpx.Client(timeout=1800) as client: - pre_upload_response = client.post( - url=f"https://generativelanguage.googleapis.com/upload/v1beta/files?key={google_api_key}", - headers={ - "X-Goog-Upload-Protocol": "resumable", - "X-Goog-Upload-Command": "start", - "X-Goog-Upload-Header-Content-Length": str(content_length), - "X-Goog-Upload-Header-Content-Type": content_type - }, - json={ - "file": { - "display_name": filename.split(".")[0], - } - }) - pre_upload_response.raise_for_status() - - upload_url = pre_upload_response.headers.get("X-Goog-Upload-Url") - - upload_response = client.post(url=upload_url, content=image, headers={ - "X-Goog-Upload-Offset": "0", - "X-Goog-Upload-Command": "upload, finalize", - "Content-Type": content_type - }) - upload_response.raise_for_status() - - return upload_response.json(), upload_response.status_code - - @SentryUtils.sentry_tracker(sentry_trace.x_trace_id, sentry_trace.x_baggage, op="make_grid_gemini", - name="将输入图拼为网格上传到Gemini网盘", fn_id=current_function_call_id()) - def _handler(google_api_key: str, - pic_info_list: List[Dict[str, str]], - image_size: int, - text_height: int, - font_size: int, - padding: int, - separator: int): - image_grid_path = f"grid_{uuid.uuid4()}.jpg" - image_grid_path = create_image_grid(pic_info_list, image_grid_path, image_size, text_height, font_size, - padding, separator) - if not image_grid_path: - raise Exception("创建图片网格失败") - - image_grid_gemini, code = upload(image_grid_path, google_api_key) - if code == 200: - image_gemini_uri = image_grid_gemini["file"]["uri"] - else: - logger.error("图片网格文件上传Gemini失败") - raise Exception("图片网格文件上传Gemini失败") - return image_gemini_uri - - return _handler(google_api_key, pic_info_list, image_size, text_height, font_size, padding, separator) - @app.function(max_containers=config.video_downloader_concurrency, timeout=130) - @modal.concurrent(max_inputs=50) - async def monitor_live_room_product_trigger(cookie: str, room_id: str, author_id: str) -> int: - def get_product_list(): - with httpx.Client(timeout=Timeout(timeout=120)) as client: - resp = client.get( - f'https://bowongai-{config.modal_environment}--{config.modal_app_name}-fastapi-webapp-tikhub.modal.run/douyin/web/fetch_live_room_product_result', - params={"cookie": cookie, "room_id": room_id, "author_id": author_id}) - resp.raise_for_status() - if resp.status_code == 200: - if resp.json()["data"]["total"] >= 0: - return 0, resp.json()["data"]["promotions"] - elif resp.json()["data"]["total"] == -1: - # 直播结束 - return 1, [] - elif resp.json()["data"]["total"] == -2: - # IP风控 - return 2, [] - # 其他错误 - return 3, [] - try: - logger.info(f"room_id {room_id} author_id {author_id} 触发监控商品...") - is_live, product_list = get_product_list() - if is_live == 1: - logger.warning(f"room_id {room_id} author_id {author_id} 直播已结束, 停止监控商品, 删除缓存") - modal_kv_product_cache.pop(room_id, raise_exception=False) - modal_kv_product_cache.batch_remove_cloudflare_kv([room_id]) - return is_live - elif is_live == 2: - logger.warning(f"room_id {room_id} author_id {author_id} 获取商品出现风控") - return is_live - elif is_live == 3: - logger.warning(f"room_id {room_id} author_id {author_id} 网络请求出现错误") - return is_live - last_cache = modal_kv_product_cache.get_cache(room_id) - if last_cache is not None: - for product in last_cache.product_list: - for new_product in product_list: - if product.title == new_product["title"]: - product_list.remove(new_product) - if len(product_list) > 0: - logger.success(f"room_id {room_id} author_id {author_id} 检测到商品变化, 增量刷新缓存") - # 最新的插入在最前 - for new_product in product_list[::-1]: - last_cache.product_list.insert(0, LiveProduct(title=new_product["title"], - leaf_category=new_product["leaf_category"], - shop_id=new_product["shop_id"], - product_id=new_product["product_id"], - cover=new_product["cover"], - detail_url=new_product["detail_url"])) - else: - logger.success(f"room_id {room_id} author_id {author_id} 新建缓存") - last_cache = LiveProductCaches(room_id=room_id, author_id=author_id, product_list=[ - LiveProduct(title=new_product["title"], - leaf_category=new_product["leaf_category"], - shop_id=new_product["shop_id"], - product_id=new_product["product_id"], - cover=new_product["cover"], - detail_url=new_product["detail_url"]) for new_product in product_list]) - last_cache.update_time = datetime.now(timezone(timedelta(hours=8))).strftime("%Y-%m-%d %H:%M:%S") - last_cache.count = len(last_cache.product_list) - modal_kv_product_cache.set_cache(last_cache) - modal_kv_product_cache.batch_update_cloudflare_kv({last_cache.room_id: last_cache.model_dump_json()}) - return is_live - except Exception as e: - logger.exception(f"room_id {room_id} author_id {author_id} 触发监控商品发生错误 {e}") - return 4 + diff --git a/src/cluster/video_apps/make_grid_upload.py b/src/cluster/video_apps/make_grid_upload.py new file mode 100644 index 0000000..e75c526 --- /dev/null +++ b/src/cluster/video_apps/make_grid_upload.py @@ -0,0 +1,301 @@ +import modal + +from ..video import downloader_image, app, config + +with downloader_image.imports(): + import os, httpx, requests + import math + import uuid + from io import BytesIO + from PIL import ImageDraw, Image, ImageFont + from typing import List, Dict + from loguru import logger + from modal import current_function_call_id + + from BowongModalFunctions.utils.SentryUtils import SentryUtils + from BowongModalFunctions.models.web_model import SentryTransactionInfo + + @app.function(cpu=(0.5, 16), max_containers=config.video_downloader_concurrency, timeout=240) + @modal.concurrent(max_inputs=200) + async def make_image_grid_upload(pic_info_list: List[Dict[str, str]], + image_size: int, + text_height: int, + font_size: int, + padding: int, + separator: int, + google_api_key: str, + sentry_trace: SentryTransactionInfo) -> str: + async def create_image_grid(image_info_list: List[Dict[str, str]], + output_path: str, + image_size: int = 450, + text_height: int = 40, + font_size: int = 18, + padding: int = 5, + separator: int = 5) -> str | None: + """ + 创建一个包含图片和文字说明的马赛克拼图,最大36张图 + + :param image_info_list: 包含图片信息的字典列表,每个字典包含 "title" 和 "cover" 键 + :param output_path: 输出图片的保存路径 + :param image_size: 单个图片网格的尺寸/像素 + :param text_height: 文本框的高度/像素 + :param font_size: 文本尺寸/像素 + :param padding: 文本距离文本框边缘距离/像素 + :param separator: 分割线宽度/像素 + :return: 图片路径 + """ + try: + cell_height = image_size + text_height + cell_width = image_size + # 提取图片路径和文字说明 + image_paths = [] + captions = [] + for info in image_info_list: + captions.append(info["title"]) + if "resize:200" in info["cover"]: + info["cover"] = info["cover"].replace("resize:200", f"resize:{image_size}") + image_paths.append(info["cover"]) + + # 检查输入 + if len(image_paths) != len(captions): + raise ValueError("图片数量与文字说明数量不匹配") + + # 处理图片(包括从网络获取) + loaded_images = [] + for path in image_paths: + if path.startswith(('http://', 'https://')): + try: + response = requests.get(path) + response.raise_for_status() + img = Image.open(BytesIO(response.content)) + loaded_images.append(img) + except requests.RequestException as e: + raise FileNotFoundError(f"无法从网络获取图片: {path}, 错误: {e}") + elif os.path.exists(path): + img = Image.open(path) + loaded_images.append(img) + else: + raise FileNotFoundError(f"找不到图片文件: {path}") + + # 计算网格大小 + n = len(image_paths) + max_cols = min(4, n) # 最大列数为6或图片数量 + min_rows = math.ceil(n / max_cols) # 最小行数 + + # 找到最接近正方形的网格布局 + best_ratio = float('inf') + best_rows, best_cols = min_rows, max_cols + + for rows in range(min_rows, 5): # 最多6行 + cols = math.ceil(n / rows) + if cols > 4: # 列数不能超过6 + continue + + ratio = abs(rows / cols - 1) # 越接近1越好 + if ratio < best_ratio: + best_ratio = ratio + best_rows, best_cols = rows, cols + + # 检查是否超出最大网格限制 + if best_rows > 4 or best_cols > 4: + raise ValueError(f"图片数量({n})过多,无法在4x4网格内合理展示") + + # 创建画布 + canvas_width = best_cols * (cell_width + separator) - separator + canvas_height = best_rows * (cell_height + separator) - separator + canvas = Image.new('RGB', (canvas_width, canvas_height), color='white') + + # 优化分割线绘制 - 先绘制所有分割线 + draw = ImageDraw.Draw(canvas) + + # 绘制垂直分割线 + for col in range(best_cols - 1): + x = (col + 1) * (cell_width + separator) - separator // 2 + draw.line([(x, 0), (x, canvas_height)], fill=(0, 0, 0), width=separator) + + # 绘制水平分割线 + for row in range(best_rows - 1): + y = (row + 1) * (cell_height + separator) - separator // 2 + draw.line([(0, y), (canvas_width, y)], fill=(0, 0, 0), width=separator) + + # 尝试加载支持中文的字体 + font = None + # 尝试常用的中文字体 + chinese_fonts = [ + "simhei.ttf", # 黑体 + "simsun.ttc", # 宋体 + "microsoftyahei.ttf", # 微软雅黑 + "arial.ttf" # 最后尝试Arial + ] + + for font_name in chinese_fonts: + try: + font = ImageFont.truetype(font_name, font_size) + break + except IOError: + continue + + # 如果所有字体都失败,使用默认字体 + if font is None: + font = ImageFont.load_default() + print("警告: 无法加载任何指定字体,使用默认字体") + + # 绘制图片和文字 + for i, (img, caption) in enumerate(zip(loaded_images, captions)): + row = i // best_cols + col = i % best_cols + + # 计算位置 + x = col * (cell_width + separator) + y = row * (cell_height + separator) + + # 保持原始长宽比 + img_width, img_height = img.size + + # 计算等比例缩放后的尺寸 + ratio = min(image_size / img_width, image_size / img_height) + new_width = int(img_width * ratio) + new_height = int(img_height * ratio) + + # 调整图片大小 + img = img.resize((new_width, new_height), Image.LANCZOS) + + # 计算图片在单元格内的居中位置(只考虑图片区域,不包括文字区域) + img_x = x + (cell_width - new_width) // 2 + + # 计算图片区域的起始位置(文字区域下方) + image_area_start_y = y + min(text_height, text_height) + + # 计算图片区域的高度(单元格高度减去文字区域高度) + image_area_height = cell_height - (image_area_start_y - y) + + # 计算图片在图片区域内的垂直居中位置 + img_y = image_area_start_y + (image_area_height - new_height) // 2 + + # 确保图片不会超出网格底部 + if img_y + new_height > y + cell_height: + img_y = y + cell_height - new_height + + # 将图片粘贴到画布上 + canvas.paste(img, (img_x, img_y)) + + # 添加文字区域 (透明背景) + draw = ImageDraw.Draw(canvas) + + # 改进的文字自动换行处理(支持中文) + max_width = cell_width - 2 * padding + lines = [] + current_line = "" + + for char in caption: + test_line = current_line + char + left, top, right, bottom = font.getbbox(test_line) + test_width = right - left + + if test_width <= max_width: + current_line = test_line + else: + lines.append(current_line) + current_line = char + + if current_line: + lines.append(current_line) + + # 计算文字总行高 + line_height = font_size # 使用固定行高 + total_text_height = len(lines) * line_height + + # 绘制多行文字 - 不绘制背景,直接绘制文字 + for j, line in enumerate(lines): + left, top, right, bottom = font.getbbox(line) + text_width = right - left + + # 水平居中 + text_x = x + padding + (cell_width - 2 * padding - text_width) // 2 + + # 垂直位置 - 从单元格顶部开始,考虑padding + text_y_offset = y + padding + j * line_height + + # 绘制文字边框 + offsets = [(-1, -1), (0, -1), (1, -1), (-1, 0), (1, 0), (-1, 1), (0, 1), (1, 1)] + for offset in offsets: + draw.text((text_x + offset[0], text_y_offset + offset[1]), line, fill=(255, 255, 255), + font=font) + + # 绘制文字 + draw.text((text_x, text_y_offset), line, fill=(0, 0, 0), font=font) + + # 保存结果 + canvas.save(output_path) + return output_path + except Exception as e: + logger.exception(f"拼图错误 {e}") + return None + + async def upload(file_path, google_api_key): + with open(file_path, "rb") as image: + image = image.read() + content_length = len(image) + if file_path.split('.')[-1] == "jpg": + content_type = f"image/jpeg" + elif file_path.split('.')[-1] == "png": + content_type = f"image/png" + elif file_path.split('.')[-1] == "gif": + content_type = f"image/gif" + elif file_path.split('.')[-1] == "webp": + content_type = f"image/webp" + else: + raise Exception(f"不支持的文件格式{file_path.split('.')[-1]}") + filename = file_path.split("\\")[-1] + logger.info(f"Uploading name = {filename}, size = {content_length}, type = {content_type} to google file") + async with httpx.AsyncClient(timeout=1800) as client: + pre_upload_response = await client.post( + url=f"https://generativelanguage.googleapis.com/upload/v1beta/files?key={google_api_key}", + headers={ + "X-Goog-Upload-Protocol": "resumable", + "X-Goog-Upload-Command": "start", + "X-Goog-Upload-Header-Content-Length": str(content_length), + "X-Goog-Upload-Header-Content-Type": content_type + }, + json={ + "file": { + "display_name": filename.split(".")[0], + } + }) + pre_upload_response.raise_for_status() + + upload_url = pre_upload_response.headers.get("X-Goog-Upload-Url") + + upload_response = await client.post(url=upload_url, content=image, headers={ + "X-Goog-Upload-Offset": "0", + "X-Goog-Upload-Command": "upload, finalize", + "Content-Type": content_type + }) + upload_response.raise_for_status() + + return upload_response.json(), upload_response.status_code + + @SentryUtils.sentry_tracker(sentry_trace.x_trace_id, sentry_trace.x_baggage, op="make_grid_gemini", + name="将输入图拼为网格上传到Gemini网盘", fn_id=current_function_call_id()) + async def _handler(google_api_key: str, + pic_info_list: List[Dict[str, str]], + image_size: int, + text_height: int, + font_size: int, + padding: int, + separator: int): + image_grid_path = f"grid_{uuid.uuid4()}.jpg" + image_grid_path = await create_image_grid(pic_info_list, image_grid_path, image_size, text_height, font_size, + padding, separator) + if not image_grid_path: + raise Exception("创建图片网格失败") + + image_grid_gemini, code = await upload(image_grid_path, google_api_key) + if code == 200: + image_gemini_uri = image_grid_gemini["file"]["uri"] + else: + logger.error("图片网格文件上传Gemini失败") + raise Exception("图片网格文件上传Gemini失败") + return image_gemini_uri + + return await _handler(google_api_key, pic_info_list, image_size, text_height, font_size, padding, separator) \ No newline at end of file diff --git a/src/cluster/video_apps/monitor_live_room_product_trigger.py b/src/cluster/video_apps/monitor_live_room_product_trigger.py new file mode 100644 index 0000000..c19890a --- /dev/null +++ b/src/cluster/video_apps/monitor_live_room_product_trigger.py @@ -0,0 +1,80 @@ +import modal + +from ..video import downloader_image, app, config, modal_kv_product_cache + +with downloader_image.imports(): + import httpx + from loguru import logger + from datetime import datetime, timedelta, timezone + from httpx import Timeout + + from BowongModalFunctions.models.web_model import LiveProduct, LiveProductCaches + + @app.function(max_containers=config.video_downloader_concurrency, timeout=130) + @modal.concurrent(max_inputs=50) + async def monitor_live_room_product_trigger(cookie: str, room_id: str, author_id: str) -> int: + def get_product_list(): + with httpx.Client(timeout=Timeout(timeout=120)) as client: + resp = client.get( + f'https://bowongai-{config.modal_environment}--{config.modal_app_name}-fastapi-webapp-tikhub.modal.run/douyin/web/fetch_live_room_product_result', + params={"cookie": cookie, "room_id": room_id, "author_id": author_id}) + resp.raise_for_status() + if resp.status_code == 200: + if resp.json()["data"]["total"] >= 0: + return 0, resp.json()["data"]["promotions"] + elif resp.json()["data"]["total"] == -1: + # 直播结束 + return 1, [] + elif resp.json()["data"]["total"] == -2: + # IP风控 + return 2, [] + # 其他错误 + return 3, [] + + try: + logger.info(f"room_id {room_id} author_id {author_id} 触发监控商品...") + is_live, product_list = get_product_list() + if is_live == 1: + logger.warning(f"room_id {room_id} author_id {author_id} 直播已结束, 停止监控商品, 删除缓存") + modal_kv_product_cache.pop(room_id, raise_exception=False) + modal_kv_product_cache.batch_remove_cloudflare_kv([room_id]) + return is_live + elif is_live == 2: + logger.warning(f"room_id {room_id} author_id {author_id} 获取商品出现风控") + return is_live + elif is_live == 3: + logger.warning(f"room_id {room_id} author_id {author_id} 网络请求出现错误") + return is_live + last_cache = modal_kv_product_cache.get_cache(room_id) + if last_cache is not None: + for product in last_cache.product_list: + for new_product in product_list: + if product.title == new_product["title"]: + product_list.remove(new_product) + if len(product_list) > 0: + logger.success(f"room_id {room_id} author_id {author_id} 检测到商品变化, 增量刷新缓存") + # 最新的插入在最前 + for new_product in product_list[::-1]: + last_cache.product_list.insert(0, LiveProduct(title=new_product["title"], + leaf_category=new_product["leaf_category"], + shop_id=new_product["shop_id"], + product_id=new_product["product_id"], + cover=new_product["cover"], + detail_url=new_product["detail_url"])) + else: + logger.success(f"room_id {room_id} author_id {author_id} 新建缓存") + last_cache = LiveProductCaches(room_id=room_id, author_id=author_id, product_list=[ + LiveProduct(title=new_product["title"], + leaf_category=new_product["leaf_category"], + shop_id=new_product["shop_id"], + product_id=new_product["product_id"], + cover=new_product["cover"], + detail_url=new_product["detail_url"]) for new_product in product_list]) + last_cache.update_time = datetime.now(timezone(timedelta(hours=8))).strftime("%Y-%m-%d %H:%M:%S") + last_cache.count = len(last_cache.product_list) + modal_kv_product_cache.set_cache(last_cache) + modal_kv_product_cache.batch_update_cloudflare_kv({last_cache.room_id: last_cache.model_dump_json()}) + return is_live + except Exception as e: + logger.exception(f"room_id {room_id} author_id {author_id} 触发监控商品发生错误 {e}") + return 4 \ No newline at end of file