KiwiMate's picture
Update handler.py
5ceedab verified
Raw
History Blame Contribute Delete
15.8 kB
import os
import glob
import io
import re
import base64
import logging
import time
from typing import Dict, List, Any, Union, Optional
import torch
from PIL import Image
from diffusers import (
AutoPipelineForText2Image,
StableDiffusionPipeline,
StableDiffusionXLPipeline,
UNet2DConditionModel,
)
from transformers import (
CLIPTextModel,
CLIPTokenizer,
CLIPTextModelWithProjection,
AutoTokenizer,
)
from safetensors import safe_open
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger("hf-endpoint-handler")
class EndpointHandler:
"""
Custom handler for Hugging Face Dedicated Inference Endpoints.
Handles Text-to-Image generation with automatic single-file detection,
Flux/SDXL/SD1.5 architecture auto-detection, step-checkpoint filtering,
missing UNet/TextEncoder base fallback, and LoRA loading.
"""
def __init__(self, path: str = ""):
logger.info(f"Initializing EndpointHandler from model directory: '{path}'")
self.device = "cuda" if torch.cuda.is_available() else "cpu"
self.dtype = torch.float16 if self.device == "cuda" else torch.float32
# 1. Try standard diffusers folder load first
try:
logger.info("Attempting standard diffusers repository load...")
self.pipeline = AutoPipelineForText2Image.from_pretrained(
path,
torch_dtype=self.dtype,
use_safetensors=True,
)
logger.info("Successfully loaded standard diffusers folder pipeline.")
except Exception as e:
logger.warning(f"Standard load failed ({e}). Running smart single-file checkpoint loader...")
self.pipeline = self._load_single_file_checkpoint(path)
# 2. Apply GPU optimizations
if self.device == "cuda":
self.pipeline.to("cuda")
if hasattr(self.pipeline, "enable_vae_slicing"):
self.pipeline.enable_vae_slicing()
if hasattr(self.pipeline, "enable_vae_tiling"):
self.pipeline.enable_vae_tiling()
logger.info(f"Pipeline successfully ready on {self.device} with class: {type(self.pipeline).__name__}")
def _select_best_checkpoint(self, path: str) -> str:
"""
Filters out training step checkpoints (_000002250, optimizer.pt, etc.)
and prioritizes merged or main model weight files.
"""
candidate_files = []
if os.path.isdir(path):
for ext in ("*.safetensors", "*.ckpt", "*.pt", "*.bin"):
candidate_files.extend(glob.glob(os.path.join(path, ext)))
candidate_files.extend(glob.glob(os.path.join(path, "**", ext), recursive=True))
elif os.path.isfile(path):
return path
if not candidate_files:
raise FileNotFoundError(f"No model checkpoint files found in directory '{path}'")
filtered_files = []
for f in candidate_files:
fname = os.path.basename(f).lower()
# Ignore optimizer files and subfolder weights
if "optimizer" in fname or "text_encoder" in f or "unet" in f or "vae" in f:
continue
# Ignore intermediate training step checkpoints (e.g. _000002250.safetensors)
if re.search(r"_\d{5,}", fname) or re.search(r"step_?\d+", fname):
logger.info(f"Filtering out intermediate step checkpoint file: {os.path.basename(f)}")
continue
filtered_files.append(f)
if not filtered_files:
logger.warning("All files matched step pattern, falling back to full list...")
filtered_files = candidate_files
# Priority 1: Merged files
merged = [f for f in filtered_files if "merged" in os.path.basename(f).lower()]
if merged:
logger.info(f"Selected merged checkpoint: {os.path.basename(merged[0])}")
return merged[0]
# Priority 2: Flux files
flux = [f for f in filtered_files if "flux" in os.path.basename(f).lower()]
if flux:
logger.info(f"Selected Flux model checkpoint: {os.path.basename(flux[0])}")
return flux[0]
logger.info(f"Selected primary checkpoint file: {os.path.basename(filtered_files[0])}")
return filtered_files[0]
def _inspect_keys(self, target_file: str) -> dict:
"""
Inspects safetensors header keys to determine architecture and component presence.
"""
is_lora = False
is_flux = False
is_sdxl = False
has_unet = False
if target_file.endswith(".safetensors"):
try:
with safe_open(target_file, framework="pt") as f:
keys = f.keys()
for k in keys:
if "lora_" in k or ".lora_down" in k or ".lora_up" in k:
is_lora = True
if "double_blocks" in k or "single_blocks" in k or "guidance_in" in k:
is_flux = True
if "conditioner.embedders" in k or "text_encoders.top" in k:
is_sdxl = True
if "model.diffusion_model" in k or "unet" in k:
has_unet = True
except Exception as e:
logger.warning(f"Failed to inspect safetensors keys: {e}")
filename_lower = os.path.basename(target_file).lower()
if "flux" in filename_lower:
is_flux = True
return {
"is_lora": is_lora,
"is_flux": is_flux,
"is_sdxl": is_sdxl,
"has_unet": has_unet,
}
def _load_single_file_checkpoint(self, path: str):
target_file = self._select_best_checkpoint(path)
info = self._inspect_keys(target_file)
logger.info(f"Checkpoint inspection result for {os.path.basename(target_file)}: {info}")
# --- CASE A: FLUX Architecture ---
if info["is_flux"]:
logger.info("Flux architecture detected. Initializing Flux pipeline...")
try:
from diffusers import FluxPipeline
if info["is_lora"]:
logger.info("Loading Flux base model and attaching LoRA weights...")
pipe = FluxPipeline.from_pretrained("black-forest-labs/FLUX.1-schnell", torch_dtype=self.dtype)
pipe.load_lora_weights(target_file)
return pipe
else:
try:
return FluxPipeline.from_single_file(target_file, torch_dtype=self.dtype)
except Exception as err:
logger.warning(f"Flux single-file load failed ({err}), loading base FLUX.1-schnell...")
pipe = FluxPipeline.from_pretrained("black-forest-labs/FLUX.1-schnell", torch_dtype=self.dtype)
if info["is_lora"] or "merged" in target_file.lower():
try:
pipe.load_lora_weights(target_file)
except Exception:
pass
return pipe
except Exception as flux_err:
logger.warning(f"Flux pipeline load failed ({flux_err}). Falling back to SDXL...")
# --- CASE B: SDXL Architecture ---
if info["is_sdxl"] or "xl" in os.path.basename(target_file).lower():
logger.info("SDXL architecture detected. Attempting SDXL single file load...")
if info["is_lora"]:
pipe = StableDiffusionXLPipeline.from_pretrained("stabilityai/stable-diffusion-xl-base-1.0", torch_dtype=self.dtype)
pipe.load_lora_weights(target_file)
return pipe
try:
return StableDiffusionXLPipeline.from_single_file(target_file, torch_dtype=self.dtype)
except Exception as err:
logger.warning(f"SDXL single-file missing components ({err}). Supplying SDXL base components...")
unet = UNet2DConditionModel.from_pretrained(
"stabilityai/stable-diffusion-xl-base-1.0", subfolder="unet", torch_dtype=self.dtype
)
text_encoder_1 = CLIPTextModel.from_pretrained("openai/clip-vit-large-patch14", torch_dtype=self.dtype)
text_encoder_2 = CLIPTextModelWithProjection.from_pretrained(
"laion/CLIP-ViT-bigG-14-laion2B-39B-b160k", torch_dtype=self.dtype
)
tokenizer_1 = CLIPTokenizer.from_pretrained("openai/clip-vit-large-patch14")
tokenizer_2 = AutoTokenizer.from_pretrained("laion/CLIP-ViT-bigG-14-laion2B-39B-b160k")
try:
return StableDiffusionXLPipeline.from_single_file(
target_file,
unet=unet,
text_encoder=text_encoder_1,
text_encoder_2=text_encoder_2,
tokenizer=tokenizer_1,
tokenizer_2=tokenizer_2,
torch_dtype=self.dtype,
)
except Exception:
pipe = StableDiffusionXLPipeline.from_pretrained("stabilityai/stable-diffusion-xl-base-1.0", torch_dtype=self.dtype)
try:
pipe.load_lora_weights(target_file)
except Exception:
pass
return pipe
# --- CASE C: Standard SD 1.5 Architecture ---
logger.info("Attempting SD 1.5 single file pipeline load...")
if info["is_lora"]:
pipe = StableDiffusionPipeline.from_pretrained("runwayml/stable-diffusion-v1-5", torch_dtype=self.dtype)
pipe.load_lora_weights(target_file)
return pipe
try:
return StableDiffusionPipeline.from_single_file(target_file, torch_dtype=self.dtype)
except Exception as err:
logger.warning(f"SD 1.5 single file missing components ({err}). Supplying base UNet and CLIP...")
unet = UNet2DConditionModel.from_pretrained("runwayml/stable-diffusion-v1-5", subfolder="unet", torch_dtype=self.dtype)
text_encoder = CLIPTextModel.from_pretrained("openai/clip-vit-large-patch14", torch_dtype=self.dtype)
tokenizer = CLIPTokenizer.from_pretrained("openai/clip-vit-large-patch14")
try:
return StableDiffusionPipeline.from_single_file(
target_file,
unet=unet,
text_encoder=text_encoder,
tokenizer=tokenizer,
torch_dtype=self.dtype,
)
except Exception:
pipe = StableDiffusionPipeline.from_pretrained("runwayml/stable-diffusion-v1-5", torch_dtype=self.dtype)
try:
pipe.load_lora_weights(target_file)
except Exception:
pass
return pipe
def _encode_image_to_base64(self, image: Image.Image, image_format: str = "PNG", quality: int = 95) -> str:
buffer = io.BytesIO()
if image_format.upper() in ["JPG", "JPEG"]:
image.save(buffer, format="JPEG", quality=quality)
else:
image.save(buffer, format="PNG")
buffer.seek(0)
return base64.b64encode(buffer.getvalue()).decode("utf-8")
def _extract_parameters(self, data: Dict[str, Any]) -> Dict[str, Any]:
inputs = data.get("inputs", "")
parameters = data.get("parameters", {})
if isinstance(inputs, dict):
prompt = inputs.get("prompt", "")
negative_prompt = inputs.get("negative_prompt", parameters.get("negative_prompt", None))
else:
prompt = str(inputs) if inputs else ""
negative_prompt = parameters.get("negative_prompt", None)
if not prompt.strip():
raise ValueError("Parameter 'prompt' (or 'inputs') must be a non-empty string.")
pipe_name = type(self.pipeline).__name__.lower()
default_dim = 1024 if ("xl" in pipe_name or "flux" in pipe_name) else 512
height = (int(parameters.get("height", default_dim)) // 8) * 8
width = (int(parameters.get("width", default_dim)) // 8) * 8
num_inference_steps = int(parameters.get("num_inference_steps", 28 if "flux" in pipe_name else 30))
guidance_scale = float(parameters.get("guidance_scale", 3.5 if "flux" in pipe_name else 7.5))
seed = parameters.get("seed", None)
num_images_per_prompt = min(max(int(parameters.get("num_images_per_prompt", 1)), 1), 4)
output_format = str(parameters.get("output_format", "pil")).lower()
image_format = str(parameters.get("image_format", "PNG")).upper()
return {
"prompt": prompt,
"negative_prompt": negative_prompt,
"height": height,
"width": width,
"num_inference_steps": num_inference_steps,
"guidance_scale": guidance_scale,
"seed": int(seed) if seed is not None else None,
"num_images_per_prompt": num_images_per_prompt,
"output_format": output_format,
"image_format": image_format,
}
def __call__(self, data: Dict[str, Any]) -> Union[List[Dict[str, Any]], Image.Image]:
start_time = time.time()
try:
params = self._extract_parameters(data)
logger.info(f"Processing prompt: '{params['prompt'][:60]}...'")
generator = None
if params["seed"] is not None:
generator = torch.Generator(device=self.device).manual_seed(params["seed"])
generation_args = {
"prompt": params["prompt"],
"negative_prompt": params["negative_prompt"],
"height": params["height"],
"width": params["width"],
"num_inference_steps": params["num_inference_steps"],
"guidance_scale": params["guidance_scale"],
"num_images_per_prompt": params["num_images_per_prompt"],
"generator": generator,
}
# Flux pipelines do not use negative_prompt
if "flux" in type(self.pipeline).__name__.lower():
generation_args.pop("negative_prompt", None)
# Filter out None values
generation_args = {k: v for k, v in generation_args.items() if v is not None}
with torch.inference_mode():
output = self.pipeline(**generation_args)
images = output.images
elapsed_seconds = round(time.time() - start_time, 3)
logger.info(f"Generated {len(images)} image(s) in {elapsed_seconds}s")
if params["output_format"] == "pil" and len(images) == 1:
return images[0]
response_payload = []
for index, img in enumerate(images):
b64_image = self._encode_image_to_base64(img, image_format=params["image_format"])
response_payload.append({
"image": b64_image,
"format": params["image_format"],
"width": img.width,
"height": img.height,
"index": index,
"execution_time": elapsed_seconds,
})
return response_payload
except Exception as err:
logger.error(f"Inference execution failed: {str(err)}", exc_info=True)
return [{"error": str(err), "status": "failed"}]