ADAM October 2026 source release: PixelRow, INRFlow, Wan Video, Oasis player and field guide
f8c73f9 verified Download adam/tools/flow_adapter.py from SyntheticMDProductions/AI_Development_Automation_Manager: direct link, hf CLI and curl.
- Browser
- Download file 15.1 kB
-
https://huggingface.co/SyntheticMDProductions/AI_Development_Automation_Manager/resolve/main/adam/tools/flow_adapter.py
- Command line
-
hf download hf://SyntheticMDProductions/AI_Development_Automation_Manager/adam/tools/flow_adapter.py
-
curl -L -o flow_adapter.py https://huggingface.co/SyntheticMDProductions/AI_Development_Automation_Manager/resolve/main/adam/tools/flow_adapter.py
15.1 kB
| """Execution bridge for the connected Rectified Flow image trainer.""" | |
| from __future__ import annotations | |
| import json | |
| import queue | |
| import subprocess | |
| import sys | |
| import threading | |
| import time | |
| import shutil | |
| from pathlib import Path | |
| from adam.config import ConfigManager | |
| from adam.executor import ToolCancelled, ToolContext, ToolExecutionError | |
| from adam.progressive_training import parse_stages, stage_batch_settings, stage_summary | |
| from adam.process_control import set_process_tree_paused, terminate_process_tree | |
| IMAGE_EXTENSIONS = {".jpg", ".jpeg", ".png", ".webp", ".bmp"} | |
| FORCE_STOP_TIMEOUT_SECONDS = 30 | |
| def _latest_preview(folder: Path) -> Path | None: | |
| try: | |
| images = [path for path in folder.rglob("*") if path.is_file() | |
| and path.suffix.lower() in IMAGE_EXTENSIONS | |
| and any(token in path.name.lower() for token in ("preview", "sample", "epoch"))] | |
| return max(images, key=lambda path: path.stat().st_mtime) if images else None | |
| except OSError: | |
| return None | |
| def _train_flow_stage( | |
| context: ToolContext, | |
| dataset_dir: str, | |
| model_name: str, | |
| epochs: int, | |
| output_dir: str, | |
| resume_from: str = "", | |
| resolution: int = 256, batch_size: int = 8, learning_rate: float = 0.0002, | |
| gradient_accumulation: int = 1, workers: int = 4, mixed_precision: str = "fp16", | |
| save_every: int = 10, preview_every: int = 10, preview_steps: int = 10, | |
| gradient_checkpointing: bool = False, preview_enabled: bool = True, | |
| preview_prompt: str = "", preview_seed: int = 123456789, | |
| ) -> dict[str, object]: | |
| """Launch the user's Flow Matching worker and relay its structured progress.""" | |
| root = Path(str(ConfigManager(context.root).get("tool_folders", {}).get("flow_trainer", ""))).expanduser() | |
| script = root / "flow_matching_app.py" | |
| dataset = Path(dataset_dir).expanduser().resolve() | |
| output = Path(output_dir).expanduser().resolve() | |
| if not script.is_file(): | |
| raise ToolExecutionError("Flow Matching flow_matching_app.py was not found. Re-scan its folder in Settings.") | |
| if not dataset.is_dir(): | |
| raise ToolExecutionError("The selected Flow Matching dataset folder no longer exists.") | |
| if sum(1 for path in dataset.iterdir() if path.is_file() and path.suffix.lower() in IMAGE_EXTENSIONS) < 2: | |
| raise ToolExecutionError("The Flow Matching dataset needs at least two image files before training can start.") | |
| if not 1 <= int(epochs) <= 100_000: | |
| raise ToolExecutionError("Epoch count must be between 1 and 100000.") | |
| if not 64 <= int(resolution) <= 512 or int(resolution) % 16 or not 1 <= int(batch_size) <= 64 or not 1e-7 <= float(learning_rate) <= 0.1 or not 1 <= int(gradient_accumulation) <= 64 or not 0 <= int(workers) <= 16 or mixed_precision not in {"fp16", "no"} or min(int(save_every), int(preview_every), int(preview_steps)) < 1: | |
| raise ToolExecutionError("Flow training options are outside ADAM's safe range.") | |
| safe_name = model_name.strip() | |
| if not safe_name or len(safe_name) > 96 or any(char in safe_name for char in "<>:\\|?*\x00"): | |
| raise ToolExecutionError("Choose a short model name without filesystem-reserved characters.") | |
| output_root = (root / "output_flow_models").resolve() | |
| try: | |
| output.relative_to(output_root) | |
| except ValueError as exc: | |
| raise ToolExecutionError("Flow Matching outputs must stay inside output_flow_models.") from exc | |
| if output.exists(): | |
| raise ToolExecutionError("The chosen Flow Matching output folder already exists; ADAM will not overwrite it.") | |
| resume = Path(resume_from).expanduser().resolve() if resume_from else None | |
| if resume: | |
| try: | |
| metadata = json.loads((resume / "flow_model_info.json").read_text(encoding="utf-8")) | |
| if metadata.get("model_type") != "rectified_flow" or not (resume / "unet" / "config.json").is_file(): | |
| raise ValueError | |
| saved_resolution = int(metadata.get("resolution", 0) or 0) | |
| except (OSError, ValueError, TypeError, json.JSONDecodeError) as exc: | |
| raise ToolExecutionError("Choose a valid completed Flow Matching model to continue.") from exc | |
| if saved_resolution and saved_resolution != int(resolution): | |
| context.log( | |
| "Resolution-change fine-tune: loading " | |
| f"{saved_resolution}px Flow weights for training at {int(resolution)}px. " | |
| "Optimizer state will start fresh." | |
| ) | |
| output.parent.mkdir(parents=True, exist_ok=True) | |
| command = [ | |
| sys.executable, str(script), "--train-worker", "--data-dir", str(dataset), | |
| "--output-dir", str(output), "--model-name", safe_name, "--epochs", str(int(epochs)), | |
| "--resolution", str(int(resolution)), "--batch-size", str(int(batch_size)), "--learning-rate", str(float(learning_rate)), | |
| "--workers", str(int(workers)), "--gradient-accumulation", str(int(gradient_accumulation)), "--mixed-precision", mixed_precision, | |
| "--save-every", str(int(save_every)), "--preview-every", str(int(preview_every) if preview_enabled else int(epochs) + 1), "--preview-steps", str(int(preview_steps)), "--tf32", | |
| ] | |
| if gradient_checkpointing: | |
| command.append("--gradient-checkpointing") | |
| if resume: | |
| command.extend(["--continue-model", str(resume)]) | |
| context.log(f"Starting real Flow Matching training. Output folder: {output}") | |
| process = subprocess.Popen(command, cwd=str(root), stdout=subprocess.PIPE, stderr=subprocess.STDOUT, | |
| text=True, encoding="utf-8", errors="replace", shell=False) | |
| lines: queue.Queue[str | None] = queue.Queue() | |
| def read_output() -> None: | |
| assert process.stdout is not None | |
| for line in process.stdout: | |
| lines.put(line.rstrip()) | |
| lines.put(None) | |
| threading.Thread(target=read_output, daemon=True).start() | |
| context.progress(1, "Starting Flow Matching trainer") | |
| stopped = False | |
| stop_requested_at: float | None = None | |
| force_stop_sent = False | |
| suspended = False | |
| stop_file = output / "stop_flow_training.flag" | |
| while True: | |
| should_pause = not context.run_event.is_set() | |
| if should_pause != suspended: | |
| if set_process_tree_paused(process, should_pause): | |
| suspended = should_pause | |
| context.log("Flow Matching trainer paused safely." if suspended else "Flow Matching trainer resumed.") | |
| if context.cancel_event.is_set() and not stopped: | |
| if suspended: | |
| set_process_tree_paused(process, False) | |
| suspended = False | |
| stop_file.touch(exist_ok=True) | |
| stopped = True | |
| stop_requested_at = time.monotonic() | |
| context.log("Safe stop requested; waiting for Flow Matching to finish its current batch.") | |
| if stop_requested_at and not force_stop_sent and time.monotonic() - stop_requested_at > FORCE_STOP_TIMEOUT_SECONDS: | |
| terminate_process_tree(process, timeout=3) | |
| force_stop_sent = True | |
| context.log("Flow Matching did not stop in time; terminating the trainer process.") | |
| try: | |
| line = lines.get(timeout=0.15) | |
| if line and line.startswith("FLOW_EVENT:"): | |
| event = json.loads(line.split(":", 1)[1]) | |
| if event.get("type") == "progress" and not stopped: | |
| current = int(event.get("epoch", 0) or 0) | |
| context.progress(max(1, min(99, round(current * 100 / int(epochs)))), | |
| f"Finished epoch {current} of {epochs}") | |
| if preview_enabled and current and current % int(preview_every) == 0: | |
| candidate = Path(str(event.get("preview_path", ""))) if event.get("preview_path") else _latest_preview(output) | |
| if candidate: | |
| context.preview(candidate, epoch=current, | |
| next_epoch=min(int(epochs), current + int(preview_every)), | |
| prompt=preview_prompt, seed=int(preview_seed), steps=int(preview_steps)) | |
| elif event.get("type") == "warning": | |
| context.log(str(event.get("message", "Flow trainer warning."))) | |
| elif line: | |
| context.log(line) | |
| except queue.Empty: | |
| pass | |
| if process.poll() is not None and lines.empty(): | |
| break | |
| if stopped: | |
| raise ToolCancelled("Flow Matching training stopped by user.") | |
| if process.returncode != 0: | |
| raise ToolExecutionError(f"Flow Matching trainer exited with code {process.returncode}. See the job log for details.") | |
| context.progress(100, "Flow Matching training completed") | |
| return {"output_folder": str(output), "model_name": safe_name, "assets": [{ | |
| "kind": "model", "name": safe_name, "path": str(output), "trainer": "flow", | |
| "dataset_path": str(dataset), "checkpoint": str(output), "epochs": int(epochs), | |
| }]} | |
| def _stage_context(context: ToolContext, *, stage_index: int, stage_count: int) -> ToolContext: | |
| def report(percent: int, message: str, **details: object) -> None: | |
| overall = round(((stage_index + max(0, min(percent, 100)) / 100) / stage_count) * 100) | |
| context.progress(overall, f"Stage {stage_index + 1}/{stage_count} · {message}", **details) | |
| return ToolContext( | |
| root=context.root, job_id=context.job_id, tool=context.tool, | |
| cancel_event=context.cancel_event, run_event=context.run_event, | |
| progress_callback=report, log_callback=context.log_callback, | |
| preview_callback=context.preview_callback, step_delay=context.step_delay, | |
| ) | |
| def _stage_output(root: Path, stage_number: int, resolution: int) -> Path: | |
| base = root / f"stage-{stage_number:02d}-{resolution}px" | |
| if not base.exists(): | |
| return base | |
| attempt = 2 | |
| while (candidate := root / f"{base.name}-retry-{attempt}").exists(): | |
| attempt += 1 | |
| return candidate | |
| def train_flow( | |
| context: ToolContext, | |
| dataset_dir: str, | |
| model_name: str, | |
| epochs: int, | |
| output_dir: str, | |
| resume_from: str = "", | |
| resolution: int = 256, batch_size: int = 8, learning_rate: float = 0.0002, | |
| gradient_accumulation: int = 1, workers: int = 4, mixed_precision: str = "fp16", | |
| save_every: int = 10, preview_every: int = 10, preview_steps: int = 10, | |
| gradient_checkpointing: bool = False, preview_enabled: bool = True, | |
| preview_prompt: str = "", preview_seed: int = 123456789, | |
| progressive_stages: list[dict[str, object]] | None = None, | |
| progressive_auto_batch: bool = True, | |
| ) -> dict[str, object]: | |
| """Train one Flow model or carry it through a saved resolution curriculum.""" | |
| if not progressive_stages: | |
| return _train_flow_stage( | |
| context, dataset_dir, model_name, epochs, output_dir, resume_from, resolution, | |
| batch_size, learning_rate, gradient_accumulation, workers, mixed_precision, | |
| save_every, preview_every, preview_steps, gradient_checkpointing, | |
| preview_enabled, preview_prompt, preview_seed, | |
| ) | |
| stages = parse_stages(progressive_stages, trainer="flow", total_epochs=epochs) | |
| public_output = Path(output_dir).expanduser().resolve() | |
| if public_output.exists(): | |
| raise ToolExecutionError("The chosen Flow Matching output folder already exists; ADAM will not overwrite it.") | |
| stage_root = public_output.parent / f".{public_output.name}.progressive" | |
| state_path = stage_root / "progressive_state.json" | |
| stage_root.mkdir(parents=True, exist_ok=True) | |
| try: | |
| state = json.loads(state_path.read_text(encoding="utf-8")) | |
| except (OSError, json.JSONDecodeError): | |
| state = {"model_name": model_name, "stages": [], "completed": []} | |
| completed = state.get("completed", []) if isinstance(state.get("completed"), list) else [] | |
| completed_by_index = { | |
| int(item.get("index")): Path(str(item.get("output"))) | |
| for item in completed if isinstance(item, dict) and str(item.get("index", "")).isdigit() | |
| } | |
| final_resolution = stages[-1].resolution | |
| prior_model = resume_from | |
| context.log( | |
| "Progressive Flow Matching schedule: " + stage_summary(stages) + ". " | |
| + ("Auto batch caps are enabled." if progressive_auto_batch else "Using the same batch settings at every stage.") | |
| ) | |
| for index, stage in enumerate(stages): | |
| completed_output = completed_by_index.get(index) | |
| if completed_output and completed_output.is_dir(): | |
| prior_model = str(completed_output) | |
| context.log(f"Stage {index + 1}/{len(stages)} already completed; using its saved weights.") | |
| continue | |
| stage_output = _stage_output(stage_root, index + 1, stage.resolution) | |
| stage_batch, stage_accumulation = stage_batch_settings( | |
| trainer="flow", stage_resolution=stage.resolution, final_resolution=final_resolution, | |
| final_batch_size=batch_size, base_accumulation=gradient_accumulation, | |
| auto_batch=bool(progressive_auto_batch), | |
| ) | |
| context.log( | |
| f"Stage {index + 1}/{len(stages)}: {stage.resolution}px for {stage.epochs} epochs; " | |
| f"batch {stage_batch}, gradient accumulation {stage_accumulation}." | |
| ) | |
| _train_flow_stage( | |
| _stage_context(context, stage_index=index, stage_count=len(stages)), | |
| dataset_dir, model_name, stage.epochs, str(stage_output), prior_model, | |
| stage.resolution, stage_batch, learning_rate, stage_accumulation, workers, | |
| mixed_precision, min(save_every, stage.epochs), min(preview_every, stage.epochs), | |
| preview_steps, gradient_checkpointing or stage.resolution >= 384, | |
| preview_enabled, preview_prompt, preview_seed, | |
| ) | |
| prior_model = str(stage_output) | |
| completed.append({"index": index, "resolution": stage.resolution, "epochs": stage.epochs, "output": prior_model}) | |
| state.update({"stages": [{"resolution": item.resolution, "epochs": item.epochs} for item in stages], "completed": completed}) | |
| state_path.write_text(json.dumps(state, indent=2), encoding="utf-8") | |
| if not prior_model or not Path(prior_model).is_dir(): | |
| raise ToolExecutionError("Progressive Flow Matching training did not produce a final stage model.") | |
| shutil.move(prior_model, public_output) | |
| context.progress(100, "Progressive Flow Matching training completed") | |
| return { | |
| "output_folder": str(public_output), "model_name": model_name, | |
| "progressive_stages": [{"resolution": stage.resolution, "epochs": stage.epochs} for stage in stages], | |
| "assets": [{ | |
| "kind": "model", "name": model_name, "path": str(public_output), "trainer": "flow", | |
| "dataset_path": str(Path(dataset_dir).expanduser().resolve()), "checkpoint": str(public_output), | |
| "epochs": sum(stage.epochs for stage in stages), | |
| }], | |
| } | |