@@ -72,7 +96,7 @@
- @for (message of messages(); track message.id) {
+ @for (message of displayedMessages(); track message.id) {
{{ message.role }}
{{ message.content }}
@@ -98,7 +122,7 @@
@if (assistantBusy()) {
assistant
- Working on the workflow...
+ {{ busyHeadline() }}
}
diff --git a/src/app/shared/flow-assistant/flow-assistant.ts b/src/app/shared/flow-assistant/flow-assistant.ts
index 060216b..bd4e4e6 100644
--- a/src/app/shared/flow-assistant/flow-assistant.ts
+++ b/src/app/shared/flow-assistant/flow-assistant.ts
@@ -1,13 +1,14 @@
import { CommonModule } from '@angular/common';
-import { Component, computed, inject, OnInit, signal } from '@angular/core';
+import { Component, computed, inject, OnDestroy, OnInit, signal } from '@angular/core';
import { FormsModule } from '@angular/forms';
+import { environment } from '@environment';
import {
+ AssistantCallPhase,
+ AssistantCallState,
AssistantChatMessage,
+ AssistantConfig,
AssistantDraftPayload,
- AssistantEditorDraft,
- AssistantFlowResponse,
- AssistantIntent,
- AssistantValidationIssue
+ AssistantSessionState
} from '@models/assistant';
import { Flow } from '@models/flow';
import { AssistantService } from '@services/assistant/assistant';
@@ -21,32 +22,81 @@ import { finalize, take } from 'rxjs';
templateUrl: './flow-assistant.html',
styleUrl: './flow-assistant.css'
})
-export class FlowAssistant implements OnInit {
+export class FlowAssistant implements OnInit, OnDestroy {
private readonly assistant = inject(AssistantService);
private readonly editorState = inject(EditorStateHolder);
private readonly authorization = inject(Authorization);
+ private pollTick: ReturnType
| null = null;
+ readonly assistantConfig = signal(null);
+ readonly sessionState = signal(null);
+ readonly currentCall = signal(null);
+ readonly localMessages = signal([]);
readonly models = signal([]);
readonly modelsLoading = signal(false);
readonly modelsError = signal(null);
- readonly selectedModel = signal('');
- readonly modelPickerOpen = signal(true);
- readonly assistantBusy = signal(false);
- readonly currentIntent = signal(null);
- readonly lastValidationErrors = signal([]);
+ readonly sessionLoading = signal(false);
readonly prompt = signal('');
+ readonly selectedModel = signal('');
+ readonly modelPickerOpen = signal(false);
readonly quickPromptsOpen = signal(true);
- readonly messages = signal([
- {
- id: crypto.randomUUID(),
- role: 'system',
- content: 'Select a model, then ask me to create, refine, fix, or explain a workflow.'
- }
- ]);
+ readonly initialSystemMessage: AssistantChatMessage = {
+ id: 'assistant-system-welcome',
+ role: 'system',
+ content: 'Select a model, then ask me to create, refine, fix, or explain a workflow.'
+ };
+
+ readonly displayedMessages = computed(() => {
+ const baseMessages = this.sessionState()?.messages?.length
+ ? this.sessionState()!.messages
+ : [this.initialSystemMessage];
+ return [...baseMessages, ...this.localMessages()];
+ });
+ readonly assistantBusy = computed(() => {
+ const status = this.currentCall()?.status;
+ return this.sessionLoading() || status === 'QUEUED' || status === 'RUNNING';
+ });
+ readonly activePhase = computed(() => {
+ const call = this.currentCall();
+ return call ? call.phase : null;
+ });
+ readonly activePhaseLabel = computed(() => {
+ const phase = this.activePhase();
+ return phase ? this.phaseText(phase) : null;
+ });
+ readonly callProgressMessage = computed(() => this.currentCall()?.progressMessage ?? '');
+ readonly busyHeadline = computed(() => {
+ const phase = this.activePhase();
+ if (!phase) {
+ return this.sessionLoading() ? 'Preparing assistant session...' : '';
+ }
+ return this.phaseText(phase);
+ });
+ readonly busyDetail = computed(() => {
+ if (this.sessionLoading()) {
+ return 'Loading assistant configuration, models, and chat session...';
+ }
+ if (!this.currentCall()) return '';
+ return this.callProgressMessage() || 'The backend may perform multiple internal steps before returning the updated conversation and flow.';
+ });
+ readonly progressSteps = [
+ { phase: 'queued', label: 'Queued...' },
+ { phase: 'routing', label: 'Routing your request...' },
+ { phase: 'planning', label: 'Planning workflow blocks...' },
+ { phase: 'configuring_blocks', label: 'Configuring blocks...' },
+ { phase: 'connecting_blocks', label: 'Connecting blocks...' },
+ { phase: 'validating', label: 'Validating flow...' },
+ { phase: 'fixing', label: 'Repairing invalid flow...' },
+ { phase: 'explaining', label: 'Explaining current flow...' }
+ ] as const;
+ readonly activeProgressIndex = computed(() => {
+ const phase = this.activePhase();
+ if (!phase) return -1;
+ return this.progressSteps.findIndex((step) => step.phase === phase);
+ });
readonly currentFlow = this.editorState.currentFlow;
- readonly draftDirty = this.editorState.isDirty;
- readonly currentDraft = computed(() => this.toAssistantDraft(this.currentFlow()));
+ readonly currentDraft = computed(() => this.sessionState()?.currentDraftFlow ?? null);
readonly starterPrompts = [
'Create a flow that classifies incoming tickets and sends urgent ones to a human',
'Modify the current flow to add a review step after the LLM block',
@@ -55,11 +105,17 @@ export class FlowAssistant implements OnInit {
];
ngOnInit(): void {
- this.loadModels();
+ this.bootstrapAssistant();
+ }
+
+ ngOnDestroy(): void {
+ this.stopPolling();
}
selectModel(model: string) {
+ if (!model || model === this.selectedModel()) return;
this.selectedModel.set(model);
+ void this.openSession(model);
}
useStarter(prompt: string) {
@@ -75,178 +131,188 @@ export class FlowAssistant implements OnInit {
}
submitPrompt() {
- const userPrompt = this.prompt().trim();
- if (!userPrompt || this.assistantBusy() || !this.selectedModel()) return;
+ const content = this.prompt().trim();
+ const sessionId = this.sessionState()?.id;
+ if (!content || this.assistantBusy() || !sessionId) return;
this.prompt.set('');
- this.pushMessage({
- role: 'user',
- content: userPrompt
- });
-
- const clarification = this.maybeClarify(userPrompt);
- if (clarification) {
- this.pushMessage({
- role: 'assistant',
- content: clarification
- });
- return;
- }
-
- const intent = this.determineIntent(userPrompt);
- this.currentIntent.set(intent);
- this.assistantBusy.set(true);
-
- if (intent === 'explain') {
- const draft = this.currentDraft();
- if (!draft) {
- this.pushMessage({
- role: 'assistant',
- content: 'Open or generate a flow first, then I can explain it in detail.',
- intent
- });
- this.assistantBusy.set(false);
- return;
+ this.localMessages.set([
+ {
+ id: crypto.randomUUID(),
+ role: 'user',
+ content
}
+ ]);
- this.assistant.explainDraft({
- userPrompt,
- model: this.selectedModel(),
- flow: draft
- }).pipe(
- take(1),
- finalize(() => this.assistantBusy.set(false))
- ).subscribe({
- next: (response) => {
- this.pushMessage({
- role: 'assistant',
- content: response.explanation || 'I analyzed the current flow.',
- intent,
- warnings: response.warnings
- });
- },
- error: (err) => this.pushAssistantError(err, intent)
- });
- return;
- }
-
- const draft = this.currentDraft();
- const maxRepairAttempts = intent === 'fix' ? 2 : 1;
- const request$ = intent === 'draft' || !draft
- ? this.assistant.createDraft({
- userPrompt,
- model: this.selectedModel(),
- maxRepairAttempts
- })
- : intent === 'fix'
- ? this.assistant.fixDraft({
- userPrompt,
- model: this.selectedModel(),
- maxRepairAttempts,
- flow: draft,
- validationErrors: this.lastValidationErrors()
- })
- : this.assistant.refineDraft({
- userPrompt,
- model: this.selectedModel(),
- maxRepairAttempts,
- flow: draft
- });
-
- request$.pipe(
- take(1),
- finalize(() => this.assistantBusy.set(false))
+ this.assistant.sendMessage(sessionId, { message: content }).pipe(
+ take(1)
).subscribe({
- next: (response) => void this.applyAssistantFlowResponse(response, intent),
- error: (err) => this.pushAssistantError(err, intent)
+ next: ({ callId }) => {
+ this.currentCall.set({
+ id: callId,
+ sessionId,
+ status: 'QUEUED',
+ phase: 'queued'
+ });
+ this.startPolling(callId, sessionId);
+ },
+ error: (err) => {
+ console.error('Assistant send message failed', err);
+ this.pushLocalAssistantMessage('The assistant request failed.');
+ }
});
}
- private loadModels() {
- this.modelsLoading.set(true);
+ private bootstrapAssistant() {
+ this.sessionLoading.set(true);
this.modelsError.set(null);
+ this.assistant.getConfig().pipe(
+ take(1)
+ ).subscribe({
+ next: (config) => {
+ this.assistantConfig.set(config);
+ this.loadModelsAndSession(config);
+ },
+ error: (err) => {
+ console.error('Assistant config loading failed', err);
+ this.sessionLoading.set(false);
+ this.modelsError.set('Unable to load assistant configuration.');
+ }
+ });
+ }
- this.assistant.listModels().pipe(
+ private loadModelsAndSession(config: AssistantConfig) {
+ this.modelsLoading.set(true);
+ const resolvedUrl = this.resolveAssistantModelsUrl(config.availableModelsRetrieverUrl);
+ this.assistant.listModels(config.availableModelsRetrieverUrl).pipe(
take(1),
finalize(() => this.modelsLoading.set(false))
).subscribe({
next: (models) => {
this.models.set(models);
- if (!this.selectedModel() && models.length) {
- this.selectedModel.set(models[0]);
+ const selectedModel = config.defaultModel || models[0] || '';
+ this.selectedModel.set(selectedModel);
+ if (!selectedModel) {
+ this.sessionLoading.set(false);
+ this.modelsError.set('No assistant model is available.');
+ return;
}
+ void this.openSession(selectedModel);
},
error: (err) => {
console.error('Assistant model loading failed', err);
- this.modelsError.set('Unable to load internal assistant models.');
+ this.sessionLoading.set(false);
+ this.modelsError.set(
+ `Unable to load assistant models from ${resolvedUrl}.`
+ );
}
});
}
- private determineIntent(prompt: string): AssistantIntent {
- const normalized = prompt.toLowerCase();
- const hasDraft = !!this.currentDraft();
+ private async openSession(model: string) {
+ this.stopPolling();
+ this.currentCall.set(null);
+ this.localMessages.set([]);
+ this.sessionLoading.set(true);
- if (/explain|what does|spiega|cosa fa|why/i.test(normalized)) return 'explain';
- if (/fix|repair|problem|invalid|error|errore|bug/i.test(normalized)) return 'fix';
- if (!hasDraft) return 'draft';
- if (/create new|new flow|from scratch|nuovo flow/i.test(normalized)) return 'draft';
- return 'refine';
+ this.assistant.createSession({ model }).pipe(
+ take(1),
+ finalize(() => this.sessionLoading.set(false))
+ ).subscribe({
+ next: (session) => {
+ this.applySessionState(session);
+ },
+ error: (err) => {
+ console.error('Assistant session creation failed', err);
+ this.pushLocalAssistantMessage('Unable to create an assistant session.');
+ }
+ });
}
- private maybeClarify(prompt: string): string | null {
- const normalized = prompt.trim();
- if (normalized.split(/\s+/).length >= 4) return null;
- if (this.currentDraft()) return null;
- return 'The request is too short to generate a useful flow. Tell me in one sentence what the workflow should do.';
+ private startPolling(callId: string, sessionId: string) {
+ this.stopPolling();
+ this.pollTick = setInterval(() => {
+ this.assistant.getCall(callId).pipe(
+ take(1)
+ ).subscribe({
+ next: (callState) => {
+ this.currentCall.set(callState);
+ if (callState.status === 'COMPLETED' || callState.status === 'FAILED') {
+ this.stopPolling();
+ void this.reloadSession(sessionId, callState.status === 'FAILED' ? callState.errorMessage : undefined);
+ }
+ },
+ error: (err) => {
+ console.error('Assistant call polling failed', err);
+ this.stopPolling();
+ this.pushLocalAssistantMessage('Polling the assistant call failed.');
+ }
+ });
+ }, 500);
}
- private async applyAssistantFlowResponse(response: AssistantFlowResponse, intent: AssistantIntent) {
- this.lastValidationErrors.set(response.validationErrors);
+ private async reloadSession(sessionId: string, failureMessage?: string) {
+ this.assistant.getSession(sessionId).pipe(
+ take(1)
+ ).subscribe({
+ next: (session) => {
+ this.applySessionState(session);
+ this.currentCall.set(null);
+ if (failureMessage) {
+ this.pushLocalAssistantMessage(failureMessage);
+ }
+ },
+ error: (err) => {
+ console.error('Assistant session refresh failed', err);
+ this.currentCall.set(null);
+ this.pushLocalAssistantMessage('Unable to refresh the assistant session.');
+ }
+ });
+ }
+
+ private applySessionState(session: AssistantSessionState) {
+ const normalizedSession = session.messages.length
+ ? session
+ : {
+ ...session,
+ messages: [this.initialSystemMessage]
+ };
+
+ this.sessionState.set(normalizedSession);
+ this.selectedModel.set(normalizedSession.selectedModel || this.selectedModel());
+ this.localMessages.set([]);
+ this.syncDraftToEditor(normalizedSession.currentDraftFlow);
+ }
+
+ private syncDraftToEditor(draft: AssistantDraftPayload | null) {
+ if (!draft) return;
const currentFlow = this.currentFlow();
- const nextFlow = this.toEditorFlow(response, currentFlow, intent);
- const isReplacingDocument = intent === 'draft' && currentFlow?.id !== nextFlow.id;
- if (isReplacingDocument) {
- const opened = await this.editorState.openDocument(nextFlow);
- if (!opened) {
- this.pushMessage({
- role: 'assistant',
- content: 'The new draft is ready, but I did not load it because the current flow has unsaved changes.',
- intent
- });
- return;
- }
+ const nextFlow = this.toEditorFlow(draft, currentFlow);
+
+ if (currentFlow?.id) {
+ this.editorState.loadAssistantFlow(nextFlow, { markDirty: true });
+ return;
}
- this.editorState.loadAssistantFlow(nextFlow, { markDirty: true });
-
- const summary = [
- response.assistantRationale || this.defaultAssistantSummary(intent, response.valid),
- response.valid ? 'The draft is valid.' : 'The draft still has validation errors.'
- ].filter(Boolean).join(' ');
-
- this.pushMessage({
- role: 'assistant',
- content: summary,
- intent,
- warnings: response.warnings,
- validationErrors: response.validationErrors
+ void this.editorState.openDocument(nextFlow, { skipDirtyCheck: false }).then((opened) => {
+ if (!opened) {
+ this.pushLocalAssistantMessage('The draft is ready, but I did not load it because the current flow has unsaved changes.');
+ return;
+ }
+ this.editorState.loadAssistantFlow(nextFlow, { markDirty: true });
});
}
- private toEditorFlow(response: AssistantFlowResponse, currentFlow: Flow | null, intent: AssistantIntent): Flow {
- const shouldReuseCurrentId = !!currentFlow && intent !== 'draft';
- const nextId = shouldReuseCurrentId
- ? currentFlow!.id
- : `${EditorStateHolder.ASSISTANT_DRAFT_PREFIX}${crypto.randomUUID()}`;
+ private toEditorFlow(draft: AssistantDraftPayload, currentFlow: Flow | null): Flow {
+ const nextId = currentFlow?.id ?? `${EditorStateHolder.ASSISTANT_DRAFT_PREFIX}${crypto.randomUUID()}`;
return {
id: nextId,
- name: response.flow.name,
- description: response.flow.description,
- data: response.flow.flow,
- status: response.valid ? 'EXECUTABLE' : 'DRAFT',
+ name: draft.name,
+ description: draft.description,
+ data: draft.flow,
+ status: 'DRAFT',
visibility: currentFlow?.visibility ?? 'PRIVATE',
author: currentFlow?.author ?? this.authorization.loggedInUser()?.username ?? 'assistant',
createdAt: currentFlow?.createdAt ?? new Date(),
@@ -256,45 +322,61 @@ export class FlowAssistant implements OnInit {
};
}
- private toAssistantDraft(flow: Flow | null): AssistantDraftPayload | null {
- if (!flow) return null;
- return {
- name: flow.name,
- description: flow.description,
- flow: flow.data
- };
- }
-
- private defaultAssistantSummary(intent: AssistantIntent, valid: boolean): string {
- if (intent === 'draft') {
- return valid
- ? 'I created a new workflow draft.'
- : 'I created an initial draft, but it still needs corrections.';
- }
- if (intent === 'fix') {
- return valid
- ? 'I fixed the current workflow.'
- : 'I tried to fix the workflow, but there are still unresolved issues.';
- }
- return 'I updated the current workflow based on your request.';
- }
-
- private pushAssistantError(err: unknown, intent: AssistantIntent) {
- console.error('Assistant request failed', err);
- this.pushMessage({
- role: 'assistant',
- content: 'The assistant request failed.',
- intent
- });
- }
-
- private pushMessage(message: Omit) {
- this.messages.update((messages) => [
- ...messages,
+ private pushLocalAssistantMessage(content: string) {
+ const filtered = this.localMessages().filter((message) => message.role !== 'assistant');
+ this.localMessages.set([
+ ...filtered,
{
- ...message,
- id: crypto.randomUUID()
+ id: crypto.randomUUID(),
+ role: 'assistant',
+ content
}
]);
}
+
+ private stopPolling() {
+ if (this.pollTick) {
+ clearInterval(this.pollTick);
+ this.pollTick = null;
+ }
+ }
+
+ private phaseText(phase: AssistantCallPhase): string {
+ switch (phase) {
+ case 'queued':
+ return 'Queued...';
+ case 'routing':
+ return 'Routing your request...';
+ case 'planning':
+ return 'Planning workflow blocks...';
+ case 'configuring_blocks':
+ return 'Configuring blocks...';
+ case 'connecting_blocks':
+ return 'Connecting blocks...';
+ case 'validating':
+ return 'Validating flow...';
+ case 'fixing':
+ return 'Repairing invalid flow...';
+ case 'explaining':
+ return 'Explaining current flow...';
+ case 'completed':
+ return 'Finalizing flow...';
+ case 'failed':
+ return 'Assistant request failed.';
+ }
+ }
+
+ private resolveAssistantModelsUrl(url: string): string {
+ if (!url) return url;
+ if (/^https?:\/\//i.test(url)) return url;
+
+ const apiBase = environment.apiUrl;
+ if (/^https?:\/\//i.test(apiBase)) {
+ return new URL(url, `${apiBase.replace(/\/+$/, '')}/`).toString();
+ }
+
+ const origin = typeof window !== 'undefined' ? window.location.origin : '';
+ const normalizedBase = apiBase.startsWith('/') ? apiBase : `/${apiBase}`;
+ return new URL(url, `${origin}${normalizedBase.replace(/\/+$/, '')}/`).toString();
+ }
}
diff --git a/src/app/shared/nodes/generic-node/generic-node.ts b/src/app/shared/nodes/generic-node/generic-node.ts
index c653718..1bbfe2d 100644
--- a/src/app/shared/nodes/generic-node/generic-node.ts
+++ b/src/app/shared/nodes/generic-node/generic-node.ts
@@ -901,6 +901,13 @@ export class GenericNodeComponent {
queueMicrotask(() => {
try {
this.cdr.detectChanges();
+ requestAnimationFrame(() => {
+ try {
+ this.rendered();
+ } catch {
+ // Node may have been removed while async refresh was running.
+ }
+ });
} catch {
// Node may have been removed while async validation was running.
}
diff --git a/src/app/shared/rete-editor/rete-editor.ts b/src/app/shared/rete-editor/rete-editor.ts
index bc95da3..1ecb17a 100644
--- a/src/app/shared/rete-editor/rete-editor.ts
+++ b/src/app/shared/rete-editor/rete-editor.ts
@@ -166,6 +166,9 @@ export class ReteEditor implements OnChanges, OnDestroy {
...movedNode.data,
position: { x: pos.x, y: pos.y }
};
+
+ // Keep socket anchors and connection paths visually in sync while dragging.
+ void rete.area.update('node', movedNode.id);
}
private markFlowChanged(rete: ReteEditorInstance, context: any, loadedFlowId: string, loadedVersion: number) {
diff --git a/src/app/utilities/rete-editor.ts b/src/app/utilities/rete-editor.ts
index 5ea85f2..f7e1e17 100644
--- a/src/app/utilities/rete-editor.ts
+++ b/src/app/utilities/rete-editor.ts
@@ -135,9 +135,56 @@ export async function addBlockToEditor(
};
const replaceWithCreatedBlock = async (createdBlock: FlowBlock) => {
if (!editor.getNode(node.id)) return;
+ 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
+ }));
const currentPosition = (node.data?.position ?? position ?? createdBlock.position) as { x: number; y: number } | undefined;
+ for (const connection of previousConnections) {
+ await editor.removeConnection(connection.id);
+ }
await editor.removeNode(node.id);
- await addBlockToEditor(editor, area, { ...createdBlock, position: currentPosition }, currentPosition);
+ const replacementNode = await addBlockToEditor(
+ editor,
+ area,
+ { ...createdBlock, position: currentPosition },
+ currentPosition
+ );
+
+ if (!replacementNode) return;
+
+ 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.source === node.id
+ ? connection.sourceOutput
+ : connection.sourceOutput;
+ const targetInput = connection.target === node.id
+ ? connection.targetInput
+ : connection.targetInput;
+
+ if (connection.source === node.id && !replacementOutputNames.has(sourceOutput)) continue;
+ if (connection.target === node.id && !replacementInputNames.has(targetInput)) continue;
+
+ await editor.addConnection(
+ new ClassicPreset.Connection(sourceNode as HFNode, sourceOutput, targetNode as HFNode, targetInput)
+ );
+ }
};
node.data = {
...cloneValue(block),
diff --git a/src/environments/environment.staging.ts b/src/environments/environment.staging.ts
index 9b8f5e9..aeae932 100644
--- a/src/environments/environment.staging.ts
+++ b/src/environments/environment.staging.ts
@@ -8,7 +8,7 @@ import { TaskExecutionsCallService } from "@services/task-executions/task-execut
export const environment = {
production: false,
apiUrl: 'http://localhost:8080',
- assistantEnabled: false,
+ assistantEnabled: true,
assistantCallService: AssistantCallService,
authorizationCallService: AuthorizationCallService,
flowsCallService: FlowsCallService,