返回 ViMax
novel2movie_pipeline.py
根目录 / pipelines / novel2movie_pipeline.py
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
1011 lines PYTHON