from __future__ import annotations import json import math import re import shutil from dataclasses import dataclass from pathlib import Path from typing import Any from adam.models import ExecutionPlan from adam.orion import apply_orion_review DEFAULT_PRESETS: dict[str, dict[str, Any]] = { "Character LoRA": { "trainer": "lora", "epochs": 100, "image_count": 60, "description": "A balanced starting point for a character or person.", }, "Style LoRA": { "trainer": "lora", "epochs": 80, "image_count": 80, "description": "A broader image set for learning a visual style.", }, "DDPM Test Run": { "trainer": "ddpm", "epochs": 25, "image_count": 40, "description": "A short run to verify the dataset and training setup.", }, "DDPM Full Run": { "trainer": "ddpm", "epochs": 100, "image_count": 100, "description": "A practical default for a full DDPM experiment.", }, "Flow Test Run": { "trainer": "flow", "epochs": 25, "image_count": 40, "description": "A short Flow Matching setup check.", }, "Oasis Pipeline Test": { "trainer": "oasis", "epochs": 3, "image_count": 500, "description": "A short action-world-model run for validating gameplay frames and controls.", "training_options": {"resolution": "256x144", "batch_size": 2, "workers": 2, "mixed_precision": "fp16", "frame_gap": 1, "chunk_size": 0, "balance_actions": True}, }, "Oasis New Game": { "trainer": "oasis", "epochs": 40, "image_count": 1000, "description": "A responsive starting recipe for training a new playable game world.", "training_options": {"resolution": "256x144", "batch_size": 2, "learning_rate": 0.00005, "workers": 4, "mixed_precision": "fp16", "frame_gap": 1, "balance_actions": True, "chunk_size": 0, "include_older_data": False}, }, "Oasis Add Actions / Poses": { "trainer": "oasis", "epochs": 30, "image_count": 1000, "description": "A conservative continuation recipe for learning new controls or poses.", "training_options": {"resolution": "256x144", "batch_size": 2, "learning_rate": 0.00002, "workers": 4, "mixed_precision": "fp16", "frame_gap": 1, "action_input_scale": 8.0, "action_contrast_weight": 0.35, "balance_actions": True}, }, "Oasis Large Dataset +10K": { "trainer": "oasis", "epochs": 45, "image_count": 10000, "description": "Uses a balanced 5,000-transition chunk so a 10K+ dataset does not create an accidental multi-day run.", "training_options": {"resolution": "256x144", "batch_size": 2, "learning_rate": 0.00002, "workers": 8, "mixed_precision": "fp16", "frame_gap": 1, "chunk_size": 5000, "chunk_mode": "balanced", "chunk_offset": 0, "replay_older_percent": 50.0, "include_older_data": True, "balance_actions": True, "recovery_minutes": 30, "tf32": True}, }, } def parse_model_batch_names(text: str) -> list[str]: """Return unique, user-ordered model subjects from a pasted line list.""" names: list[str] = [] seen: set[str] = set() for raw in text.splitlines(): name = re.sub(r"^\s*(?:[-*•]|\d+[.)])\s*", "", raw).strip() key = re.sub(r"\s+", " ", name).casefold() if name and key not in seen: names.append(re.sub(r"\s+", " ", name)) seen.add(key) return names def build_dataset_collection_request( subject: str, *, image_count: int = 100, collection_mode: str = "target", ) -> str: """Build the dataset-only first phase used by a saved model batch.""" subject = subject.strip() if not subject: raise ValueError("Dataset collection requires a subject.") if collection_mode == "all_available": return ( f"Collect a dataset of {subject} with as many available images as Bing " "returns (up to 5,000)." ) return f"Collect a dataset of {image_count} images of {subject}." @dataclass(slots=True) class PreflightItem: level: str message: str @dataclass(slots=True) class DatasetMatch: status: str dataset_name: str = "" score: float = 0.0 _DATASET_NAME_NOISE = { "dataset", "datasets", "image", "images", "picture", "pictures", "photo", "photos", "collection", "collected", } def _dataset_name_tokens(value: object) -> set[str]: words = re.findall(r"[a-z0-9]+", str(value).casefold()) return { word[:-1] if word.endswith("s") and len(word) > 3 else word for word in words if word not in _DATASET_NAME_NOISE } def suggest_existing_dataset(state: dict[str, object], datasets: list[Any]) -> DatasetMatch: """Safely match one batch model to a registered dataset by its human name.""" queries = [ _dataset_name_tokens(state.get("model_name", "")), _dataset_name_tokens(state.get("subject", "")), ] queries = [query for query in queries if query] if not queries: return DatasetMatch("unmatched") scored: list[tuple[float, Any]] = [] for asset in datasets: name = str(getattr(asset, "name", "")) path = Path(str(getattr(asset, "path", ""))) tokens = _dataset_name_tokens(name) if not name or not tokens or not path.is_dir(): continue score = 0.0 for query in queries: overlap = len(query & tokens) / len(query) extra_penalty = min(0.20, len(tokens - query) * 0.08) score = max(score, overlap - extra_penalty) if score >= 0.80: scored.append((score, asset)) if not scored: return DatasetMatch("unmatched") scored.sort(key=lambda item: (-item[0], len(str(getattr(item[1], "name", ""))))) best_score, best = scored[0] if len(scored) > 1 and best_score - scored[1][0] < 0.10: return DatasetMatch("ambiguous", score=best_score) return DatasetMatch("matched", str(getattr(best, "name", "")), best_score) def combine_training_plans(plans: list[ExecutionPlan]) -> ExecutionPlan: """Combine independently validated model plans into one sequential job.""" usable = [plan for plan in plans if plan.steps] if not usable: raise ValueError("A training batch needs at least one actionable model plan.") if len(usable) == 1: return usable[0] summaries = [ f"{index}. {plan.project_name}: {plan.summary.splitlines()[0]}" for index, plan in enumerate(usable, 1) ] reasons = [plan.confirmation_reason for plan in usable if plan.confirmation_reason] has_training = any( step.tool_id.endswith("_trainer") for plan in usable for step in plan.steps ) batch_kind = "training" if has_training else "dataset collection" return ExecutionPlan( request="\n\n".join(plan.request for plan in usable), summary=( f"Sequential {batch_kind} batch with {len(usable)} items. ADAM will finish " "each item before starting the next; a failed step stops the batch.\n\n" + "\n".join(summaries) ), steps=[step for plan in usable for step in plan.steps], requires_confirmation=any(plan.requires_confirmation for plan in usable), confirmation_reason="; ".join(dict.fromkeys(reasons)) or ( "This batch contains multiple model workflows. Review every model and its " "output path before starting." ), project_name=( f"Training batch ({len(usable)} models)" if has_training else f"Dataset collection batch ({len(usable)} datasets)" ), ) def estimate_plan(plan: Any) -> list[PreflightItem]: """Add deliberately conservative, clearly labelled planning estimates.""" estimates: list[PreflightItem] = [] for step in plan.steps: if not step.tool_id.endswith("_trainer"): continue if step.tool_id == "wan_video_trainer": from adam.video_lora import clips_in count = len(clips_in(Path(str(step.arguments.get("dataset_dir", ""))))) estimates.append(PreflightItem("estimate", f"Wan: {count} clips; frame buckets {step.arguments.get('target_frames', '25,49')}. Video duration and memory cost depend on frame count and block swapping; image-training time estimates do not apply.")) continue if step.tool_id == "oasis_trainer": from adam.oasis_dataset import inspect_oasis_dataset epochs = max(1, int(step.arguments.get("epochs", 1) or 1)) gap = max(1, int(step.arguments.get("frame_gap", 1) or 1)) report = inspect_oasis_dataset(step.arguments.get("dataset_dir", ""), frame_gap=gap) if report.ok and report.valid_transitions: batch = max(1, int(step.arguments.get("batch_size", 1) or 1)) accumulation = max(1, int(step.arguments.get("gradient_accumulation", 1) or 1)) chunk_size = max(0, int(step.arguments.get("chunk_size", 0) or 0)) transitions_per_epoch = min(report.valid_transitions, chunk_size) if chunk_size else report.valid_transitions updates_per_epoch = math.ceil(transitions_per_epoch / batch / accumulation) updates = updates_per_epoch * epochs pace = f"; {report.capture_fps:g} FPS capture → {report.native_ai_fps:g} native AI FPS" if report.capture_fps and report.native_ai_fps else "" estimates.extend([ PreflightItem("estimate", f"Oasis workload: {updates:,} optimizer steps ({epochs} epochs × {updates_per_epoch:,} steps){pace}"), PreflightItem("estimate", f"Oasis dataset: {report.valid_transitions:,} valid transitions; {transitions_per_epoch:,} used per epoch"), PreflightItem("estimate", "Suggested capacity: about 12 GB VRAM. Run the Oasis speed test before treating duration estimates as reliable."), ]) continue trainer = step.tool_id.removesuffix("_trainer") epochs = max(1, int(step.arguments.get("epochs", 1) or 1)) raw_dataset = str(step.arguments.get("dataset_dir", "") or "").strip() dataset = Path(raw_dataset).expanduser() image_count = 0 if raw_dataset and dataset.is_dir(): try: image_count = sum( 1 for path in dataset.rglob("*") if path.is_file() and path.suffix.casefold() in {".jpg", ".jpeg", ".png", ".webp", ".bmp"} ) except OSError: image_count = 0 image_count = image_count or 60 workload = epochs * image_count seconds_per_image_epoch = { "lora": 0.12, "ddpm": 0.07, "flow": 0.10, "oasis": 0.18, }.get(trainer, 0.10) center_minutes = max(1, int(workload * seconds_per_image_epoch / 60)) low = max(1, center_minutes // 2) high = max(low + 1, center_minutes * 3) checkpoint_gb = { "lora": 0.25, "ddpm": 1.0, "flow": 1.0, "oasis": 1.0, }.get(trainer, 0.75) checkpoint_count = max(1, min(20, epochs // 25 + 1)) disk_gb = checkpoint_gb * checkpoint_count typical_vram = {"lora": 8, "ddpm": 6, "flow": 8, "oasis": 12}.get(trainer, 8) estimates.extend( [ PreflightItem( "estimate", f"Estimated workload: {workload:,} image-epochs " f"({epochs:,} epochs × about {image_count:,} images)", ), PreflightItem( "estimate", f"Rough duration: {low}–{high} minutes; model size, resolution, " "batch size, and GPU can change this substantially", ), PreflightItem( "estimate", f"Suggested capacity: about {typical_vram} GB VRAM and " f"{disk_gb:.1f} GB free for checkpoints", ), ] ) return estimates def presets_from_config(config: Any) -> dict[str, dict[str, Any]]: presets = {name: dict(values) for name, values in DEFAULT_PRESETS.items()} stored = config.get("training_presets", {}) if isinstance(stored, dict): for name, values in stored.items(): if isinstance(name, str) and isinstance(values, dict): presets[name] = dict(values) return presets def build_training_request( *, trainer: str, subject: str, dataset_name: str, create_dataset: bool, epochs: int, image_count: int, model_name: str, collection_mode: str = "target", training_options: dict[str, Any] | None = None, ) -> str: subject = subject.strip() dataset_name = dataset_name.strip() model_name = model_name.strip() or subject or dataset_name trainer_label = {"lora": "LoRA", "ddpm": "DDPM", "flow": "Flow Matching", "oasis": "Oasis Action World Model"}.get( trainer, trainer.replace("_", " ").title(), ) if create_dataset: collection_phrase = ( "as many available images as Bing returns (up to 5,000)" if collection_mode == "all_available" else f"up to {image_count} images" ) if trainer == "lora": request = ( f"Create and train a LoRA of {subject} for {epochs} epochs " f"using {collection_phrase}. Name the model {model_name}." ) elif trainer == "ddpm": request = ( f"Grab a dataset of {subject} off the internet with {collection_phrase}, " f"name the model {model_name}, train it on a DDPM for {epochs} epochs, " "and save it to the DDPM output." ) else: request = ( f"Collect a dataset of {collection_phrase} of {subject}. Then train the " f"{subject} dataset with Flow Matching for {epochs} epochs and name the model {model_name}." ) else: request = ( f"From the {dataset_name} dataset, train a {trainer_label} model for {epochs} epochs. " f"Name the model {model_name}." ) if training_options: request += " [ADAM_TRAINING_OPTIONS:" + json.dumps(training_options, sort_keys=True) + "]" request += " [ADAM_TRAINER:" + trainer + "]" return request def build_fine_tune_request( *, model_name: str, trainer: str, epochs: int, output_model_name: str = "", dataset_mode: str = "original", dataset_name: str = "", new_subject: str = "", image_count: int = 60, training_options: dict[str, Any] | None = None, ) -> str: """Build the explicit continuation request used by the Fine-Tune assistant.""" labels = {"lora": "LoRA", "ddpm": "DDPM", "flow": "Flow Matching", "oasis": "Oasis Action World Model"} if trainer not in labels: raise ValueError("Fine-tuning requires a supported trainer.") if not model_name.strip(): raise ValueError("Fine-tuning requires a model name.") if epochs < 1: raise ValueError("Fine-tuning requires at least one additional epoch.") if dataset_mode not in {"original", "existing", "new"}: raise ValueError("Fine-tuning requires a valid dataset choice.") payload = { "model_name": model_name.strip(), "output_model_name": output_model_name.strip(), "trainer": trainer, "epochs": epochs, "dataset_mode": dataset_mode, "dataset_name": dataset_name.strip(), "new_subject": new_subject.strip(), "image_count": max(10, min(int(image_count), 5000)), "training_options": dict(training_options or {}), } return ( f"Fine-tune {model_name.strip()} for {epochs} epochs with {labels[trainer]}. " "[ADAM_FINE_TUNE:" + json.dumps(payload, sort_keys=True) + "]" ) def inspect_plan(plan: Any, config: Any) -> list[PreflightItem]: items: list[PreflightItem] = [] folders = config.get("tool_folders", {}) folders = folders if isinstance(folders, dict) else {} checked_tools: set[str] = set() checked_paths: set[str] = set() for step in plan.steps: if step.tool_id in {"wan_video_trainer", "wan_video_generator", "wan_video_dataset"}: from adam.video_lora import trainer_root, setup_errors, dataset_errors errors = setup_errors(trainer_root(config.root, config)) if step.tool_id == "wan_video_trainer": errors += dataset_errors(Path(str(step.arguments.get("dataset_dir", ""))), str(step.arguments.get("trigger_word", "subject_token"))) items.extend(PreflightItem("warning", error) for error in errors[:12]) if not errors: items.append(PreflightItem("ready", "Wan Video LoRA environment and required files are available")) if step.tool_id != "wan_video_trainer": continue output = Path(str(step.arguments.get("output_dir", ""))) while not output.exists() and output.parent != output: output = output.parent try: free = shutil.disk_usage(output).free / (1024 ** 3) items.append(PreflightItem("ready" if free >= 10 else "warning", f"Video output/cache drive: {free:.1f} GB free")) except OSError: items.append(PreflightItem("warning", "Could not check video output drive space")) continue if step.tool_id == "oasis_trainer": from adam.oasis_dataset import oasis_pace pace = oasis_pace( step.arguments.get("dataset_dir", ""), frame_gap=int(step.arguments.get("frame_gap", 1) or 1), ) capture_fps = pace["capture_fps"] native_fps = pace["native_ai_fps"] recommended_gap = pace["recommended_frame_gap"] if isinstance(capture_fps, (int, float)) and isinstance(native_fps, (int, float)): message = ( f"Oasis pacing: {float(capture_fps):g} FPS capture → " f"{float(native_fps):g} native AI FPS at prediction gap " f"{int(step.arguments.get('frame_gap', 1) or 1)}" ) if isinstance(recommended_gap, int) and recommended_gap != int(step.arguments.get("frame_gap", 1) or 1): items.append(PreflightItem("warning", message + f"; dataset recommends gap {recommended_gap}")) else: items.append(PreflightItem("ready", message)) if step.tool_id.endswith("_trainer") or step.tool_id == "dataset_collector": if step.tool_id not in checked_tools: raw_folder = str(folders.get(step.tool_id, "") or "").strip() configured = Path(raw_folder).expanduser() if raw_folder and configured.is_dir(): items.append(PreflightItem("ready", f"{step.title}: connected")) else: items.append(PreflightItem("warning", f"{step.title}: program folder is not connected")) checked_tools.add(step.tool_id) dataset = str(step.arguments.get("dataset_dir", "")) if dataset and dataset not in checked_paths: if Path(dataset).is_dir(): image_count = sum( 1 for path in Path(dataset).iterdir() if path.suffix.casefold() in {".jpg", ".jpeg", ".png", ".webp", ".bmp"} ) detail = f"{image_count} images found" if image_count else "folder found; no top-level images detected" items.append(PreflightItem("ready" if image_count else "warning", f"Dataset: {detail}")) elif not any( prior.tool_id == "dataset_collector" and prior.arguments.get("output_dir") == dataset for prior in plan.steps ): items.append(PreflightItem("warning", "Dataset folder does not exist yet")) checked_paths.add(dataset) base_model = str(step.arguments.get("base_model", "")) if step.tool_id == "lora_trainer": items.append( PreflightItem( "ready" if base_model and Path(base_model).is_file() else "warning", "LoRA base model is available" if base_model and Path(base_model).is_file() else "LoRA base model still needs to be selected", ) ) output = str(step.arguments.get("output_dir", "")) if output: probe = Path(output) while not probe.exists() and probe.parent != probe: probe = probe.parent try: free_gb = shutil.disk_usage(probe).free / (1024 ** 3) items.append( PreflightItem( "ready" if free_gb >= 10 else "warning", f"Output drive has {free_gb:.1f} GB free", ) ) except OSError: items.append(PreflightItem("warning", "Output drive space could not be checked")) return items def append_preflight_summary(plan: Any, config: Any) -> None: if not plan.steps: return if "Pre-flight:" not in plan.summary: items = inspect_plan(plan, config) + estimate_plan(plan) if items: lines = [ ( "Ready" if item.level == "ready" else "Estimate" if item.level == "estimate" else "Check" ) + f": {item.message}" for item in items ] plan.summary += "\n\nPre-flight:\n" + "\n".join(f"• {line}" for line in lines) if not getattr(plan, "orion_review", None): apply_orion_review(plan) def completion_recommendation(plan: Any) -> str: tools = {step.tool_id for step in plan.steps} if "wan_video_trainer" in tools: return "Generate a short clip in Video LoRA and review motion and subject consistency before a longer run." if "oasis_trainer" in tools: return ( "Recommended next step: launch the Oasis player with the saved checkpoint and a starting frame from the same game, then record more control-balanced gameplay before long training." ) if "lora_trainer" in tools or "ddpm_trainer" in tools or "flow_trainer" in tools: return ( "Recommended next step: generate a few preview images and compare them with " "the training dataset. If the subject is weak, improve the dataset before adding epochs." ) if "dataset_collector" in tools: return ( "Recommended next step: review the images and captions, remove weak or duplicate " "examples, then open the Model Creation Assistant to start a short test run." ) return ""