| 1 | # TODO: NOT IMPLEMENTED YET |
| 2 | |
| 3 | import os |
| 4 | import shutil |
| 5 | import yaml |
| 6 | import json |
| 7 | import importlib |
| 8 | import asyncio |
| 9 | from typing import Any, Callable, List, Dict |
| 10 | from langchain.embeddings import CacheBackedEmbeddings |
| 11 | from langchain.storage import LocalFileStore |
| 12 | from langchain_text_splitters import RecursiveCharacterTextSplitter |
| 13 | from langchain_community.vectorstores import FAISS |
| 14 | from PIL import Image |
| 15 | |
| 16 | from interfaces import ( |
| 17 | Event, |
| 18 | Scene, |
| 19 | CharacterInScene, |
| 20 | CharacterInNovel, |
| 21 | CharacterInEvent, |
| 22 | ) |
| 23 | from tenacity import retry |
| 24 | |
| 25 | from utils.text import safe_path_component |
| 26 | |
| 27 | |
| 28 | |
| 29 | def _pipeline_print(quiet: bool, message: str) -> None: |
| 30 | if not quiet: |
| 31 | print(message) |
| 32 | |
| 33 | |
| 34 | def _emit_text_plan_progress(progress, stage: str, message: str, metadata: dict | None = None) -> None: |
| 35 | if progress is not None: |
| 36 | progress(stage, message, metadata or {}) |
| 37 | |
| 38 | |
| 39 | def _event_file_index(path: str) -> int: |
| 40 | return int(os.path.basename(path).split("_")[1].split(".")[0]) |
| 41 | |
| 42 | |
| 43 | def _scene_file_index(path: str) -> int: |
| 44 | return int(os.path.basename(path).split("_")[1].split(".")[0]) |
| 45 | |
| 46 | class Novel2MoviePipeline: |
| 47 | def __init__( |
| 48 | self, |
| 49 | novel_compressor: Any, |
| 50 | event_extractor: Any, |
| 51 | embeddings: Any, |
| 52 | rerank_model: Any, |
| 53 | scene_extractor: Any, |
| 54 | global_information_planner: Any, |
| 55 | image_generator: Any, |
| 56 | rewriter: Any, |
| 57 | script2video_pipeline: Any, |
| 58 | working_dir: str, |
| 59 | ): |
| 60 | self.novel_compressor = novel_compressor |
| 61 | self.event_extractor = event_extractor |
| 62 | self.embeddings = embeddings |
| 63 | self.rerank_model = rerank_model |
| 64 | self.scene_extractor = scene_extractor |
| 65 | self.global_information_planner = global_information_planner |
| 66 | self.image_generator = image_generator |
| 67 | self.rewriter = rewriter |
| 68 | self.script2video_pipeline = script2video_pipeline |
| 69 | self.working_dir = working_dir |
| 70 | os.makedirs(self.working_dir, exist_ok=True) |
| 71 | |
| 72 | |
| 73 | async def plan_text_artifacts( |
| 74 | self, |
| 75 | novel_text: str, |
| 76 | user_requirement: str = "", |
| 77 | style: str = "", |
| 78 | progress: Callable[[str, str, Dict[str, Any] | None], None] | None = None, |
| 79 | quiet: bool = False, |
| 80 | ) -> dict[str, Any]: |
| 81 | """Generate structured text artifacts for novel adaptation only. |
| 82 | |
| 83 | This helper intentionally stops before character portrait generation, |
| 84 | scene video generation, and final concatenation so the agent loop can |
| 85 | pause after the novel planning stage. |
| 86 | """ |
| 87 | del user_requirement, style |
| 88 | |
| 89 | _emit_text_plan_progress(progress, "save_novel", "Saving and splitting novel text") |
| 90 | working_dir_novel = os.path.join(self.working_dir, "novel") |
| 91 | os.makedirs(working_dir_novel, exist_ok=True) |
| 92 | with open(os.path.join(working_dir_novel, "novel.txt"), "w", encoding="utf-8") as f: |
| 93 | f.write(novel_text) |
| 94 | |
| 95 | novel_chunks = self.novel_compressor.split(novel_text) |
| 96 | for idx, novel_chunk in enumerate(novel_chunks): |
| 97 | with open(os.path.join(working_dir_novel, f"novel_chunk_{idx}.txt"), "w", encoding="utf-8") as f: |
| 98 | f.write(novel_chunk) |
| 99 | _pipeline_print(quiet, f"Split novel into {len(novel_chunks)} chunks.") |
| 100 | |
| 101 | _emit_text_plan_progress(progress, "compress_novel", "Compressing novel chunks", {"chunk_count": len(novel_chunks)}) |
| 102 | compressed_novel_chunks: list[str | None] = [None] * len(novel_chunks) |
| 103 | unfinished_pairs = [] |
| 104 | for index, novel_chunk in enumerate(novel_chunks): |
| 105 | path = os.path.join(working_dir_novel, f"novel_chunk_{index}_compressed.txt") |
| 106 | if os.path.exists(path): |
| 107 | compressed_novel_chunks[index] = open(path, "r", encoding="utf-8").read() |
| 108 | else: |
| 109 | unfinished_pairs.append((index, novel_chunk)) |
| 110 | if unfinished_pairs: |
| 111 | sem = asyncio.Semaphore(5) |
| 112 | outputs = await asyncio.gather(*[ |
| 113 | self.novel_compressor.compress_single_novel_chunk(sem, index, novel_chunk) |
| 114 | for index, novel_chunk in unfinished_pairs |
| 115 | ]) |
| 116 | for index, compressed in outputs: |
| 117 | path = os.path.join(working_dir_novel, f"novel_chunk_{index}_compressed.txt") |
| 118 | with open(path, "w", encoding="utf-8") as f: |
| 119 | f.write(compressed) |
| 120 | compressed_novel_chunks[index] = compressed |
| 121 | |
| 122 | compressed_path = os.path.join(working_dir_novel, "novel_compressed.txt") |
| 123 | if os.path.exists(compressed_path): |
| 124 | compressed_novel = open(compressed_path, "r", encoding="utf-8").read() |
| 125 | else: |
| 126 | compressed_novel = self.novel_compressor.aggregate([chunk or "" for chunk in compressed_novel_chunks]) |
| 127 | with open(compressed_path, "w", encoding="utf-8") as f: |
| 128 | f.write(compressed_novel) |
| 129 | |
| 130 | _emit_text_plan_progress(progress, "extract_events", "Extracting events from compressed novel") |
| 131 | working_dir_events = os.path.join(self.working_dir, "events") |
| 132 | os.makedirs(working_dir_events, exist_ok=True) |
| 133 | extracted_events: list[Event] = [] |
| 134 | event_files = [ |
| 135 | os.path.join(working_dir_events, fname) |
| 136 | for fname in os.listdir(working_dir_events) |
| 137 | if fname.startswith("event_") and fname.endswith(".json") |
| 138 | ] |
| 139 | for event_path in sorted(event_files, key=_event_file_index): |
| 140 | with open(event_path, "r", encoding="utf-8") as f: |
| 141 | extracted_events.append(Event.model_validate(json.load(f))) |
| 142 | while len(extracted_events) == 0 or not extracted_events[-1].is_last: |
| 143 | _ensure_extraction_cap(len(extracted_events), MAX_EXTRACTED_EVENTS, "events") |
| 144 | next_event = self.event_extractor.extract_next_event( |
| 145 | novel_text=compressed_novel, |
| 146 | extracted_events=extracted_events, |
| 147 | ) |
| 148 | event_path = os.path.join(working_dir_events, f"event_{len(extracted_events)}.json") |
| 149 | with open(event_path, "w", encoding="utf-8") as f: |
| 150 | json.dump(next_event.model_dump(), f, ensure_ascii=False, indent=4) |
| 151 | extracted_events.append(next_event) |
| 152 | |
| 153 | _emit_text_plan_progress(progress, "retrieve_chunks", "Retrieving relevant chunks for events", {"event_count": len(extracted_events)}) |
| 154 | working_dir_knowledge_base = os.path.join(self.working_dir, "knowledge_base") |
| 155 | working_dir_retrieve = os.path.join(self.working_dir, "relevant_chunks") |
| 156 | os.makedirs(working_dir_knowledge_base, exist_ok=True) |
| 157 | os.makedirs(working_dir_retrieve, exist_ok=True) |
| 158 | embeddings = CacheBackedEmbeddings.from_bytes_store( |
| 159 | underlying_embeddings=self.embeddings, |
| 160 | document_embedding_cache=LocalFileStore(root_path=working_dir_knowledge_base), |
| 161 | namespace=getattr(self.embeddings, "model", "default"), |
| 162 | key_encoder="sha256", |
| 163 | ) |
| 164 | novel_splitter = RecursiveCharacterTextSplitter(chunk_size=512, chunk_overlap=128) |
| 165 | knowledge_chunks = novel_splitter.split_text(novel_text) |
| 166 | knowledge_base = FAISS.from_texts(texts=knowledge_chunks, embedding=embeddings) |
| 167 | event_idx_to_relevant_chunk_score_dict: dict[int, dict[str, float]] = {} |
| 168 | |
| 169 | async def retrieve_relevant_chunks(sem, event: Event): |
| 170 | async with sem: |
| 171 | relevant: dict[str, float] = {} |
| 172 | for process in event.process_chain: |
| 173 | chunks = knowledge_base.similarity_search(process, k=10) |
| 174 | chunk_texts = [chunk.page_content for chunk in chunks if chunk.page_content not in relevant] |
| 175 | if not chunk_texts: |
| 176 | continue |
| 177 | chunk_score_pairs = await self.rerank_model(documents=chunk_texts, query=process, top_n=10) |
| 178 | for chunk, score in chunk_score_pairs: |
| 179 | if score >= 0.7: |
| 180 | relevant[chunk] = relevant.get(chunk, 0.0) + score |
| 181 | return event.index, relevant |
| 182 | |
| 183 | retrieve_tasks = [] |
| 184 | retrieve_sem = asyncio.Semaphore(10) |
| 185 | for event in extracted_events: |
| 186 | chunks_dir = os.path.join(working_dir_retrieve, f"event_{event.index}") |
| 187 | if os.path.exists(chunks_dir) and os.listdir(chunks_dir): |
| 188 | relevant = {} |
| 189 | for chunk_fname in os.listdir(chunks_dir): |
| 190 | chunk_path = os.path.join(chunks_dir, chunk_fname) |
| 191 | score = float(chunk_fname.split("-score_")[1].split(".txt")[0]) |
| 192 | with open(chunk_path, "r", encoding="utf-8") as f: |
| 193 | relevant[f.read()] = score |
| 194 | event_idx_to_relevant_chunk_score_dict[event.index] = relevant |
| 195 | else: |
| 196 | retrieve_tasks.append(retrieve_relevant_chunks(retrieve_sem, event)) |
| 197 | if retrieve_tasks: |
| 198 | for event_index, relevant in await asyncio.gather(*retrieve_tasks): |
| 199 | chunks_dir = os.path.join(working_dir_retrieve, f"event_{event_index}") |
| 200 | os.makedirs(chunks_dir, exist_ok=True) |
| 201 | for idx, (chunk, score) in enumerate(relevant.items()): |
| 202 | with open(os.path.join(chunks_dir, f"chunk_{idx}-score_{score:.2f}.txt"), "w", encoding="utf-8") as f: |
| 203 | f.write(chunk) |
| 204 | event_idx_to_relevant_chunk_score_dict[event_index] = relevant |
| 205 | |
| 206 | _emit_text_plan_progress(progress, "extract_scenes", "Extracting screenplay scenes", {"event_count": len(extracted_events)}) |
| 207 | working_dir_scenes = os.path.join(self.working_dir, "scenes") |
| 208 | os.makedirs(working_dir_scenes, exist_ok=True) |
| 209 | event_idx_to_scenes: dict[int, list[Scene]] = {event.index: [] for event in extracted_events} |
| 210 | unfinished_events: list[Event] = [] |
| 211 | for event in extracted_events: |
| 212 | scenes_dir = os.path.join(working_dir_scenes, f"event_{event.index}") |
| 213 | if os.path.exists(scenes_dir): |
| 214 | scene_files = [ |
| 215 | os.path.join(scenes_dir, fname) |
| 216 | for fname in os.listdir(scenes_dir) |
| 217 | if fname.startswith("scene_") and fname.endswith(".json") |
| 218 | ] |
| 219 | for scene_path in sorted(scene_files, key=_scene_file_index): |
| 220 | with open(scene_path, "r", encoding="utf-8") as f: |
| 221 | event_idx_to_scenes[event.index].append(Scene.model_validate(json.load(f))) |
| 222 | if not event_idx_to_scenes[event.index] or not event_idx_to_scenes[event.index][-1].is_last: |
| 223 | unfinished_events.append(event) |
| 224 | |
| 225 | async def extract_scenes_for_event(sem, event: Event, previous_scenes: list[Scene]): |
| 226 | async with sem: |
| 227 | scenes_dir = os.path.join(working_dir_scenes, f"event_{event.index}") |
| 228 | os.makedirs(scenes_dir, exist_ok=True) |
| 229 | while len(previous_scenes) == 0 or not previous_scenes[-1].is_last: |
| 230 | _ensure_extraction_cap(len(previous_scenes), MAX_SCENES_PER_EVENT, "scenes") |
| 231 | next_scene = await self.scene_extractor.get_next_scene( |
| 232 | relevant_chunks=list(event_idx_to_relevant_chunk_score_dict.get(event.index, {}).keys()), |
| 233 | event=event, |
| 234 | previous_scenes=previous_scenes, |
| 235 | ) |
| 236 | scene_path = os.path.join(scenes_dir, f"scene_{len(previous_scenes)}.json") |
| 237 | with open(scene_path, "w", encoding="utf-8") as f: |
| 238 | json.dump(next_scene.model_dump(), f, ensure_ascii=False, indent=4) |
| 239 | previous_scenes.append(next_scene) |
| 240 | return event.index, previous_scenes |
| 241 | |
| 242 | if unfinished_events: |
| 243 | sem = asyncio.Semaphore(8) |
| 244 | scene_outputs = await asyncio.gather(*[ |
| 245 | extract_scenes_for_event(sem, event, event_idx_to_scenes[event.index]) |
| 246 | for event in unfinished_events |
| 247 | ]) |
| 248 | for event_index, scenes in scene_outputs: |
| 249 | event_idx_to_scenes[event_index] = scenes |
| 250 | |
| 251 | _emit_text_plan_progress(progress, "merge_characters", "Merging scene characters into novel-level characters", {"event_count": len(extracted_events)}) |
| 252 | working_dir_global = os.path.join(self.working_dir, "global_information") |
| 253 | working_dir_characters = os.path.join(working_dir_global, "characters") |
| 254 | os.makedirs(working_dir_characters, exist_ok=True) |
| 255 | event_idx_to_characters_in_event: dict[int, list[CharacterInEvent]] = {} |
| 256 | |
| 257 | async def merge_event_characters(sem, event: Event): |
| 258 | async with sem: |
| 259 | characters = await self.global_information_planner.merge_characters_across_scenes_in_event( |
| 260 | event_idx=event.index, |
| 261 | scenes=event_idx_to_scenes[event.index], |
| 262 | ) |
| 263 | path = os.path.join(working_dir_characters, "event_level", f"event_{event.index}_characters.json") |
| 264 | os.makedirs(os.path.dirname(path), exist_ok=True) |
| 265 | with open(path, "w", encoding="utf-8") as f: |
| 266 | json.dump([char.model_dump() for char in characters], f, ensure_ascii=False, indent=4) |
| 267 | return event.index, characters |
| 268 | |
| 269 | merge_tasks = [] |
| 270 | merge_sem = asyncio.Semaphore(8) |
| 271 | for event in extracted_events: |
| 272 | path = os.path.join(working_dir_characters, "event_level", f"event_{event.index}_characters.json") |
| 273 | if os.path.exists(path): |
| 274 | with open(path, "r", encoding="utf-8") as f: |
| 275 | event_idx_to_characters_in_event[event.index] = [CharacterInEvent.model_validate(item) for item in json.load(f)] |
| 276 | else: |
| 277 | merge_tasks.append(merge_event_characters(merge_sem, event)) |
| 278 | if merge_tasks: |
| 279 | for event_index, characters in await asyncio.gather(*merge_tasks): |
| 280 | event_idx_to_characters_in_event[event_index] = characters |
| 281 | |
| 282 | working_dir_novel_chars = os.path.join(working_dir_characters, "novel_level") |
| 283 | os.makedirs(working_dir_novel_chars, exist_ok=True) |
| 284 | existing_files = [fname for fname in os.listdir(working_dir_novel_chars) if fname.startswith("novel_characters_after_event_") and fname.endswith(".json")] |
| 285 | if existing_files: |
| 286 | latest = max(existing_files, key=lambda fname: int(fname.split("_")[-1].split(".json")[0])) |
| 287 | start_event_idx = int(latest.split("_")[-1].split(".json")[0]) + 1 |
| 288 | with open(os.path.join(working_dir_novel_chars, latest), "r", encoding="utf-8") as f: |
| 289 | characters_in_novel = [CharacterInNovel.model_validate(item) for item in json.load(f)] |
| 290 | else: |
| 291 | start_event_idx = 0 |
| 292 | characters_in_novel = [] |
| 293 | for event in extracted_events[start_event_idx:]: |
| 294 | characters_in_novel = self.global_information_planner.merge_characters_to_existing_characters_in_novel( |
| 295 | event_idx=event.index, |
| 296 | existing_characters_in_novel=characters_in_novel, |
| 297 | characters_in_event=event_idx_to_characters_in_event[event.index], |
| 298 | ) |
| 299 | path = os.path.join(working_dir_novel_chars, f"novel_characters_after_event_{event.index}.json") |
| 300 | with open(path, "w", encoding="utf-8") as f: |
| 301 | json.dump([char.model_dump() for char in characters_in_novel], f, ensure_ascii=False, indent=4) |
| 302 | |
| 303 | _emit_text_plan_progress(progress, "completed", "Novel structured text planning complete", {"event_count": len(extracted_events)}) |
| 304 | return { |
| 305 | "compressed_novel": compressed_novel, |
| 306 | "events": extracted_events, |
| 307 | "scenes": event_idx_to_scenes, |
| 308 | "characters_in_novel": characters_in_novel, |
| 309 | } |
| 310 | |
| 311 | |
| 312 | async def render_video_artifacts( |
| 313 | self, |
| 314 | style: str, |
| 315 | user_requirement: str = "", |
| 316 | progress: Callable[[str, str, Dict[str, Any] | None], None] | None = None, |
| 317 | quiet: bool = False, |
| 318 | ) -> dict[str, Any]: |
| 319 | """Render portraits and per-scene videos from existing novel planning artifacts. |
| 320 | |
| 321 | This helper assumes plan_text_artifacts has already completed. It does not |
| 322 | re-run compression, event extraction, RAG retrieval, scene extraction, or |
| 323 | character merging. |
| 324 | """ |
| 325 | del user_requirement |
| 326 | |
| 327 | _emit_text_plan_progress(progress, "novel_render_load", "Loading novel structured text artifacts") |
| 328 | working_dir_events = os.path.join(self.working_dir, "events") |
| 329 | working_dir_scenes = os.path.join(self.working_dir, "scenes") |
| 330 | working_dir_characters = os.path.join(self.working_dir, "global_information", "characters") |
| 331 | event_level_dir = os.path.join(working_dir_characters, "event_level") |
| 332 | novel_level_dir = os.path.join(working_dir_characters, "novel_level") |
| 333 | |
| 334 | if not os.path.isdir(working_dir_events): |
| 335 | raise RuntimeError("novel2video/events is missing; run vimax_novel_planning first") |
| 336 | if not os.path.isdir(working_dir_scenes): |
| 337 | raise RuntimeError("novel2video/scenes is missing; run vimax_novel_planning first") |
| 338 | if not os.path.isdir(event_level_dir) or not os.path.isdir(novel_level_dir): |
| 339 | raise RuntimeError("novel2video/global_information/characters is missing; run vimax_novel_planning first") |
| 340 | |
| 341 | event_files = [ |
| 342 | os.path.join(working_dir_events, fname) |
| 343 | for fname in os.listdir(working_dir_events) |
| 344 | if fname.startswith("event_") and fname.endswith(".json") |
| 345 | ] |
| 346 | extracted_events = [] |
| 347 | for event_path in sorted(event_files, key=_event_file_index): |
| 348 | with open(event_path, "r", encoding="utf-8") as f: |
| 349 | extracted_events.append(Event.model_validate(json.load(f))) |
| 350 | if not extracted_events: |
| 351 | raise RuntimeError("novel2video/events has no event_*.json files") |
| 352 | |
| 353 | event_idx_to_scenes: dict[int, list[Scene]] = {} |
| 354 | for event in extracted_events: |
| 355 | scenes_dir = os.path.join(working_dir_scenes, f"event_{event.index}") |
| 356 | if not os.path.isdir(scenes_dir): |
| 357 | raise RuntimeError(f"novel2video/scenes/event_{event.index} is missing") |
| 358 | scene_files = [ |
| 359 | os.path.join(scenes_dir, fname) |
| 360 | for fname in os.listdir(scenes_dir) |
| 361 | if fname.startswith("scene_") and fname.endswith(".json") |
| 362 | ] |
| 363 | scenes = [] |
| 364 | for scene_path in sorted(scene_files, key=_scene_file_index): |
| 365 | with open(scene_path, "r", encoding="utf-8") as f: |
| 366 | scenes.append(Scene.model_validate(json.load(f))) |
| 367 | if not scenes: |
| 368 | raise RuntimeError(f"novel2video/scenes/event_{event.index} has no scene_*.json files") |
| 369 | event_idx_to_scenes[event.index] = scenes |
| 370 | |
| 371 | event_idx_to_characters_in_event: dict[int, list[CharacterInEvent]] = {} |
| 372 | for event in extracted_events: |
| 373 | path = os.path.join(event_level_dir, f"event_{event.index}_characters.json") |
| 374 | if not os.path.exists(path): |
| 375 | raise RuntimeError(f"novel2video/global_information/characters/event_level/event_{event.index}_characters.json is missing") |
| 376 | with open(path, "r", encoding="utf-8") as f: |
| 377 | event_idx_to_characters_in_event[event.index] = [CharacterInEvent.model_validate(item) for item in json.load(f)] |
| 378 | |
| 379 | novel_files = [fname for fname in os.listdir(novel_level_dir) if fname.startswith("novel_characters_after_event_") and fname.endswith(".json")] |
| 380 | if not novel_files: |
| 381 | raise RuntimeError("novel2video/global_information/characters/novel_level has no novel characters file") |
| 382 | latest_novel_file = max(novel_files, key=lambda fname: int(fname.split("_")[-1].split(".json")[0])) |
| 383 | with open(os.path.join(novel_level_dir, latest_novel_file), "r", encoding="utf-8") as f: |
| 384 | characters_in_novel = [CharacterInNovel.model_validate(item) for item in json.load(f)] |
| 385 | |
| 386 | _emit_text_plan_progress(progress, "novel_portraits_start", "Generating novel character portraits", {"character_count": len(characters_in_novel)}) |
| 387 | working_dir_character_portrait = os.path.join(self.working_dir, "character_portraits") |
| 388 | base_character_portrait_dir = os.path.join(working_dir_character_portrait, "base") |
| 389 | os.makedirs(base_character_portrait_dir, exist_ok=True) |
| 390 | |
| 391 | async def generate_base_portrait(sem, character: CharacterInNovel): |
| 392 | async with sem: |
| 393 | image_path = os.path.join(base_character_portrait_dir, f"character_{character.index}_{safe_path_component(character.identifier_in_novel)}.png") |
| 394 | if os.path.exists(image_path): |
| 395 | return image_path |
| 396 | prompt = f"Generate a full-body, front-view portrait based on the following description, in the style of {style}:" |
| 397 | prompt += f"\nCharacter Identifier: {character.identifier_in_novel}" |
| 398 | prompt += f"\nFeatures: {character.static_features}" |
| 399 | prompt += "\nThe character should be centered in the image, occupying most of the frame. Gazing straight ahead. Standing with arms relaxed at sides. Natural expression. The background should be plain white." |
| 400 | image = await self.image_generator.generate_single_image(prompt=prompt, size="512x512") |
| 401 | image.save(image_path) |
| 402 | return image_path |
| 403 | |
| 404 | sem = asyncio.Semaphore(5) |
| 405 | await asyncio.gather(*[generate_base_portrait(sem, character) for character in characters_in_novel]) |
| 406 | _emit_text_plan_progress(progress, "novel_portraits_base_done", "Base character portraits ready", {"character_count": len(characters_in_novel)}) |
| 407 | |
| 408 | async def generate_scene_portrait(sem, base_character_image_path: str, character: CharacterInScene, event_idx: int, scene_idx: int): |
| 409 | async with sem: |
| 410 | image_path = os.path.join(working_dir_character_portrait, f"event_{event_idx}", f"scene_{scene_idx}", f"character_{character.idx}_{safe_path_component(character.identifier_in_scene)}.png") |
| 411 | os.makedirs(os.path.dirname(image_path), exist_ok=True) |
| 412 | if os.path.exists(image_path): |
| 413 | return image_path |
| 414 | if not character.is_visible or character.dynamic_features is None: |
| 415 | shutil.copy(base_character_image_path, image_path) |
| 416 | return image_path |
| 417 | prompt = f"Generate a full-body, front-view portrait based on the provided base image. Modify the base image according to the following dynamic features, in the style of {style}. Keep the character's identity consistent with the base image:" |
| 418 | prompt += f"\nCharacter Identifier: {character.identifier_in_scene}" |
| 419 | prompt += f"\nDynamic Features: {character.dynamic_features}" |
| 420 | prompt += "\nThe character should be centered in the image, occupying most of the frame. Gazing straight ahead. Standing with arms relaxed at sides. Natural expression. The background should be plain white." |
| 421 | prompt = await self.rewriter(prompt) |
| 422 | image = await self.image_generator.generate_single_image(prompt=prompt, reference_image_paths=[base_character_image_path], size="512x512") |
| 423 | image.save(image_path) |
| 424 | return image_path |
| 425 | |
| 426 | _emit_text_plan_progress(progress, "novel_portraits_scene_start", "Generating scene character portraits") |
| 427 | scene_portrait_tasks = [] |
| 428 | sem = asyncio.Semaphore(3) |
| 429 | for character in characters_in_novel: |
| 430 | base_path = os.path.join(base_character_portrait_dir, f"character_{character.index}_{safe_path_component(character.identifier_in_novel)}.png") |
| 431 | for event_idx, identifier_in_event in character.active_events.items(): |
| 432 | event_characters = event_idx_to_characters_in_event[int(event_idx)] |
| 433 | character_in_event = [char for char in event_characters if char.identifier_in_event == identifier_in_event][0] |
| 434 | for scene_idx, identifier_in_scene in character_in_event.active_scenes.items(): |
| 435 | scene = event_idx_to_scenes[int(event_idx)][int(scene_idx)] |
| 436 | character_in_scene = [char for char in scene.characters if char.identifier_in_scene == identifier_in_scene][0] |
| 437 | scene_portrait_tasks.append(generate_scene_portrait(sem, base_path, character_in_scene, int(event_idx), int(scene_idx))) |
| 438 | if scene_portrait_tasks: |
| 439 | await asyncio.gather(*scene_portrait_tasks) |
| 440 | _emit_text_plan_progress(progress, "novel_portraits_done", "Scene character portraits ready") |
| 441 | |
| 442 | working_dir_scene_videos = os.path.join(self.working_dir, "videos") |
| 443 | os.makedirs(working_dir_scene_videos, exist_ok=True) |
| 444 | scene_video_dirs: list[str] = [] |
| 445 | for event in extracted_events: |
| 446 | for scene in event_idx_to_scenes[event.index]: |
| 447 | scene_video_dir = os.path.join(working_dir_scene_videos, f"event_{event.index}", f"scene_{scene.idx}") |
| 448 | os.makedirs(scene_video_dir, exist_ok=True) |
| 449 | self.script2video_pipeline.working_dir = scene_video_dir |
| 450 | character_portraits_registry = {} |
| 451 | for character in scene.characters: |
| 452 | character_portraits_registry[character.identifier_in_scene] = { |
| 453 | "portrait": { |
| 454 | "path": os.path.join(working_dir_character_portrait, f"event_{event.index}", f"scene_{scene.idx}", f"character_{character.idx}_{safe_path_component(character.identifier_in_scene)}.png"), |
| 455 | "description": f"A portrait of {character.identifier_in_scene}", |
| 456 | } |
| 457 | } |
| 458 | _emit_text_plan_progress(progress, "novel_scene_render_start", "Rendering novel scene video", {"event_idx": event.index, "scene_idx": scene.idx}) |
| 459 | await self.script2video_pipeline( |
| 460 | script=scene.script, |
| 461 | user_requirement="", |
| 462 | style=style or "realistic movie style", |
| 463 | characters=scene.characters, |
| 464 | character_portraits_registry=character_portraits_registry, |
| 465 | quiet=quiet, |
| 466 | progress=progress, |
| 467 | ) |
| 468 | scene_video_dirs.append(scene_video_dir) |
| 469 | _emit_text_plan_progress(progress, "novel_scene_render_done", "Rendered novel scene video", {"event_idx": event.index, "scene_idx": scene.idx, "path": scene_video_dir}) |
| 470 | |
| 471 | _emit_text_plan_progress(progress, "novel_render_completed", "Novel scene render complete", {"scene_count": len(scene_video_dirs)}) |
| 472 | return { |
| 473 | "character_portraits_dir": working_dir_character_portrait, |
| 474 | "scene_videos_dir": working_dir_scene_videos, |
| 475 | "scene_video_dirs": scene_video_dirs, |
| 476 | "scene_count": len(scene_video_dirs), |
| 477 | } |
| 478 | |
| 479 | async def __call__( |
| 480 | self, |
| 481 | novel_text: str, |
| 482 | style: str, |
| 483 | ): |
| 484 | print("🎬 Novel to Movie Pipeline Started".center(80, "=")) |
| 485 | |
| 486 | # Step 1: Compress the novel text |
| 487 | print() |
| 488 | print("📋 Step 1: Compress the novel text".center(80, "-")) |
| 489 | |
| 490 | working_dir_novel_compressor = os.path.join(self.working_dir, "novel") |
| 491 | os.makedirs(working_dir_novel_compressor, exist_ok=True) |
| 492 | with open(os.path.join(working_dir_novel_compressor, "novel.txt"), "w", encoding="utf-8") as f: |
| 493 | f.write(novel_text) |
| 494 | print(f"🗂️ Working directory: {working_dir_novel_compressor}") |
| 495 | |
| 496 | print("🔖 Splitting the novel into chunks...") |
| 497 | novel_chunks = self.novel_compressor.split(novel_text) |
| 498 | for idx, novel_chunk in enumerate(novel_chunks): |
| 499 | with open(os.path.join(working_dir_novel_compressor, f"novel_chunk_{idx}.txt"), "w", encoding="utf-8") as f: |
| 500 | f.write(novel_chunk) |
| 501 | print(f"🔖 Split the novel into {len(novel_chunks)} chunks, all saved to {working_dir_novel_compressor}.") |
| 502 | |
| 503 | |
| 504 | print() |
| 505 | print("🔖 Compressing the novel chunks...") |
| 506 | compressed_novel_chunks = [None] * len(novel_chunks) |
| 507 | index_chunk_pairs_unfinished = [] |
| 508 | for index, novel_chunk in enumerate(novel_chunks): |
| 509 | path = os.path.join(working_dir_novel_compressor, f"novel_chunk_{index}_compressed.txt") |
| 510 | if os.path.exists(path): |
| 511 | compressed_novel_chunks[index] = open(path, "r", encoding="utf-8").read() |
| 512 | print(f"⏭️ Skipping compression for chunk {index} as it already exists.") |
| 513 | else: |
| 514 | index_chunk_pairs_unfinished.append((index, novel_chunk)) |
| 515 | |
| 516 | sem = asyncio.Semaphore(5) |
| 517 | tasks = [ |
| 518 | self.novel_compressor.compress_single_novel_chunk(sem, index, novel_chunk) |
| 519 | for index, novel_chunk in index_chunk_pairs_unfinished |
| 520 | ] |
| 521 | task_outputs = await asyncio.gather(*tasks) |
| 522 | for index, novel_chunk_compressed in task_outputs: |
| 523 | save_path = os.path.join(working_dir_novel_compressor, f"novel_chunk_{index}_compressed.txt") |
| 524 | with open(save_path, "w", encoding="utf-8") as f: |
| 525 | f.write(novel_chunk_compressed) |
| 526 | print(f"✅ Compressed chunk {index}, saved to {save_path}") |
| 527 | compressed_novel_chunks[index] = novel_chunk_compressed |
| 528 | print("🔖 Compressed all novel chunks.") |
| 529 | |
| 530 | |
| 531 | print() |
| 532 | print("🔖 Merging the compressed novel chunks...") |
| 533 | path = os.path.join(working_dir_novel_compressor, "novel_compressed.txt") |
| 534 | if os.path.exists(path): |
| 535 | compressed_novel = open(path, "r", encoding="utf-8").read() |
| 536 | print(f"⏭️ Skipping merging as {path} already exists.") |
| 537 | else: |
| 538 | compressed_novel = self.novel_compressor.aggregate(compressed_novel_chunks) |
| 539 | with open(path, "w", encoding="utf-8") as f: |
| 540 | f.write(compressed_novel) |
| 541 | print(f"✅ Merged the compressed novel chunks, saved to {path}") |
| 542 | print(f"🔖 Merging completed.") |
| 543 | |
| 544 | # summary |
| 545 | print() |
| 546 | print("📌 Summary:") |
| 547 | print(f"📌 Before Compression: {len(novel_text)} characters") |
| 548 | print(f"📌 After Compression: {len(compressed_novel)} characters") |
| 549 | print(f"📌 Compression Ratio: {len(compressed_novel) / len(novel_text):.2%}") |
| 550 | |
| 551 | print("📋 Step 1: Compress the novel text".center(80, "-")) |
| 552 | |
| 553 | |
| 554 | # Step 2: Extract events from the compressed novel |
| 555 | print() |
| 556 | print("📋 Step 2: Extract events from the compressed novel".center(80, "-")) |
| 557 | working_dir_event_extractor = os.path.join(self.working_dir, "events") |
| 558 | os.makedirs(working_dir_event_extractor, exist_ok=True) |
| 559 | print(f"🗂️ Working directory: {working_dir_event_extractor}") |
| 560 | |
| 561 | extracted_events = [] |
| 562 | for event_json_fname in sorted(os.listdir(working_dir_event_extractor), key=lambda x: int(x.split('_')[1].split('.')[0])): |
| 563 | event_json_path = os.path.join(working_dir_event_extractor, event_json_fname) |
| 564 | if os.path.exists(event_json_path): |
| 565 | with open(event_json_path, "r", encoding="utf-8") as f: |
| 566 | event_data = json.load(f) |
| 567 | event: Event = Event.model_validate(event_data) |
| 568 | extracted_events.append(event) |
| 569 | |
| 570 | if len(extracted_events) > 0: |
| 571 | if extracted_events[-1].is_last: |
| 572 | print(f"⏭️ Skipping event extraction as all events already exist in {working_dir_event_extractor}.") |
| 573 | else: |
| 574 | print(f"🔖 Continuing event extraction from {len(extracted_events)} existing events...") |
| 575 | else: |
| 576 | print("🔖 Starting event extraction ...") |
| 577 | |
| 578 | while len(extracted_events) == 0 or not extracted_events[-1].is_last: |
| 579 | next_event = self.event_extractor.extract_next_event( |
| 580 | novel_text=compressed_novel, |
| 581 | extracted_events=extracted_events, |
| 582 | ) |
| 583 | event_json_path = os.path.join(working_dir_event_extractor, f"event_{len(extracted_events)}.json") |
| 584 | with open(event_json_path, "w", encoding="utf-8") as f: |
| 585 | json.dump(next_event.model_dump(), f, ensure_ascii=False, indent=4) |
| 586 | print(f"✅ Extracted event {next_event.index}, saved to {event_json_path}") |
| 587 | |
| 588 | extracted_events.append(next_event) |
| 589 | |
| 590 | # summary |
| 591 | print() |
| 592 | print("📌 Summary:") |
| 593 | print(f"📌 Extracted a total of {len(extracted_events)} events.") |
| 594 | |
| 595 | print("📋 Step 2: Extract events from the compressed novel".center(80, "-")) |
| 596 | |
| 597 | |
| 598 | # Step 3: Extract relevant chunks for each event |
| 599 | print() |
| 600 | print("📋 Step 3: Retrieve relevant chunks for each event".center(80, "-")) |
| 601 | working_dir_knowledge_base = os.path.join(self.working_dir, "knowledge_base") |
| 602 | working_dir_retrieve = os.path.join(self.working_dir, "relevant_chunks") |
| 603 | os.makedirs(working_dir_knowledge_base, exist_ok=True) |
| 604 | os.makedirs(working_dir_retrieve, exist_ok=True) |
| 605 | print(f"🗂️ Working directory: {working_dir_knowledge_base} and {working_dir_retrieve}") |
| 606 | |
| 607 | print("🔖 Constructing knowledge base from the raw novel text...") |
| 608 | embeddings = CacheBackedEmbeddings.from_bytes_store( |
| 609 | underlying_embeddings=self.embeddings, |
| 610 | document_embedding_cache=LocalFileStore( |
| 611 | root_path=working_dir_knowledge_base, |
| 612 | ), |
| 613 | namespace=self.embeddings.model, |
| 614 | key_encoder="sha256", |
| 615 | ) |
| 616 | novel_splitter = RecursiveCharacterTextSplitter( |
| 617 | chunk_size=512, |
| 618 | chunk_overlap=128, |
| 619 | ) |
| 620 | novel_chunks = novel_splitter.split_text(novel_text) |
| 621 | knowledge_base = FAISS.from_texts(texts=novel_chunks, embedding=embeddings) |
| 622 | print(f"🔖 Constructed knowledge base with {len(novel_chunks)} chunks, saved to {working_dir_knowledge_base}") |
| 623 | |
| 624 | |
| 625 | print("🔖 Retrieving relevant chunks for each event...") |
| 626 | async def retrieve_relevant_chunks(sem, knowledge_base, event): |
| 627 | async with sem: |
| 628 | relevant_chunk_score_dict = {} |
| 629 | for process in event.process_chain: |
| 630 | chunks = knowledge_base.similarity_search(process, k=10) |
| 631 | chunks = [chunk.page_content for chunk in chunks if chunk.page_content not in relevant_chunk_score_dict] |
| 632 | |
| 633 | chunk_score_pairs = await self.rerank_model( |
| 634 | documents=chunks, |
| 635 | query=process, |
| 636 | top_n=10, |
| 637 | ) |
| 638 | |
| 639 | threshold = 0.7 |
| 640 | for chunk, score in chunk_score_pairs: |
| 641 | if score >= threshold: |
| 642 | if chunk not in relevant_chunk_score_dict: |
| 643 | relevant_chunk_score_dict[chunk] = score |
| 644 | else: |
| 645 | relevant_chunk_score_dict[chunk] += score |
| 646 | |
| 647 | return event.index, relevant_chunk_score_dict |
| 648 | |
| 649 | event_idx_to_relevant_chunk_score_dict = {} |
| 650 | |
| 651 | sem = asyncio.Semaphore(10) |
| 652 | tasks = [] |
| 653 | for event in extracted_events: |
| 654 | chunks_dir = os.path.join(working_dir_retrieve, f"event_{event.index}") |
| 655 | if os.path.exists(chunks_dir) and len(os.listdir(chunks_dir)) > 0: |
| 656 | relevant_chunk_score_dict = {} |
| 657 | for chunk_fname in os.listdir(chunks_dir): |
| 658 | chunk_path = os.path.join(chunks_dir, chunk_fname) |
| 659 | score = float(chunk_fname.split('-score_')[1].split('.txt')[0]) |
| 660 | with open(chunk_path, "r", encoding="utf-8") as f: |
| 661 | chunk = f.read() |
| 662 | relevant_chunk_score_dict[chunk] = score |
| 663 | event_idx_to_relevant_chunk_score_dict[event.index] = relevant_chunk_score_dict |
| 664 | print(f"⏭️ Skipping retrieval for event {event.index} as it already exists.") |
| 665 | else: |
| 666 | tasks.append(retrieve_relevant_chunks(sem, knowledge_base, event)) |
| 667 | |
| 668 | if len(tasks) > 0: |
| 669 | for task in asyncio.as_completed(tasks): |
| 670 | event_index, relevant_chunk_score_dict = await task |
| 671 | chunks_dir = os.path.join(working_dir_retrieve, f"event_{event_index}") |
| 672 | os.makedirs(chunks_dir, exist_ok=True) |
| 673 | for idx, (chunk, score) in enumerate(relevant_chunk_score_dict.items()): |
| 674 | chunk_path = os.path.join(chunks_dir, f"chunk_{idx}-score_{score:.2f}.txt") |
| 675 | with open(chunk_path, "w", encoding="utf-8") as f: |
| 676 | f.write(chunk) |
| 677 | event_idx_to_relevant_chunk_score_dict[event_index] = relevant_chunk_score_dict |
| 678 | print(f"✅ Retrieved {len(relevant_chunk_score_dict)} relevant chunks for event {event_index}, saved to {chunks_dir}") |
| 679 | |
| 680 | print("🔖 Retrieved relevant chunks for all events.") |
| 681 | print("📋 Step 3: Retrieve relevant chunks for each event".center(80, "-")) |
| 682 | |
| 683 | |
| 684 | |
| 685 | # Step 4: Extract scenes for each event, design the script for each scene |
| 686 | print() |
| 687 | print("📋 Step 4: Extract scenes for each event, design the script for each scene".center(80, "-")) |
| 688 | working_dir_scene_extractor = os.path.join(self.working_dir, "scenes") |
| 689 | os.makedirs(working_dir_scene_extractor, exist_ok=True) |
| 690 | print(f"🗂️ Working directory: {working_dir_scene_extractor}") |
| 691 | |
| 692 | |
| 693 | unfinished_event_indices = [] |
| 694 | event_idx_to_scenes = {event.index: [] for event in extracted_events} |
| 695 | for event in extracted_events: |
| 696 | scenes_dir = os.path.join(working_dir_scene_extractor, f"event_{event.index}") |
| 697 | if os.path.exists(scenes_dir): |
| 698 | for scene_json_fname in sorted(os.listdir(scenes_dir), key=lambda x: int(x.split('_')[1].split('.')[0])): |
| 699 | scene_json_path = os.path.join(scenes_dir, scene_json_fname) |
| 700 | with open(scene_json_path, "r", encoding="utf-8") as f: |
| 701 | scene_data = json.load(f) |
| 702 | scene = Scene.model_validate(scene_data) |
| 703 | event_idx_to_scenes[event.index].append(scene) |
| 704 | |
| 705 | if len(event_idx_to_scenes[event.index]) > 0 and event_idx_to_scenes[event.index][-1].is_last: |
| 706 | print(f"⏭️ Skipping scene extraction for event {event.index} as all scenes already exist in {scenes_dir}.") |
| 707 | else: |
| 708 | unfinished_event_indices.append(event.index) |
| 709 | |
| 710 | if len(unfinished_event_indices) > 0: |
| 711 | if len(unfinished_event_indices) == len(extracted_events): |
| 712 | print(f"🔖 Starting scene extraction for all events...") |
| 713 | else: |
| 714 | print(f"🔖 Continuing scene extraction for events: {unfinished_event_indices}") |
| 715 | |
| 716 | |
| 717 | async def extract_scenes_for_event(sem, relevant_chunks, event, previous_scenes): |
| 718 | async with sem: |
| 719 | os.makedirs(os.path.join(working_dir_scene_extractor, f"event_{event.index}"), exist_ok=True) |
| 720 | |
| 721 | while len(previous_scenes) == 0 or not previous_scenes[-1].is_last: |
| 722 | next_scene = await self.scene_extractor.get_next_scene( |
| 723 | relevant_chunks=relevant_chunks, |
| 724 | event=event, |
| 725 | previous_scenes=previous_scenes, |
| 726 | ) |
| 727 | scene_json_path = os.path.join(working_dir_scene_extractor, f"event_{event.index}", f"scene_{len(previous_scenes)}.json") |
| 728 | with open(scene_json_path, "w", encoding="utf-8") as f: |
| 729 | json.dump(next_scene.model_dump(), f, ensure_ascii=False, indent=4) |
| 730 | print(f"✔️ Extracted scene {next_scene.idx} for event {event.index}, saved to {scene_json_path}") |
| 731 | previous_scenes.append(next_scene) |
| 732 | |
| 733 | print(f"✅ Extracted all {len(previous_scenes)} scenes for event {event.index}.") |
| 734 | return event.index, previous_scenes |
| 735 | |
| 736 | |
| 737 | sem = asyncio.Semaphore(8) |
| 738 | for event_index in unfinished_event_indices: |
| 739 | relevant_chunks = list(event_idx_to_relevant_chunk_score_dict[event_index].keys()) |
| 740 | tasks.append(extract_scenes_for_event(sem, relevant_chunks, extracted_events[event_index], event_idx_to_scenes[event_index])) |
| 741 | |
| 742 | task_outputs = await asyncio.gather(*tasks) |
| 743 | for event_index, previous_scenes in task_outputs: |
| 744 | event_idx_to_scenes[event_index] = previous_scenes |
| 745 | |
| 746 | print("🔖 Extracted scenes for all events.") |
| 747 | print("📋 Step 4: Extract scenes for each event, design the script for each scene".center(80, "-")) |
| 748 | |
| 749 | |
| 750 | |
| 751 | # Step 5: Merge characters from scene-level to event-level, then to novel-level |
| 752 | print() |
| 753 | print("📋 Step 5: Merge characters from scene-level to novel-level".center(80, "-")) |
| 754 | working_dir_global_information_planner = os.path.join(self.working_dir, "global_information") |
| 755 | os.makedirs(working_dir_global_information_planner, exist_ok=True) |
| 756 | print(f"🗂️ Working directory: {working_dir_global_information_planner}") |
| 757 | |
| 758 | # Step 5.1: Merge characters from scene-level to event-level |
| 759 | print("🔖 Merging characters across scenes in each event...") |
| 760 | working_dir_characters = os.path.join(working_dir_global_information_planner, "characters") |
| 761 | os.makedirs(working_dir_characters, exist_ok=True) |
| 762 | |
| 763 | async def merge_characters_across_scenes_in_event(sem, event_idx, scenes): |
| 764 | async with sem: |
| 765 | merged_characters = await self.global_information_planner.merge_characters_across_scenes_in_event( |
| 766 | event_idx=event_idx, |
| 767 | scenes=scenes, |
| 768 | ) |
| 769 | path = os.path.join(working_dir_characters, "event_level", f"event_{event_idx}_characters.json") |
| 770 | os.makedirs(os.path.dirname(path), exist_ok=True) |
| 771 | with open(path, "w", encoding="utf-8") as f: |
| 772 | json.dump([char.model_dump() for char in merged_characters], f, ensure_ascii=False, indent=4) |
| 773 | print(f"✅ Merged characters for event {event_idx}, saved to {path}") |
| 774 | |
| 775 | return event_idx, merged_characters |
| 776 | |
| 777 | |
| 778 | event_idx_to_characters_in_event = {} |
| 779 | |
| 780 | sem = asyncio.Semaphore(8) |
| 781 | tasks = [] |
| 782 | for event in extracted_events: |
| 783 | path = os.path.join(working_dir_characters, "event_level", f"event_{event.index}_characters.json") |
| 784 | if os.path.exists(path): |
| 785 | with open(path, "r", encoding="utf-8") as f: |
| 786 | character_data = json.load(f) |
| 787 | characters = [CharacterInEvent.model_validate(char) for char in character_data] |
| 788 | event_idx_to_characters_in_event[event.index] = characters |
| 789 | print(f"⏭️ Skipping character merging for event {event.index} as it already exists.") |
| 790 | else: |
| 791 | tasks.append(merge_characters_across_scenes_in_event(sem, event.index, event_idx_to_scenes[event.index])) |
| 792 | |
| 793 | task_outputs = await asyncio.gather(*tasks) |
| 794 | for event_index, merged_characters in task_outputs: |
| 795 | event_idx_to_characters_in_event[event_index] = merged_characters |
| 796 | |
| 797 | print("🔖 Merged characters across scenes in each event.") |
| 798 | |
| 799 | # Step 5.2: Merge characters from event-level to novel-level |
| 800 | print("🔖 Merging characters across events in the novel...") |
| 801 | |
| 802 | working_dir_characters_novel = os.path.join(working_dir_characters, f"novel_level") |
| 803 | os.makedirs(working_dir_characters_novel, exist_ok=True) |
| 804 | |
| 805 | fnames = os.listdir(working_dir_characters_novel) |
| 806 | existing_characters_in_novel = [] |
| 807 | if len(fnames) > 0: |
| 808 | fname = max(fnames, key=lambda x: int(x.split('_')[-1].split('.json')[0])) |
| 809 | start_event_idx = int(fname.split('_')[-1].split('.json')[0]) + 1 |
| 810 | path = os.path.join(working_dir_characters_novel, fname) |
| 811 | with open(path, "r", encoding="utf-8") as f: |
| 812 | character_data = json.load(f) |
| 813 | existing_characters_in_novel = [CharacterInNovel.model_validate(char) for char in character_data] |
| 814 | |
| 815 | if start_event_idx == len(extracted_events): |
| 816 | print(f"⏭️ Skipping merging as all events already merged to novel-level in {working_dir_characters_novel}.") |
| 817 | else: |
| 818 | print(f"🔖 Continuing merging from event {start_event_idx}, currently {len(existing_characters_in_novel)} characters in novel.") |
| 819 | |
| 820 | else: |
| 821 | existing_characters_in_novel = [] |
| 822 | start_event_idx = 0 |
| 823 | |
| 824 | for event in extracted_events[start_event_idx:]: |
| 825 | characters_in_event = event_idx_to_characters_in_event[event.index] |
| 826 | path = os.path.join(working_dir_characters_novel, f"novel_characters_after_event_{event.index}.json") |
| 827 | existing_characters_in_novel = self.global_information_planner.merge_characters_to_existing_characters_in_novel( |
| 828 | event_idx=event.index, |
| 829 | existing_characters_in_novel=existing_characters_in_novel, |
| 830 | characters_in_event=characters_in_event, |
| 831 | ) |
| 832 | with open(path, "w", encoding="utf-8") as f: |
| 833 | json.dump([char.model_dump() for char in existing_characters_in_novel], f, ensure_ascii=False, indent=4) |
| 834 | print(f"✅ Merged characters from event {event.index} to novel-level, now {len(existing_characters_in_novel)} characters in novel, saved to {path}") |
| 835 | |
| 836 | print("🔖 Merged characters across events in the novel.") |
| 837 | |
| 838 | characters_in_novel = existing_characters_in_novel |
| 839 | |
| 840 | print("📋 Step 5: Merge characters from scene-level to novel-level".center(80, "-")) |
| 841 | |
| 842 | |
| 843 | |
| 844 | |
| 845 | # Step 6: Generate the portrait for all characters in the novel |
| 846 | print() |
| 847 | print("📋 Step 6: Generate the reference images for all characters in the specific scene") |
| 848 | |
| 849 | working_dir_character_portrait = os.path.join(self.working_dir, "character_portraits") |
| 850 | os.makedirs(working_dir_character_portrait, exist_ok=True) |
| 851 | print(f"🗂️ Working directory: {working_dir_character_portrait}") |
| 852 | |
| 853 | print("🔖 Generating character portraits based on static features ...") |
| 854 | base_character_portrait_dir = os.path.join(working_dir_character_portrait, "base") |
| 855 | os.makedirs(base_character_portrait_dir, exist_ok=True) |
| 856 | |
| 857 | async def generate_portrait_for_character(sem, character: CharacterInNovel): |
| 858 | async with sem: |
| 859 | image_path = os.path.join(base_character_portrait_dir, f"character_{character.index}_{safe_path_component(character.identifier_in_novel)}.png") |
| 860 | |
| 861 | if os.path.exists(image_path): |
| 862 | print(f"⏭️ Skipping portrait generation for character {character.idx} as it already exists.") |
| 863 | return |
| 864 | |
| 865 | prompt = f"Generate a full-body, front-view portrait based on the following description, in the style of {style}:" |
| 866 | prompt += f"\nCharacter Identifier: {character.identifier_in_novel}" |
| 867 | prompt += f"\nFeatures: {character.static_features}" |
| 868 | prompt += f"\nThe character should be centered in the image, occupying most of the frame. Gazing straight ahead. Standing with arms relaxed at sides. Natural expression. The background should be plain white." |
| 869 | |
| 870 | image = await self.image_generator.generate_single_image( |
| 871 | prompt=prompt, |
| 872 | size="512x512", |
| 873 | ) |
| 874 | image.save(image_path) |
| 875 | print(f"✅ Generated portrait for character {character.index} ({character.identifier_in_novel}), saved to {image_path}") |
| 876 | |
| 877 | |
| 878 | sem = asyncio.Semaphore(5) |
| 879 | tasks = [ |
| 880 | generate_portrait_for_character(sem, character) |
| 881 | for character in characters_in_novel |
| 882 | ] |
| 883 | |
| 884 | await asyncio.gather(*tasks) |
| 885 | print("🔖 Generated character portraits based on static features.") |
| 886 | |
| 887 | |
| 888 | print("🔖 Generating character portraits based on dynamic features in the specific scene") |
| 889 | |
| 890 | async def generate_portrait_for_character_in_scene( |
| 891 | sem, |
| 892 | base_character_image_path: str, |
| 893 | character: CharacterInScene, |
| 894 | event_idx: int, |
| 895 | scene_idx: int, |
| 896 | ): |
| 897 | async with sem: |
| 898 | image_path = os.path.join( |
| 899 | working_dir_character_portrait, |
| 900 | f"event_{event_idx}", |
| 901 | f"scene_{scene_idx}", |
| 902 | f"character_{character.idx}_{character.identifier_in_scene}.png", |
| 903 | ) |
| 904 | os.makedirs(os.path.dirname(image_path), exist_ok=True) |
| 905 | |
| 906 | if os.path.exists(image_path): |
| 907 | print(f"⏭️ Skipping portrait generation for event {event_idx}, scene {scene_idx}, character {character.idx} as it already exists.") |
| 908 | return |
| 909 | |
| 910 | if not character.is_visible: |
| 911 | shutil.copy(base_character_image_path, image_path) |
| 912 | print(f"⏭️ For event {event_idx}, scene {scene_idx}, character {character.idx} ({character.identifier_in_scene}) is not visible, copied base portrait to {image_path}") |
| 913 | return |
| 914 | |
| 915 | if character.dynamic_features is None: |
| 916 | shutil.copy(base_character_image_path, image_path) |
| 917 | print(f"⏭️ For event {event_idx}, scene {scene_idx}, character {character.idx} ({character.identifier_in_scene}) has no dynamic features, copied base portrait to {image_path}") |
| 918 | return |
| 919 | |
| 920 | prompt = f"Generate a full-body, front-view portrait based on the provided base image. Modify the base image according to the following dynamic features, in the style of {style}. Keep the character's identity consistent with the base image:" |
| 921 | prompt += f"\nCharacter Identifier: {character.identifier_in_scene}" |
| 922 | prompt += f"\nDynamic Features: {character.dynamic_features}" |
| 923 | prompt += f"\nThe character should be centered in the image, occupying most of the frame. Gazing straight ahead. Standing with arms relaxed at sides. Natural expression. The background should be plain white." |
| 924 | |
| 925 | prompt = await self.rewriter(prompt) |
| 926 | |
| 927 | |
| 928 | image = await self.image_generator.generate_single_image( |
| 929 | prompt=prompt, |
| 930 | reference_image_paths=[base_character_image_path], |
| 931 | size="512x512", |
| 932 | ) |
| 933 | image.save(image_path) |
| 934 | print(f"✅ For event {event_idx}, scene {scene_idx}, generated portrait for character {character.idx} ({character.identifier_in_scene}), saved to {image_path}") |
| 935 | |
| 936 | |
| 937 | sem = asyncio.Semaphore(3) |
| 938 | tasks = [] |
| 939 | for character in characters_in_novel: |
| 940 | character_base_image_path = os.path.join(base_character_portrait_dir, f"character_{character.index}_{safe_path_component(character.identifier_in_novel)}.png") |
| 941 | for event_idx, identifier_in_event in character.active_events.items(): |
| 942 | characters_in_event: List[CharacterInEvent] = event_idx_to_characters_in_event[event_idx] |
| 943 | character_in_event = [char for char in characters_in_event if char.identifier_in_event == identifier_in_event][0] # TODO: 这里的数据结构没有做好,居然还要遍历查找。。。 |
| 944 | for scene_idx, identifier_in_scene in character_in_event.active_scenes.items(): |
| 945 | scene = event_idx_to_scenes[event_idx][scene_idx] |
| 946 | character_in_scene: CharacterInScene = [char for char in scene.characters if char.identifier_in_scene == identifier_in_scene][0] # TODO: 这里的数据结构也没有做好 |
| 947 | tasks.append( |
| 948 | generate_portrait_for_character_in_scene( |
| 949 | sem, |
| 950 | character_base_image_path, |
| 951 | character_in_scene, |
| 952 | event_idx, |
| 953 | scene_idx, |
| 954 | ) |
| 955 | ) |
| 956 | await asyncio.gather(*tasks) |
| 957 | print("🔖 Generated character portraits based on dynamic features in the specific scene") |
| 958 | |
| 959 | print("📋 Step 6: Generate the reference images for all characters in the specific scene".center(80, "-")) |
| 960 | |
| 961 | |
| 962 | |
| 963 | # Step 7: Generate video for each scene |
| 964 | print("📋 Step 7: Generate the video for each scene".center(80, "-")) |
| 965 | working_dir_scene_videos = os.path.join(self.working_dir, "videos") |
| 966 | os.makedirs(working_dir_scene_videos, exist_ok=True) |
| 967 | |
| 968 | for event in extracted_events: |
| 969 | scenes: List[Scene] = event_idx_to_scenes[event.index] |
| 970 | for scene in scenes: |
| 971 | scene_video_dir = os.path.join(working_dir_scene_videos, f"event_{event.index}", f"scene_{scene.idx}") |
| 972 | os.makedirs(scene_video_dir, exist_ok=True) |
| 973 | |
| 974 | self.script2video_pipeline.working_dir = scene_video_dir |
| 975 | script = scene.script |
| 976 | style = "realistic movie style" |
| 977 | character_registry = {} |
| 978 | for character in scene.characters: |
| 979 | character_registry[character.identifier_in_scene] = [ |
| 980 | { |
| 981 | "path": os.path.join( |
| 982 | working_dir_character_portrait, |
| 983 | f"event_{event.index}", |
| 984 | f"scene_{scene.idx}", |
| 985 | f"character_{character.idx}_{character.identifier_in_scene}.png", |
| 986 | ), |
| 987 | "description": f"A portrait of {character.identifier_in_scene}", |
| 988 | } |
| 989 | ] |
| 990 | await self.script2video_pipeline( |
| 991 | script=script, |
| 992 | style=style, |
| 993 | character_registry=character_registry |
| 994 | ) |
| 995 | print(f"✅ Generated video for event {event.index}, scene {scene.idx}, saved to {scene_video_dir}") |
| 996 | print("📋 Step 7: Generate the video for each scene".center(80, "-")) |
| 997 | |
| 998 | |
| 999 | # is_last flags are asserted by the LLM only; cap the extraction loops so a |
| 1000 | # model that never sets one cannot spend tokens forever. |
| 1001 | MAX_EXTRACTED_EVENTS = 50 |
| 1002 | MAX_SCENES_PER_EVENT = 30 |
| 1003 | |
| 1004 | |
| 1005 | def _ensure_extraction_cap(count: int, cap: int, what: str) -> None: |
| 1006 | if count >= cap: |
| 1007 | raise RuntimeError( |
| 1008 | f"Extraction reached {count} {what} without an is_last marker (cap: {cap}); " |
| 1009 | "aborting to avoid unbounded LLM calls." |
| 1010 | ) |
| 1011 |