diff --git a/facefusion/ffmpeg_builder.py b/facefusion/ffmpeg_builder.py index a818e225..57aca55d 100644 --- a/facefusion/ffmpeg_builder.py +++ b/facefusion/ffmpeg_builder.py @@ -201,7 +201,6 @@ def set_audio_volume(audio_volume : int) -> List[Command]: return [ '-filter:a', 'volume=' + str(audio_volume / 100) ] -#todo: needs review - [encoding] [critical: low] explicit ffmpeg thread cap def set_thread_count(thread_count : int) -> List[Command]: return [ '-threads', str(thread_count) ] diff --git a/facefusion/workflows/core.py b/facefusion/workflows/core.py index cba1d56f..274efd42 100644 --- a/facefusion/workflows/core.py +++ b/facefusion/workflows/core.py @@ -1,18 +1,15 @@ -from concurrent.futures import ThreadPoolExecutor, as_completed from typing import List import numpy -from tqdm import tqdm -import facefusion.workflows.image_to_video as image_to_video -from facefusion import logger, process_manager, state_manager, translator, video_manager +from facefusion import logger, process_manager, state_manager, translator from facefusion.audio import create_empty_audio_frame, get_audio_frame, get_voice_frame from facefusion.common_helper import get_first from facefusion.filesystem import filter_audio_paths from facefusion.processors.core import get_processors_modules -from facefusion.temp_helper import clear_temp_directory, create_temp_directory, resolve_temp_frame_set +from facefusion.temp_helper import clear_temp_directory, create_temp_directory from facefusion.types import AudioFrame, ErrorCode, VisionFrame -from facefusion.vision import conditional_merge_vision_mask, detect_video_resolution, extract_vision_mask, read_static_image, read_static_images, read_static_video_frame, restrict_trim_frame, restrict_video_fps, restrict_video_resolution, scale_resolution +from facefusion.vision import conditional_merge_vision_mask, extract_vision_mask, read_static_image, read_static_images, read_static_video_frame, restrict_trim_frame, restrict_video_fps, select_video_frames def is_process_stopping() -> bool: @@ -34,14 +31,12 @@ def clear() -> ErrorCode: return 0 -#todo: copy - renamed workflow to workflow_mode def conditional_get_reference_vision_frame() -> VisionFrame: if state_manager.get_item('workflow_mode') == 'image-to-video': return read_static_video_frame(state_manager.get_item('target_path'), state_manager.get_item('reference_frame_number')) return read_static_image(state_manager.get_item('target_path')) -#todo: copy - renamed workflow to workflow_mode, video only, trim offset def conditional_get_source_audio_frame(frame_number : int) -> AudioFrame: if state_manager.get_item('workflow_mode') == 'image-to-video': trim_frame_start, _ = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) @@ -55,7 +50,6 @@ def conditional_get_source_audio_frame(frame_number : int) -> AudioFrame: return create_empty_audio_frame() -#todo: copy - renamed workflow to workflow_mode, video only, trim offset def conditional_get_source_voice_frame(frame_number : int) -> AudioFrame: if state_manager.get_item('workflow_mode') == 'image-to-video': trim_frame_start, _ = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) @@ -69,8 +63,13 @@ def conditional_get_source_voice_frame(frame_number : int) -> AudioFrame: return create_empty_audio_frame() -#todo: needs review - [workflow] [critical: medium] copy of process_temp_frame, receives vision frames instead of paths -def process_target_frame(frame_number : int, target_vision_frames : List[VisionFrame], temp_vision_frame : VisionFrame) -> VisionFrame: +def conditional_get_target_vision_frames(frame_number : int) -> List[VisionFrame]: + if state_manager.get_item('workflow_mode') == 'image-to-video': + return select_video_frames(state_manager.get_item('target_path'), frame_number, state_manager.get_item('target_frame_amount')) + return [ read_static_image(state_manager.get_item('target_path')) ] + + +def process_temp_frame(target_vision_frames : List[VisionFrame], temp_vision_frame : VisionFrame, frame_number : int) -> VisionFrame: reference_vision_frame = conditional_get_reference_vision_frame() source_vision_frames = read_static_images(state_manager.get_item('source_paths')) source_audio_frame = conditional_get_source_audio_frame(frame_number) @@ -90,84 +89,3 @@ def process_target_frame(frame_number : int, target_vision_frames : List[VisionF }) return conditional_merge_vision_mask(temp_vision_frame, temp_vision_mask) - - -#todo: needs review - [workflow] [critical: medium] copy of process_frames, uses temp_frame_set and warms the reference frame -def process_disk_frames() -> ErrorCode: - temp_frame_set = resolve_temp_frame_set(state_manager.get_item('target_path')) - - if temp_frame_set: - with tqdm(total = len(temp_frame_set), desc = translator.get('processing'), unit = 'frame', ascii = ' =', disable = state_manager.get_item('log_level') in [ 'warn', 'error' ]) as progress: - progress.set_postfix(execution_providers = state_manager.get_item('execution_providers')) - - read_static_video_frame(state_manager.get_item('target_path'), state_manager.get_item('reference_frame_number')) - - with ThreadPoolExecutor(max_workers = state_manager.get_item('execution_thread_count')) as executor: - futures = [] - - for frame_number, temp_frame_path in temp_frame_set.items(): - future = executor.submit(image_to_video.process_disk_frame, temp_frame_path, frame_number, temp_frame_set) - futures.append(future) - - for future in as_completed(futures): - if is_process_stopping(): - for pending_future in futures: - pending_future.cancel() - - if not future.cancelled(): - future.result() - progress.update() - - for processor_module in get_processors_modules(state_manager.get_item('processors')): - processor_module.post_process() - - if is_process_stopping(): - return 4 - else: - logger.error(translator.get('temp_frames_not_found'), __name__) - return 1 - return 0 - - -#todo: needs review - [streaming] [critical: high] ordered submit and drain keeps writes in order, failed close stops the process manager -def process_stream_frames() -> ErrorCode: - trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) - output_video_resolution = scale_resolution(detect_video_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_video_scale')) - temp_video_resolution = restrict_video_resolution(state_manager.get_item('target_path'), output_video_resolution) - temp_video_fps = restrict_video_fps(state_manager.get_item('target_path'), state_manager.get_item('output_video_fps')) - frame_range = range(trim_frame_start, trim_frame_end) - - if frame_range: - video_writer = video_manager.get_writer(state_manager.get_item('target_path'), temp_video_fps, temp_video_resolution, output_video_resolution, state_manager.get_item('output_video_fps')) - frame_look_ahead = image_to_video.calculate_frame_look_ahead(temp_video_resolution) - - with tqdm(total = len(frame_range), desc = translator.get('processing'), unit = 'frame', ascii = ' =', disable = state_manager.get_item('log_level') in [ 'warn', 'error' ]) as progress: - progress.set_postfix(execution_providers = state_manager.get_item('execution_providers')) - - read_static_video_frame(state_manager.get_item('target_path'), state_manager.get_item('reference_frame_number')) - - with ThreadPoolExecutor(max_workers = state_manager.get_item('execution_thread_count')) as executor: - futures = [] - - for frame_number in frame_range: - future = executor.submit(image_to_video.process_stream_frame, frame_number, temp_video_resolution) - futures.append(future) - - if len(futures) > frame_look_ahead: - image_to_video.write_stream_frame(video_writer, futures, progress) - - for _ in list(futures): - image_to_video.write_stream_frame(video_writer, futures, progress) - - if not video_manager.close_video_writer(video_writer): - process_manager.stop() - - for processor_module in get_processors_modules(state_manager.get_item('processors')): - processor_module.post_process() - - if is_process_stopping(): - return 4 - else: - logger.error(translator.get('temp_frames_not_found'), __name__) - return 1 - return 0 diff --git a/facefusion/workflows/image_to_image.py b/facefusion/workflows/image_to_image.py index 51a5cac8..f7300808 100644 --- a/facefusion/workflows/image_to_image.py +++ b/facefusion/workflows/image_to_image.py @@ -1,13 +1,9 @@ from functools import partial -from facefusion import content_analyser, ffmpeg, logger, process_manager, state_manager, translator -from facefusion.filesystem import is_image -from facefusion.processors.core import get_processors_modules -from facefusion.temp_helper import get_temp_file_path -from facefusion.time_helper import calculate_end_time +from facefusion import process_manager from facefusion.types import ErrorCode -from facefusion.vision import detect_image_resolution, pack_resolution, read_static_image, restrict_image_resolution, scale_resolution, write_image -from facefusion.workflows.core import clear, is_process_stopping, process_target_frame, setup +from facefusion.workflows.core import clear, setup +from facefusion.workflows.to_image import analyse_image, finalize_image, prepare_image, process_image def process(start_time : float) -> ErrorCode: @@ -32,55 +28,3 @@ def process(start_time : float) -> ErrorCode: process_manager.end() return 0 - - -def analyse_image() -> ErrorCode: - if content_analyser.analyse_image(state_manager.get_item('target_path')): - return 3 - return 0 - - -def prepare_image() -> ErrorCode: - output_image_resolution = scale_resolution(detect_image_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_image_scale')) - temp_image_resolution = restrict_image_resolution(state_manager.get_item('target_path'), output_image_resolution) - - logger.info(translator.get('copying_image').format(resolution = pack_resolution(temp_image_resolution)), __name__) - if ffmpeg.copy_image(state_manager.get_item('target_path'), temp_image_resolution): - logger.debug(translator.get('copying_image_succeeded'), __name__) - else: - logger.error(translator.get('copying_image_failed'), __name__) - process_manager.end() - return 1 - return 0 - - -def process_image() -> ErrorCode: - temp_image_path = get_temp_file_path(state_manager.get_item('target_path')) - target_vision_frame = read_static_image(state_manager.get_item('target_path')) - temp_vision_frame = read_static_image(temp_image_path, 'rgba') - temp_vision_frame = process_target_frame(0, [ target_vision_frame ], temp_vision_frame) - write_image(temp_image_path, temp_vision_frame) - - for processor_module in get_processors_modules(state_manager.get_item('processors')): - processor_module.post_process() - - if is_process_stopping(): - return 4 - return 0 - - -def finalize_image(start_time : float) -> ErrorCode: - output_image_resolution = scale_resolution(detect_image_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_image_scale')) - - logger.info(translator.get('finalizing_image').format(resolution = pack_resolution(output_image_resolution)), __name__) - if ffmpeg.finalize_image(state_manager.get_item('target_path'), state_manager.get_item('output_path'), output_image_resolution): - logger.debug(translator.get('finalizing_image_succeeded'), __name__) - else: - logger.warn(translator.get('finalizing_image_skipped'), __name__) - - if is_image(state_manager.get_item('output_path')): - logger.info(translator.get('processing_image_succeeded').format(seconds = calculate_end_time(start_time)), __name__) - else: - logger.error(translator.get('processing_image_failed'), __name__) - return 1 - return 0 diff --git a/facefusion/workflows/image_to_video.py b/facefusion/workflows/image_to_video.py index b0ad862b..6cf8c728 100644 --- a/facefusion/workflows/image_to_video.py +++ b/facefusion/workflows/image_to_video.py @@ -1,46 +1,35 @@ -from concurrent.futures import Future from functools import partial -from typing import List, Tuple -import cv2 -import numpy -from tqdm import tqdm - -import facefusion.workflows.core as core -from facefusion import content_analyser, ffmpeg, logger, process_manager, state_manager, translator, video_manager -from facefusion.common_helper import get_first, get_middle -from facefusion.filesystem import filter_audio_paths, is_video -from facefusion.temp_helper import move_temp_file -from facefusion.time_helper import calculate_end_time -from facefusion.types import ErrorCode, FrameSet, Resolution, VideoWriter, VisionFrame -from facefusion.vision import detect_video_resolution, pack_resolution, read_static_image, restrict_trim_frame, restrict_video_fps, restrict_video_resolution, scale_resolution, select_video_frames, write_image +from facefusion import process_manager, state_manager +from facefusion.types import ErrorCode +from facefusion.workflows.core import clear, setup +from facefusion.workflows.to_video import analyse_video, extract_frames, finalize_video, merge_frames, process_disk_frames, process_stream_frames, restore_audio -#todo: copy - task list adds workflow_strategy disk and stream branching def process(start_time : float) -> ErrorCode: tasks =\ [ analyse_video, - core.clear, - core.setup + clear, + setup ] if state_manager.get_item('workflow_strategy') == 'disk': tasks.extend( [ extract_frames, - core.process_disk_frames, + process_disk_frames, merge_frames ]) if state_manager.get_item('workflow_strategy') == 'stream': - tasks.append(core.process_stream_frames) + tasks.append(process_stream_frames) tasks.extend( [ restore_audio, partial(finalize_video, start_time), - core.clear + clear ]) process_manager.start() @@ -53,148 +42,3 @@ def process(start_time : float) -> ErrorCode: process_manager.end() return 0 - - -#todo: copy - restrict_trim_video_frame renamed to restrict_trim_frame -def analyse_video() -> ErrorCode: - trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) - - if content_analyser.analyse_video(state_manager.get_item('target_path'), trim_frame_start, trim_frame_end): - return 3 - return 0 - - -#todo: copy - renamed from create_temp_frames, ffmpeg.extract_frames without output_path -def extract_frames() -> ErrorCode: - trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) - output_video_resolution = scale_resolution(detect_video_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_video_scale')) - temp_video_resolution = restrict_video_resolution(state_manager.get_item('target_path'), output_video_resolution) - temp_video_fps = restrict_video_fps(state_manager.get_item('target_path'), state_manager.get_item('output_video_fps')) - logger.info(translator.get('extracting_frames').format(resolution=pack_resolution(temp_video_resolution), fps=temp_video_fps), __name__) - - if ffmpeg.extract_frames(state_manager.get_item('target_path'), temp_video_resolution, temp_video_fps, trim_frame_start, trim_frame_end): - logger.debug(translator.get('extracting_frames_succeeded'), __name__) - else: - if core.is_process_stopping(): - return 4 - logger.error(translator.get('extracting_frames_failed'), __name__) - return 1 - return 0 - - -#todo: copy - dropped conditional resolution and fps helpers -def merge_frames() -> ErrorCode: - trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) - output_video_resolution = scale_resolution(detect_video_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_video_scale')) - temp_video_fps = restrict_video_fps(state_manager.get_item('target_path'), state_manager.get_item('output_video_fps')) - - logger.info(translator.get('merging_video').format(resolution = pack_resolution(output_video_resolution), fps = state_manager.get_item('output_video_fps')), __name__) - if ffmpeg.merge_video(state_manager.get_item('target_path'), temp_video_fps, output_video_resolution, state_manager.get_item('output_video_fps'), trim_frame_start, trim_frame_end): - logger.debug(translator.get('merging_video_succeeded'), __name__) - else: - if core.is_process_stopping(): - return 4 - logger.error(translator.get('merging_video_failed'), __name__) - return 1 - return 0 - - -#todo: copy - clear_video_pool inline, move_temp_file takes target_path -def restore_audio() -> ErrorCode: - trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) - - if state_manager.get_item('output_audio_volume') == 0: - logger.info(translator.get('skipping_audio'), __name__) - move_temp_file(state_manager.get_item('target_path'), state_manager.get_item('output_path')) - else: - source_audio_path = get_first(filter_audio_paths(state_manager.get_item('source_paths'))) - if source_audio_path: - if ffmpeg.replace_audio(state_manager.get_item('target_path'), source_audio_path, state_manager.get_item('output_path')): - video_manager.clear_video_pool() - logger.debug(translator.get('replacing_audio_succeeded'), __name__) - else: - video_manager.clear_video_pool() - if core.is_process_stopping(): - return 4 - logger.warn(translator.get('replacing_audio_skipped'), __name__) - move_temp_file(state_manager.get_item('target_path'), state_manager.get_item('output_path')) - else: - if ffmpeg.restore_audio(state_manager.get_item('target_path'), state_manager.get_item('output_path'), trim_frame_start, trim_frame_end): - video_manager.clear_video_pool() - logger.debug(translator.get('restoring_audio_succeeded'), __name__) - else: - video_manager.clear_video_pool() - if core.is_process_stopping(): - return 4 - logger.warn(translator.get('restoring_audio_skipped'), __name__) - move_temp_file(state_manager.get_item('target_path'), state_manager.get_item('output_path')) - return 0 - - -#todo: needs review - [correctness] [critical: high] missing neighbor frames resolve to none paths and rely on read_static_image returning none -def resolve_temp_vision_frames(frame_number : int, frame_amount : int, temp_frame_set : FrameSet) -> List[VisionFrame]: - temp_vision_frames = [] - frame_range = range(frame_number - frame_amount, frame_number + frame_amount + 1) - - for temp_frame_number in frame_range: - temp_vision_frames.append(read_static_image(temp_frame_set.get(temp_frame_number))) - - return temp_vision_frames - - -#todo: needs review - [workflow] [critical: low] disk variant reads temp frame and neighbors from disk and writes back in place -def process_disk_frame(temp_frame_path : str, frame_number : int, temp_frame_set : FrameSet) -> bool: - target_vision_frames = resolve_temp_vision_frames(frame_number, state_manager.get_item('target_frame_amount'), temp_frame_set) - temp_vision_frame = read_static_image(temp_frame_path, 'rgba') - temp_vision_frame = core.process_target_frame(frame_number, target_vision_frames, temp_vision_frame) - return write_image(temp_frame_path, temp_vision_frame) - - -#todo: needs review - [streaming] [critical: high] middle frame resized to temp resolution before processing, channels follow temp_pixel_format to match the writer pipe -def process_stream_frame(frame_number : int, temp_video_resolution : Resolution) -> Tuple[int, VisionFrame]: - target_vision_frames = select_video_frames(state_manager.get_item('target_path'), frame_number, state_manager.get_item('target_frame_amount')) - target_vision_frame = get_middle(target_vision_frames) - temp_vision_frame = target_vision_frame.copy() - - if not (target_vision_frame.shape[1], target_vision_frame.shape[0]) == temp_video_resolution: - temp_vision_frame = cv2.resize(target_vision_frame, temp_video_resolution) - temp_vision_frame = core.process_target_frame(frame_number, target_vision_frames, temp_vision_frame) - - if state_manager.get_item('temp_pixel_format') == 'bgra': - temp_vision_frame = cv2.cvtColor(temp_vision_frame, cv2.COLOR_BGR2BGRA) - - if state_manager.get_item('temp_pixel_format') == 'bgr24': - temp_vision_frame = temp_vision_frame[:, :, :3] - - return frame_number, numpy.ascontiguousarray(temp_vision_frame) - - -#todo: needs review - [memory] [critical: high] frame look ahead bound by a hardcoded 3gb budget -def calculate_frame_look_ahead(temp_video_resolution : Resolution) -> int: - width, height = temp_video_resolution - frame_memory_budget = 3 * 1024 ** 3 - frame_memory_usage = width * height * 4 * 6 - return min(state_manager.get_item('execution_thread_count') * 2, max(2, frame_memory_budget // frame_memory_usage)) - - -def write_stream_frame(video_writer : VideoWriter, futures : List[Future[Tuple[int, VisionFrame]]], progress : tqdm) -> None: - if core.is_process_stopping(): - for pending_future in futures: - pending_future.cancel() - - future = futures.pop(0) - - if not future.cancelled(): - _, temp_vision_frame = future.result() - video_manager.write_video_frame(video_writer, temp_vision_frame) - progress.update() - - -#todo: copy -def finalize_video(start_time : float) -> ErrorCode: - if is_video(state_manager.get_item('output_path')): - logger.info(translator.get('processing_video_succeeded').format(seconds = calculate_end_time(start_time)), __name__) - else: - logger.error(translator.get('processing_video_failed'), __name__) - return 1 - return 0 diff --git a/facefusion/workflows/to_image.py b/facefusion/workflows/to_image.py new file mode 100644 index 00000000..e9a70c5e --- /dev/null +++ b/facefusion/workflows/to_image.py @@ -0,0 +1,60 @@ +from facefusion import content_analyser, ffmpeg, logger, process_manager, state_manager, translator +from facefusion.filesystem import is_image +from facefusion.processors.core import get_processors_modules +from facefusion.temp_helper import get_temp_file_path +from facefusion.time_helper import calculate_end_time +from facefusion.types import ErrorCode +from facefusion.vision import detect_image_resolution, pack_resolution, read_static_image, restrict_image_resolution, scale_resolution, write_image +from facefusion.workflows.core import conditional_get_target_vision_frames, is_process_stopping, process_temp_frame + + +def analyse_image() -> ErrorCode: + if content_analyser.analyse_image(state_manager.get_item('target_path')): + return 3 + return 0 + + +def prepare_image() -> ErrorCode: + output_image_resolution = scale_resolution(detect_image_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_image_scale')) + temp_image_resolution = restrict_image_resolution(state_manager.get_item('target_path'), output_image_resolution) + + logger.info(translator.get('copying_image').format(resolution = pack_resolution(temp_image_resolution)), __name__) + if ffmpeg.copy_image(state_manager.get_item('target_path'), temp_image_resolution): + logger.debug(translator.get('copying_image_succeeded'), __name__) + else: + logger.error(translator.get('copying_image_failed'), __name__) + process_manager.end() + return 1 + return 0 + + +def process_image() -> ErrorCode: + temp_image_path = get_temp_file_path(state_manager.get_item('target_path')) + target_vision_frames = conditional_get_target_vision_frames(0) + temp_vision_frame = read_static_image(temp_image_path, 'rgba') + temp_vision_frame = process_temp_frame(target_vision_frames, temp_vision_frame, 0) + write_image(temp_image_path, temp_vision_frame) + + for processor_module in get_processors_modules(state_manager.get_item('processors')): + processor_module.post_process() + + if is_process_stopping(): + return 4 + return 0 + + +def finalize_image(start_time : float) -> ErrorCode: + output_image_resolution = scale_resolution(detect_image_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_image_scale')) + + logger.info(translator.get('finalizing_image').format(resolution = pack_resolution(output_image_resolution)), __name__) + if ffmpeg.finalize_image(state_manager.get_item('target_path'), state_manager.get_item('output_path'), output_image_resolution): + logger.debug(translator.get('finalizing_image_succeeded'), __name__) + else: + logger.warn(translator.get('finalizing_image_skipped'), __name__) + + if is_image(state_manager.get_item('output_path')): + logger.info(translator.get('processing_image_succeeded').format(seconds = calculate_end_time(start_time)), __name__) + else: + logger.error(translator.get('processing_image_failed'), __name__) + return 1 + return 0 diff --git a/facefusion/workflows/to_video.py b/facefusion/workflows/to_video.py new file mode 100644 index 00000000..e9ca0c20 --- /dev/null +++ b/facefusion/workflows/to_video.py @@ -0,0 +1,205 @@ +from concurrent.futures import ThreadPoolExecutor, as_completed +from typing import Tuple + +import cv2 +import numpy +from tqdm import tqdm + +from facefusion import content_analyser, ffmpeg, logger, process_manager, state_manager, translator, video_manager +from facefusion.common_helper import get_first, get_middle +from facefusion.filesystem import filter_audio_paths, is_video +from facefusion.processors.core import get_processors_modules +from facefusion.temp_helper import move_temp_file, resolve_temp_frame_set +from facefusion.time_helper import calculate_end_time +from facefusion.types import ErrorCode, Resolution, VisionFrame +from facefusion.vision import detect_video_resolution, pack_resolution, read_static_image, read_static_video_frame, restrict_trim_frame, restrict_video_fps, restrict_video_resolution, scale_resolution, select_video_frames, write_image +from facefusion.workflows.core import conditional_get_target_vision_frames, is_process_stopping, process_temp_frame + + +def analyse_video() -> ErrorCode: + trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) + + if content_analyser.analyse_video(state_manager.get_item('target_path'), trim_frame_start, trim_frame_end): + return 3 + return 0 + + +def extract_frames() -> ErrorCode: + trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) + output_video_resolution = scale_resolution(detect_video_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_video_scale')) + temp_video_resolution = restrict_video_resolution(state_manager.get_item('target_path'), output_video_resolution) + temp_video_fps = restrict_video_fps(state_manager.get_item('target_path'), state_manager.get_item('output_video_fps')) + logger.info(translator.get('extracting_frames').format(resolution=pack_resolution(temp_video_resolution), fps=temp_video_fps), __name__) + + if ffmpeg.extract_frames(state_manager.get_item('target_path'), temp_video_resolution, temp_video_fps, trim_frame_start, trim_frame_end): + logger.debug(translator.get('extracting_frames_succeeded'), __name__) + else: + if is_process_stopping(): + return 4 + logger.error(translator.get('extracting_frames_failed'), __name__) + return 1 + return 0 + + +def process_disk_frame(temp_frame_path : str, frame_number : int) -> bool: + target_vision_frames = conditional_get_target_vision_frames(frame_number) + temp_vision_frame = read_static_image(temp_frame_path, 'rgba') + temp_vision_frame = process_temp_frame(target_vision_frames, temp_vision_frame, frame_number) + return write_image(temp_frame_path, temp_vision_frame) + + +def process_disk_frames() -> ErrorCode: + temp_frame_set = resolve_temp_frame_set(state_manager.get_item('target_path')) + + if temp_frame_set: + with tqdm(total = len(temp_frame_set), desc = translator.get('processing'), unit = 'frame', ascii = ' =', disable = state_manager.get_item('log_level') in [ 'warn', 'error' ]) as progress: + progress.set_postfix(execution_providers = state_manager.get_item('execution_providers')) + + read_static_video_frame(state_manager.get_item('target_path'), state_manager.get_item('reference_frame_number')) + + with ThreadPoolExecutor(max_workers = state_manager.get_item('execution_thread_count')) as executor: + futures = [] + + for frame_number, temp_frame_path in temp_frame_set.items(): + future = executor.submit(process_disk_frame, temp_frame_path, frame_number) + futures.append(future) + + for future in as_completed(futures): + if is_process_stopping(): + for pending_future in futures: + pending_future.cancel() + + if not future.cancelled(): + future.result() + progress.update() + + for processor_module in get_processors_modules(state_manager.get_item('processors')): + processor_module.post_process() + + if is_process_stopping(): + return 4 + else: + logger.error(translator.get('temp_frames_not_found'), __name__) + return 1 + return 0 + + +def process_stream_frame(frame_number : int, temp_video_resolution : Resolution) -> Tuple[int, VisionFrame]: + target_vision_frames = select_video_frames(state_manager.get_item('target_path'), frame_number, state_manager.get_item('target_frame_amount')) + target_vision_frame = get_middle(target_vision_frames) + temp_vision_frame = target_vision_frame.copy() + + if not (target_vision_frame.shape[1], target_vision_frame.shape[0]) == temp_video_resolution: + temp_vision_frame = cv2.resize(target_vision_frame, temp_video_resolution) + + temp_vision_frame = process_temp_frame(target_vision_frames, temp_vision_frame, frame_number) + + if state_manager.get_item('temp_pixel_format') == 'bgra': + temp_vision_frame = cv2.cvtColor(temp_vision_frame, cv2.COLOR_BGR2BGRA) + + if state_manager.get_item('temp_pixel_format') == 'bgr24': + temp_vision_frame = temp_vision_frame[:, :, :3] + + return frame_number, numpy.ascontiguousarray(temp_vision_frame) + + +def process_stream_frames() -> ErrorCode: + trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) + output_video_resolution = scale_resolution(detect_video_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_video_scale')) + temp_video_resolution = restrict_video_resolution(state_manager.get_item('target_path'), output_video_resolution) + temp_video_fps = restrict_video_fps(state_manager.get_item('target_path'), state_manager.get_item('output_video_fps')) + temp_frame_range = range(trim_frame_start, trim_frame_end) + + if temp_frame_range: + video_writer = video_manager.get_writer(state_manager.get_item('target_path'), temp_video_fps, temp_video_resolution, output_video_resolution, state_manager.get_item('output_video_fps')) + + with tqdm(total = len(temp_frame_range), desc = translator.get('processing'), unit = 'frame', ascii = ' =', disable = state_manager.get_item('log_level') in [ 'warn', 'error' ]) as progress: + progress.set_postfix(execution_providers = state_manager.get_item('execution_providers')) + + read_static_video_frame(state_manager.get_item('target_path'), state_manager.get_item('reference_frame_number')) + + with ThreadPoolExecutor(max_workers = state_manager.get_item('execution_thread_count')) as executor: + futures = [] + + for frame_number in temp_frame_range: + future = executor.submit(process_stream_frame, frame_number, temp_video_resolution) + futures.append(future) + + for future in futures: + if is_process_stopping(): + for pending_future in futures: + pending_future.cancel() + + if not future.cancelled(): + _, temp_vision_frame = future.result() + video_manager.write_video_frame(video_writer, temp_vision_frame) + progress.update() + + if not video_manager.close_video_writer(video_writer): + process_manager.stop() + + for processor_module in get_processors_modules(state_manager.get_item('processors')): + processor_module.post_process() + + if is_process_stopping(): + return 4 + else: + logger.error(translator.get('temp_frames_not_found'), __name__) + return 1 + return 0 + + +def merge_frames() -> ErrorCode: + trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) + output_video_resolution = scale_resolution(detect_video_resolution(state_manager.get_item('target_path')), state_manager.get_item('output_video_scale')) + temp_video_fps = restrict_video_fps(state_manager.get_item('target_path'), state_manager.get_item('output_video_fps')) + + logger.info(translator.get('merging_video').format(resolution = pack_resolution(output_video_resolution), fps = state_manager.get_item('output_video_fps')), __name__) + if ffmpeg.merge_video(state_manager.get_item('target_path'), temp_video_fps, output_video_resolution, state_manager.get_item('output_video_fps'), trim_frame_start, trim_frame_end): + logger.debug(translator.get('merging_video_succeeded'), __name__) + else: + if is_process_stopping(): + return 4 + logger.error(translator.get('merging_video_failed'), __name__) + return 1 + return 0 + + +def restore_audio() -> ErrorCode: + trim_frame_start, trim_frame_end = restrict_trim_frame(state_manager.get_item('target_path'), state_manager.get_item('trim_frame_start'), state_manager.get_item('trim_frame_end')) + + if state_manager.get_item('output_audio_volume') == 0: + logger.info(translator.get('skipping_audio'), __name__) + move_temp_file(state_manager.get_item('target_path'), state_manager.get_item('output_path')) + else: + source_audio_path = get_first(filter_audio_paths(state_manager.get_item('source_paths'))) + if source_audio_path: + if ffmpeg.replace_audio(state_manager.get_item('target_path'), source_audio_path, state_manager.get_item('output_path')): + video_manager.clear_video_pool() + logger.debug(translator.get('replacing_audio_succeeded'), __name__) + else: + video_manager.clear_video_pool() + if is_process_stopping(): + return 4 + logger.warn(translator.get('replacing_audio_skipped'), __name__) + move_temp_file(state_manager.get_item('target_path'), state_manager.get_item('output_path')) + else: + if ffmpeg.restore_audio(state_manager.get_item('target_path'), state_manager.get_item('output_path'), trim_frame_start, trim_frame_end): + video_manager.clear_video_pool() + logger.debug(translator.get('restoring_audio_succeeded'), __name__) + else: + video_manager.clear_video_pool() + if is_process_stopping(): + return 4 + logger.warn(translator.get('restoring_audio_skipped'), __name__) + move_temp_file(state_manager.get_item('target_path'), state_manager.get_item('output_path')) + return 0 + + +def finalize_video(start_time : float) -> ErrorCode: + if is_video(state_manager.get_item('output_path')): + logger.info(translator.get('processing_video_succeeded').format(seconds = calculate_end_time(start_time)), __name__) + else: + logger.error(translator.get('processing_video_failed'), __name__) + return 1 + return 0