882 lines
34 KiB
TypeScript
882 lines
34 KiB
TypeScript
// SPDX-FileCopyrightText: 2025-2026 Lucio Lelii <lucio.lelii@isti.cnr.it> - ISTI-CNR
|
|
// SPDX-License-Identifier: AGPL-3.0-or-later
|
|
// Attribution term under AGPL-3.0 section 7(b): see LICENSE-ADDENDUM.
|
|
|
|
import { Injector } from "@angular/core";
|
|
import { NodeEditor, ClassicPreset } from "rete";
|
|
import { AreaPlugin, AreaExtensions } from "rete-area-plugin";
|
|
import {
|
|
ConnectionPlugin,
|
|
Presets as ConnectionPresets
|
|
} from "rete-connection-plugin";
|
|
import { AngularPlugin, Presets, AngularArea2D } from "rete-angular-plugin/21";
|
|
import { HFNode, HFSchemes } from "@models/nodes";
|
|
import {
|
|
areFlowValueKindsCompatible,
|
|
DEFAULT_NODE_CAPABILITIES,
|
|
FlowBlock,
|
|
FlowData,
|
|
FlowGlobalInput,
|
|
FlowLane,
|
|
FlowLoopEdgeSettings,
|
|
FLOW_DEPENDANT_PORT_KEY,
|
|
FLOW_DEPENDENCY_PORT_KEY,
|
|
FLOW_DEPENDENCY_SOCKET_TYPE,
|
|
FlowNode,
|
|
NodeTypeCapabilities,
|
|
normalizeFlowPortValueKinds
|
|
} from "@models/flow";
|
|
import { BlocksService } from "@services/blocks/blocks";
|
|
import { ContainersService } from "@services/containers/containers";
|
|
import { EditorStateHolder } from "@stores/flow-editor";
|
|
import { ContainerNodeComponent } from "@shared/nodes/container-node/container-node";
|
|
import { GenericNodeComponent } from "@shared/nodes/generic-node/generic-node";
|
|
import { TaskStepNodeComponent } from "@shared/nodes/task-step-node/task-step-node";
|
|
import { CustomSocket } from "@shared/custom-socket/custom-socket";
|
|
import { CustomConnectionComponent } from "@shared/custom-connection/custom-connection";
|
|
import { deleteSchemaValueByPath, setSchemaValueByPath } from "@shared/nodes/schema-driven-fields";
|
|
import { firstValueFrom } from "rxjs";
|
|
import { findLoopBackEdgeIds } from "./flow-loops";
|
|
|
|
type AreaExtra = AngularArea2D<HFSchemes>;
|
|
const editorSockets = new WeakMap<NodeEditor<HFSchemes>, Map<string, ClassicPreset.Socket>>();
|
|
const editorRuntime = new WeakMap<NodeEditor<HFSchemes>, ReteRuntimeContext>();
|
|
const areaProgrammaticTranslations = new WeakMap<AreaPlugin<HFSchemes, AreaExtra>, Set<string>>();
|
|
|
|
export const RETE_ZOOM_RANGE = {
|
|
min: 0.1,
|
|
max: 2.4
|
|
} as const;
|
|
|
|
export type ReteEditorInstance = {
|
|
editor: NodeEditor<HFSchemes>;
|
|
area: AreaPlugin<HFSchemes, AreaExtra>;
|
|
};
|
|
|
|
type GraphConnectionKind = "data" | "dependency";
|
|
|
|
export type ReteRuntimeContext = {
|
|
blocksService: BlocksService;
|
|
containersService: ContainersService;
|
|
flowState: EditorStateHolder;
|
|
readonly: boolean;
|
|
globalInputs: FlowGlobalInput[];
|
|
lanes: FlowLane[];
|
|
/**
|
|
* Set while connections are being put back programmatically - loading a flow, replacing a node -
|
|
* when what is added is already a decided graph and nothing it adds should displace anything.
|
|
*/
|
|
restoringConnections?: number;
|
|
};
|
|
|
|
/**
|
|
* What the editor keeps on a Rete connection beyond Rete's own fields. `loop` travels to and from
|
|
* the flow; `__loopBack` and `__readonly` are worked out here for the connection component to draw.
|
|
*/
|
|
export type LoopAwareConnection = HFSchemes['Connection'] & {
|
|
loop?: FlowLoopEdgeSettings;
|
|
__loopBack?: boolean;
|
|
__readonly?: boolean;
|
|
/** Canvas y above both ends' nodes, where the way back is drawn so it crosses neither. */
|
|
__loopTop?: number;
|
|
};
|
|
|
|
export async function createEditor(
|
|
container: HTMLElement,
|
|
injector: Injector,
|
|
flowData: FlowData,
|
|
options?: { nodeView?: "editor" | "execution"; readonly?: boolean }
|
|
): Promise<ReteEditorInstance> {
|
|
|
|
const editor = new NodeEditor<HFSchemes>();
|
|
const area = new AreaPlugin<HFSchemes, AreaExtra>(container);
|
|
const connection = new ConnectionPlugin<HFSchemes, AreaExtra>();
|
|
const render = new AngularPlugin<HFSchemes, AreaExtra>({ injector });
|
|
const nodeView = options?.nodeView ?? "editor";
|
|
const readonly = options?.readonly === true;
|
|
const programmaticTranslations = new Set<string>();
|
|
areaProgrammaticTranslations.set(area, programmaticTranslations);
|
|
const runtime: ReteRuntimeContext = {
|
|
blocksService: injector.get(BlocksService),
|
|
containersService: injector.get(ContainersService),
|
|
flowState: injector.get(EditorStateHolder),
|
|
readonly,
|
|
globalInputs: cloneValue(flowData.globalInputs ?? []),
|
|
lanes: cloneValue(flowData.lanes ?? [])
|
|
};
|
|
editorRuntime.set(editor, runtime);
|
|
|
|
render.addPreset(
|
|
Presets.classic.setup({
|
|
customize: {
|
|
node(context: any) {
|
|
if (nodeView === "execution") return TaskStepNodeComponent;
|
|
const nodeFamily = context?.payload?.data?.nodeFamily;
|
|
return nodeFamily === "container" ? ContainerNodeComponent : GenericNodeComponent;
|
|
},
|
|
connection() {
|
|
return CustomConnectionComponent;
|
|
},
|
|
socket(context: any) {
|
|
// rete-angular passes only `payload` to the socket component.
|
|
// Build a per-render payload copy to avoid mutating shared socket objects.
|
|
const socketPayload = context?.payload;
|
|
const socketSide = context?.side === "output" ? "output" : "input";
|
|
context.payload = {
|
|
...(socketPayload ?? {}),
|
|
__hfSide: socketSide
|
|
};
|
|
return CustomSocket;
|
|
}
|
|
},
|
|
})
|
|
);
|
|
editor.addPipe(async (context) => {
|
|
if (context.type !== "connectioncreate") return context;
|
|
|
|
const sourceNode = editor.getNode(context.data.source) as HFNode | undefined;
|
|
const targetNode = editor.getNode(context.data.target) as HFNode | undefined;
|
|
const sourceCapabilities = resolveNodeCapabilities(runtime, sourceNode);
|
|
const targetCapabilities = resolveNodeCapabilities(runtime, targetNode);
|
|
|
|
const connectionKind = getGraphConnectionKind(context.data.sourceOutput, context.data.targetInput);
|
|
if (connectionKind === "dependency") {
|
|
if (context.data.source === context.data.target) return;
|
|
if (!sourceCapabilities.canHaveDependentNodes || !targetCapabilities.canDependOnOtherNodes) return;
|
|
return context;
|
|
}
|
|
|
|
if (!sourceCapabilities.allowsOutgoingConnections || !targetCapabilities.allowsIncomingConnections) return;
|
|
|
|
const sourcePort = resolveNodePort(sourceNode, "output", context.data.sourceOutput);
|
|
const targetPort = resolveNodePort(targetNode, "input", context.data.targetInput);
|
|
|
|
if (!sourcePort || !targetPort) return;
|
|
|
|
const compatible = areFlowValueKindsCompatible(
|
|
normalizeFlowPortValueKinds(sourcePort),
|
|
normalizeFlowPortValueKinds(targetPort)
|
|
);
|
|
if (!compatible) return undefined;
|
|
|
|
if (!runtime.restoringConnections) {
|
|
const leadsBackRoundALoop = await makeRoomOnInput(editor, runtime, context.data);
|
|
// Drawn closing a cycle from a router, this is the connection its author means to go round
|
|
// by. Saying so settles the rare shape where the graph alone could not tell which one it is.
|
|
if (leadsBackRoundALoop) {
|
|
(context.data as LoopAwareConnection).loop ??= {};
|
|
}
|
|
}
|
|
return context;
|
|
});
|
|
|
|
editor.addPipe((context) => {
|
|
if (context.type === "connectioncreated" || context.type === "connectionremoved") {
|
|
void refreshLoopMarkers(editor, area, runtime);
|
|
}
|
|
return context;
|
|
});
|
|
|
|
area.addPipe((context: any) => {
|
|
if (context?.type === 'nodetranslated') {
|
|
const nodeId = String(context.data?.id ?? '');
|
|
for (const connection of editor.getConnections() as LoopAwareConnection[]) {
|
|
if (!connection.__loopBack || (connection.source !== nodeId && connection.target !== nodeId)) continue;
|
|
connection.__loopTop = loopTop(area, connection);
|
|
void area.update('connection', connection.id);
|
|
}
|
|
}
|
|
return context;
|
|
});
|
|
|
|
connection.addPreset(ConnectionPresets.classic.setup());
|
|
|
|
AreaExtensions.simpleNodesOrder(area);
|
|
|
|
area.addPipe((context: any) => {
|
|
if (readonly && nodeView !== "execution" && context?.type === 'nodetranslate') {
|
|
const nodeId = String(context?.data?.id ?? '');
|
|
if (!programmaticTranslations.has(nodeId)) return;
|
|
}
|
|
return context;
|
|
});
|
|
|
|
editor.use(area);
|
|
area.use(connection);
|
|
area.use(render);
|
|
|
|
AreaExtensions.simpleNodesOrder(area);
|
|
AreaExtensions.restrictor(area, {
|
|
scaling: RETE_ZOOM_RANGE
|
|
});
|
|
|
|
if (flowData)
|
|
await loadFlowData(editor, area, flowData, runtime);
|
|
await refreshLoopMarkers(editor, area, runtime);
|
|
|
|
AreaExtensions.zoomAt(area, editor.getNodes());
|
|
return { editor, area };
|
|
}
|
|
|
|
export function exportGraph(editor: NodeEditor<HFSchemes>) {
|
|
const runtime = editorRuntime.get(editor);
|
|
const nodeIdToBlockId = new Map<string, string>();
|
|
const nodes: FlowNode[] = editor.getNodes().map((node) => {
|
|
const blockData = node.data;
|
|
const blockRecord = blockData as unknown as Record<string, unknown> | undefined;
|
|
const blockId = blockData?.id ?? node.id;
|
|
nodeIdToBlockId.set(node.id, blockId);
|
|
|
|
const inputs = cloneValue(blockData?.inputs ?? []);
|
|
const outputs = cloneValue(blockData?.outputs ?? []);
|
|
|
|
const biasAnnotationsProperty = typeof blockRecord?.['__biasAnnotationsProperty'] === 'string'
|
|
? String(blockRecord['__biasAnnotationsProperty'])
|
|
: 'biasAnnotations';
|
|
|
|
return {
|
|
id: blockId,
|
|
name: blockData?.name ?? node.label,
|
|
position: blockData?.position,
|
|
inputs,
|
|
outputs,
|
|
specificConfiguration: cloneValue(blockData?.specificConfiguration ?? {}),
|
|
capabilities: cloneValue(blockData?.capabilities),
|
|
...(blockData?.nodeFamily === 'container' ? {} : {
|
|
[biasAnnotationsProperty]: cloneValue(blockRecord?.[biasAnnotationsProperty] ?? [])
|
|
}),
|
|
typeName: blockData?.typeName ?? "LLMBlock",
|
|
nodeFamily: blockData?.nodeFamily === 'container' ? 'container' : 'block',
|
|
laneId: typeof blockRecord?.['laneId'] === 'string' ? blockRecord['laneId'] : null
|
|
};
|
|
});
|
|
|
|
const allConnections = editor.getConnections().map((c) => {
|
|
const loop = (c as LoopAwareConnection).loop;
|
|
return {
|
|
id: String(c.id),
|
|
sourceId: nodeIdToBlockId.get(c.source) ?? c.source,
|
|
sourceName: c.sourceOutput,
|
|
targetId: nodeIdToBlockId.get(c.target) ?? c.target,
|
|
targetName: c.targetInput,
|
|
...(loop ? { loop: cloneValue(loop) } : {}),
|
|
kind: getGraphConnectionKind(c.sourceOutput, c.targetInput)
|
|
};
|
|
});
|
|
|
|
return {
|
|
blocks: nodes.filter((node): node is FlowBlock => node.nodeFamily === 'block'),
|
|
containers: nodes.filter((node) => node.nodeFamily === 'container'),
|
|
connections: allConnections
|
|
.filter((connection) => connection.kind === 'data')
|
|
.map(({ kind, ...connection }) => connection),
|
|
dependencies: allConnections
|
|
.filter((connection) => connection.kind === 'dependency')
|
|
.map(({ sourceId, targetId }) => ({ sourceId, targetId })),
|
|
globalInputs: cloneValue(runtime?.globalInputs ?? []),
|
|
lanes: cloneValue(runtime?.lanes ?? [])
|
|
};
|
|
}
|
|
|
|
export function setEditorGlobalInputs(editor: NodeEditor<HFSchemes>, globalInputs: FlowGlobalInput[]) {
|
|
const runtime = editorRuntime.get(editor);
|
|
if (!runtime) return;
|
|
runtime.globalInputs = cloneValue(globalInputs ?? []);
|
|
}
|
|
|
|
export function setEditorLanes(editor: NodeEditor<HFSchemes>, lanes: FlowLane[]) {
|
|
const runtime = editorRuntime.get(editor);
|
|
if (!runtime) return;
|
|
runtime.lanes = cloneValue(lanes ?? []);
|
|
}
|
|
|
|
export function getEditorLanes(editor: NodeEditor<HFSchemes>): FlowLane[] {
|
|
return editorRuntime.get(editor)?.lanes ?? [];
|
|
}
|
|
|
|
export function isProgrammaticNodeTranslation(area: AreaPlugin<HFSchemes, AreaExtra>, nodeId: string): boolean {
|
|
return areaProgrammaticTranslations.get(area)?.has(nodeId) === true;
|
|
}
|
|
|
|
export async function addBlockToEditor(
|
|
editor: NodeEditor<HFSchemes>,
|
|
area: AreaPlugin<HFSchemes, AreaExtra>,
|
|
block: FlowNode,
|
|
position?: { x: number; y: number },
|
|
runtime?: ReteRuntimeContext
|
|
) {
|
|
const resolvedRuntime = runtime ?? editorRuntime.get(editor);
|
|
const node = new ClassicPreset.Node(block.typeName) as HFNode;
|
|
const removeNode = async () => {
|
|
if (!editor.getNode(node.id)) return;
|
|
const relatedConnectionIds = editor.getConnections()
|
|
.filter((connection) => connection.source === node.id || connection.target === node.id)
|
|
.map((connection) => connection.id);
|
|
|
|
for (const connectionId of relatedConnectionIds) {
|
|
await editor.removeConnection(connectionId);
|
|
}
|
|
|
|
await editor.removeNode(node.id);
|
|
};
|
|
const clearContainerSubflow = async (targetPath = "subFlow") => {
|
|
const currentNode = editor.getNode(node.id) as HFNode | undefined;
|
|
if (!currentNode?.data) return;
|
|
const nextConfiguration = {
|
|
...cloneValue(currentNode.data.specificConfiguration ?? {})
|
|
};
|
|
deleteSchemaValueByPath(nextConfiguration as Record<string, unknown>, targetPath);
|
|
const replacement = {
|
|
...cloneValue(currentNode.data),
|
|
inputs: [],
|
|
outputs: [],
|
|
specificConfiguration: nextConfiguration
|
|
};
|
|
const replaceNode = currentNode.data['replaceWithCreatedNode'];
|
|
if (typeof replaceNode === 'function') {
|
|
await replaceNode(replacement);
|
|
} else {
|
|
currentNode.data = {
|
|
...currentNode.data,
|
|
inputs: [],
|
|
outputs: [],
|
|
specificConfiguration: nextConfiguration,
|
|
__containerValidationErrors: [],
|
|
__containerAssignmentError: null,
|
|
__containerAssigning: false
|
|
};
|
|
await area.update("node", node.id);
|
|
}
|
|
if (resolvedRuntime) {
|
|
resolvedRuntime.flowState.updateData(exportGraph(editor));
|
|
}
|
|
};
|
|
const applyContainerSubflow = async (
|
|
candidateSubFlow: FlowData,
|
|
options?: { selectedIds?: Set<string>; targetPath?: string; validationUrl?: string | null; preValidate?: boolean; source?: 'drag' | 'import' }
|
|
) => {
|
|
if (!resolvedRuntime) return;
|
|
|
|
const liveNode = editor.getNode(node.id) as HFNode | undefined;
|
|
if (!liveNode?.data) return;
|
|
|
|
liveNode.data = {
|
|
...liveNode.data,
|
|
__containerAssigning: true,
|
|
__containerAssignmentError: null,
|
|
__containerValidationErrors: []
|
|
};
|
|
await area.update("node", node.id);
|
|
try {
|
|
const currentLiveNode = editor.getNode(node.id) as HFNode | undefined;
|
|
if (!currentLiveNode?.data) return;
|
|
|
|
if (options?.preValidate) {
|
|
const validation = await firstValueFrom(
|
|
resolvedRuntime.containersService.validateContainerSubflow(candidateSubFlow, options?.validationUrl)
|
|
);
|
|
|
|
if (!validation.valid) {
|
|
currentLiveNode.data = {
|
|
...currentLiveNode.data,
|
|
__containerAssigning: false,
|
|
__containerAssignmentError: validation.errors[0]?.message ?? "Selected subflow is not valid",
|
|
__containerValidationErrors: validation.errors
|
|
};
|
|
await area.update("node", node.id);
|
|
return;
|
|
}
|
|
}
|
|
|
|
const currentConfiguration = cloneValue(currentLiveNode.data.specificConfiguration ?? {}) as Record<string, unknown>;
|
|
const targetPath = options?.targetPath ?? 'subFlow';
|
|
const nextConfiguration: Record<string, unknown> = {
|
|
...currentConfiguration,
|
|
name: String(currentConfiguration['name'] ?? currentLiveNode.data['name'] ?? 'Container')
|
|
};
|
|
setSchemaValueByPath(nextConfiguration, targetPath, candidateSubFlow);
|
|
const nextPosition = cloneValue(currentLiveNode.data['position'] ?? null);
|
|
|
|
const containerId = String(currentLiveNode.data['id'] ?? '');
|
|
const replacementFromServer = await firstValueFrom(
|
|
resolvedRuntime.containersService.createContainer(containerId, {
|
|
...nextConfiguration,
|
|
position: nextPosition,
|
|
typeName: currentLiveNode.data['typeName']
|
|
})
|
|
);
|
|
const selectedIds = options?.selectedIds ?? new Set<string>();
|
|
const selectedNodeIds = editor.getNodes()
|
|
.filter((candidate) => selectedIds.has(String(candidate.data?.id ?? "")))
|
|
.map((candidate) => candidate.id)
|
|
.filter((candidateId) => candidateId !== node.id);
|
|
|
|
for (const selectedNodeId of selectedNodeIds) {
|
|
if (!editor.getNode(selectedNodeId)) continue;
|
|
|
|
const relatedConnectionIds = editor.getConnections()
|
|
.filter((connection) => connection.source === selectedNodeId || connection.target === selectedNodeId)
|
|
.map((connection) => connection.id);
|
|
|
|
for (const connectionId of relatedConnectionIds) {
|
|
await editor.removeConnection(connectionId);
|
|
}
|
|
|
|
await editor.removeNode(selectedNodeId);
|
|
}
|
|
|
|
const replacement = {
|
|
...cloneValue(currentLiveNode.data),
|
|
...cloneValue(replacementFromServer),
|
|
position: cloneValue(currentLiveNode.data['position'] ?? replacementFromServer.position),
|
|
specificConfiguration: nextConfiguration,
|
|
__containerAssigning: false,
|
|
__containerAssignmentError: null,
|
|
__containerValidationErrors: []
|
|
};
|
|
const replaceNode = currentLiveNode.data['replaceWithCreatedNode'];
|
|
if (typeof replaceNode === 'function') {
|
|
await replaceNode(replacement);
|
|
} else {
|
|
currentLiveNode.data = replacement;
|
|
await area.update("node", node.id);
|
|
}
|
|
|
|
resolvedRuntime.flowState.clearBlockSelection();
|
|
resolvedRuntime.flowState.updateData(exportGraph(editor));
|
|
} catch (error) {
|
|
const currentLiveNode = editor.getNode(node.id) as HFNode | undefined;
|
|
if (!currentLiveNode?.data) return;
|
|
|
|
currentLiveNode.data = {
|
|
...currentLiveNode.data,
|
|
__containerAssigning: false,
|
|
__containerValidationErrors: [],
|
|
__containerAssignmentError: error instanceof Error
|
|
? error.message
|
|
: "Container update failed"
|
|
};
|
|
await area.update("node", node.id);
|
|
}
|
|
};
|
|
const assignSelectedBlocksToContainer = async (
|
|
selectedBlockIds?: string[],
|
|
targetPath = 'subFlow',
|
|
validationUrl?: string | null
|
|
) => {
|
|
if (!resolvedRuntime) return;
|
|
|
|
const currentNode = editor.getNode(node.id) as HFNode | undefined;
|
|
const containerBlockId = currentNode?.data?.id;
|
|
if (!currentNode?.data || typeof containerBlockId !== "string") return;
|
|
|
|
const selection = Array.from(new Set((selectedBlockIds ?? resolvedRuntime.flowState.selectedBlockIds())
|
|
.filter((id) => typeof id === "string" && id.length > 0)
|
|
.filter((id) => id !== containerBlockId)));
|
|
|
|
if (!selection.length) {
|
|
currentNode.data = {
|
|
...currentNode.data,
|
|
__containerAssignmentError: "Select one or more nodes before dropping them into the container.",
|
|
__containerValidationErrors: []
|
|
};
|
|
await area.update("node", node.id);
|
|
return;
|
|
}
|
|
|
|
const currentFlow = exportGraph(editor);
|
|
const selectedBlocks = currentFlow.blocks.filter((candidate) => selection.includes(candidate.id));
|
|
const selectedContainers = currentFlow.containers.filter((candidate) => selection.includes(candidate.id));
|
|
const selectedIds = new Set([
|
|
...selectedBlocks.map((candidate) => candidate.id),
|
|
...selectedContainers.map((candidate) => candidate.id)
|
|
]);
|
|
const candidateSubFlow: FlowData = {
|
|
blocks: cloneValue(selectedBlocks),
|
|
containers: cloneValue(selectedContainers),
|
|
connections: cloneValue(
|
|
currentFlow.connections.filter((connection) =>
|
|
selectedIds.has(connection.sourceId) && selectedIds.has(connection.targetId)
|
|
)
|
|
),
|
|
dependencies: cloneValue(
|
|
(currentFlow.dependencies ?? []).filter((dependency) =>
|
|
selectedIds.has(dependency.sourceId) && selectedIds.has(dependency.targetId)
|
|
)
|
|
)
|
|
};
|
|
await applyContainerSubflow(candidateSubFlow, { selectedIds, targetPath, validationUrl, preValidate: true });
|
|
};
|
|
const assignImportedSubflow = async (subFlow: FlowData, targetPath = 'subFlow', validationUrl?: string | null) => {
|
|
await applyContainerSubflow(cloneValue(subFlow), { targetPath, validationUrl, preValidate: true, source: 'import' });
|
|
};
|
|
const replaceWithCreatedNode = async (createdBlock: FlowNode) => {
|
|
if (!editor.getNode(node.id)) return;
|
|
const createdBlockData = createdBlock as FlowNode & Record<string, unknown>;
|
|
const previousConnections = editor.getConnections()
|
|
.filter((connection) => connection.source === node.id || connection.target === node.id)
|
|
.map((connection) => ({
|
|
id: connection.id,
|
|
source: connection.source,
|
|
sourceOutput: connection.sourceOutput,
|
|
target: connection.target,
|
|
targetInput: connection.targetInput,
|
|
loop: (connection as LoopAwareConnection).loop
|
|
}));
|
|
const currentPosition = (node.data?.position ?? position ?? createdBlock.position) as { x: number; y: number } | undefined;
|
|
const replacementNode = await addBlockToEditor(
|
|
editor,
|
|
area,
|
|
{
|
|
...createdBlockData,
|
|
position: currentPosition,
|
|
__focusOpen: node.data?.['__focusOpen'] === true || createdBlockData['__focusOpen'] === true
|
|
} as FlowNode & Record<string, unknown>,
|
|
currentPosition,
|
|
resolvedRuntime
|
|
);
|
|
|
|
if (!replacementNode) return;
|
|
|
|
for (const connection of previousConnections) {
|
|
await editor.removeConnection(connection.id);
|
|
}
|
|
await editor.removeNode(node.id);
|
|
|
|
const replacementOutputNames = new Set(Object.keys(replacementNode.outputs));
|
|
const replacementInputNames = new Set(Object.keys(replacementNode.inputs));
|
|
|
|
for (const connection of previousConnections) {
|
|
const sourceNode = connection.source === node.id
|
|
? replacementNode
|
|
: editor.getNode(connection.source);
|
|
const targetNode = connection.target === node.id
|
|
? replacementNode
|
|
: editor.getNode(connection.target);
|
|
|
|
if (!sourceNode || !targetNode) continue;
|
|
|
|
const sourceOutput = connection.sourceOutput;
|
|
const targetInput = connection.targetInput;
|
|
|
|
if (connection.source === node.id && !replacementOutputNames.has(sourceOutput)) continue;
|
|
if (connection.target === node.id && !replacementInputNames.has(targetInput)) continue;
|
|
|
|
try {
|
|
const restored = new ClassicPreset.Connection(sourceNode as HFNode, sourceOutput, targetNode as HFNode, targetInput) as LoopAwareConnection;
|
|
restored.id = connection.id;
|
|
if (connection.loop) restored.loop = cloneValue(connection.loop);
|
|
await withRestoredConnections(resolvedRuntime, () => editor.addConnection(restored));
|
|
} catch (error) {
|
|
console.warn('Failed to restore connection after node replacement', {
|
|
connection,
|
|
error
|
|
});
|
|
}
|
|
}
|
|
};
|
|
const cloneNode = async () => {
|
|
if (resolvedRuntime?.readonly) return;
|
|
|
|
const currentNode = editor.getNode(node.id) as HFNode | undefined;
|
|
if (!currentNode?.data) return;
|
|
|
|
const sourceData = currentNode.data as Record<string, unknown>;
|
|
const sourcePosition = sourceData['position'] as { x: number; y: number } | undefined;
|
|
const nextPosition = sourcePosition
|
|
? { x: sourcePosition.x + 48, y: sourcePosition.y + 48 }
|
|
: { x: 168, y: 148 };
|
|
|
|
const clonedNode = {
|
|
...cloneValue(block),
|
|
...cloneValue(sourceData),
|
|
id: globalThis.crypto?.randomUUID?.() ?? `${Date.now()}`,
|
|
position: nextPosition,
|
|
inputs: cloneValue((sourceData['inputs'] as FlowNode['inputs'] | undefined) ?? block.inputs ?? []),
|
|
outputs: cloneValue((sourceData['outputs'] as FlowNode['outputs'] | undefined) ?? block.outputs ?? []),
|
|
specificConfiguration: cloneValue((sourceData['specificConfiguration'] as FlowNode['specificConfiguration'] | undefined) ?? block.specificConfiguration ?? {}),
|
|
nodeFamily: sourceData['nodeFamily'] === 'container' ? 'container' : 'block',
|
|
__needsServerCreate: sourceData['nodeFamily'] === 'container' ? false : true,
|
|
__createdOnServer: false,
|
|
__isCreatingOnServer: false,
|
|
__focusOpen: false,
|
|
__updateBlockError: null
|
|
} as FlowNode & Record<string, unknown>;
|
|
|
|
await addBlockToEditor(editor, area, clonedNode, nextPosition, resolvedRuntime);
|
|
};
|
|
node.data = {
|
|
...cloneValue(block),
|
|
position: position ?? block.position,
|
|
__readonly: resolvedRuntime?.readonly === true,
|
|
deleteNode: removeNode,
|
|
cloneNode,
|
|
replaceWithCreatedNode,
|
|
assignSelectedBlocksToContainer,
|
|
assignImportedSubflow,
|
|
clearContainerSubflow,
|
|
__containerValidationErrors: [],
|
|
__containerAssignmentError: null,
|
|
__containerAssigning: false
|
|
};
|
|
|
|
const capabilities = resolveNodeCapabilities(resolvedRuntime, node);
|
|
if (capabilities.canHaveDependentNodes) {
|
|
node.addOutput(FLOW_DEPENDANT_PORT_KEY, new ClassicPreset.Output(getSocket(editor, FLOW_DEPENDENCY_SOCKET_TYPE)));
|
|
}
|
|
if (capabilities.canDependOnOtherNodes) {
|
|
node.addInput(FLOW_DEPENDENCY_PORT_KEY, new ClassicPreset.Input(getSocket(editor, FLOW_DEPENDENCY_SOCKET_TYPE), undefined, true));
|
|
}
|
|
|
|
for (const output of block.outputs ?? []) {
|
|
node.addOutput(output.name, new ClassicPreset.Output(getSocket(editor, output.type ?? "ANY")));
|
|
}
|
|
|
|
for (const input of block.inputs ?? []) {
|
|
// Multiple so that Rete does not drop the input's connection itself before asking: the entry of
|
|
// a loop takes a second one, from the way back. makeRoomOnInput keeps every other input to one.
|
|
node.addInput(input.name, new ClassicPreset.Input(getSocket(editor, input.type ?? "ANY"), undefined, true));
|
|
}
|
|
|
|
await editor.addNode(node);
|
|
|
|
const targetPosition = position ?? block.position;
|
|
if (targetPosition) {
|
|
node.data = {
|
|
...node.data,
|
|
position: { x: targetPosition.x, y: targetPosition.y }
|
|
};
|
|
await applyNodePosition(area, node.id, targetPosition);
|
|
}
|
|
|
|
return node;
|
|
}
|
|
|
|
async function loadFlowData(
|
|
editor: NodeEditor<HFSchemes>,
|
|
area: AreaPlugin<HFSchemes, AreaExtra>,
|
|
flowData: FlowData,
|
|
runtime?: ReteRuntimeContext
|
|
) {
|
|
const topLevelNodes = [...(flowData.blocks ?? []), ...(flowData.containers ?? [])];
|
|
if (!topLevelNodes.length) return;
|
|
|
|
const nodeMapping = new Map<string, any>();
|
|
|
|
for (const [index, block] of topLevelNodes.entries()) {
|
|
const fallbackPosition = block.position ?? {
|
|
x: 120 + (index % 3) * 340,
|
|
y: 100 + Math.floor(index / 3) * 220
|
|
};
|
|
const node = await addBlockToEditor(editor, area, block, fallbackPosition, runtime);
|
|
nodeMapping.set(block.id, node.id);
|
|
}
|
|
|
|
await withRestoredConnections(runtime, async () => {
|
|
for (const c of flowData.connections ?? []) {
|
|
if (!nodeMapping.has(c.sourceId) || !nodeMapping.has(c.targetId)) continue;
|
|
|
|
const sourceNode = editor.getNode(nodeMapping.get(c.sourceId)) as any;
|
|
const targetNode = editor.getNode(nodeMapping.get(c.targetId)) as any;
|
|
|
|
const connection = new ClassicPreset.Connection(sourceNode, c.sourceName, targetNode, c.targetName) as LoopAwareConnection;
|
|
// The saved id, not a fresh one: settings and selections refer to a connection by it, and a
|
|
// new id on every load would detach them.
|
|
if (c.id) connection.id = c.id;
|
|
if (c.loop) connection.loop = cloneValue(c.loop);
|
|
await editor.addConnection(connection);
|
|
}
|
|
|
|
for (const dependency of flowData.dependencies ?? []) {
|
|
if (!nodeMapping.has(dependency.sourceId) || !nodeMapping.has(dependency.targetId)) continue;
|
|
|
|
const sourceNode = editor.getNode(nodeMapping.get(dependency.sourceId)) as any;
|
|
const targetNode = editor.getNode(nodeMapping.get(dependency.targetId)) as any;
|
|
|
|
await editor.addConnection(
|
|
new ClassicPreset.Connection(sourceNode, FLOW_DEPENDANT_PORT_KEY, targetNode, FLOW_DEPENDENCY_PORT_KEY)
|
|
);
|
|
}
|
|
});
|
|
}
|
|
|
|
async function withRestoredConnections<T>(runtime: ReteRuntimeContext | undefined, restore: () => Promise<T>): Promise<T> {
|
|
if (!runtime) return restore();
|
|
runtime.restoringConnections = (runtime.restoringConnections ?? 0) + 1;
|
|
try {
|
|
return await restore();
|
|
} finally {
|
|
runtime.restoringConnections -= 1;
|
|
}
|
|
}
|
|
|
|
type NewConnection = { source: string; sourceOutput: string; target: string; targetInput: string };
|
|
|
|
/**
|
|
* Keeps an input to one connection, as it always was, with one exception: the entry of a loop,
|
|
* which takes one connection from before the loop and one leading back from the router that
|
|
* decides whether to go round. A connection "leads back" here when it comes from a node that routes
|
|
* exclusively and its target can already reach that node - it closes a cycle.
|
|
*
|
|
* So an existing connection stays when exactly one of it and the new one leads back; otherwise the
|
|
* new one replaces it, as dropping onto a taken input always did.
|
|
*/
|
|
async function makeRoomOnInput(editor: NodeEditor<HFSchemes>, runtime: ReteRuntimeContext, created: NewConnection) {
|
|
const newLeadsBack = leadsBack(editor, runtime, created, null);
|
|
const existing = editor.getConnections().filter((connection) =>
|
|
connection.target === created.target && connection.targetInput === created.targetInput);
|
|
for (const connection of existing) {
|
|
const existingLeadsBack = leadsBack(editor, runtime, connection, connection.id);
|
|
if (newLeadsBack !== existingLeadsBack) continue;
|
|
await editor.removeConnection(connection.id);
|
|
}
|
|
return newLeadsBack;
|
|
}
|
|
|
|
function leadsBack(editor: NodeEditor<HFSchemes>, runtime: ReteRuntimeContext, connection: NewConnection,
|
|
ignoredConnectionId: string | null): boolean {
|
|
const source = editor.getNode(connection.source) as HFNode | undefined;
|
|
if (resolveNodeCapabilities(runtime, source).routesExclusively !== true) return false;
|
|
const outgoing = new Map<string, string[]>();
|
|
for (const candidate of editor.getConnections()) {
|
|
if (candidate.id === ignoredConnectionId) continue;
|
|
outgoing.set(candidate.source, [...(outgoing.get(candidate.source) ?? []), candidate.target]);
|
|
}
|
|
const seen = new Set<string>();
|
|
const pending = [connection.target];
|
|
while (pending.length) {
|
|
const next = pending.pop()!;
|
|
if (next === connection.source) return true;
|
|
if (seen.has(next)) continue;
|
|
seen.add(next);
|
|
pending.push(...(outgoing.get(next) ?? []));
|
|
}
|
|
return false;
|
|
}
|
|
|
|
/**
|
|
* Marks which connections lead back round a loop, for the connection component to draw as such.
|
|
* Worked out with the same rule the server uses, from the editor's graph as it now stands.
|
|
*/
|
|
export async function refreshLoopMarkers(
|
|
editor: NodeEditor<HFSchemes>,
|
|
area: AreaPlugin<HFSchemes, AreaExtra>,
|
|
runtime?: ReteRuntimeContext
|
|
) {
|
|
const resolvedRuntime = runtime ?? editorRuntime.get(editor);
|
|
const connections = editor.getConnections() as LoopAwareConnection[];
|
|
const data = connections.filter((c) => getGraphConnectionKind(c.sourceOutput, c.targetInput) === 'data');
|
|
const backEdges = findLoopBackEdgeIds({
|
|
nodes: editor.getNodes().map((node) => ({
|
|
id: node.id,
|
|
routesExclusively: resolveNodeCapabilities(resolvedRuntime, node as HFNode).routesExclusively === true
|
|
})),
|
|
connections: data.map((c) => ({ id: c.id, sourceId: c.source, targetId: c.target, hasLoopSettings: c.loop != null })),
|
|
dependencies: connections
|
|
.filter((c) => getGraphConnectionKind(c.sourceOutput, c.targetInput) === 'dependency')
|
|
.map((c) => ({ sourceId: c.source, targetId: c.target }))
|
|
});
|
|
// Every connection, dependencies included: in a read-only graph none of them can be selected.
|
|
for (const connection of connections) {
|
|
const loopBack = backEdges.has(connection.id);
|
|
const readonly = resolvedRuntime?.readonly === true;
|
|
const top = loopBack ? loopTop(area, connection) : undefined;
|
|
if (connection.__loopBack === loopBack && connection.__readonly === readonly && connection.__loopTop === top) continue;
|
|
connection.__loopBack = loopBack;
|
|
connection.__readonly = readonly;
|
|
connection.__loopTop = top;
|
|
await area.update('connection', connection.id);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Above the tops of the nodes at both ends. Their tops, not their bottoms: a node's height changes
|
|
* as it is expanded or edited, its position only when it is moved.
|
|
*/
|
|
function loopTop(area: AreaPlugin<HFSchemes, AreaExtra>, connection: LoopAwareConnection): number | undefined {
|
|
const tops = [connection.source, connection.target]
|
|
.map((id) => area.nodeViews.get(id)?.position.y)
|
|
.filter((y): y is number => typeof y === 'number');
|
|
return tops.length ? Math.min(...tops) - 48 : undefined;
|
|
}
|
|
|
|
/**
|
|
* Sets a loop connection's iteration limit; null goes back to the default. The settings stay on the
|
|
* connection either way, so it remains the one its author marked as leading back.
|
|
*/
|
|
export async function setLoopMaxIterations(
|
|
editor: NodeEditor<HFSchemes>,
|
|
area: AreaPlugin<HFSchemes, AreaExtra>,
|
|
connectionId: string,
|
|
maxIterations: number | null
|
|
) {
|
|
const connection = editor.getConnections().find((c) => String(c.id) === connectionId) as LoopAwareConnection | undefined;
|
|
if (!connection) return false;
|
|
connection.loop = maxIterations == null ? {} : { maxIterations };
|
|
await area.update('connection', connection.id);
|
|
return true;
|
|
}
|
|
|
|
function getSocket(editor: NodeEditor<HFSchemes>, type: string) {
|
|
if (!editorSockets.has(editor)) {
|
|
editorSockets.set(editor, new Map<string, ClassicPreset.Socket>());
|
|
}
|
|
const map = editorSockets.get(editor)!;
|
|
if (!map.has(type)) {
|
|
const socket = new ClassicPreset.Socket(type) as ClassicPreset.Socket & { __hfKind?: GraphConnectionKind };
|
|
socket.__hfKind = type === FLOW_DEPENDENCY_SOCKET_TYPE ? 'dependency' : 'data';
|
|
map.set(type, socket);
|
|
}
|
|
return map.get(type)!;
|
|
}
|
|
|
|
export function resolveNodeCapabilities(runtime: ReteRuntimeContext | undefined, node: HFNode | undefined): NodeTypeCapabilities {
|
|
if (node?.data?.capabilities) return node.data.capabilities;
|
|
|
|
const typeName = node?.data?.typeName;
|
|
if (!runtime || typeof typeName !== "string" || !typeName) return DEFAULT_NODE_CAPABILITIES;
|
|
|
|
const descriptor = node?.data?.nodeFamily === "container"
|
|
? runtime.containersService.peekContainerType(typeName)
|
|
: runtime.blocksService.peekBlockType(typeName);
|
|
|
|
return descriptor?.capabilities ?? DEFAULT_NODE_CAPABILITIES;
|
|
}
|
|
|
|
function resolveNodePort(node: HFNode | undefined, kind: "input" | "output", portName: string) {
|
|
const ports = node?.data?.[kind === "input" ? "inputs" : "outputs"];
|
|
if (!Array.isArray(ports)) return null;
|
|
return ports.find((port) => port?.name === portName) ?? null;
|
|
}
|
|
|
|
async function applyNodePosition(
|
|
area: AreaPlugin<HFSchemes, AreaExtra>,
|
|
nodeId: string,
|
|
position: { x: number; y: number }
|
|
) {
|
|
const programmaticTranslations = areaProgrammaticTranslations.get(area);
|
|
if (programmaticTranslations) programmaticTranslations.add(nodeId);
|
|
try {
|
|
await area.translate(nodeId, position);
|
|
} finally {
|
|
programmaticTranslations?.delete(nodeId);
|
|
}
|
|
}
|
|
|
|
function cloneValue<T>(value: T): T {
|
|
if (typeof globalThis.structuredClone === "function") {
|
|
try {
|
|
return globalThis.structuredClone(value);
|
|
} catch {
|
|
// Functions and runtime socket objects are not cloneable; fall back to JSON-safe clone.
|
|
}
|
|
}
|
|
return JSON.parse(JSON.stringify(value)) as T;
|
|
}
|
|
|
|
function getGraphConnectionKind(sourceOutput: string, targetInput: string): GraphConnectionKind {
|
|
return sourceOutput === FLOW_DEPENDANT_PORT_KEY && targetInput === FLOW_DEPENDENCY_PORT_KEY
|
|
? 'dependency'
|
|
: 'data';
|
|
}
|