add sse support

This commit is contained in:
qnsh
2026-08-28 15:16:26 +08:00
parent b5e602133f
commit a131426c6f
10 changed files with 1128 additions and 33 deletions
+2
View File
@@ -33,6 +33,8 @@ The ModelScope image generation interface only requires you to fill in the corre
The node supports text, image, and video inputs. Video is sent directly using the current protocol's content format; if the target API or protocol does not support video, its error is surfaced as an explicit video-unsupported message. File input is not implemented in this version. The `OpenAI Text Advanced Options` node accepts a JSON object for protocol/API-specific parameters, for example `{"temperature":0.7,"max_output_tokens":4096}`. JSON options cannot override request fields such as `model`, `messages`, `input`, `instructions`, `stream`, `api_key`, `base_url`, or `timeout`.
When `stream` is enabled, the model response is streamed to the client as it is generated using server-sent events (SSE). Streaming and `skill_options` cannot currently be enabled together.
### skills
The `OpenAI Text Skill Options` node discovers local skills organized around `SKILL.md`. Select a Skill and connect the node to `OpenAI Text API` through `skill_options`. Configure Skill locations with `skills.paths`; relative paths resolve from this plugin directory, and the default location is `skills/`. Set `skills.allow_call` to `true` to enable Skill calls.
+2
View File
@@ -33,6 +33,8 @@ git clone https://github.com/ycyy/ComfyUI-YCYY-API.git
节点支持文本、图像和视频输入。视频会按当前协议的格式直接发送;如果目标 API 或协议不支持视频,接口错误会转换为明确的视频不支持提示。`files` 输入当前版本暂不支持。`OpenAI 文本高级选项(JSON)` 节点接受协议或 API 特有参数,例如 `{"temperature":0.7,"max_output_tokens":4096}`。JSON 参数不能覆盖 `model`、`messages`、`input`、`instructions`、`stream`、`api_key`、`base_url` 或 `timeout` 等请求字段。
启用 `stream` 后,模型响应会在生成过程中通过服务器发送事件(SSE)流式传输到客户端。当前流式模式不能与 `skill_options` 同时启用。
### skills
`OpenAI 文本 Skill 选项` 节点可以发现以 `SKILL.md` 组织的本地 Skill。选择 Skill 后,通过 `skill_options` 连接到 `OpenAI 文本 API` 节点。使用 `skills.paths` 配置 Skill 位置;相对路径以本插件目录为基准,默认位置为 `skills/`。将 `skills.allow_call` 设置为 `true` 后即可启用 Skill 调用。
+2
View File
@@ -19,6 +19,7 @@ from .openai.openai_text_node import *
from .options.openai_text_advanced_options_node import *
from .options.openai_text_skill_options_node import *
from .images.image_compare import ImageCompare
from .text.preview_api_result_node import PreviewAPIResult
WEB_DIRECTORY = "./web/js"
@@ -44,6 +45,7 @@ class APIExtension(ComfyExtension):
OpenAITextAPI,
OpenAITextAdvancedOptions,
OpenAITextSkillOptions,
PreviewAPIResult,
]
+11
View File
@@ -9,6 +9,7 @@
"user_prompt": {"name": "user_prompt"},
"persist_context": {"name": "persist_context", "tooltip": "Persist chat context between calls"},
"clear_history": {"name": "clear_history", "tooltip": "Clear the current conversation history"},
"stream": {"name": "stream", "tooltip": "If true, the model response is streamed to the client as it is generated using server-sent events (SSE)."},
"images": {"name": "images", "tooltip": "Optional image input"},
"videos": {"name": "videos", "tooltip": "Optional video input; unsupported APIs return an error"},
"config_options": {"name": "config_options", "tooltip": "Configuration override"},
@@ -22,6 +23,16 @@
"2": {"name": "Skill Trace"}
}
},
"YCYY_Preview_API_Result": {
"display_name": "Preview API Result",
"description": "Render text as a Markdown preview with copy support.",
"inputs": {
"source": {"name": "source", "tooltip": "Markdown or plain text to preview"}
},
"outputs": {
"0": {"name": "text"}
}
},
"YCYY_OpenAI_Text_Skill_Options": {
"display_name": "OpenAI Text Skill Options",
"description": "Load a local SKILL.md package for progressive use by OpenAI Text API.",
+11
View File
@@ -9,6 +9,7 @@
"user_prompt": {"name": "用户提示词"},
"persist_context": {"name": "保持上下文", "tooltip": "在多次调用之间保持对话上下文"},
"clear_history": {"name": "清除历史", "tooltip": "清除当前会话历史"},
"stream": {"name": "流式输出", "tooltip": "设为 true 时,模型响应会在生成过程中通过服务器发送事件(SSE)流式传输到客户端。"},
"images": {"name": "图像", "tooltip": "可选图像输入"},
"videos": {"name": "视频", "tooltip": "可选视频输入;不支持视频的接口会返回错误"},
"config_options": {"name": "配置选项", "tooltip": "配置覆盖"},
@@ -22,6 +23,16 @@
"2": {"name": "Skill 调用记录"}
}
},
"YCYY_Preview_API_Result": {
"display_name": "Preview API Result",
"description": "将文本渲染为 Markdown 网页预览,并支持复制。",
"inputs": {
"source": {"name": "文本", "tooltip": "要预览的 Markdown 或纯文本"}
},
"outputs": {
"0": {"name": "文本"}
}
},
"YCYY_OpenAI_Text_Skill_Options": {
"display_name": "OpenAI 文本 Skill 选项",
"description": "加载本地 SKILL.md,并供 OpenAI 文本 API 渐进式读取引用资料。",
+120 -4
View File
@@ -1,5 +1,7 @@
import json
import hashlib
import re
from uuid import uuid4
from aiohttp import web
from server import PromptServer
from comfy_api.latest import io
@@ -9,12 +11,77 @@ from ..utils.image_utils import tensor_to_base64_string
from ..utils.request_utils import (
get_proxy_config,
post_openai_json,
post_openai_stream,
resolve_endpoint,
video_to_data_uri,
)
from ..utils.skill_utils import SkillRequestContext
STREAM_EVENT = "ycyy_openai_text_stream"
def _safe_stream_error(exc):
detail = str(exc).replace("\r", " ").replace("\n", " ")
detail = re.sub(r"(?i)bearer\s+[a-z0-9._~+/-]+", "Bearer <redacted>", detail)
detail = re.sub(r"(?i)(api[_-]?key[=:]\s*)[^\s&]+", r"\1<redacted>", detail)
detail = re.sub(r"\bsk-[a-zA-Z0-9_-]{8,}\b", "sk-<redacted>", detail)
return detail[:300]
class _TextStreamSink:
"""Route text deltas to the executing ComfyUI client immediately."""
def __init__(self, node_id):
try:
from comfy_execution.utils import get_executing_context
context = get_executing_context()
except (ImportError, RuntimeError):
context = None
server = PromptServer.instance
context_node_id = getattr(context, "node_id", None)
effective_node_id = node_id if node_id is not None else context_node_id
self.node_id = str(effective_node_id) if effective_node_id is not None else ""
self.prompt_id = getattr(context, "prompt_id", None)
self.run_id = uuid4().hex
self.client_id = getattr(server, "client_id", None)
self.seq = 0
self.activities = set()
def _send(self, phase, **extra):
PromptServer.instance.send_sync(
STREAM_EVENT,
{
"node_id": self.node_id,
"prompt_id": self.prompt_id,
"run_id": self.run_id,
"seq": self.seq,
"phase": phase,
**extra,
},
self.client_id,
)
self.seq += 1
def start(self):
self._send("start")
def delta(self, value):
if value:
self._send("delta", delta=value)
def activity(self, kind):
if kind and kind not in self.activities:
self.activities.add(kind)
self._send("activity", activity=kind)
def end(self, value):
self._send("end", text=value)
def error(self, exc):
self._send("error", message=_safe_stream_error(exc))
@PromptServer.instance.routes.get("/ycyy/openai/apis/all")
async def get_all_openai_apis(request):
try:
@@ -43,6 +110,21 @@ class OpenAITextAPI(io.ComfyNode):
_conversation_history = {}
_max_history_items = 40
@classmethod
def _runtime_node_id(cls, explicit_id=None):
"""Resolve the client graph node id for ComfyUI V3 and direct tests."""
if explicit_id is not None:
return explicit_id
hidden_id = getattr(getattr(cls, "hidden", None), "unique_id", None)
if hidden_id is not None:
return hidden_id
try:
from comfy_execution.utils import get_executing_context
context = get_executing_context()
except (ImportError, RuntimeError):
context = None
return getattr(context, "node_id", None)
@classmethod
def define_schema(cls) -> io.Schema:
apis = get_openai_apis()
@@ -62,6 +144,14 @@ class OpenAITextAPI(io.ComfyNode):
io.String.Input(id="user_prompt", multiline=True),
io.Boolean.Input(id="persist_context", default=True),
io.Boolean.Input(id="clear_history", default=False),
io.Boolean.Input(
id="stream",
default=False,
tooltip=(
"If true, the model response is streamed to the client as it is "
"generated using server-sent events (SSE)."
),
),
io.Image.Input("images", optional=True, tooltip="Optional image input"),
io.Video.Input("videos", optional=True, tooltip="Optional video input"),
io.AnyType.Input(id="config_options", optional=True),
@@ -190,7 +280,7 @@ class OpenAITextAPI(io.ComfyNode):
@classmethod
def execute(cls, api_name, model, system_prompt, user_prompt, persist_context,
clear_history=False, images=None, videos=None, config_options=None,
clear_history=False, stream=False, images=None, videos=None, config_options=None,
proxy_options=None, advanced_options=None, skill_options=None,
unique_id=None) -> io.NodeOutput:
if not user_prompt or not user_prompt.strip():
@@ -199,9 +289,15 @@ class OpenAITextAPI(io.ComfyNode):
base_url, api_key, timeout, protocol = cls._resolve_api_settings(api, config_options)
endpoint = resolve_endpoint(base_url, protocol)
skill = SkillRequestContext.create(skill_options, protocol)
stream_enabled = bool(stream)
runtime_node_id = cls._runtime_node_id(unique_id)
if stream_enabled and skill.enabled:
raise ValueError(
"stream cannot be enabled together with skill_options in this version"
)
session_prompt = skill.session_discriminator(system_prompt)
key = cls._session_key(
unique_id, endpoint, protocol, model, session_prompt,
runtime_node_id, endpoint, protocol, model, session_prompt,
)
if clear_history:
cls._conversation_history.pop(key, None)
@@ -226,7 +322,7 @@ class OpenAITextAPI(io.ComfyNode):
request_history.append(user_message)
else:
history.append(user_message)
payload = {"model": model, "messages": request_history, "stream": False}
payload = {"model": model, "messages": request_history, "stream": stream_enabled}
else:
history = list(cls._conversation_history.get(key, [])) if persist_context and not skill.enabled else []
content = [{"type": "input_text", "text": user_prompt}] + _image_parts(images, protocol)
@@ -238,7 +334,7 @@ class OpenAITextAPI(io.ComfyNode):
instructions = system_prompt
else:
instructions = None
payload = {"model": model, "input": history, "stream": False}
payload = {"model": model, "input": history, "stream": stream_enabled}
if instructions:
payload["instructions"] = instructions
skill.validate_advanced_options(advanced_options)
@@ -259,6 +355,26 @@ class OpenAITextAPI(io.ComfyNode):
key,
persist_context=persist_context,
)
elif stream_enabled:
sink = _TextStreamSink(runtime_node_id)
sink.start()
try:
result = post_openai_stream(
endpoint,
headers,
payload,
timeout,
proxies,
protocol,
sink.delta,
final_response_parser=cls._parse_responses,
on_activity=sink.activity,
)
except Exception as exc:
sink.error(exc)
raise
else:
sink.end(result)
else:
data = post_openai_json(endpoint, headers, payload, timeout, proxies)
result = cls._parse_completions(data) if protocol == "openai-completions" else cls._parse_responses(data)
+29
View File
@@ -0,0 +1,29 @@
from comfy_api.latest import io, ui
class PreviewAPIResult(io.ComfyNode):
"""Render a string in the frontend Markdown preview widget."""
@classmethod
def define_schema(cls) -> io.Schema:
return io.Schema(
node_id="YCYY_Preview_API_Result",
display_name="Preview API Result",
category="YCYY/API/utils",
inputs=[
io.String.Input(
id="source",
force_input=True,
multiline=True,
tooltip="Markdown or plain text to preview",
),
],
outputs=[io.String.Output(id="text", display_name="text")],
is_output_node=True,
description="Render text as a Markdown preview with copy support.",
)
@classmethod
def execute(cls, source) -> io.NodeOutput:
text = source if isinstance(source, str) else ("" if source is None else str(source))
return io.NodeOutput(text, ui=ui.PreviewText(text))
+191 -29
View File
@@ -18,40 +18,46 @@ class FunctionToolsRejected(RuntimeError):
"""The provider explicitly rejected the requested function tool schema."""
def _raise_openai_http_error(response, payload):
"""Raise the same stable error classes for JSON and streaming requests."""
if 200 <= response.status_code < 300:
return
detail = response.text[:1000]
lower_detail = detail.lower()
choice_rejected = (
response.status_code in (400, 404, 415, 422)
and "tool_choice" in payload
and any(marker in lower_detail for marker in ("tool_choice", "tool choice"))
and any(
marker in lower_detail
for marker in (
"unsupported", "unknown", "unrecognized", "invalid",
"not support", "not allowed", "not permitted", "extra field",
)
)
)
if choice_rejected:
raise ToolChoiceRejected(
f"API request rejected tool_choice ({response.status_code}): {detail}"
)
tool_error = response.status_code in (400, 404, 415, 422) and any(
marker in lower_detail
for marker in ("tool", "function", "unsupported", "unknown field")
)
if payload.get("tools") and tool_error:
raise FunctionToolsRejected(
f"API request failed ({response.status_code}); this API/model may not support "
f"the requested Skill tools: {detail}"
)
raise RuntimeError(f"API request failed ({response.status_code}): {detail}")
def post_openai_json(endpoint, headers, payload, timeout, proxies):
"""POST an OpenAI-compatible JSON request with stable error classification."""
response = requests.post(
endpoint, headers=headers, json=payload, timeout=timeout, proxies=proxies
)
if response.status_code < 200 or response.status_code >= 300:
detail = response.text[:1000]
lower_detail = detail.lower()
choice_rejected = (
response.status_code in (400, 404, 415, 422)
and "tool_choice" in payload
and any(marker in lower_detail for marker in ("tool_choice", "tool choice"))
and any(
marker in lower_detail
for marker in (
"unsupported", "unknown", "unrecognized", "invalid",
"not support", "not allowed", "not permitted", "extra field",
)
)
)
if choice_rejected:
raise ToolChoiceRejected(
f"API request rejected tool_choice ({response.status_code}): {detail}"
)
tool_error = response.status_code in (400, 404, 415, 422) and any(
marker in lower_detail
for marker in ("tool", "function", "unsupported", "unknown field")
)
if payload.get("tools") and tool_error:
raise FunctionToolsRejected(
f"API request failed ({response.status_code}); this API/model may not support "
f"the requested Skill tools: {detail}"
)
raise RuntimeError(f"API request failed ({response.status_code}): {detail}")
_raise_openai_http_error(response, payload)
if not response.text.strip():
raise ValueError("API returned an empty response")
try:
@@ -60,6 +66,162 @@ def post_openai_json(endpoint, headers, payload, timeout, proxies):
raise ValueError("API returned invalid JSON") from exc
def _iter_sse_data(response):
"""Yield complete SSE data payloads, independent of HTTP chunk boundaries."""
data_lines = []
# The SSE specification mandates UTF-8. Requests may otherwise infer
# ISO-8859-1 for text/event-stream responses without a charset.
response.encoding = "utf-8"
# requests defaults to 512-byte buffering here. Small token events can
# otherwise sit in that buffer until generation is nearly complete,
# making a real SSE response appear non-streaming in the UI.
for line in response.iter_lines(chunk_size=1, decode_unicode=True):
if line is None:
continue
if isinstance(line, bytes):
line = line.decode(response.encoding or "utf-8", errors="replace")
line = line.rstrip("\r")
if line == "":
if data_lines:
yield "\n".join(data_lines)
data_lines.clear()
continue
if line.startswith(":"):
continue
if line.startswith("data:"):
data_lines.append(line[5:].lstrip())
if data_lines:
yield "\n".join(data_lines)
def _completion_stream_event(event):
"""Return the provider-neutral meaning of a Chat Completions event."""
choices = event.get("choices")
if not isinstance(choices, list) or not choices or not isinstance(choices[0], dict):
return None, None, False
choice = choices[0]
delta = choice.get("delta")
delta = delta if isinstance(delta, dict) else {}
content = delta.get("content")
reasoning = delta.get("reasoning_content")
return (
content if isinstance(content, str) and content else None,
"reasoning" if isinstance(reasoning, str) and reasoning else None,
choice.get("finish_reason") is not None,
)
def post_openai_stream(
endpoint,
headers,
payload,
timeout,
proxies,
protocol,
on_delta,
final_response_parser=None,
on_activity=None,
):
"""POST an OpenAI-compatible SSE request and return its complete text."""
if protocol not in ("openai-completions", "openai-responses"):
raise ValueError(f"Unsupported streaming protocol: {protocol}")
chunks = []
completed_response = None
reasoning_notified = False
completed = False
stream_headers = {
**headers,
"Accept": "text/event-stream",
"Cache-Control": "no-cache",
# Some OpenAI-compatible gateways gzip SSE responses and only flush
# compressed blocks occasionally. Identity encoding keeps token-sized
# events observable as soon as the provider sends them.
"Accept-Encoding": "identity",
}
with requests.post(
endpoint,
headers=stream_headers,
json=payload,
timeout=timeout,
proxies=proxies,
stream=True,
) as response:
_raise_openai_http_error(response, payload)
content_type = response.headers.get("Content-Type", "").lower()
if "text/event-stream" not in content_type:
raise ValueError("Streaming API returned a non-SSE response")
for raw_data in _iter_sse_data(response):
if raw_data.strip() == "[DONE]":
completed = True
break
try:
event = json.loads(raw_data)
except json.JSONDecodeError as exc:
raise ValueError("Streaming API returned invalid SSE JSON") from exc
if not isinstance(event, dict):
raise ValueError("Streaming API returned a non-object SSE event")
if event.get("error"):
raise ValueError(f"Streaming API failed: {event['error']}")
if protocol == "openai-completions":
text, activity, event_completed = _completion_stream_event(event)
completed = completed or event_completed
else:
event_type = event.get("type")
text = (
event.get("delta")
if event_type == "response.output_text.delta"
else None
)
activity = (
"reasoning"
if event_type == "response.reasoning_text.delta"
and isinstance(event.get("delta"), str)
and event.get("delta")
else None
)
if event_type == "response.completed":
completed_response = event.get("response")
completed = True
if event_type in ("error", "response.failed", "response.incomplete"):
detail = event.get("error") or event.get("response") or event
raise ValueError(f"Streaming API failed: {detail}")
if activity == "reasoning" and not reasoning_notified:
reasoning_notified = True
if on_activity is not None:
on_activity("reasoning")
if isinstance(text, str) and text:
chunks.append(text)
on_delta(text)
if not completed:
raise ValueError("Streaming API ended before a completion marker")
result = "".join(chunks)
if (
protocol == "openai-responses"
and isinstance(completed_response, dict)
and final_response_parser is not None
):
try:
final_result = final_response_parser(completed_response)
except (TypeError, ValueError):
if not result:
raise
else:
if isinstance(final_result, str) and final_result.strip():
result = final_result
if not isinstance(result, str) or not result.strip():
if reasoning_notified:
raise ValueError("Streaming API returned no formal text content")
raise ValueError("Streaming API returned no text content")
return result
def parse_json_options(options_json):
if not options_json or not str(options_json).strip():
return {}
File diff suppressed because one or more lines are too long
+680
View File
@@ -0,0 +1,680 @@
import { app } from "../../scripts/app.js";
import { api } from "../../scripts/api.js";
import "./marked.umd.js";
const NODE_CLASS = "YCYY_Preview_API_Result";
const OPENAI_NODE_CLASS = "YCYY_OpenAI_Text_API";
const STREAM_EVENT = "ycyy_openai_text_stream";
const LOCALE_SETTING = "Comfy.Locale";
const MIN_WIDTH = 360;
const MIN_HEIGHT = 260;
const MAX_TYPING_UNITS_PER_FRAME = 18;
const RENDER_INTERVAL_MS = 40;
const MESSAGES = {
en: {
empty: "Connect an API result and run the workflow",
ready: "Ready",
waiting: "Waiting for response…",
reasoning: "Thinking…",
generating: "Generating…",
displaying: "Displaying result…",
complete: "Complete",
copy: "Copy",
copied: "Copied",
copyFailed: "Copy failed",
success: "Success",
copiedDetail: "Text copied to clipboard",
error: "Error",
},
zh: {
empty: "连接 API 结果并运行工作流",
ready: "就绪",
waiting: "等待响应…",
reasoning: "正在思考…",
generating: "正在生成…",
displaying: "正在显示结果…",
complete: "完成",
copy: "复制",
copied: "已复制",
copyFailed: "复制失败",
success: "成功",
copiedDetail: "文本已复制到剪贴板",
error: "错误",
},
};
const ALLOWED_TAGS = new Set([
"P", "BR", "H1", "H2", "H3", "H4", "H5", "H6",
"UL", "OL", "LI", "BLOCKQUOTE", "PRE", "CODE", "HR",
"STRONG", "EM", "DEL", "A", "TABLE", "THEAD", "TBODY",
"TR", "TH", "TD",
]);
const DROP_CONTENT_TAGS = new Set([
"SCRIPT", "STYLE", "IFRAME", "OBJECT", "EMBED", "FORM",
"SVG", "MATH", "META", "LINK", "BASE", "IMG", "VIDEO", "AUDIO",
]);
let currentLocale = "en";
const graphemeSegmenter = globalThis.Intl?.Segmenter
? new Intl.Segmenter(undefined, { granularity: "grapheme" })
: null;
const reducedMotion = globalThis.matchMedia?.("(prefers-reduced-motion: reduce)");
function normalizeLocale(locale) {
return String(locale || "en").toLowerCase().startsWith("zh") ? "zh" : "en";
}
function message(key) {
return MESSAGES[normalizeLocale(currentLocale)]?.[key] ?? MESSAGES.en[key] ?? key;
}
function nodeClass(node) {
return node?.constructor?.comfyClass || node?.comfyClass || node?.type;
}
function getLink(graph, linkId) {
return graph?.links?.get?.(linkId) ?? graph?.links?.[linkId] ?? null;
}
function safeHref(value) {
try {
const url = new URL(value, window.location.href);
return ["http:", "https:", "mailto:"].includes(url.protocol) ? value : null;
} catch {
return null;
}
}
function sanitizeMarkedHtml(html) {
const template = document.createElement("template");
template.innerHTML = html;
for (const element of [...template.content.querySelectorAll("*")]) {
if (DROP_CONTENT_TAGS.has(element.tagName)) {
element.remove();
continue;
}
if (!ALLOWED_TAGS.has(element.tagName)) {
element.replaceWith(...element.childNodes);
continue;
}
const href = element.tagName === "A"
? safeHref(element.getAttribute("href") || "")
: null;
const title = element.getAttribute("title");
for (const attribute of [...element.attributes]) {
element.removeAttribute(attribute.name);
}
if (href) {
element.setAttribute("href", href);
element.setAttribute("target", "_blank");
element.setAttribute("rel", "noopener noreferrer");
}
if (title) element.setAttribute("title", title);
}
return template.innerHTML;
}
function markedApi() {
return globalThis.marked;
}
function isNearBottom(element) {
return element.scrollHeight - element.scrollTop - element.clientHeight < 48;
}
function splitTextUnits(value) {
const text = String(value ?? "");
if (!text) return [];
if (graphemeSegmenter) {
return [...graphemeSegmenter.segment(text)].map(part => part.segment);
}
return Array.from(text);
}
function typingRate(state) {
const queueLength = state.pendingUnits.length;
if (state.sourceEnded) {
if (queueLength > 500) return 720;
if (queueLength > 200) return 480;
return 260;
}
if (queueLength > 500) return 260;
if (queueLength > 200) return 180;
if (queueLength > 80) return 110;
return 48;
}
function renderNow(state) {
const followTail = isNearBottom(state.content);
const raw = state.displayedText || "";
state.host.dataset.hasContent = raw ? "true" : "false";
if (!raw) {
state.content.replaceChildren();
if (!["waiting", "reasoning", "error"].includes(state.statusKey)) {
const placeholder = document.createElement("div");
placeholder.className = "empty";
placeholder.textContent = message("empty");
state.content.append(placeholder);
}
} else {
const parser = markedApi();
if (!parser?.parse) {
state.content.textContent = raw;
setStatus(state, "error", "marked.umd.js unavailable");
} else {
const html = parser.parse(raw, { gfm: true, breaks: true });
state.content.innerHTML = sanitizeMarkedHtml(html);
}
}
if (followTail) state.content.scrollTop = state.content.scrollHeight;
}
function scheduleRender(state, immediate = false) {
if (immediate) {
clearTimeout(state.renderTimer);
state.renderTimer = null;
state.lastRenderTime = performance.now();
renderNow(state);
return;
}
if (state.renderTimer !== null) return;
const elapsed = performance.now() - state.lastRenderTime;
const delay = Math.max(0, RENDER_INTERVAL_MS - elapsed);
state.renderTimer = setTimeout(() => {
state.renderTimer = null;
state.lastRenderTime = performance.now();
renderNow(state);
}, delay);
}
function setStatus(state, key, detail = "") {
state.statusKey = key;
state.statusDetail = detail;
const visible = ["waiting", "reasoning", "generating", "displaying", "error"].includes(key);
state.host.dataset.state = key;
state.status.textContent = key === "error" && detail
? `${message("error")}: ${detail}`
: message(key);
state.status.hidden = !visible;
}
function updateCopyAvailability(state) {
state.copyButton.disabled = !(state.finalText || state.receivedText || state.displayedText);
}
function setCopyState(state, copyState) {
state.copyButton.dataset.copyState = copyState;
const label = copyState === "copied"
? message("copied")
: (copyState === "failed" ? message("copyFailed") : message("copy"));
state.copyButton.title = label;
state.copyButton.setAttribute("aria-label", label);
}
function showCopyToast(success) {
const toast = app.extensionManager?.toast;
if (!toast?.add) return;
toast.add(success
? {
severity: "success",
summary: message("success"),
detail: message("copiedDetail"),
life: 3000,
}
: {
severity: "error",
summary: message("error"),
detail: message("copyFailed"),
});
}
async function copyRawText(state) {
const text = state.sourceEnded ? state.finalText : (state.receivedText || state.displayedText);
if (!text) return;
try {
if (navigator.clipboard?.writeText && window.isSecureContext) {
await navigator.clipboard.writeText(text);
} else {
const textarea = document.createElement("textarea");
textarea.value = text;
textarea.style.position = "fixed";
textarea.style.opacity = "0";
document.body.append(textarea);
try {
textarea.select();
if (!document.execCommand("copy")) throw new Error("copy command failed");
} finally {
textarea.remove();
}
}
setCopyState(state, "copied");
showCopyToast(true);
} catch (error) {
console.error("[YCYY] Failed to copy preview text:", error);
setCopyState(state, "failed");
showCopyToast(false);
} finally {
clearTimeout(state.copyTimer);
state.copyTimer = setTimeout(() => {
setCopyState(state, "idle");
}, 1200);
}
}
function stopTyping(state) {
if (state.typingFrame !== null) cancelAnimationFrame(state.typingFrame);
state.typingFrame = null;
state.lastTypingTime = 0;
state.typingBudget = 0;
}
function finishTypingIfReady(state) {
if (!state.sourceEnded || state.pendingUnits.length > 0) return false;
// Normally the queue already produced finalText. This assignment also
// reconciles unusual providers whose final response differs from deltas.
if (state.displayedText !== state.finalText) {
state.displayedText = state.finalText;
}
state.receivedText = state.finalText;
state.streaming = false;
state.activeRunId = null;
state.promptId = null;
if (state.terminalError) {
setStatus(state, "error", state.terminalError);
} else {
setStatus(state, "complete");
}
updateCopyAvailability(state);
scheduleRender(state, true);
return true;
}
function typingTick(state, timestamp) {
state.typingFrame = null;
if (!state.streaming) return;
if (!state.lastTypingTime) state.lastTypingTime = timestamp;
const elapsed = Math.min(250, Math.max(0, timestamp - state.lastTypingTime));
state.lastTypingTime = timestamp;
state.typingBudget += elapsed * typingRate(state) / 1000;
const available = Math.floor(state.typingBudget);
const count = Math.min(
available,
MAX_TYPING_UNITS_PER_FRAME,
state.pendingUnits.length,
);
if (count > 0) {
state.displayedText += state.pendingUnits.splice(0, count).join("");
state.typingBudget -= count;
scheduleRender(state);
}
if (state.pendingUnits.length > 0) {
state.typingFrame = requestAnimationFrame(time => typingTick(state, time));
} else {
state.lastTypingTime = 0;
state.typingBudget = 0;
finishTypingIfReady(state);
}
}
function startTyping(state) {
if (state.typingFrame !== null) return;
if (!state.pendingUnits.length) {
finishTypingIfReady(state);
return;
}
state.typingFrame = requestAnimationFrame(time => typingTick(state, time));
}
function enqueueDelta(state, value) {
const delta = String(value ?? "");
if (!delta) return;
state.receivedText += delta;
updateCopyAvailability(state);
if (reducedMotion?.matches) {
state.displayedText = state.receivedText;
state.pendingUnits = [];
scheduleRender(state);
return;
}
state.pendingUnits.push(...splitTextUnits(delta));
startTyping(state);
}
function enqueueFinalText(state, value) {
const finalText = String(value ?? state.receivedText);
state.sourceEnded = true;
state.finalText = finalText;
state.receivedText = finalText;
updateCopyAvailability(state);
if (reducedMotion?.matches) {
state.displayedText = finalText;
state.pendingUnits = [];
finishTypingIfReady(state);
return;
}
// Rebuild the unplayed tail from the authoritative final response. This
// also prevents duplicate characters when both `end` and `onExecuted`
// deliver the same full text before the animation has caught up.
if (finalText.startsWith(state.displayedText)) {
state.pendingUnits = splitTextUnits(finalText.slice(state.displayedText.length));
} else {
const shown = splitTextUnits(state.displayedText);
const final = splitTextUnits(finalText);
let common = 0;
while (common < shown.length && common < final.length && shown[common] === final[common]) {
common += 1;
}
state.displayedText = shown.slice(0, common).join("");
state.pendingUnits = final.slice(common);
}
startTyping(state);
}
function createPreview(node) {
if (node._ycyyApiResultPreview) return node._ycyyApiResultPreview;
const root = document.createElement("div");
root.style.width = "100%";
root.style.height = "100%";
root.style.minHeight = "210px";
root.style.boxSizing = "border-box";
const shadow = root.attachShadow({ mode: "open" });
shadow.innerHTML = `
<style>
:host { display: block; width: 100%; height: 100%; }
.preview {
position: relative; height: 100%; min-height: 200px; overflow: hidden;
box-sizing: border-box; border: 1px solid var(--border-color, #484848);
border-radius: 8px; background: var(--comfy-input-bg, #181818);
color: var(--input-text, var(--fg-color, #ddd));
font: 13px/1.55 system-ui, -apple-system, "Segoe UI", sans-serif;
}
.content {
width: 100%; height: 100%; overflow: auto; box-sizing: border-box;
padding: 12px 14px 22px; overflow-wrap: anywhere;
scrollbar-color: color-mix(in srgb, currentColor 35%, transparent) transparent;
}
.empty { height: 100%; display: grid; place-items: center; text-align: center; opacity: .42; }
.copy {
position: absolute; z-index: 2; top: 6px; right: 8px; width: 28px; height: 28px;
display: grid; place-items: center; padding: 0; border: 0; border-radius: 6px;
background: color-mix(in srgb, var(--comfy-input-bg, #181818) 88%, currentColor 12%);
color: inherit; cursor: pointer; opacity: 0; transition: opacity .12s, background .12s;
}
.preview:hover .copy, .copy:focus-visible { opacity: .82; }
.copy:hover { opacity: 1; background: color-mix(in srgb, var(--comfy-input-bg, #181818) 76%, currentColor 24%); }
.copy:focus-visible { outline: 1px solid #60a5fa; outline-offset: 1px; }
.copy svg { width: 16px; height: 16px; fill: none; stroke: currentColor; stroke-width: 1.8; stroke-linecap: round; stroke-linejoin: round; }
.copy .check-icon { display: none; }
.copy[data-copy-state="copied"] { color: #69d69b; opacity: 1; }
.copy[data-copy-state="copied"] .copy-icon { display: none; }
.copy[data-copy-state="copied"] .check-icon { display: block; }
.copy[data-copy-state="failed"] { color: #f87171; opacity: 1; }
.copy:disabled { display: none; }
.status {
position: absolute; z-index: 2; left: 10px; bottom: 7px; max-width: calc(100% - 20px);
overflow: hidden; text-overflow: ellipsis; white-space: nowrap; pointer-events: none;
padding: 2px 7px; border-radius: 999px; background: color-mix(in srgb, var(--comfy-input-bg, #181818) 84%, currentColor 16%);
color: inherit; font-size: 11px; opacity: .72;
}
.preview[data-state="waiting"] .status,
.preview[data-state="reasoning"] .status,
.preview[data-state="error"][data-has-content="false"] .status {
left: 50%; bottom: 50%; max-width: calc(100% - 32px);
transform: translate(-50%, 50%); padding: 5px 10px;
}
.preview[data-state="error"] .status { color: #fca5a5; opacity: .95; }
h1, h2, h3, h4, h5, h6 { margin: 1.1em 0 .55em; line-height: 1.25; }
h1 { font-size: 1.55em; } h2 { font-size: 1.35em; } h3 { font-size: 1.18em; }
p, ul, ol, blockquote, pre, table { margin: .65em 0; }
ul, ol { padding-left: 1.7em; } li + li { margin-top: .2em; }
blockquote { margin-left: 0; padding: .1em 1em; border-left: 3px solid #7c8cf8; opacity: .82; }
code { border-radius: 4px; padding: .12em .35em; background: color-mix(in srgb, currentColor 10%, transparent); font-family: ui-monospace, SFMono-Regular, Consolas, monospace; }
pre { overflow: auto; padding: 12px; border-radius: 7px; background: color-mix(in srgb, currentColor 8%, transparent); }
pre code { padding: 0; background: transparent; white-space: pre; }
table { display: block; max-width: 100%; overflow-x: auto; border-collapse: collapse; }
th, td { padding: 6px 9px; border: 1px solid color-mix(in srgb, currentColor 18%, transparent); }
th { background: color-mix(in srgb, currentColor 7%, transparent); }
a { color: #7c8cf8; } hr { border: 0; border-top: 1px solid color-mix(in srgb, currentColor 18%, transparent); }
.preview[data-state="generating"] .content::after,
.preview[data-state="displaying"] .content::after { content: ""; display: inline-block; width: 6px; height: 1em; margin-left: 3px; vertical-align: -.12em; background: currentColor; animation: blink .85s steps(1) infinite; }
@keyframes blink { 50% { opacity: 0; } }
@media (prefers-reduced-motion: reduce) {
*, *::before, *::after { animation: none !important; transition: none !important; }
.preview[data-state="generating"] .content::after,
.preview[data-state="displaying"] .content::after { display: none; }
}
</style>
<div class="preview" data-state="idle" data-has-content="false">
<button class="copy" type="button" data-copy-state="idle">
<svg class="copy-icon" viewBox="0 0 24 24" aria-hidden="true"><rect x="9" y="9" width="11" height="11" rx="2"></rect><path d="M15 9V6a2 2 0 0 0-2-2H6a2 2 0 0 0-2 2v7a2 2 0 0 0 2 2h3"></path></svg>
<svg class="check-icon" viewBox="0 0 24 24" aria-hidden="true"><path d="m5 12 4 4L19 6"></path></svg>
</button>
<div class="content"></div>
<span class="status" hidden aria-live="polite"></span>
</div>`;
const state = {
displayedText: "",
receivedText: "",
finalText: "",
pendingUnits: [],
sourceEnded: false,
terminalError: "",
activeRunId: null,
promptId: null,
lastSeq: -1,
streaming: false,
typingFrame: null,
lastTypingTime: 0,
typingBudget: 0,
renderTimer: null,
lastRenderTime: 0,
copyTimer: null,
statusKey: "ready",
statusDetail: "",
root,
host: shadow.querySelector(".preview"),
content: shadow.querySelector(".content"),
status: shadow.querySelector(".status"),
copyButton: shadow.querySelector("button"),
};
setStatus(state, "ready");
setCopyState(state, "idle");
updateCopyAvailability(state);
state.copyButton.addEventListener("click", () => copyRawText(state));
root.addEventListener("pointerdown", event => event.stopPropagation());
root.addEventListener("wheel", event => event.stopPropagation(), { passive: true });
node.addDOMWidget("preview_api_result", "div", root, {
serialize: false,
hideOnZoom: false,
getMinHeight: () => 210,
});
node._ycyyApiResultPreview = state;
node.size ||= [MIN_WIDTH, MIN_HEIGHT];
node.size[0] = Math.max(node.size[0] || 0, MIN_WIDTH);
node.size[1] = Math.max(node.size[1] || 0, MIN_HEIGHT);
scheduleRender(state);
return state;
}
function disposePreview(node) {
const state = node?._ycyyApiResultPreview;
if (!state) return;
stopTyping(state);
clearTimeout(state.renderTimer);
clearTimeout(state.copyTimer);
node._ycyyApiResultPreview = null;
}
function setFinalText(node, value) {
const state = createPreview(node);
const text = String(value ?? "");
if (state.streaming || state.typingFrame !== null || state.pendingUnits.length > 0) {
enqueueFinalText(state, text);
return;
}
const unchanged = state.displayedText === text && !state.streaming;
stopTyping(state);
state.displayedText = text;
state.receivedText = text;
state.finalText = text;
state.pendingUnits = [];
state.sourceEnded = true;
state.terminalError = "";
state.streaming = false;
state.activeRunId = null;
state.promptId = null;
setStatus(state, "complete");
updateCopyAvailability(state);
if (!unchanged) scheduleRender(state, true);
}
function connectedPreviewNodes(sourceNodeId) {
const graph = app.graph;
const source = graph?.getNodeById?.(sourceNodeId)
?? graph?.getNodeById?.(Number(sourceNodeId));
if (!source || nodeClass(source) !== OPENAI_NODE_CLASS) return [];
const resultOutput = source.outputs?.[0];
const targets = [];
for (const linkId of resultOutput?.links || []) {
const link = getLink(graph, linkId);
if (!link || link.origin_slot !== 0 || link.target_slot !== 0) continue;
const target = graph.getNodeById(link.target_id);
if (target && nodeClass(target) === NODE_CLASS) targets.push(target);
}
return targets;
}
function receiveStreamEvent(event) {
const data = event.detail || {};
for (const node of connectedPreviewNodes(data.node_id)) {
const state = createPreview(node);
const sequence = Number(data.seq);
if (data.phase === "start") {
stopTyping(state);
state.activeRunId = data.run_id;
state.promptId = data.prompt_id ?? null;
state.lastSeq = Number.isFinite(sequence) ? sequence : -1;
state.displayedText = "";
state.receivedText = "";
state.finalText = "";
state.pendingUnits = [];
state.sourceEnded = false;
state.terminalError = "";
state.streaming = true;
setStatus(state, "waiting");
updateCopyAvailability(state);
scheduleRender(state, true);
continue;
}
if (state.activeRunId !== data.run_id) continue;
if (Number.isFinite(sequence) && sequence <= state.lastSeq) continue;
if (Number.isFinite(sequence)) state.lastSeq = sequence;
if (data.phase === "activity") {
if (data.activity === "reasoning" && !state.receivedText) {
setStatus(state, "reasoning");
scheduleRender(state, true);
}
} else if (data.phase === "delta") {
if (!state.receivedText) setStatus(state, "generating");
enqueueDelta(state, data.delta);
} else if (data.phase === "end") {
if (state.pendingUnits.length || state.displayedText !== String(data.text ?? state.receivedText)) {
setStatus(state, "displaying");
}
enqueueFinalText(state, data.text ?? state.receivedText);
} else if (data.phase === "error") {
state.sourceEnded = true;
state.finalText = state.receivedText;
state.terminalError = data.message || "unknown";
setStatus(state, "error", state.terminalError);
updateCopyAvailability(state);
scheduleRender(state, true);
startTyping(state);
}
}
}
function refreshLocale() {
currentLocale = app.ui?.settings?.getSettingValue?.(LOCALE_SETTING) || "en";
for (const node of app.graph?._nodes || []) {
const state = node?._ycyyApiResultPreview;
if (!state) continue;
setCopyState(state, state.copyButton.dataset.copyState || "idle");
setStatus(state, state.statusKey, state.statusDetail);
scheduleRender(state);
}
}
app.registerExtension({
name: "YCYY.Preview.API.Result",
async setup() {
refreshLocale();
app.ui?.settings?.addEventListener?.(`${LOCALE_SETTING}.change`, refreshLocale);
api.addEventListener(STREAM_EVENT, receiveStreamEvent);
},
async beforeRegisterNodeDef(nodeType, nodeData) {
if (nodeData.name !== NODE_CLASS && nodeType.comfyClass !== NODE_CLASS) return;
const originalCreated = nodeType.prototype.onNodeCreated;
nodeType.prototype.onNodeCreated = function () {
const result = originalCreated?.apply(this, arguments);
createPreview(this);
return result;
};
const originalConfigured = nodeType.prototype.onConfigure;
nodeType.prototype.onConfigure = function () {
const result = originalConfigured?.apply(this, arguments);
setTimeout(() => createPreview(this), 0);
return result;
};
const originalExecuted = nodeType.prototype.onExecuted;
nodeType.prototype.onExecuted = function (executionMessage) {
originalExecuted?.apply(this, arguments);
const value = Array.isArray(executionMessage?.text)
? executionMessage.text[0]
: (executionMessage?.text ?? "");
setFinalText(this, value);
};
const originalRemoved = nodeType.prototype.onRemoved;
nodeType.prototype.onRemoved = function () {
disposePreview(this);
return originalRemoved?.apply(this, arguments);
};
},
async nodeCreated(node) {
if (nodeClass(node) === NODE_CLASS) createPreview(node);
},
async onNodeOutputsUpdated(nodeOutputs) {
for (const [nodeId, output] of Object.entries(nodeOutputs || {})) {
const node = app.graph?.getNodeById?.(nodeId)
?? app.graph?.getNodeById?.(Number(nodeId));
if (node && nodeClass(node) === NODE_CLASS) {
const value = Array.isArray(output?.text)
? output.text[0]
: (output?.text ?? "");
setFinalText(node, value);
}
}
},
});