录制HLS时 结果复制到S3挂载点增加一个fallback处理
* 录制HLS时 结果复制到S3挂载点增加一个fallback处理 * fix : TikHub的Logger换为Loguru * fix : 修复直播切分视频接口 * Merge remote-tracking branch 'origin/feature/modal-cluster' into feature/modal-cluster * 合并分支 * fix : hls录制缓存读写问题 * fix : 增加s3分片上传接口 * fix 处理jpg图片的codec_name应为mjpg * fix 获取媒体metadata的类型报错,modal client升级到1.0.2 * fix 一些webhook相关的bug,添加了更详细的接口说明 * 更新直播录制为可跳转时间的hls流接口,用于替代掉腾讯VOD+火山云拉流转推,片段保证每片开头为关键帧,时间精度准确到毫秒级 * - KVCache类改为可拓展,基于环境变量设置KV space - test 环境配置与CF测试环境对齐 * 对齐预发环境 --------- Merge request URL: https://g-ldyi2063.coding.net/p/dev/d/modalDeploy/git/merge/4784?initial=true Co-authored-by: 康宇佳,shuohigh@gmail.com
This commit is contained in:
@@ -1,4 +1,5 @@
|
||||
import asyncio
|
||||
import datetime
|
||||
import os
|
||||
from typing import Annotated, Optional
|
||||
|
||||
@@ -9,6 +10,9 @@ import sentry_sdk
|
||||
from fastapi import APIRouter, Depends, UploadFile, HTTPException, File, Form
|
||||
from fastapi.responses import JSONResponse, RedirectResponse
|
||||
from starlette import status
|
||||
import boto3
|
||||
from botocore.config import Config
|
||||
|
||||
from ..config import WorkerConfig
|
||||
from ..middleware.authorization import verify_token
|
||||
from ..models.media_model import (MediaSources,
|
||||
@@ -17,20 +21,36 @@ from ..models.media_model import (MediaSources,
|
||||
MediaCacheStatus,
|
||||
DownloadResult,
|
||||
UploadResultResponse,
|
||||
UploadBase64Request
|
||||
UploadBase64Request, UploadPresignRequest, UploadPresignResponse,
|
||||
UploadMultipartPresignRequest, UploadMultipartPresignResponse
|
||||
)
|
||||
from ..models.web_model import SentryTransactionInfo
|
||||
from ..utils.KVCache import KVCache
|
||||
from ..models.web_model import SentryTransactionInfo, MonitorLiveRoomProductRequest, ModalTaskResponse, \
|
||||
LiveRoomProductCachesResponse
|
||||
from ..utils.KVCache import MediaSourceKVCache, LiveProductKVCache
|
||||
from ..utils.SentryUtils import SentryUtils
|
||||
|
||||
config = WorkerConfig()
|
||||
|
||||
router = APIRouter(prefix="/cache")
|
||||
modal_kv_cache = KVCache(kv_name=config.modal_kv_name, environment=config.modal_environment)
|
||||
client = boto3.client("s3",
|
||||
aws_access_key_id=os.environ.get("AWS_ACCESS_KEY_ID"),
|
||||
aws_secret_access_key=os.environ.get("AWS_SECRET_ACCESS_KEY"),
|
||||
region_name=config.S3_region,
|
||||
endpoint_url="https://s3-accelerate.amazonaws.com",
|
||||
config=Config(
|
||||
s3={'addressing_style': 'virtual'},
|
||||
signature_version='s3v4', )
|
||||
)
|
||||
|
||||
router = APIRouter(prefix="/cache", tags=['缓存'], )
|
||||
|
||||
modal_kv_cache = MediaSourceKVCache(kv_name=config.modal_kv_name,
|
||||
environment=config.modal_environment)
|
||||
|
||||
modal_kv_product_cache = LiveProductKVCache(kv_name=config.modal_product_kv_name,
|
||||
environment=config.modal_environment)
|
||||
|
||||
|
||||
@router.post("/",
|
||||
tags=["缓存"],
|
||||
summary="缓存视频文件",
|
||||
description="异步缓存视频文件到S3存储桶和Modal Dict(KV)",
|
||||
dependencies=[Depends(verify_token)])
|
||||
@@ -91,14 +111,19 @@ async def cache(medias: MediaSources) -> CacheResult:
|
||||
async with asyncio.TaskGroup() as group:
|
||||
tasks = [group.create_task(cache_handler(media)) for media in medias.inputs]
|
||||
|
||||
cache_task_result = [task.result() for task in tasks]
|
||||
cache_task_result_dict = {}
|
||||
cache_task_result_list = []
|
||||
|
||||
KVCache.batch_update_cloudflare_kv(cache_task_result)
|
||||
return CacheResult(caches={media.urn: media for media in cache_task_result})
|
||||
for task in tasks:
|
||||
result = task.result()
|
||||
cache_task_result_dict[result.urn] = result.model_dump_json()
|
||||
cache_task_result_list.append(result)
|
||||
|
||||
modal_kv_cache.batch_update_cloudflare_kv(cache_task_result_dict)
|
||||
return CacheResult(caches={media.urn: media for media in cache_task_result_list})
|
||||
|
||||
|
||||
@router.delete("/",
|
||||
tags=["缓存"],
|
||||
summary="清除指定的所有缓存",
|
||||
description="清除指定的所有缓存(包括KV记录和S3存储文件)",
|
||||
dependencies=[Depends(verify_token)])
|
||||
@@ -119,12 +144,11 @@ async def purge_media_kv_file(medias: MediaSources):
|
||||
tasks = [group.create_task(purge_handle(media)) for media in medias.inputs]
|
||||
|
||||
keys = [task.result() for task in tasks]
|
||||
KVCache.batch_remove_cloudflare_kv(keys)
|
||||
modal_kv_cache.batch_remove_cloudflare_kv(keys)
|
||||
return JSONResponse(content={"success": True, "keys": keys})
|
||||
|
||||
|
||||
@router.post("/download",
|
||||
tags=["缓存"],
|
||||
summary="批量获取下载地址",
|
||||
description="获取已缓存的视频下载地址",
|
||||
dependencies=[Depends(verify_token)])
|
||||
@@ -138,7 +162,6 @@ async def download_caches(medias: MediaSources) -> DownloadResult:
|
||||
|
||||
|
||||
@router.get("/download",
|
||||
tags=["缓存"],
|
||||
summary="下载已缓存的视频",
|
||||
description="通过CDN下载已缓存的视频文件")
|
||||
@sentry_sdk.trace
|
||||
@@ -149,7 +172,6 @@ async def download_cache(media: str) -> RedirectResponse:
|
||||
|
||||
|
||||
@router.delete("/kv",
|
||||
tags=["缓存"],
|
||||
summary="清除KV记录",
|
||||
description="清除当前环境下KV缓存过的所有数据(S3存储桶内的文件会保留)",
|
||||
dependencies=[Depends(verify_token)])
|
||||
@@ -163,7 +185,6 @@ async def purge_kv_all():
|
||||
|
||||
|
||||
@router.post("/kv",
|
||||
tags=["缓存"],
|
||||
summary="删除对应的KV记录",
|
||||
description="删除请求中对应的视频缓存记录",
|
||||
dependencies=[Depends(verify_token)])
|
||||
@@ -172,14 +193,13 @@ async def purge_kv(medias: MediaSources):
|
||||
for media in medias.inputs:
|
||||
modal_kv_cache.pop(media.urn)
|
||||
keys = [media.urn for media in medias.inputs]
|
||||
KVCache.batch_remove_cloudflare_kv(keys)
|
||||
modal_kv_cache.batch_remove_cloudflare_kv(keys)
|
||||
return JSONResponse(content={"success": True, "keys": keys})
|
||||
except Exception as e:
|
||||
return JSONResponse(content={"success": False, "error": str(e)})
|
||||
|
||||
|
||||
@router.post("/media",
|
||||
tags=["缓存"],
|
||||
summary="清除指定的所有缓存",
|
||||
description="清除指定的所有缓存(包括KV记录和S3存储文件), 将要被淘汰,使用DELETE /cache/替代",
|
||||
deprecated=True,
|
||||
@@ -201,12 +221,11 @@ async def purge_media(medias: MediaSources):
|
||||
tasks = [group.create_task(purge_handle(media)) for media in medias.inputs]
|
||||
|
||||
keys = [task.result() for task in tasks]
|
||||
KVCache.batch_remove_cloudflare_kv(keys)
|
||||
modal_kv_cache.batch_remove_cloudflare_kv(keys)
|
||||
return JSONResponse(content={"success": True, "keys": keys})
|
||||
|
||||
|
||||
@router.post("/upload-s3",
|
||||
tags=['缓存'],
|
||||
summary="上传文件到S3",
|
||||
description="上传文件到S3的文件必须小于200M",
|
||||
dependencies=[Depends(verify_token)])
|
||||
@@ -231,7 +250,6 @@ async def s3_upload(file: Annotated[UploadFile, File(description="上传的文
|
||||
|
||||
|
||||
@router.post('/upload-s3-b64',
|
||||
tags=['缓存'],
|
||||
summary="基于Base64格式上传文件到S3",
|
||||
description="上传文件到S3当文件必须小于200M",
|
||||
dependencies=[Depends(verify_token)])
|
||||
@@ -252,3 +270,100 @@ async def s3_upload_base64(body: UploadBase64Request) -> UploadResultResponse:
|
||||
media_source.status = MediaCacheStatus.ready
|
||||
media_source.downloader_id = fn_id
|
||||
return UploadResultResponse(media=media_source)
|
||||
|
||||
|
||||
@router.post("/monitor_live_room_product_trigger",
|
||||
summary="触发监控直播间商品信息并缓存",
|
||||
description="触发监控直播间商品信息并缓存, 如果直播结束清除缓存, 触发间隔请控制在60s以上",
|
||||
dependencies=[Depends(verify_token)])
|
||||
async def monitor_live_room_product(body: MonitorLiveRoomProductRequest) -> LiveRoomProductCachesResponse:
|
||||
fn = modal.Function.from_name(config.modal_app_name, "monitor_live_room_product_trigger",
|
||||
environment_name=config.modal_environment)
|
||||
status = await fn.remote.aio(body.cookie, body.room_id, body.author_id)
|
||||
if status == 0:
|
||||
product_list = modal_kv_product_cache.get_cache(body.room_id)
|
||||
return LiveRoomProductCachesResponse(status=status, cache_json=product_list.model_dump_json())
|
||||
elif status == 1:
|
||||
return LiveRoomProductCachesResponse(status=status, message="直播已结束")
|
||||
elif status == 2:
|
||||
return LiveRoomProductCachesResponse(status=status, message="分配到风控IP, 请稍后重试")
|
||||
elif status == 3:
|
||||
return LiveRoomProductCachesResponse(status=status, message="请求Tikhub API出现错误")
|
||||
else:
|
||||
return LiveRoomProductCachesResponse(status=4, message="内部错误")
|
||||
|
||||
|
||||
@router.post('/upload-s3/simple/presign',
|
||||
summary="S3简单上传预签名",
|
||||
description="利用S3就近接入点上传",
|
||||
dependencies=[Depends(verify_token)])
|
||||
async def s3_presign_upload(body: UploadPresignRequest) -> UploadPresignResponse:
|
||||
expires_in = 3600
|
||||
expired_at = datetime.datetime.now() + datetime.timedelta(seconds=expires_in)
|
||||
signed_url = client.generate_presigned_url("put_object",
|
||||
Params={
|
||||
'Bucket': config.S3_bucket_name,
|
||||
'Key': f"upload/{body.key}",
|
||||
"ContentType": body.content_type,
|
||||
}, ExpiresIn=expires_in, )
|
||||
return UploadPresignResponse(url=signed_url,
|
||||
urn=f"s3://{config.S3_region}/{config.S3_bucket_name}/upload/{body.key}",
|
||||
expired_at=expired_at)
|
||||
|
||||
|
||||
@router.post("/upload-s3/multipart/presign",
|
||||
summary="S3分片上传预签名",
|
||||
description="""
|
||||
1. 本地按文件总大小分Chunk大小按urls链接内的顺序通过HTTP PUT请求上传文件分片; 并将上传完成后获得的返回头ETag值记录与PartNumber对应, PartNumber对应使用url在urls内的顺位, 以1开始
|
||||
\n2. 所有分片上传完成后使用XML格式拼装出用于确认上传的body; 并通过HTTP POST complete_url确认上传, ContentType 需要确保为application/xml\n\n
|
||||
<CompleteMultipartUpload>
|
||||
<Part>
|
||||
<PartNumber>1</PartNumber>
|
||||
<ETag>"60575364b098a1a48765a28c3a48e0ef"</ETag>
|
||||
</Part>
|
||||
<Part>
|
||||
<PartNumber>2</PartNumber>
|
||||
<ETag>"a38691c31fd242faee5533c65b4501d7"</ETag>
|
||||
</Part>
|
||||
<Part>
|
||||
<PartNumber>3</PartNumber>
|
||||
<ETag>"7a7510cc83f98feea28f319567e4cf66"</ETag>
|
||||
</Part>
|
||||
</CompleteMultipartUpload>
|
||||
\n3. 如无法确认上传,可使用HTTP GET list_url debug当前分片上传状态,确认成功后list_url无法返回有效数据
|
||||
""", dependencies=[Depends(verify_token)])
|
||||
async def s3_presign_upload_multipart(body: UploadMultipartPresignRequest) -> UploadMultipartPresignResponse:
|
||||
chunk_count = body.parts_count
|
||||
multipart_upload_response = client.create_multipart_upload(Bucket=config.S3_bucket_name, Key=body.key,
|
||||
ContentType=body.content_type, )
|
||||
upload_id = multipart_upload_response.get("UploadId")
|
||||
signed_urls = []
|
||||
expires_in = 3600
|
||||
expired_at = datetime.datetime.now() + datetime.timedelta(seconds=expires_in)
|
||||
for i in range(chunk_count):
|
||||
signed_url = client.generate_presigned_url("upload_part",
|
||||
Params={
|
||||
'Bucket': config.S3_bucket_name,
|
||||
'Key': body.key,
|
||||
'PartNumber': i + 1,
|
||||
'UploadId': upload_id,
|
||||
}, ExpiresIn=expires_in)
|
||||
signed_urls.append(signed_url)
|
||||
|
||||
signed_completed_url = client.generate_presigned_url("complete_multipart_upload",
|
||||
Params={
|
||||
'Bucket': config.S3_bucket_name,
|
||||
'Key': body.key,
|
||||
'UploadId': upload_id,
|
||||
}, ExpiresIn=expires_in)
|
||||
signed_list_url = client.generate_presigned_url("list_parts",
|
||||
Params={
|
||||
'Bucket': config.S3_bucket_name,
|
||||
'Key': body.key,
|
||||
'UploadId': upload_id,
|
||||
}, ExpiresIn=expires_in)
|
||||
return UploadMultipartPresignResponse(urls=signed_urls,
|
||||
urn=f"s3://{config.S3_region}/{config.S3_bucket_name}/upload/{body.key}",
|
||||
expired_at=expired_at,
|
||||
complete_url=signed_completed_url,
|
||||
list_url=signed_list_url)
|
||||
|
||||
Reference in New Issue
Block a user