- BrainTools - https://www.braintools.ru -
Некоторое время назад я написал первую версию генератора коротких видео для Shorts, Reels и TikTok. На вход подавалась тема — например, «Сюжет “1984” за одну минуту», — а на выходе получался вертикальный ролик со сценарием, изображениями, анимацией, озвучкой и субтитрами.
Первая версия доказала, что идея работает. Но она же показала, что сгенерировать ролик один раз и построить надёжный production pipeline — совершенно разные задачи.
Если генерация падала на восьмой сцене из десяти, перезапуск мог повторить уже оплаченные запросы. Если процесс завершался между получением файла и сохранением состояния, было непонятно, считать сцену готовой или нет. А обычный asyncio.gather() давал параллелизм, но не отвечал на главный вопрос: что делать с частично успешным результатом?
Поэтому v2 я переписал почти с нуля. Clean Architecture осталась, но вокруг неё появились DDD, агрегат проекта, явная машина состояний, атомарные checkpoint’ы, content-addressed артефакты и возможность продолжить работу после сбоя.
В этой статье разберу не столько генеративные модели, сколько инженерную часть системы, в которой один API-вызов может длиться несколько минут и стоить реальных денег.
В v1 пайплайн выглядел примерно так:
script = await screenwriter.generate(topic)
script = await art_director.enhance(script)
script = await motion_director.enhance(script)
await asyncio.gather(
generate_images(script),
generate_speech(script),
)
await generate_videos(script)
return await editor.render(script)
Для прототипа это хороший код: последовательность понятна, внешние сервисы скрыты за интерфейсами, независимая работа запускается параллельно.
Но по мере развития проекта проявились проблемы.
Пока Python-процесс жив, он знает, какие сцены готовы. После Ctrl+C, сетевой ошибки [1] или падения машины это знание теряется.
Ретрай всего пайплайна мог снова вызвать LLM, генератор изображений или image-to-video API. Для локальной функции это просто повторный вызов. Для генеративного сервиса — новые минуты ожидания, новый результат и новая сумма в биллинге.
Допустим, из десяти изображений девять сгенерировались, а одно упало по timeout. Бизнес-состояние в этот момент не равно ни «успех», ни «ошибка». Это частично завершённая стадия, и её нужно уметь сохранить и продолжить.
Одни и те же модели постепенно использовались как ответы API, внутренние сущности и формат хранения на диске. В результате изменение JSON провайдера начинало влиять на бизнес-логику.
Главный вывод после v1 оказался таким: меняются не только провайдеры. Меняется и само состояние долгоживущего процесса. Значит, оно должно быть частью предметной модели.
Shorts Maker v2 — это один bounded context: производство короткого видео.
Он начинается с темы и заканчивается готовым видеофайлом. Публикация в соцсети, аналитика, аккаунты пользователей и биллинг провайдеров пока находятся за его границами.
Основные стадии:
topic
|
v
plan -------> images -------> videos ---+
| |
+----------> speech -------------------+--> render
Внутри стадии изображения, озвучка и видео генерируются отдельно для каждой сцены с ограниченной конкурентностью. Каждая успешно завершённая сцена сразу сохраняется.
Текущие реальные адаптеры:
CometAPI — структурированный production plan через LLM и генерация изображений;
Kie/Grok — image-to-video;
ImgBB — временный публичный мост для передачи исходного изображения в Kie;
Yandex SpeechKit — озвучка;
MoviePy и FFmpeg — финальный монтаж, аудио и субтитры.
Есть и полностью автономный профиль fake. Он проходит тот же application flow, но не вызывает платные API.
Код v2 находится в src/shorts_maker и разделён на четыре слоя:
presentation ─────> application ─────> domain
^
|
infrastructure ────────────+
Здесь находятся Project, сцены, персонажи, value objects, инварианты и переходы между стадиями. Домен ничего не знает о Pydantic, HTTP, MoviePy, файловой системе и CLI.
Здесь живут use cases: создать проект, построить план, сгенерировать изображения, речь и видео, выполнить рендер. Этот же слой объявляет исходящие порты.
Реализации портов: HTTP-клиенты провайдеров, JSON-репозиторий, хранилище медиа и MoviePy-рендерер.
CLI на Typer и Rich. Он разбирает аргументы, вызывает use case и показывает результат, но не принимает бизнес-решений.
Единственное место, которому разрешено знать обо всех слоях, — bootstrap.py. Это composition root, где порты связываются с конкретными адаптерами.
Направление зависимостей проверяется отдельным архитектурным тестом: он разбирает импорты через ast и падает, если, например, domain начал импортировать infrastructure.
FORBIDDEN_PREFIXES = {
"domain": (
"shorts_maker.application",
"shorts_maker.infrastructure",
"shorts_maker.presentation",
"shorts_maker.bootstrap",
),
"application": (
"shorts_maker.infrastructure",
"shorts_maker.presentation",
"shorts_maker.bootstrap",
),
}
Такие тесты кажутся избыточными, пока проект небольшой. Но правило, существующее только в architecture.md, рано или поздно нарушается. Исполняемое правило живёт дольше документации.
В центре домена находится агрегат Project. Сцена не сохраняется отдельно и не меняет свой production status самостоятельно: все переходы проходят через корень агрегата.
Домен написан на обычных dataclass, а не на Pydantic:
@dataclass(slots=True)
class Project:
id: ProjectId
topic: str
title: str
style_hint: str | None
characters: tuple[Character, ...]
scenes: list[Scene]
stages: dict[Stage, StageProgress]
revision: int
created_at: datetime
updated_at: datetime
Pydantic в v2 остался, но работает на границах: проверяет настройки, JSON-манифест и ответы внешних API. Это важное разделение. Формат ответа провайдера — не моя предметная модель.
Агрегат защищает инварианты:
тема не может быть пустой;
номера сцен уникальны, идут по порядку и начинаются с нуля;
в ролике не больше трёх главных персонажей;
сцена не может ссылаться на персонажа, которого нет в character bible;
стадию нельзя начать, пока не завершены её зависимости;
стадия изображений не может прикрепить видео или аудио;
время обновления проекта не может двигаться назад.
Например, начало стадии — это не присваивание строки в сервисе, а доменная операция:
def start_stage(self, stage: Stage, *, resume_interrupted: bool = False) -> bool:
progress = self.stages[stage]
if progress.status is StageStatus.COMPLETED:
return False
missing = [
dependency
for dependency in STAGE_DEPENDENCIES[stage]
if self.stages[dependency].status is not StageStatus.COMPLETED
]
if missing:
raise MissingStageDependencyError(...)
progress.status = StageStatus.RUNNING
progress.attempts += 1
self._touch()
return True
Если метод вернул False, use case знает, что стадия уже завершена, и не вызывает провайдера повторно.
У каждой стадии есть один из пяти статусов:
class StageStatus(StrEnum):
PENDING = "pending"
RUNNING = "running"
COMPLETED = "completed"
PARTIAL = "partial"
FAILED = "failed"
PARTIAL — один из наиболее полезных статусов во всём проекте.
Если две сцены из трёх готовы, стадия изображений сохраняется как частичная. Следующий запуск выбирает только сцены без нужного артефакта:
def pending_scenes(self, stage: Stage) -> tuple[Scene, ...]:
kind = SCENE_STAGE_ARTIFACT[stage]
return tuple(
scene for scene in self.scenes
if kind not in scene.artifacts
)
Это делает pipeline идемпотентным на уровне бизнес-операции. Команда может запускаться повторно, но уже готовая работа не повторяется.
Сохранённый статус running тоже можно забрать после прерванного процесса. Новый attempt переводит стадию в работу и продолжает с отсутствующих результатов.
В v1 основным инструментом был asyncio.gather(). В v2 задачи по-прежнему выполняются параллельно, но orchestration устроен иначе.
Упрощённый фрагмент генерации изображений:
pending = project.pending_scenes(Stage.IMAGES)
semaphore = asyncio.Semaphore(max_concurrency)
async def generate(scene_id: SceneId):
async with semaphore:
return await image_generator.generate(...)
tasks = [asyncio.create_task(generate(scene.id)) for scene in pending]
for future in asyncio.as_completed(tasks):
result = await future
if isinstance(result, SceneSuccess):
project.attach_scene_artifact(
stage=Stage.IMAGES,
scene_id=result.scene_id,
artifact=result.value,
)
else:
project.record_scene_failure(...)
await repository.save(project, expected_revision=expected_revision)
Здесь есть три принципиальных момента.
Во-первых, Semaphore ограничивает число одновременных запросов. Без него легко упереться в rate limit или случайно отправить провайдеру десятки платных задач.
Во-вторых, worker возвращает значение и не мутирует агрегат. Результаты применяются последовательно в application loop, поэтому внутри процесса нет гонки за общим объектом Project.
В-третьих, после каждого результата сохраняется checkpoint. Если процесс упадёт после седьмой сцены, первые семь останутся частью состояния проекта.
Это чуть многословнее одного gather(), зато семантика частичного выполнения становится явной.
Для v2 я сознательно не стал сразу поднимать PostgreSQL. Один проект — один агрегат, поэтому его состояние хранится в одном schema-versioned файле:
output/
└── the-trial-a1b2c3d4/
├── project.json
└── assets/
├── scene-000/
├── scene-001/
└── final/
В project.json находятся production plan, статусы стадий, число попыток, ошибки и ссылки на артефакты. Текущая версия схемы — 2.
Сохранение устроено так:
Захватывается локальный async lock.
Затем cross-process file lock.
Проверяется ожидаемая ревизия агрегата.
Новый JSON записывается во временный файл.
Данные сбрасываются через flush() и fsync().
os.replace() атомарно заменяет старый манифест.
with temporary.open("w", encoding="utf-8") as stream:
stream.write(content)
stream.flush()
os.fsync(stream.fileno())
os.replace(temporary, path)
Поле revision работает как optimistic lock. Use case загружает, например, ревизию 12 и при сохранении сообщает: «я изменял именно версию 12». Если другой процесс уже записал версию 13, репозиторий отклонит устаревшее обновление вместо тихой потери данных.
if persisted.revision != expected_revision:
raise ConcurrentProjectUpdateError(...)
Для локального приложения этого достаточно. Если появятся несколько машин и распределённая очередь, consistency boundary останется прежней, а файловый адаптер можно будет заменить базой данных.
Медиафайлы не кладутся под именами вроде scene_1_final_final_v2.png. Для каждого результата вычисляется SHA-256, а часть checksum входит в storage key:
checksum = hashlib.sha256(payload.content).hexdigest()
key = (
f"{project_id}/assets/{owner}/"
f"{kind.value}-{checksum[:12]}.{payload.extension}"
)
Из этого следуют полезные свойства:
одинаковое содержимое получает предсказуемое имя;
конкурирующие попытки не перезаписывают разные результаты одним путём;
manifest ссылается на конкретный immutable-результат;
при чтении можно повторно проверить SHA-256 и обнаружить повреждение файла.
Если попытка проиграла optimistic lock, созданный ею файл может остаться orphaned. Но он не попадёт в winning manifest и не испортит состояние проекта. Позже такие файлы можно безопасно чистить отдельным garbage collector’ом.
В первой версии интерфейсы были построены на ABC. В v2 я перешёл на структурную типизацию через Protocol:
class ImageGenerator(Protocol):
async def generate(self, request: ImageRequest) -> MediaPayload: ...
class VideoGenerator(Protocol):
async def generate(self, request: VideoRequest) -> MediaPayload: ...
class ProjectRepository(Protocol):
async def get(self, project_id: ProjectId) -> Project: ...
async def save(
self,
project: Project,
*,
expected_revision: int,
) -> None: ...
Адаптеру не нужно наследоваться от общего класса. Достаточно соблюдать контракт, который проверяет mypy.
Composition root выбирает реализации в одном месте:
if settings.profile == "fake":
planner = FakeScriptPlanner()
image_generator = FakeImageGenerator()
video_generator = FakeVideoGenerator()
speech_synthesizer = FakeSpeechSynthesizer()
renderer = FakeRenderer()
else:
planner = CometScriptPlanner(...)
image_generator = CometImageGenerator(...)
video_generator = KieVideoGenerator(...)
speech_synthesizer = YandexSpeechSynthesizer(...)
renderer = MoviePyRenderer(...)
Application layer при этом не знает, какой профиль запущен.
Провайдеры генерации отличаются не только URL и форматом JSON. У них разные схемы polling, ограничения, статусы ошибок и таймауты.
Общий HTTP-слой повторяет transport errors, HTTP 429 и выбранные 5xx с ограниченным exponential backoff:
_RETRYABLE_STATUSES = {408, 409, 425, 429, 500, 502, 503, 504}
delay = policy.initial_delay
for attempt in range(1, policy.attempts + 1):
...
await asyncio.sleep(delay)
delay = min(delay * 2, policy.maximum_delay)
Polling генерации имеет жёсткий deadline. Иначе одна зависшая задача способна удерживать весь pipeline бесконечно.
При этом ошибки конфигурации обрабатываются отдельно. Если не задан SHORTS_KIE_API_KEY, это не «неудача сцены» и не повод сделать четыре ретрая. CLI сразу сообщает, какой секрет отсутствует, а стадия остаётся доступной для продолжения после исправления .env.
Полностью автоматическая генерация удобна до первой сцены с шестью пальцами, сломанной перспективой или внезапно изменившимся персонажем.
В v2 появился флаг:
shorts-maker run <project-id> --review-images
После генерации CLI открывает изображение и предлагает:
принять результат;
перегенерировать с тем же prompt;
отредактировать prompt и повторить попытку.
Сам review тоже является портом ImageReviewer. В автоматическом режиме используется AutoApproveImageReviewer, а в интерактивном — реализация для консоли. Поэтому human approval встроен в use case, но UI не протекает в бизнес-логику.
Число попыток ограничено настройкой SHORTS_IMAGE_REVIEW_ATTEMPTS: пользовательский цикл тоже не должен быть бесконечным.
В v1 был булев флаг --video-image-only. Со временем стало ясно, что двух состояний недостаточно.
В v2 у video adapter есть четыре режима:
class VideoPromptMode(StrEnum):
COMBINED = "combined" # image prompt + motion prompt
MOTION = "motion" # только движение
IMAGE = "image" # только визуальное описание
NONE = "none" # только исходная картинка
Запуск выглядит так:
shorts-maker run <project-id> --video-prompt motion
Это небольшой пример полезной эволюции модели: вместо флага с неочевидной семантикой появился тип, перечисляющий все допустимые политики.
Для полного запуска достаточно одной команды:
uv run shorts-maker make "The Trial" --style "retro pixel art" --profile real
Можно разделить создание и выполнение:
uv run shorts-maker new "1984" --style "Soviet constructivism"
uv run shorts-maker run <project-id> --until images --profile real
uv run shorts-maker status <project-id>
uv run shorts-maker run <project-id> --profile real
--until images полезен не только для отладки. Можно сначала построить план и визуалы, просмотреть их, а уже потом запускать наиболее дорогую стадию image-to-video.
Для автономной проверки есть fake profile:
uv run shorts-maker make "The Trial" --style "pixel art" --profile fake
Он создаёт детерминированные тестовые артефакты и позволяет проверить orchestration, persistence и resume без API-ключей.
Платные вызовы не входят в обычный test suite. Вместо этого тесты разделены по уровням:
domain tests проверяют инварианты и переходы состояний;
architecture tests контролируют направление импортов;
application и integration tests запускают use cases с fake-портами;
contract tests проверяют HTTP-адаптеры через httpx.MockTransport;
smoke test действительно собирает короткий медиафайл через MoviePy и FFmpeg.
Один из наиболее ценных integration-сценариев специально ломает генерацию изображения для одной сцены. Первый запуск должен завершить стадию как partial, а второй — запросить только отсутствующую сцену.
Отдельный тест загружает один проект через два экземпляра репозитория и проверяет, что устаревший агрегат не может затереть более новую ревизию.
Для такого проекта тесты на happy path важны, но тесты на повторный запуск и частичный сбой важнее.
Python 3.12+;
asyncio — конкурентное выполнение I/O-задач;
dataclasses — чистая доменная модель;
Pydantic 2 и pydantic-settings — валидация внешних данных и конфигурации;
httpx — асинхронные provider adapters;
Typer и Rich — CLI;
filelock — межпроцессная блокировка manifest;
Pillow — подготовка изображений;
MoviePy и FFmpeg — финальный рендер;
pytest, mypy и Ruff — тесты и статический контроль.
LangChain, который использовался в первой версии, в v2 больше не является обязательной частью ядра. Для одного структурированного LLM-вызова прямой адаптер через HTTP оказался проще и прозрачнее.
Интерфейс VideoGenerator позволяет заменить Kling на другой сервис. Но он не отвечает на вопросы, что делать после timeout, когда сохранять результат и можно ли безопасно повторить команду. Для этого нужны явная модель состояния и семантика выполнения.
Повторять [2] производство целиком слишком дорого. В моём случае естественной единицей стала одна сцена внутри одной стадии.
Запустить десять coroutine просто. Гораздо сложнее решить, кто и в каком порядке меняет aggregate, когда пишется checkpoint и что произойдёт при двух процессах.
Один локальный JSON может быть разумным решением, если есть schema version, atomic replace, fsync, file lock и optimistic revision check. База данных нужна по требованиям масштаба, а не для красоты диаграммы.
Он позволяет воспроизводимо проверить весь application flow. Без него любая правка orchestration либо требует денег, либо остаётся непроверенной.
v2 пока остаётся альфа-версией. Ближайшие направления развития:
мультиязычная озвучка и отдельный финальный ролик для каждой locale;
web-интерфейс с просмотром и редактированием production plan;
ручное подтверждение не только изображений, но и сценария;
garbage collection orphaned-артефактов;
метрики стоимости и времени по провайдерам;
распределённые workers, когда локального процесса станет недостаточно;
публикация готовых роликов в Shorts, Reels и TikTok как отдельный bounded context.
Если вы строили похожие AI-pipeline’ы, интересно сравнить подходы: какую единицу повтора вы выбрали, как храните частичный прогресс и где проводите границу между автоматикой и human in the loop?
Автор: Ykrops
Источник [3]
Сайт-источник BrainTools: https://www.braintools.ru
Путь до страницы источника: https://www.braintools.ru/article/33677
URLs in this post:
[1] ошибки: http://www.braintools.ru/article/4192
[2] Повторять: http://www.braintools.ru/article/4012
[3] Источник: https://habr.com/ru/articles/1064134/?utm_source=habrahabr&utm_medium=rss&utm_campaign=1064134
Нажмите здесь для печати.