Author SHA1 Message Date
Sebastian Monroy 08b953f40f Merge branch 'develop' 2026-02-17 15:41:59 +00:00
Sebastian Monroy 207de4487c fix(worker_consumer): propagate submit failures to retry/DLQ and emit failed status
- raise a dedicated JobSubmissionError when ComfyUI prompt submission fails
- re-raise submission failures in process_message so consume loop does not delete SQS messages
- emit failed status updates on submit-time errors using status_policy.fail_status fallback
- remove failed-status dedupe guard so each retry attempt remains observable
- align parse-error logging with retry behavior ("Message will be retried")
2026-02-17 15:40:02 +00:00
Sebastian Monroy c36846c008 Merge branch 'main' into develop 2026-02-17 12:47:45 +00:00
Sebastian Monroy ab7792c8ed add NilorUserInput_Seed node with robust randomization (#14)
Add NilorUserInput_Seed with sentinel-based randomization, modulo normalization for invalid values, cache-busting via IS_CHANGED, and isolated RNG state to avoid global random pollution.
2026-02-16 15:21:22 +00:00
Sebastian Monroy aa594e80d6 Merge branch 'develop' 2026-01-16 18:15:07 +00:00
3 changed files with 175 additions and 88 deletions
+50 -86
View File
@@ -8,15 +8,13 @@ import imageio.v2 as imageio
import mimetypes
import boto3
import json
import tempfile
import os
from .logger import logger
from .config.config import load_nilor_nodes_config
# Load shared configuration once
_CFG = load_nilor_nodes_config()
# --- Node Categories ---
category = "Nilor Nodes 👺"
subcategories = {
@@ -56,10 +54,10 @@ class MediaStreamInput:
CATEGORY = category + subcategories["streaming"]
def download(
self,
presigned_download_url: str,
format: str,
input_name: str = "default_input",
self,
presigned_download_url: str,
format: str,
input_name: str = "default_input",
):
logger.info(
f"ℹ️\u2009 Nilor-Nodes: MediaStreamInput: Downloading from {presigned_download_url} for input '{input_name}' with format '{format}'"
@@ -96,35 +94,19 @@ class MediaStreamInput:
return self._process_image_batch(asset_responses)
# --- Single-file download ---
if format == "video":
# Stream video to temp file to avoid loading entire video into RAM
temp_file = tempfile.NamedTemporaryFile(delete=False, suffix=".mp4")
try:
logger.info(
f"ℹ️\u2009 Nilor-Nodes (MediaStreamInput): Streaming video to temp file: {temp_file.name}")
with requests.get(presigned_download_url, timeout=180, stream=True) as response:
response.raise_for_status()
for chunk in response.iter_content(chunk_size=8192):
temp_file.write(chunk)
temp_file.close()
return self._process_video(temp_file.name)
finally:
# Clean up temp file
if os.path.exists(temp_file.name):
os.unlink(temp_file.name)
else:
# For images, load into memory (they're small)
response = requests.get(presigned_download_url, timeout=180)
response.raise_for_status()
media_bytes = response.content
response = requests.get(presigned_download_url, timeout=180)
response.raise_for_status()
media_bytes = response.content
if format == "image":
return self._process_image(media_bytes)
else:
# Should not happen if UI choices are respected
raise ValueError(
f"🛑\u2009 Nilor-Nodes (MediaStreamInput): Unsupported format '{format}' for single media download."
)
if format == "video":
return self._process_video(media_bytes)
elif format == "image":
return self._process_image(media_bytes)
else:
# Should not happen if UI choices are respected
raise ValueError(
f"🛑\u2009 Nilor-Nodes (MediaStreamInput): Unsupported format '{format}' for single media download."
)
except requests.RequestException as e:
logger.error(
@@ -174,45 +156,27 @@ class MediaStreamInput:
logger.info("✅ Nilor-Nodes (MediaStreamInput): Image processing successful.")
return (image_tensor,)
def _process_video(self, video_path):
logger.info(f"ℹ️\u2009 Nilor-Nodes (MediaStreamInput): Processing video from {video_path}...")
# Open video to get metadata first
with imageio.get_reader(video_path, format="mp4") as reader:
# Get video metadata
metadata = reader.get_meta_data()
num_frames = reader.count_frames()
if num_frames == 0:
raise ValueError(
"🛑\u2009 Nilor-Nodes (MediaStreamInput): No frames could be read from the video."
)
# Read first frame to get dimensions
first_frame = reader.get_data(0)
height, width = first_frame.shape[:2]
logger.info(
f"ℹ️\u2009 Nilor-Nodes (MediaStreamInput): Video has {num_frames} frames at {width}x{height}"
)
# Pre-allocate tensor for all frames (N, H, W, 3)
video_tensor = torch.empty((num_frames, height, width, 3), dtype=torch.float32)
# Process first frame (already read for dimensions)
pil_image = Image.fromarray(first_frame).convert("RGB")
numpy_image = np.array(pil_image).astype(np.float32) / 255.0
video_tensor[0] = torch.from_numpy(numpy_image)
# Read remaining frames by explicit index to avoid iterator position ambiguity
for i in range(1, num_frames):
frame = reader.get_data(i)
def _process_video(self, video_bytes):
logger.info("ℹ️\u2009 Nilor-Nodes (MediaStreamInput): Processing as video...")
frames = []
with imageio.get_reader(io.BytesIO(video_bytes), format="mp4") as reader:
for frame in reader:
# Convert frame to RGB PIL Image and then to tensor
pil_image = Image.fromarray(frame).convert("RGB")
numpy_image = np.array(pil_image).astype(np.float32) / 255.0
video_tensor[i] = torch.from_numpy(numpy_image)
tensor_frame = torch.from_numpy(numpy_image)
frames.append(tensor_frame)
logger.info(
f"✅ Nilor-Nodes (MediaStreamInput): Video processing successful. Tensor shape: {video_tensor.shape}"
if not frames:
raise ValueError(
"🛑\u2009 Nilor-Nodes (MediaStreamInput): No frames could be read from the video."
)
# Stack frames into a single tensor (batch of images)
video_tensor = torch.stack(frames)
logging.info(
f"✅ Nilor-Nodes (MediaStreamInput): Video processing successful. Image Shape: {video_tensor.shape}"
)
return (video_tensor,)
@@ -276,21 +240,21 @@ class MediaStreamOutput:
CATEGORY = category + subcategories["streaming"]
def upload_and_notify(
self,
images,
format,
content_id,
venue,
canvas,
scene,
presigned_upload_url,
job_completions_queue_url,
output_object_keys,
framerate,
output_name: str = "default_output",
prompt=None,
extra_pnginfo=None,
job_type: str | None = None,
self,
images,
format,
content_id,
venue,
canvas,
scene,
presigned_upload_url,
job_completions_queue_url,
output_object_keys,
framerate,
output_name: str = "default_output",
prompt=None,
extra_pnginfo=None,
job_type: str | None = None,
):
if not content_id:
raise ValueError(
+82
View File
@@ -3,6 +3,9 @@ subcategories = {
"io": "/IO",
}
import random
from datetime import datetime
from .controllers import CONTROLLER_HOOK
@@ -50,6 +53,83 @@ class NilorUserInput_Int:
return (value, None)
class NilorUserInput_Seed:
MAX_COMFYUI_SEED = 1125899906842624
SEED_RANDOM_STATE = None
@classmethod
def _ensure_seed_random_state(cls):
if cls.SEED_RANDOM_STATE is not None:
return
initial_random_state = random.getstate()
random.seed(datetime.now().timestamp())
cls.SEED_RANDOM_STATE = random.getstate()
random.setstate(initial_random_state)
@classmethod
def generate_random_seed(cls):
cls._ensure_seed_random_state()
prev_random_state = random.getstate()
random.setstate(cls.SEED_RANDOM_STATE)
seed = random.randint(0, cls.MAX_COMFYUI_SEED)
cls.SEED_RANDOM_STATE = random.getstate()
random.setstate(prev_random_state)
return seed
@classmethod
def resolve_seed(cls, value):
if value in (None, 0, -1):
return cls.generate_random_seed()
try:
return int(value) % (cls.MAX_COMFYUI_SEED + 1)
except (TypeError, ValueError):
return cls.generate_random_seed()
@classmethod
def INPUT_TYPES(cls):
return {
"required": {
"input_name": (
"STRING",
{"default": "my_seed_input", "multiline": False},
),
"value": (
"INT",
{
"default": -1,
"min": -1,
"max": cls.MAX_COMFYUI_SEED,
},
),
},
"hidden": {
"prompt": "PROMPT",
"extra_pnginfo": "EXTRA_PNGINFO",
"unique_id": "UNIQUE_ID",
},
}
RETURN_TYPES = ("INT", CONTROLLER_HOOK)
RETURN_NAMES = ("seed", "_controller_hook")
FUNCTION = "get_value"
CATEGORY = category + subcategories["io"]
@classmethod
def IS_CHANGED(
cls, input_name, value, prompt=None, extra_pnginfo=None, unique_id=None
):
# Force node re-execution while using randomize sentinel values.
return cls.resolve_seed(value)
def get_value(
self, input_name, value, prompt=None, extra_pnginfo=None, unique_id=None
):
value = self.resolve_seed(value)
return (value, None)
class NilorUserInput_Float:
@classmethod
def INPUT_TYPES(cls):
@@ -97,6 +177,7 @@ class NilorUserInput_Boolean:
NODE_CLASS_MAPPINGS = {
"NilorUserInput_String": NilorUserInput_String,
"NilorUserInput_Int": NilorUserInput_Int,
"NilorUserInput_Seed": NilorUserInput_Seed,
"NilorUserInput_Float": NilorUserInput_Float,
"NilorUserInput_Boolean": NilorUserInput_Boolean,
}
@@ -104,6 +185,7 @@ NODE_CLASS_MAPPINGS = {
NODE_DISPLAY_NAME_MAPPINGS = {
"NilorUserInput_String": "👺 User Input (String)",
"NilorUserInput_Int": "👺 User Input (Int)",
"NilorUserInput_Seed": "👺 User Input (Seed)",
"NilorUserInput_Float": "👺 User Input (Float)",
"NilorUserInput_Boolean": "👺 User Input (Boolean)",
}
+43 -2
View File
@@ -29,6 +29,10 @@ from .config.config import load_nilor_nodes_config, NilorNodesConfig
_CFG: NilorNodesConfig = load_nilor_nodes_config()
class JobSubmissionError(Exception):
"""Raised when a job cannot be submitted to local ComfyUI."""
class WorkerConsumer:
def __init__(self, cfg: NilorNodesConfig):
self.session = get_session()
@@ -397,7 +401,19 @@ class WorkerConsumer:
return
# Submit to ComfyUI
await self._submit_job_to_comfyui(content_id, job_payload)
try:
await self._submit_job_to_comfyui(content_id, job_payload)
except JobSubmissionError as e:
logger.error(
f"🛑\u2009 Nilor-Nodes (worker_consumer): Submission failed for content_id {content_id}: {e}. Message will be retried/DLQ'd."
)
await self._emit_failed_status_for_submission_error(
content_id=content_id,
job_payload=job_payload,
error_message=str(e),
)
# Re-raise so consume_loop does not delete the message.
raise
# Cache context for subsequent status updates
try:
@@ -418,6 +434,28 @@ class WorkerConsumer:
# Re-raise to prevent deletion from queue if we want SQS to handle retry
raise
async def _emit_failed_status_for_submission_error(
self, content_id, job_payload, error_message: str
):
"""Best-effort failed status emission for submit-time errors."""
policy = job_payload.get("status_policy") or {}
fail_status = policy.get("fail_status", "failed")
await self._send_status_update(
content_id,
fail_status,
job_payload.get("venue"),
job_payload.get("canvas"),
job_payload.get("scene"),
job_payload.get("job_type"),
)
logger.info(
"ℹ️\u2009 Nilor-Nodes (worker_consumer): Emitted failed status '%s' for content_id %s after submit error: %s",
fail_status,
content_id,
error_message,
)
async def _submit_job_to_comfyui(self, content_id, workflow_data):
"""Submits a single job to the ComfyUI API."""
try:
@@ -484,15 +522,18 @@ class WorkerConsumer:
logger.error(
f"🛑\u2009 Nilor-Nodes (worker_consumer): Failed to submit job to ComfyUI: {e}. Message will be retried."
)
raise JobSubmissionError(str(e)) from e
except (json.JSONDecodeError, KeyError) as e:
logger.error(
f"🛑\u2009 Nilor-Nodes (worker_consumer): Failed to parse ComfyUI response: {e}. Discarding malformed response."
f"🛑\u2009 Nilor-Nodes (worker_consumer): Failed to parse ComfyUI response: {e}. Message will be retried."
)
raise JobSubmissionError(f"Malformed ComfyUI response: {e}") from e
except Exception as e:
logger.error(
f"🛑\u2009 Nilor-Nodes (worker_consumer): An unexpected error occurred while submitting job to ComfyUI: {e}",
exc_info=True,
)
raise JobSubmissionError(str(e)) from e
async def _send_status_update(
self, content_id, status, venue=None, canvas=None, scene=None, job_type=None