新增产出视频的metadata数据和文件字节大小,通过Webhook方式返回调用者

This commit is contained in:
shuohigh@gmail.com
2025-05-23 15:52:03 +08:00
parent dc7b140c2d
commit 008626e0df
5 changed files with 195 additions and 115 deletions

View File

@@ -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

View File

@@ -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="媒体元数据")

View File

@@ -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

View File

@@ -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

View File

@@ -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