diff --git a/src/BowongModalFunctions/models/ffmpeg_worker_model.py b/src/BowongModalFunctions/models/ffmpeg_worker_model.py index c053199..b4d5ef7 100644 --- a/src/BowongModalFunctions/models/ffmpeg_worker_model.py +++ b/src/BowongModalFunctions/models/ffmpeg_worker_model.py @@ -3,9 +3,12 @@ from pydantic import BaseModel, Field, computed_field, field_validator from pydantic.json_schema import JsonSchemaValue from ..utils.TimeUtils import TimeDelta + class FFMpegSliceSegment(BaseModel): - start: TimeDelta = Field(description="视频切割的开始时间点秒数, 可为浮点小数(精确到小数点后3位,毫秒级)或者为标准格式的时间戳") - end: TimeDelta = Field(description="视频切割的结束时间点秒数, 可为浮点小数(精确到小数点后3位,毫秒级)或者标准格式的时间戳") + start: TimeDelta = Field( + description="视频切割的开始时间点秒数, 可为浮点小数(精确到小数点后3位,毫秒级)或者为标准格式的时间戳") + end: TimeDelta = Field( + description="视频切割的结束时间点秒数, 可为浮点小数(精确到小数点后3位,毫秒级)或者标准格式的时间戳") @computed_field @property diff --git a/src/BowongModalFunctions/models/web_model.py b/src/BowongModalFunctions/models/web_model.py index b599893..619415d 100644 --- a/src/BowongModalFunctions/models/web_model.py +++ b/src/BowongModalFunctions/models/web_model.py @@ -1,11 +1,12 @@ import uuid from enum import Enum -from typing import List, Union, Optional +from typing import List, Union, Optional, Dict, Any import pydantic from pydantic import BaseModel, Field, field_validator, ConfigDict, HttpUrl from .ffmpeg_worker_model import FFMpegSliceSegment from .media_model import MediaSource, MediaSources +from ..utils.VideoUtils import VideoMetadata class TaskStatus(str, Enum): @@ -52,8 +53,7 @@ class ModalTaskResponse(BaseModel): class WebhookNotify(BaseModel): endpoint: HttpUrl = Field(description="Webhook回调端点", examples=["https://webhook.example.com"]) method: WebhookMethodEnum = Field(description="Webhook回调请求方法", examples=["get", "post"]) - authentication: Optional[str] = Field(description="Webhook authentication回调请求降权token,如不需要鉴权则不填", - examples=["Bearer 123456"], default=None) + headers: Optional[Dict[str, str]] = Field(description="Webhook回调附带的Headers", default=None) class BaseFFMPEGTaskRequest(BaseModel): @@ -65,7 +65,7 @@ class BaseFFMPEGTaskStatusResponse(BaseModel): status: TaskStatus = Field(description="任务运行状态") error: Optional[str] = Field(description="任务错误原因", default=None) code: Optional[int] = Field(description="任务错误原因代码", default=None) - results: Optional[List[str]] = Field(description="任务运行结果", default=None) + results: Optional[List[Any]] = Field(description="任务运行结果", default=None) model_config = ConfigDict(extra='ignore') @@ -347,3 +347,9 @@ class ComfyTaskRequest(BaseFFMPEGTaskRequest): return v else: raise TypeError(v) + + +class FFMPEGResult(BaseModel): + urn: str = Field(description="FFMPEG任务结果urn") + content_length: int = Field(description="媒体资源文件字节大小(Byte)") + metadata: VideoMetadata = Field(description="媒体元数据") diff --git a/src/BowongModalFunctions/utils/SentryUtils.py b/src/BowongModalFunctions/utils/SentryUtils.py index f5b1817..192d0f9 100644 --- a/src/BowongModalFunctions/utils/SentryUtils.py +++ b/src/BowongModalFunctions/utils/SentryUtils.py @@ -1,4 +1,5 @@ -from typing import Optional, Tuple, Any +import os +from typing import Optional, Tuple, Any, List import httpx import psutil import sentry_sdk @@ -6,7 +7,7 @@ from loguru import logger import functools from BowongModalFunctions.models.web_model import WebhookNotify, WebhookMethodEnum, BaseFFMPEGTaskStatusResponse, \ - TaskStatus, ErrorCode + TaskStatus, ErrorCode, FFMPEGResult class SentryUtils: @@ -24,6 +25,7 @@ class SentryUtils: def decorator(func): @functools.wraps(func) def wrapper(*args, **kwargs): + logger.info(f"sentry-trace={sentry_trace_id}, baggage={sentry_baggage}") if sentry_trace_id and sentry_baggage: transaction = sentry_sdk.continue_trace(environ_or_headers={"sentry-trace": sentry_trace_id, "baggage": sentry_baggage, }) @@ -34,10 +36,27 @@ class SentryUtils: cpu_freq, cpu_count, mem = SentryUtils.capture_hardware_info() total_mb = mem.total >> 20 # available_mb = mem.available >> 20 - sentry_sdk.set_tag('fn.id', fn_id) - span.set_data('cpu.count', cpu_count) - span.set_data('cpu.frequency', f"{cpu_freq.current / 1000:.2f} GHz") - span.set_data('memory.available', f"{total_mb} Mb") + sentry_sdk.set_tags({ + 'fn.id': fn_id, + 'x_trace.id': sentry_trace_id, + 'x_trace.baggage': sentry_baggage, + 'cpu.count': cpu_count, + 'cpu.frequency': cpu_freq.current / 1000, + 'cpu.frequency.format': f"{cpu_freq.current / 1000:.2f} GHz", + 'memory.available': total_mb, + 'memory.available.format': f"{total_mb} Mb", + 'cloud_resource': { + "cloud.provider": 'Modal', + "type": "cloud_resource" + }, + 'modal': { + 'modal.region': os.environ.get('MODAL_REGION', 'unknown'), + 'modal.provider': os.environ.get('MODAL_CLOUD_PROVIDER', 'unknown'), + 'modal.task.id': os.environ.get('MODAL_TASK_ID', 'unknown'), + 'modal.identity.token': os.environ.get('MODAL_IDENTITY_TOKEN', 'unknown'), + 'modal.image.id': os.environ.get('MODAL_IMAGE_ID', 'unknown'), + } + }) result = func(*args, **kwargs) return result @@ -67,21 +86,33 @@ class SentryUtils: code = ErrorCode.SYSTEM_ERROR.value logger.info(f"webhook = {webhook}") if webhook: + logger.info(f"result = {result}") + if isinstance(result, str): + results = [result] + elif isinstance(result, FFMPEGResult): + results = [result] + elif isinstance(result, List): + results = result + else: + results = [result] with httpx.Client() as client: match webhook.method: case WebhookMethodEnum.POST: response = client.post(url=webhook.endpoint.__str__(), + headers=webhook.headers, json=BaseFFMPEGTaskStatusResponse( - taskId=func_id, status=status, error=error, - code=code, - results=[result] if isinstance(result, str) else result, + taskId=func_id, status=status, + error=error, code=code, + results=results, ).model_dump()) case WebhookMethodEnum.GET: - response = client.get(url=webhook.endpoint.__str__(), params=BaseFFMPEGTaskStatusResponse( - taskId=func_id, status=status, error=error, - code=code, - results=[result] if isinstance(result, str) else result, - ).model_dump()) + response = client.get(url=webhook.endpoint.__str__(), + headers=webhook.headers, + params=BaseFFMPEGTaskStatusResponse( + taskId=func_id, status=status, error=error, + code=code, + results=results, + ).model_dump()) logger.info(f"webhook response = {response}") response.raise_for_status() return result diff --git a/src/BowongModalFunctions/utils/VideoUtils.py b/src/BowongModalFunctions/utils/VideoUtils.py index b45096c..1aee8c5 100644 --- a/src/BowongModalFunctions/utils/VideoUtils.py +++ b/src/BowongModalFunctions/utils/VideoUtils.py @@ -34,6 +34,7 @@ class MediaStream(BaseModel): class AudioStream(MediaStream): + stream_type: str = "audio" sample_rate: str # bit_depth : str # sample_fmt: str @@ -42,6 +43,7 @@ class AudioStream(MediaStream): class VideoStream(MediaStream): + stream_type: str = "video" width: int height: int bit_rate: int @@ -82,6 +84,22 @@ class VideoUtils: video_metadata = VideoMetadata.model_validate_json(ffprobe.execute()) return video_metadata.streams[0] + @staticmethod + def ffprobe_media_metadata(media_path: str) -> VideoMetadata: + ffprobe = FFmpeg(executable="ffprobe").input( + media_path, print_format="json", show_streams=None + ) + video_metadata = VideoMetadata.model_validate_json(ffprobe.execute()) + return video_metadata + + @staticmethod + def ffprobe_video_duration(media_path: str) -> TimeDelta: + ffprobe_cmd = VideoUtils.ffmpeg_init(use_ffprobe=True) + ffprobe_cmd.input(media_path, print_format="json", show_streams=None) + metadata_json = ffprobe_cmd.execute() + metadata = VideoMetadata.model_validate_json(metadata_json) + return TimeDelta(seconds=metadata.streams[0].duration) + @staticmethod def ffprobe_audio_duration(media_path: str) -> TimeDelta: ffprobe_cmd = VideoUtils.ffmpeg_init(use_ffprobe=True) @@ -120,7 +138,7 @@ class VideoUtils: @staticmethod def noise_reduce(media_path: str, noise_sample_path: Optional[str] = None, - output_path: Optional[str] = None) -> str: + output_path: Optional[str] = None) -> Tuple[str, VideoStream]: samplerate = 44100 with AudioFile(media_path).resampled_to(float(samplerate)) as f: audio = f.read(f.frames) @@ -276,7 +294,7 @@ class VideoUtils: @staticmethod async def ffmpeg_slice_media(media_path: str, media_markers: List[FFMpegSliceSegment], - output_path: Optional[str] = None) -> List[str]: + output_path: Optional[str] = None) -> List[Tuple[str, VideoMetadata]]: """ 使用本地视频文件按时间段切割出分段视频 :param media_path: 本地视频路径 @@ -288,7 +306,7 @@ class VideoUtils: ffmpeg_cmd = VideoUtils.async_ffmpeg_init() ffmpeg_cmd.input(media_path) filter_complex: List[str] = [] - outputs: List[str] = [] + outputs: List[Tuple[str, VideoMetadata]] = [] if not output_path: output_path = FileUtils.file_path_extend(media_path, "slice") @@ -312,7 +330,6 @@ class VideoUtils: raise ValueError( f"第{i}个切割点结束点{marker.end.total_seconds()}s超出视频时长[0-{video_metadata.duration}s]范围") segment_output_path = FileUtils.file_path_extend(output_path, str(i)) - outputs.append(segment_output_path) ffmpeg_cmd.output(segment_output_path, map=[f"[cut{i}]", f"[acut{i}]"], reset_timestamps="1", @@ -324,6 +341,8 @@ class VideoUtils: crf=16, r=30, ) + video_metadata = VideoUtils.ffprobe_media_metadata(segment_output_path) + outputs.append((segment_output_path, video_metadata)) await ffmpeg_cmd.execute() return outputs @@ -331,13 +350,13 @@ class VideoUtils: @staticmethod async def ffmpeg_slice_stream_media(media_path: str, media_markers: List[FFMpegSliceSegment], - output_path: Optional[str] = None) -> List[str]: + output_path: Optional[str] = None) -> List[Tuple[str, VideoMetadata]]: """ 按时间分段切割HLS视频流 :param media_path: hls manifest URL :param media_markers: 分段起始结束时间标记 :param output_path: 最终输出文件路径, 片段会根据指定路径附加_1.mp4, _2.mp4等片段编号 - :return: 输出片段的本地路径 + :return: 输出片段的本地路径, 输出片段时长 """ import m3u8 playlist = m3u8.load(media_path) @@ -350,7 +369,7 @@ class VideoUtils: reconnect_streamed="1", reconnect_delay_max="5") filter_complex: List[str] = [] - outputs: List[str] = [] + outputs: List[Tuple[str, VideoMetadata]] = [] if not output_path: output_path = FileUtils.file_path_extend(media_path, "slice") @@ -373,7 +392,6 @@ class VideoUtils: ffmpeg_cmd.option('filter_complex', ';'.join(filter_complex)) for i, marker in enumerate(media_markers): output_filepath = FileUtils.file_path_extend(output_path, str(i)) - outputs.append(output_filepath) ffmpeg_cmd.output(output_filepath, map=[f"[cut{i}]", f"[acut{i}]"], reset_timestamps="1", @@ -384,6 +402,8 @@ class VideoUtils: acodec="aac", crf=16, r=30, ) + video_metadata = VideoUtils.ffprobe_media_metadata(output_filepath) + outputs.append((output_filepath, video_metadata)) await ffmpeg_cmd.execute() return outputs @@ -392,14 +412,14 @@ class VideoUtils: async def ffmpeg_concat_medias(media_paths: List[str], target_width: int = 1080, target_height: int = 1920, - output_path: Optional[str] = None) -> str: + output_path: Optional[str] = None) -> Tuple[str, VideoMetadata]: """ 将待处理的视频合并为一个视频 :param media_paths: 待合并的多个视频文件路径 :param target_width: 输出的视频分辨率宽 :param target_height: 输出的视频分辨率高 :param output_path: 指定输出视频路径 - :return: + :return: 最终合并结果路径,最终合并结果时长 """ total_videos = len(media_paths) @@ -451,15 +471,17 @@ class VideoUtils: }, ) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod - async def ffmpeg_extract_audio_async(media_path: str, output_path: Optional[str] = None) -> str: + async def ffmpeg_extract_audio_async(media_path: str, output_path: Optional[str] = None) -> Tuple[ + str, VideoMetadata]: """ 提取源视频的音频 :param media_path: 待处理的源视频 :param output_path: 指定输出的音频文件路径(可选) - :return: 最终输出音频文件路径 + :return: 最终输出音频文件路径,音频文件时长 """ if not output_path: output_path = FileUtils.file_path_change_extension(output_path, 'wav') @@ -484,11 +506,12 @@ class VideoUtils: ar=44100, ac=1) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod async def ffmpeg_mix_bgm(origin_audio_path: str, bgm_audio_path: str, video_volume: float = 1.4, - music_volume: float = 0.1, output_path: Optional[str] = None) -> str: + music_volume: float = 0.1, output_path: Optional[str] = None) -> Tuple[str, VideoMetadata]: """ 给待处理视频混合BGM :param origin_audio_path: 待处理的源视频 @@ -496,7 +519,7 @@ class VideoUtils: :param video_volume: 最终输出视频的音量系数 :param music_volume: BGM在源视频音量内占比的音量系数 :param output_path: 指定最终输出的视频路径(可选) - :return: 最终输出视频文件路径 + :return: 最终输出视频文件路径,最终输出视频时长 """ if not output_path: output_path = FileUtils.file_path_extend(origin_audio_path, "bgm") @@ -522,14 +545,15 @@ class VideoUtils: ac=2, # 音频通道数 ) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod async def ffmpeg_mix_bgm_with_noise_reduce(media_path: str, bgm_audio_path: str, video_volume: float = 1.4, music_volume: float = 0.1, noise_sample_path: Optional[str] = None, - output_path: Optional[str] = None) -> str: + output_path: Optional[str] = None) -> Tuple[str, VideoMetadata]: """ 先对待处理的视频音轨降噪,再将降噪后的结果添加BGM,最终输出降噪过且混合BGM的视频; 由于最终视频画面和音轨是同步混合+合成视频,所以处理速度会比分步降噪, 加BGM快; @@ -539,7 +563,7 @@ class VideoUtils: :param music_volume: 最终输出的BGM音量系数 :param noise_sample_path: 降噪使用的噪音样本,如不指定将使用源视频的前2秒作为样本(可选) :param output_path: 指定输出视频的路径(可选) - :return: 最终输出视频的路径 + :return: 最终输出视频的路径, 最终输出视频时长 """ if not output_path: @@ -571,16 +595,18 @@ class VideoUtils: ar=48000, ab='192k', ac=2, ) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod - async def ffmpeg_overlay_gif(media_path: str, overlay_gif_path: str, output_path: Optional[str] = None) -> str: + async def ffmpeg_overlay_gif(media_path: str, overlay_gif_path: str, output_path: Optional[str] = None) -> Tuple[ + str, VideoMetadata]: """ 将GIF特效叠加到视频上,如果视频较长则循环播放GIF :param media_path: 输入视频路径 :param overlay_gif_path: GIF特效文件路径 :param output_path: 指定输出路径 - :return: 输出视频路径 + :return: 输出视频路径, 最终输出视频时长 """ if not output_path: output_path = FileUtils.file_path_extend(media_path, "overlay") @@ -604,18 +630,19 @@ class VideoUtils: r=video_metadata.video_frame_rate, # 帧率 ) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod async def ffmpeg_zoom_loop(media_path: str, duration: float = 6.0, zoom: float = 0.1, - output_path: Optional[str] = None) -> str: + output_path: Optional[str] = None) -> Tuple[str, VideoMetadata]: """ 视频放大缩小循环特效 :param media_path: 待处理的视频文件路径 :param duration: 视频特效循环时间长度 :param zoom: 视频特效放大缩小系数 :param output_path: 指定输出视频地址(可选) - :return: 最终输出视频地址 + :return: 最终输出视频地址, 最终输出视频时长 """ if not output_path: @@ -640,13 +667,13 @@ class VideoUtils: r=video_metadata.video_frame_rate, # 帧率 ) await ffmpeg_cmd.execute() - - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod async def ffmpeg_corner_mirror(media_path: str, mirror_scale_down_size: int = 6, mirror_from_right: bool = True, mirror_position: tuple[float, float] = (40, 40), - output_path: Optional[str] = None) -> str: + output_path: Optional[str] = None) -> Tuple[str, VideoMetadata]: """ 对源视频添加镜像小窗特效 :param media_path: 待处理的源视频 @@ -654,7 +681,7 @@ class VideoUtils: :param mirror_from_right: 小窗原点是否使用右下角 :param mirror_position: 小窗基于原点坐标轴的偏移量 :param output_path: 指定的输出视频路径(可选) - :return: 返回最终输出视频的路径 + :return: 返回最终输出视频的路径, 最终输出视频时长 """ if not output_path: output_path = FileUtils.file_path_extend(media_path, 'mir') @@ -685,18 +712,19 @@ class VideoUtils: r=video_metadata.video_frame_rate # 帧率 ) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod async def ffmpeg_subtitle_apply(media_path: str, subtitle_path: str, - font_dir: str, output_path: Optional[str] = None) -> str: + font_dir: str, output_path: Optional[str] = None) -> Tuple[str, VideoMetadata]: """ 给视频画面叠加字幕,需要确保字幕文件为ass字幕,并且subtitle文件内设置的字体存在与font_dir文件夹内 :param media_path: 待处理的源视频 :param subtitle_path: ass字幕文件路径 :param font_dir: 字体文件目录路径 :param output_path: 指定输出文件路径(可选) - :return: 返回最终输出视频路径 + :return: 返回最终输出视频路径, 最终输出视频时长 """ if not output_path: output_path = FileUtils.file_path_extend(media_path, 'sub') @@ -713,16 +741,18 @@ class VideoUtils: acodec="copy", ) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata @staticmethod - async def ffmpeg_fill_longest(video_path: str, audio_path: str, output_path: Optional[str] = None) -> str: + async def ffmpeg_fill_longest(video_path: str, audio_path: str, output_path: Optional[str] = None) -> Tuple[ + str, VideoMetadata]: """ 用视频循环对齐音频时长,如果短于音频时长则循环填满音频时长,如短于音频时长则裁剪结尾 :param video_path: 使用的视频文件路径 :param audio_path: 匹配的音频文件路径 :param output_path: 指定输出文件地址 - :return: 最终输出的文件路径 + :return: 最终输出的文件路径, 最终输出视频详细信息 """ video_metadata = VideoUtils.ffprobe_video_format(video_path) audio_duration = VideoUtils.ffprobe_audio_duration(audio_path) @@ -746,4 +776,5 @@ class VideoUtils: shortest=None, ) await ffmpeg_cmd.execute() - return output_path + video_metadata = VideoUtils.ffprobe_media_metadata(output_path) + return output_path, video_metadata diff --git a/src/cluster/ffmpeg_app.py b/src/cluster/ffmpeg_app.py index f106e34..8ad6a78 100644 --- a/src/cluster/ffmpeg_app.py +++ b/src/cluster/ffmpeg_app.py @@ -17,7 +17,7 @@ app = modal.App( with ffmpeg_worker_image.imports(): import shutil, os, backoff, sentry_sdk - from typing import List, Optional, Tuple + from typing import List, Optional, Tuple, Dict, Any from loguru import logger from modal import current_function_call_id from ffmpeg.asyncio import FFmpeg @@ -26,7 +26,7 @@ with ffmpeg_worker_image.imports(): from BowongModalFunctions.utils.VideoUtils import VideoUtils from BowongModalFunctions.models.ffmpeg_worker_model import FFMpegSliceSegment from BowongModalFunctions.models.media_model import MediaSources, MediaSource, MediaProtocol - from BowongModalFunctions.models.web_model import SentryTransactionInfo, WebhookNotify + from BowongModalFunctions.models.web_model import SentryTransactionInfo, WebhookNotify, FFMPEGResult from BowongModalFunctions.config import WorkerConfig config = WorkerConfig() @@ -120,19 +120,20 @@ with ffmpeg_worker_image.imports(): sentry_trace_id=sentry_trace.x_trace_id if sentry_trace else None, sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) - async def ffmpeg_process(media_sources: MediaSources, output_filepath: str) -> str: + async def ffmpeg_process(media_sources: MediaSources, output_filepath: str) -> FFMPEGResult: input_videos = [f"{s3_mount}/{media_source.cache_filepath}" for media_source in media_sources.inputs] - local_output_path = await VideoUtils.ffmpeg_concat_medias(media_paths=input_videos, - output_path=output_filepath) + local_output_path, metadata = await VideoUtils.ffmpeg_concat_medias(media_paths=input_videos, + output_path=output_filepath) s3_outputs = local_copy_to_s3([local_output_path]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, + content_length=os.path.getsize(local_output_path), ) output_path = f"{output_path_prefix}/{config.modal_environment}/concat/outputs/{fn_id}/output.mp4" s3_output = await ffmpeg_process(media_sources=medias, output_filepath=output_path) if not sentry_trace: sentry_trace = SentryTransactionInfo(x_trace_id=sentry_sdk.get_traceparent(), x_baggage=sentry_sdk.get_baggage()) - return s3_output, sentry_trace + return s3_output.urn, sentry_trace @app.function( @@ -160,7 +161,7 @@ with ffmpeg_worker_image.imports(): sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) async def ffmpeg_slice_process(media_source: MediaSource, media_markers: List[FFMpegSliceSegment], - fn_id: str) -> List[str]: + fn_id: str) -> List[FFMPEGResult]: cache_filepath = f"{s3_mount}/{media_source.cache_filepath}" logger.info(f"从{media_source.urn}切割") for marker in media_markers: @@ -168,21 +169,21 @@ with ffmpeg_worker_image.imports(): segments = await VideoUtils.ffmpeg_slice_media(media_path=cache_filepath, media_markers=media_markers, output_path=f"{output_path_prefix}/{config.modal_environment}/slice/outputs/{fn_id}/output.mp4") - s3_outputs = local_copy_to_s3(segments) - return s3_outputs + return [FFMPEGResult(urn=local_copy_to_s3([segment[0]])[0], metadata=segment[1], + content_length=os.path.getsize(segment[0])) for segment in segments] @SentryUtils.sentry_tracker(name="视频切割任务", op="ffmpeg.slice", fn_id=fn_id, sentry_trace_id=sentry_trace.x_trace_id if sentry_trace else None, sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) async def ffmpeg_hls_slice_process(media_source: MediaSource, media_markers: List[FFMpegSliceSegment], - fn_id: str) -> List[str]: + fn_id: str) -> List[FFMPEGResult]: hls_m3u8_url = media_source.path segments = await VideoUtils.ffmpeg_slice_stream_media(media_path=hls_m3u8_url, media_markers=media_markers, output_path=f"{output_path_prefix}/{config.modal_environment}/slice/outputs/{fn_id}/output.mp4") - s3_outputs = local_copy_to_s3(segments) - return s3_outputs + return [FFMPEGResult(urn=local_copy_to_s3([segment[0]])[0], metadata=segment[1], + content_length=os.path.getsize(segment[0])) for segment in segments] match media.protocol: case MediaProtocol.hls: @@ -194,7 +195,7 @@ with ffmpeg_worker_image.imports(): media_markers=markers, fn_id=fn_id) - return outputs, sentry_trace + return [result.urn for result in outputs], sentry_trace @app.function(timeout=600, cloud="aws", @@ -219,19 +220,19 @@ with ffmpeg_worker_image.imports(): sentry_trace_id=sentry_trace.x_trace_id if sentry_trace else None, sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) - async def ffmpeg_process(media: MediaSource, fn_id: str) -> str: + async def ffmpeg_process(media: MediaSource, fn_id: str) -> FFMPEGResult: cache_filepath = f"{s3_mount}/{media.cache_filepath}" output_path = f"{output_path_prefix}/{config.modal_environment}/extract_audio/outputs/{fn_id}/output.wav" - output_path = await VideoUtils.ffmpeg_extract_audio_async(cache_filepath, output_path) + output_path, metadata = await VideoUtils.ffmpeg_extract_audio_async(cache_filepath, output_path) s3_outputs = local_copy_to_s3([output_path]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, content_length=os.path.getsize(output_path), ) match media_source.protocol: case MediaProtocol.hls: return None, sentry_trace case _: output = await ffmpeg_process(media_source, fn_id=fn_id) - return output, sentry_trace + return output.urn, sentry_trace @app.function(timeout=600, cloud="aws", @@ -258,24 +259,25 @@ with ffmpeg_worker_image.imports(): @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) async def ffmpeg_process(media: MediaSource, func_id: str, mirror_scale_down_size: int = 6, mirror_from_right: bool = True, - mirror_position: tuple[float, float] = (40, 40)) -> str: + mirror_position: tuple[float, float] = (40, 40)) -> FFMPEGResult: media_filepath = f"{s3_mount}/{media.cache_filepath}" output_path = f"{output_path_prefix}/{config.modal_environment}/corner_mirror/outputs/{func_id}/output.mp4" - local_output_filepath = await VideoUtils.ffmpeg_corner_mirror(media_path=media_filepath, - output_path=output_path, - mirror_from_right=mirror_from_right, - mirror_position=mirror_position, - mirror_scale_down_size=mirror_scale_down_size) + local_output_filepath, metadata = await VideoUtils.ffmpeg_corner_mirror(media_path=media_filepath, + output_path=output_path, + mirror_from_right=mirror_from_right, + mirror_position=mirror_position, + mirror_scale_down_size=mirror_scale_down_size) s3_outputs = local_copy_to_s3([local_output_filepath]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, + content_length=os.path.getsize(local_output_filepath), ) result = await ffmpeg_process(media=media, mirror_scale_down_size=mirror_scale_down_size, func_id=fn_id, mirror_from_right=mirror_from_right, mirror_position=mirror_position) if not sentry_trace: sentry_trace = SentryTransactionInfo(x_trace_id=sentry_sdk.get_traceparent(), x_baggage=sentry_sdk.get_baggage()) - return result, sentry_trace + return result.urn, sentry_trace @app.function(timeout=600, cloud="aws", @@ -297,21 +299,22 @@ with ffmpeg_worker_image.imports(): sentry_trace_id=sentry_trace.x_trace_id if sentry_trace else None, sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) - async def ffmpeg_process(media: MediaSource, func_id: str, gif: MediaSource) -> str: + async def ffmpeg_process(media: MediaSource, func_id: str, gif: MediaSource) -> FFMPEGResult: media_filepath = f"{s3_mount}/{media.cache_filepath}" gif_filepath = f"{s3_mount}/{gif.cache_filepath}" output_path = f"{output_path_prefix}/{config.modal_environment}/overlay/outputs/{func_id}/output.mp4" - local_output_filepath = await VideoUtils.ffmpeg_overlay_gif(media_path=media_filepath, - output_path=output_path, - overlay_gif_path=gif_filepath) + local_output_filepath, metadata = await VideoUtils.ffmpeg_overlay_gif(media_path=media_filepath, + output_path=output_path, + overlay_gif_path=gif_filepath) s3_outputs = local_copy_to_s3([local_output_filepath]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, + content_length=os.path.getsize(local_output_filepath), ) result = await ffmpeg_process(media=media, func_id=fn_id, gif=gif) if not sentry_trace: sentry_trace = SentryTransactionInfo(x_trace_id=sentry_sdk.get_traceparent(), x_baggage=sentry_sdk.get_baggage()) - return result, sentry_trace + return result.urn, sentry_trace @app.function(timeout=600, cloud="aws", @@ -333,21 +336,23 @@ with ffmpeg_worker_image.imports(): sentry_trace_id=sentry_trace.x_trace_id if sentry_trace else None, sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) - async def ffmpeg_process(media: MediaSource, func_id: str, duration: int = 6, zoom: float = 0.1) -> str: + async def ffmpeg_process(media: MediaSource, func_id: str, duration: int = 6, + zoom: float = 0.1) -> FFMPEGResult: media_filepath = f"{s3_mount}/{media.cache_filepath}" output_path = f"{output_path_prefix}/{config.modal_environment}/zoom_loop/outputs/{func_id}/output.mp4" - local_output_filepath = await VideoUtils.ffmpeg_zoom_loop(media_path=media_filepath, - output_path=output_path, - duration=duration, - zoom=zoom) + local_output_filepath, metadata = await VideoUtils.ffmpeg_zoom_loop(media_path=media_filepath, + output_path=output_path, + duration=duration, + zoom=zoom) s3_outputs = local_copy_to_s3([local_output_filepath]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, + content_length=os.path.getsize(local_output_filepath), ) result = await ffmpeg_process(media=media, duration=duration, zoom=zoom, func_id=fn_id) if not sentry_trace: sentry_trace = SentryTransactionInfo(x_trace_id=sentry_sdk.get_traceparent(), x_baggage=sentry_sdk.get_baggage()) - return result, sentry_trace + return result.urn, sentry_trace @app.function(timeout=600, cloud="aws", @@ -371,20 +376,23 @@ with ffmpeg_worker_image.imports(): sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) async def ffmpeg_process(video: MediaSource, bgm: MediaSource, func_id: str, - noise_sample: Optional[MediaSource] = None) -> str: + noise_sample: Optional[MediaSource] = None) -> FFMPEGResult: local_input_filepath = f"{s3_mount}/{video.cache_filepath}" bgm_filepath = f"{s3_mount}/{bgm.cache_filepath}" noise_sample_path = f"{s3_mount}/{noise_sample.cache_filepath}" if noise_sample else None output_path = f"{output_path_prefix}/{config.modal_environment}/bgm_nosie_reduce/outputs/{func_id}/output.mp4" - local_output_filepath = await VideoUtils.ffmpeg_mix_bgm_with_noise_reduce(media_path=local_input_filepath, - bgm_audio_path=bgm_filepath, - video_volume=video_volume, - music_volume=music_volume, - noise_sample_path=noise_sample_path, - output_path=output_path) + local_output_filepath, metadata = await VideoUtils.ffmpeg_mix_bgm_with_noise_reduce( + media_path=local_input_filepath, + bgm_audio_path=bgm_filepath, + video_volume=video_volume, + music_volume=music_volume, + noise_sample_path=noise_sample_path, + output_path=output_path + ) s3_outputs = local_copy_to_s3([local_output_filepath]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, + content_length=os.path.getsize(local_output_filepath), ) result = await ffmpeg_process(video=media, bgm=bgm, video_volume=video_volume, music_volume=music_volume, noise_sample=noise_sample) @@ -392,7 +400,7 @@ with ffmpeg_worker_image.imports(): if not sentry_trace: sentry_trace = SentryTransactionInfo(x_trace_id=sentry_sdk.get_traceparent(), x_baggage=sentry_sdk.get_baggage()) - return result, sentry_trace + return result.urn, sentry_trace @app.function(timeout=600, cloud="aws", @@ -415,7 +423,7 @@ with ffmpeg_worker_image.imports(): sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) async def ffmpeg_process(video: MediaSource, subtitle: MediaSource, - fonts: List[MediaSource], func_id: str) -> str: + fonts: List[MediaSource], func_id: str) -> FFMPEGResult: media_path = f"{s3_mount}/{video.cache_filepath}" subtitle_path = f"{s3_mount}/{subtitle.cache_filepath}" output_path = f"{output_path_prefix}/{config.modal_environment}/subtitle_apply/{func_id}/output.mp4" @@ -425,16 +433,17 @@ with ffmpeg_worker_image.imports(): for font in local_fonts: font_filename = os.path.basename(font) shutil.copy(font, f"{font_dir}/{font_filename}") - local_output = await VideoUtils.ffmpeg_subtitle_apply(media_path=media_path, subtitle_path=subtitle_path, - font_dir=font_dir, output_path=output_path) + local_output, metadata = await VideoUtils.ffmpeg_subtitle_apply(media_path=media_path, + subtitle_path=subtitle_path, + font_dir=font_dir, output_path=output_path) s3_outputs = local_copy_to_s3([local_output]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, content_length=os.path.getsize(local_output), ) result = await ffmpeg_process(video=media, subtitle=subtitle, fonts=fonts, func_id=fn_id) if not sentry_trace: sentry_trace = SentryTransactionInfo(x_trace_id=sentry_sdk.get_traceparent(), x_baggage=sentry_sdk.get_baggage()) - return result, sentry_trace + return result.urn, sentry_trace @app.function(timeout=600, cloud="aws", @@ -456,15 +465,15 @@ with ffmpeg_worker_image.imports(): sentry_trace_id=sentry_trace.x_trace_id if sentry_trace else None, sentry_baggage=sentry_trace.x_baggage if sentry_trace else None) @SentryUtils.webhook_handler(webhook=webhook, func_id=fn_id) - async def ffmpeg_process(media: MediaSource, audio: MediaSource, func_id: str) -> str: + async def ffmpeg_process(media: MediaSource, audio: MediaSource, func_id: str) -> FFMPEGResult: video_path = f"{s3_mount}/{media.cache_filepath}" audio_path = f"{s3_mount}/{audio.cache_filepath}" - local_output = await VideoUtils.ffmpeg_fill_longest(video_path=video_path, - audio_path=audio_path, - output_path=f"{output_path_prefix}/{config.modal_environment}/loop_fill/{func_id}/output.mp4") + local_output, metadata = await VideoUtils.ffmpeg_fill_longest(video_path=video_path, + audio_path=audio_path, + output_path=f"{output_path_prefix}/{config.modal_environment}/loop_fill/{func_id}/output.mp4") s3_outputs = local_copy_to_s3([local_output]) - return s3_outputs[0] + return FFMPEGResult(urn=s3_outputs[0], metadata=metadata, content_length=os.path.getsize(local_output), ) output = await ffmpeg_process(media=media, audio=audio, func_id=fn_id) @@ -472,4 +481,4 @@ with ffmpeg_worker_image.imports(): sentry_trace = SentryTransactionInfo(x_trace_id=sentry_sdk.get_traceparent(), x_baggage=sentry_sdk.get_baggage()) - return output, sentry_trace + return output.urn, sentry_trace