Update DistributedValue to track enabled workers
This commit is contained in:
+65
-5
@@ -5,13 +5,22 @@ const NODE_CLASS = "DistributedValue";
|
||||
const CONVERTED_WIDGET = "converted-widget";
|
||||
const DYNAMIC_DEFAULT_WIDGET = "_dv_default";
|
||||
const DYNAMIC_WORKER_WIDGET_PREFIX = "_dv_worker_";
|
||||
const WORKERS_CHANGED_EVENT = "distributed:workers-changed";
|
||||
|
||||
const trackedNodes = new Set();
|
||||
let workersChangedListenerAttached = false;
|
||||
|
||||
function filterEnabledWorkers(workers) {
|
||||
if (!Array.isArray(workers)) return [];
|
||||
return workers.filter((worker) => Boolean(worker?.enabled));
|
||||
}
|
||||
|
||||
async function fetchWorkers() {
|
||||
try {
|
||||
const resp = await fetch(ENDPOINTS.CONFIG);
|
||||
if (!resp.ok) return [];
|
||||
const config = await resp.json();
|
||||
return config.workers || [];
|
||||
return filterEnabledWorkers(config.workers);
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
@@ -201,13 +210,21 @@ function serializeWorkerStoreFromWidgets(node, inputType, comboOptions) {
|
||||
if (inputType === "COMBO" && Array.isArray(comboOptions)) {
|
||||
nextStore._options = comboOptions;
|
||||
}
|
||||
const valuesByWorkerId = {};
|
||||
|
||||
for (const widget of getDynamicWorkerWidgets(node)) {
|
||||
const key = widget.name.slice(DYNAMIC_WORKER_WIDGET_PREFIX.length);
|
||||
if (widget.value !== "" && widget.value !== null && widget.value !== undefined) {
|
||||
nextStore[key] = String(widget.value);
|
||||
const value = String(widget.value);
|
||||
nextStore[key] = value;
|
||||
if (widget._dvWorkerId) {
|
||||
valuesByWorkerId[widget._dvWorkerId] = value;
|
||||
}
|
||||
}
|
||||
}
|
||||
if (Object.keys(valuesByWorkerId).length) {
|
||||
nextStore._by_worker_id = valuesByWorkerId;
|
||||
}
|
||||
|
||||
writeWorkerStore(node, nextStore);
|
||||
}
|
||||
@@ -277,8 +294,11 @@ function createDynamicDefaultWidget(node, inputType, comboOptions) {
|
||||
widget.label = "default_value";
|
||||
}
|
||||
|
||||
function getWorkerInitialValue(store, key, inputType, comboOptions) {
|
||||
const saved = store[key];
|
||||
function getWorkerInitialValue(store, key, workerId, inputType, comboOptions) {
|
||||
const byWorkerId = store?._by_worker_id;
|
||||
const saved = (byWorkerId && workerId && byWorkerId[workerId] != null)
|
||||
? byWorkerId[workerId]
|
||||
: store[key];
|
||||
if (saved == null) {
|
||||
if (inputType === "INT" || inputType === "FLOAT") return 0;
|
||||
if (inputType === "COMBO" && Array.isArray(comboOptions) && comboOptions.length) {
|
||||
@@ -305,7 +325,7 @@ function createWorkerWidgets(node, workers, inputType, comboOptions) {
|
||||
const worker = workers[i];
|
||||
const label = worker.name || worker.id || `Worker ${key}`;
|
||||
const widgetName = `${DYNAMIC_WORKER_WIDGET_PREFIX}${key}`;
|
||||
const initial = getWorkerInitialValue(store, key, inputType, comboOptions);
|
||||
const initial = getWorkerInitialValue(store, key, worker.id, inputType, comboOptions);
|
||||
let widget;
|
||||
|
||||
if (inputType === "COMBO" && Array.isArray(comboOptions) && comboOptions.length) {
|
||||
@@ -355,6 +375,7 @@ function createWorkerWidgets(node, workers, inputType, comboOptions) {
|
||||
}
|
||||
|
||||
widget.label = label;
|
||||
widget._dvWorkerId = worker.id;
|
||||
}
|
||||
|
||||
serializeWorkerStoreFromWidgets(node, inputType, comboOptions);
|
||||
@@ -382,12 +403,43 @@ function rebuildWidgets(node) {
|
||||
if (node.setDirtyCanvas) node.setDirtyCanvas(true, true);
|
||||
}
|
||||
|
||||
function refreshNodeWorkers(node, workers) {
|
||||
if (!node || !node.graph) return;
|
||||
node._dvWorkers = workers;
|
||||
rebuildWidgets(node);
|
||||
}
|
||||
|
||||
async function refreshTrackedNodes(workers = null) {
|
||||
const nextWorkers = workers || (await fetchWorkers());
|
||||
for (const node of trackedNodes) {
|
||||
refreshNodeWorkers(node, nextWorkers);
|
||||
}
|
||||
}
|
||||
|
||||
function attachWorkersChangedListener() {
|
||||
if (workersChangedListenerAttached) return;
|
||||
if (typeof window === "undefined" || typeof window.addEventListener !== "function") return;
|
||||
|
||||
window.addEventListener(WORKERS_CHANGED_EVENT, (event) => {
|
||||
const changedWorkers = filterEnabledWorkers(event?.detail?.workers);
|
||||
if (changedWorkers.length > 0 || Array.isArray(event?.detail?.workers)) {
|
||||
void refreshTrackedNodes(changedWorkers);
|
||||
return;
|
||||
}
|
||||
void refreshTrackedNodes();
|
||||
});
|
||||
|
||||
workersChangedListenerAttached = true;
|
||||
}
|
||||
|
||||
app.registerExtension({
|
||||
name: "Distributed.DistributedValue",
|
||||
async nodeCreated(node) {
|
||||
if (node.comfyClass !== NODE_CLASS) return;
|
||||
|
||||
try {
|
||||
attachWorkersChangedListener();
|
||||
trackedNodes.add(node);
|
||||
node._dvWorkers = await fetchWorkers();
|
||||
rebuildWidgets(node);
|
||||
|
||||
@@ -407,6 +459,14 @@ app.registerExtension({
|
||||
setTimeout(() => rebuildWidgets(this), 20);
|
||||
return result;
|
||||
};
|
||||
|
||||
const originalOnRemoved = node.onRemoved;
|
||||
node.onRemoved = function () {
|
||||
trackedNodes.delete(this);
|
||||
if (originalOnRemoved) {
|
||||
return originalOnRemoved.call(this);
|
||||
}
|
||||
};
|
||||
} catch (error) {
|
||||
console.error("Error in DistributedValue extension:", error);
|
||||
}
|
||||
|
||||
+13
@@ -13,6 +13,8 @@ import { checkAllWorkerStatuses, checkWorkerStatus, loadManagedWorkers } from '.
|
||||
import { detectMasterIP } from './masterDetection.js';
|
||||
import { parseHostInput, getMasterUrl as buildMasterUrl } from './urlUtils.js';
|
||||
|
||||
const WORKERS_CHANGED_EVENT = "distributed:workers-changed";
|
||||
|
||||
class DistributedExtension {
|
||||
constructor() {
|
||||
this.config = null;
|
||||
@@ -150,12 +152,22 @@ class DistributedExtension {
|
||||
this.state.updateWorker(w.id, { enabled: w.enabled });
|
||||
});
|
||||
}
|
||||
this._emitWorkersChanged();
|
||||
} catch (error) {
|
||||
this.log("Failed to load config: " + error.message, "error");
|
||||
this.config = { workers: [], settings: { has_auto_populated_workers: false } };
|
||||
}
|
||||
}
|
||||
|
||||
_emitWorkersChanged() {
|
||||
if (typeof window === "undefined" || typeof window.dispatchEvent !== "function") {
|
||||
return;
|
||||
}
|
||||
window.dispatchEvent(new CustomEvent(WORKERS_CHANGED_EVENT, {
|
||||
detail: { workers: this.config?.workers || [] },
|
||||
}));
|
||||
}
|
||||
|
||||
_applyMasterHost(host) {
|
||||
if (!host || !this.config) return;
|
||||
if (!this.config.master) this.config.master = {};
|
||||
@@ -187,6 +199,7 @@ class DistributedExtension {
|
||||
if (worker) {
|
||||
worker.enabled = enabled;
|
||||
this.state.updateWorker(workerId, { enabled });
|
||||
this._emitWorkersChanged();
|
||||
|
||||
// Immediately update status dot based on enabled state
|
||||
const statusDot = document.getElementById(`status-${workerId}`);
|
||||
|
||||
@@ -3,6 +3,17 @@ import { generateUUID } from './constants.js';
|
||||
import { parseHostInput } from './urlUtils.js';
|
||||
import { toggleWorkerExpanded } from './workerLifecycle.js';
|
||||
|
||||
const WORKERS_CHANGED_EVENT = "distributed:workers-changed";
|
||||
|
||||
function emitWorkersChanged(extension) {
|
||||
if (typeof window === "undefined" || typeof window.dispatchEvent !== "function") {
|
||||
return;
|
||||
}
|
||||
window.dispatchEvent(new CustomEvent(WORKERS_CHANGED_EVENT, {
|
||||
detail: { workers: extension.config?.workers || [] },
|
||||
}));
|
||||
}
|
||||
|
||||
export function isRemoteWorker(extension, worker) {
|
||||
const workerType = String(worker?.type || "").toLowerCase();
|
||||
|
||||
@@ -158,6 +169,7 @@ export async function saveWorkerSettings(extension, workerId) {
|
||||
|
||||
// Sync to state
|
||||
extension.state.updateWorker(workerId, { enabled: nextEnabled });
|
||||
emitWorkersChanged(extension);
|
||||
|
||||
extension.app.extensionManager.toast.add({
|
||||
severity: "success",
|
||||
@@ -225,6 +237,7 @@ export async function deleteWorker(extension, workerId) {
|
||||
if (index !== -1) {
|
||||
extension.config.workers.splice(index, 1);
|
||||
}
|
||||
emitWorkersChanged(extension);
|
||||
|
||||
extension.app.extensionManager.toast.add({
|
||||
severity: "success",
|
||||
@@ -325,6 +338,7 @@ export async function addNewWorker(extension) {
|
||||
|
||||
// Sync to state
|
||||
extension.state.updateWorker(newId, { enabled: newWorker.enabled });
|
||||
emitWorkersChanged(extension);
|
||||
|
||||
extension.app.extensionManager.toast.add({
|
||||
severity: fallbackToRemote ? "warn" : "success",
|
||||
|
||||
Reference in New Issue
Block a user