Привет, Хабр!
Некоторое время назад я написал первую версию генератора коротких видео для 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, сетевой ошибки или падения машины это знание теряется.
Повторный запуск не был безопасным
Ретрай всего пайплайна мог снова вызвать LLM, генератор изображений или image-to-video API. Для локальной функции это просто повторный вызов. Для генеративного сервиса — новые минуты ожидания, новый результат и новая сумма в биллинге.
gather() скрывал частичный успех
Допустим, из десяти изображений девять сгенерировались, а одно упало по timeout. Бизнес-состояние в этот момент не равно ни «успех», ни «ошибка». Это частично завершённая стадия, и её нужно уметь сохранить и продолжить.
Pydantic-модели стали выполнять слишком много ролей
Одни и те же модели постепенно использовались как ответы API, внутренние сущности и формат хранения на диске. В результате изменение JSON провайдера начинало влиять на бизнес-логику.
Главный вывод после v1 оказался таким: меняются не только провайдеры. Меняется и само состояние долгоживущего процесса. Значит, оно должно быть частью предметной модели.
Что представляет собой v2
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 ────────────+
Domain
Здесь находятся Project, сцены, персонажи, value objects, инварианты и переходы между стадиями. Домен ничего не знает о Pydantic, HTTP, MoviePy, файловой системе и CLI.
Application
Здесь живут use cases: создать проект, построить план, сгенерировать изображения, речь и видео, выполнить рендер. Этот же слой объявляет исходящие порты.
Infrastructure
Реализации портов: HTTP-клиенты провайдеров, JSON-репозиторий, хранилище медиа и MoviePy-рендерер.
Presentation
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 как агрегат
В центре домена находится агрегат 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 знает, что стадия уже завершена, и не вызывает провайдера повторно.
Не просто success или failed
У каждой стадии есть один из пяти статусов:
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 переводит стадию в работу и продолжает с отсутствующих результатов.
Параллелизм с checkpoint после каждой сцены
В 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(), зато семантика частичного выполнения становится явной.
Один JSON-манифест как consistency boundary
Для 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 при этом не знает, какой профиль запущен.
Внешний API всегда ненадёжен
Провайдеры генерации отличаются не только 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.
Human in the loop для изображений
Полностью автоматическая генерация удобна до первой сцены с шестью пальцами, сломанной перспективой или внезапно изменившимся персонажем.
В v2 появился флаг:
shorts-maker run <project-id> --review-images
После генерации CLI открывает изображение и предлагает:
принять результат;
перегенерировать с тем же prompt;
отредактировать prompt и повторить попытку.
Сам review тоже является портом ImageReviewer. В автоматическом режиме используется AutoApproveImageReviewer, а в интерактивном — реализация для консоли. Поэтому human approval встроен в use case, но UI не протекает в бизнес-логику.
Число попыток ограничено настройкой SHORTS_IMAGE_REVIEW_ATTEMPTS: пользовательский цикл тоже не должен быть бесконечным.
Четыре режима prompt для image-to-video
В 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
Это небольшой пример полезной эволюции модели: вместо флага с неочевидной семантикой появился тип, перечисляющий все допустимые политики.
CLI: создать, продолжить, проверить
Для полного запуска достаточно одной команды:
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-ключей.
Как тестируется pipeline с дорогими 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 важны, но тесты на повторный запуск и частичный сбой важнее.
Стек v2
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 оказался проще и прозрачнее.
Что я понял после переписывания
Clean Architecture сама по себе не даёт надёжность
Интерфейс VideoGenerator позволяет заменить Kling на другой сервис. Но он не отвечает на вопросы, что делать после timeout, когда сохранять результат и можно ли безопасно повторить команду. Для этого нужны явная модель состояния и семантика выполнения.
Единица восстановления должна быть меньше всего pipeline
Повторять производство целиком слишком дорого. В моём случае естественной единицей стала одна сцена внутри одной стадии.
Конкурентность и консистентность нужно проектировать вместе
Запустить десять coroutine просто. Гораздо сложнее решить, кто и в каком порядке меняет aggregate, когда пишется checkpoint и что произойдёт при двух процессах.
Файловое хранилище не обязательно означает ненадёжность
Один локальный JSON может быть разумным решением, если есть schema version, atomic replace, fsync, file lock и optimistic revision check. База данных нужна по требованиям масштаба, а не для красоты диаграммы.
Fake provider — часть архитектуры, а не временная заглушка
Он позволяет воспроизводимо проверить весь application flow. Без него любая правка orchestration либо требует денег, либо остаётся непроверенной.
Что дальше
v2 пока остаётся альфа-версией. Ближайшие направления развития:
мультиязычная озвучка и отдельный финальный ролик для каждой locale;
web-интерфейс с просмотром и редактированием production plan;
ручное подтверждение не только изображений, но и сценария;
garbage collection orphaned-артефактов;
метрики стоимости и времени по провайдерам;
распределённые workers, когда локального процесса станет недостаточно;
публикация готовых роликов в Shorts, Reels и TikTok как отдельный bounded context.
Если вы строили похожие AI-pipeline’ы, интересно сравнить подходы: какую единицу повтора вы выбрали, как храните частичный прогресс и где проводите границу между автоматикой и human in the loop?