Migrate flow assistant chat to the async session call flow

Replace the direct draft/refine/fix/explain HTTP calls (which no longer exist
on the backend) with submitMessage + polling on getCall, wiring the
currentCall/progress-phase UI that was already scaffolded but never
connected. Intent (draft vs refine vs fix vs explain) is now inferred
server-side instead of guessed client-side. Updates the fake service and
specs to match.
This commit is contained in:
Lucio Lelii 2026-09-02 13:32:21 +02:00
parent 3b48e71b82
commit 38b01b2bee
8 changed files with 338 additions and 169 deletions

View File

@ -89,12 +89,14 @@ export type AssistantSessionRequest = {
llmSelection?: AssistantLlmSelection;
};
export type AssistantFlowRequest = {
userPrompt: string;
llmSelection?: AssistantLlmSelection;
export type AssistantSessionMessageRequest = {
message: string;
flow?: AssistantDraftPayload;
validationErrors?: AssistantValidationIssue[];
maxRepairAttempts?: number;
};
export type AssistantCallAccepted = {
sessionId: string;
callId: string;
};
export type AssistantFlowActionResult = {
@ -129,6 +131,7 @@ export type AssistantCallState = {
intent?: AssistantIntent | null;
errorMessage?: string;
flowResult?: AssistantFlowResult | null;
actionResult?: AssistantFlowActionResult | null;
};
export type AssistantSessionState = {

View File

@ -1,7 +1,8 @@
import {
AssistantCallAccepted,
AssistantCallState,
AssistantConfig,
AssistantFlowActionResult,
AssistantFlowRequest,
AssistantSessionMessageRequest,
AssistantSessionRequest,
AssistantSessionState
} from '@models/assistant';
@ -16,12 +17,9 @@ export abstract class AssistantCallServiceBase {
abstract createSession(request: AssistantSessionRequest): Observable<AssistantSessionState>;
abstract draft(request: AssistantFlowRequest): Observable<AssistantFlowActionResult>;
abstract submitMessage(sessionId: string, request: AssistantSessionMessageRequest): Observable<AssistantCallAccepted>;
abstract refine(request: AssistantFlowRequest): Observable<AssistantFlowActionResult>;
abstract fix(request: AssistantFlowRequest): Observable<AssistantFlowActionResult>;
abstract explain(request: AssistantFlowRequest): Observable<AssistantFlowActionResult>;
abstract getCall(callId: string): Observable<AssistantCallState>;
abstract cancelCall(callId: string): Observable<AssistantCallState>;
}

View File

@ -1,3 +1,4 @@
import { firstValueFrom } from 'rxjs';
import { AssistantCallServiceFake } from './assistant-call.fake';
describe('AssistantCallServiceFake', () => {
@ -7,17 +8,31 @@ describe('AssistantCallServiceFake', () => {
service = new AssistantCallServiceFake();
});
it('uses the selected model only when a custom selection is supplied', async () => {
let selectedModel = '';
service.draft({
userPrompt: 'Create a flow',
llmSelection: { provider: 'OpenAI', model: 'custom-model' }
}).subscribe((result) => {
selectedModel = String(
(result.flow?.flow.blocks[1]?.specificConfiguration as Record<string, any>)['llmDescriptor']?.model
);
});
it('drafts a new flow when the session has no flow yet', async () => {
const accepted = await firstValueFrom(service.submitMessage('session-1', { message: 'Create a flow' }));
const call = await firstValueFrom(service.getCall(accepted.callId));
expect(selectedModel).toBe('custom-model');
expect(call.status).toBe('COMPLETED');
expect(call.intent).toBe('draft');
expect(call.actionResult?.flow?.flow.blocks.length).toBeGreaterThan(0);
});
it('refines the attached flow when one is present', async () => {
const draftAccepted = await firstValueFrom(service.submitMessage('session-1', { message: 'Create a flow' }));
const draftCall = await firstValueFrom(service.getCall(draftAccepted.callId));
const flow = draftCall.actionResult!.flow!;
const refineAccepted = await firstValueFrom(service.submitMessage('session-1', { message: 'Add a review step', flow }));
const refineCall = await firstValueFrom(service.getCall(refineAccepted.callId));
expect(refineCall.intent).toBe('refine');
expect(refineCall.actionResult?.flow?.flow.blocks.some((block) => block.id === 'assistant-extra-review')).toBe(true);
});
it('reports a cancelled call', async () => {
const accepted = await firstValueFrom(service.submitMessage('session-1', { message: 'Create a flow' }));
const cancelled = await firstValueFrom(service.cancelCall(accepted.callId));
expect(cancelled.status).toBe('CANCELLED');
});
});

View File

@ -1,8 +1,11 @@
import {
AssistantCallAccepted,
AssistantCallState,
AssistantConfig,
AssistantDraftPayload,
AssistantFlowActionResult,
AssistantFlowRequest,
AssistantIntent,
AssistantSessionMessageRequest,
AssistantSessionRequest,
AssistantSessionState
} from '@models/assistant';
@ -13,6 +16,7 @@ import { AssistantCallServiceBase } from './assistant-call.base';
export class AssistantCallServiceFake extends AssistantCallServiceBase {
private readonly models = ['llama3.1:8b', 'qwen2.5:7b', 'mistral:7b'];
private readonly providers = ['InternalOllama', 'OpenAI'];
private readonly calls = new Map<string, { sessionId: string; intent: AssistantIntent; result: AssistantFlowActionResult }>();
override getConfig(): Observable<AssistantConfig> {
return of({
@ -51,50 +55,105 @@ export class AssistantCallServiceFake extends AssistantCallServiceBase {
return of(structuredClone(session));
}
override draft(request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
const model = request.llmSelection?.model ?? this.models[0];
override submitMessage(sessionId: string, request: AssistantSessionMessageRequest): Observable<AssistantCallAccepted> {
const intent = this.inferIntent(request.message, request.flow);
const result = this.buildActionResult(intent, request);
const callId = crypto.randomUUID();
this.calls.set(callId, { sessionId, intent, result });
return of({ sessionId, callId });
}
override getCall(callId: string): Observable<AssistantCallState> {
const call = this.calls.get(callId);
if (!call) {
return of({
id: callId,
sessionId: '',
status: 'FAILED',
phase: 'failed',
errorMessage: 'Assistant call not found',
intent: null,
flowResult: null,
actionResult: null
});
}
return of({
flow: {
name: 'Ticket classification with urgent review',
description: `Draft generated from prompt: ${request.userPrompt}`,
flow: buildTicketFlow(model)
},
valid: true,
validationErrors: [],
warnings: ['Fake assistant response'],
message: 'I created a new workflow draft.'
id: callId,
sessionId: call.sessionId,
status: 'COMPLETED',
phase: call.intent === 'explain' ? 'explaining' : 'completed',
progressMessage: 'Assistant request completed',
intent: call.intent,
flowResult: call.result.flow,
actionResult: call.result
});
}
override refine(request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
const flow = request.flow ? structuredClone(request.flow) : null;
if (flow) addHumanReviewTail(flow.flow);
override cancelCall(callId: string): Observable<AssistantCallState> {
const call = this.calls.get(callId);
return of({
flow,
valid: true,
validationErrors: [],
warnings: ['Fake assistant response'],
message: 'I updated the current workflow based on your request.'
id: callId,
sessionId: call?.sessionId ?? '',
status: 'CANCELLED',
phase: 'cancelled',
progressMessage: 'Assistant request cancelled',
intent: call?.intent ?? null,
flowResult: null,
actionResult: null
});
}
override fix(request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
return of({
flow: request.flow ? structuredClone(request.flow) : null,
valid: true,
validationErrors: [],
warnings: ['Fake assistant response'],
message: 'I fixed the current workflow.'
});
private inferIntent(message: string, flow: AssistantDraftPayload | undefined): AssistantIntent {
const normalized = message.toLowerCase();
if (!flow) return 'draft';
if (/\b(fix|repair|invalid|error|broken)\b/.test(normalized)) return 'fix';
if (/\b(explain|what does|why|describe)\b/.test(normalized)) return 'explain';
return 'refine';
}
override explain(_request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
return of({
flow: null,
validationErrors: [],
warnings: [],
message: 'This workflow classifies incoming tickets and routes urgent cases to a human reviewer.'
});
private buildActionResult(intent: AssistantIntent, request: AssistantSessionMessageRequest): AssistantFlowActionResult {
switch (intent) {
case 'draft': {
const model = this.models[0];
return {
flow: {
name: 'Ticket classification with urgent review',
description: `Draft generated from prompt: ${request.message}`,
flow: buildTicketFlow(model)
},
valid: true,
validationErrors: [],
warnings: ['Fake assistant response'],
message: 'I created a new workflow draft.'
};
}
case 'refine': {
const flow = request.flow ? structuredClone(request.flow) : null;
if (flow) addHumanReviewTail(flow.flow);
return {
flow,
valid: true,
validationErrors: [],
warnings: ['Fake assistant response'],
message: 'I updated the current workflow based on your request.'
};
}
case 'fix':
return {
flow: request.flow ? structuredClone(request.flow) : null,
valid: true,
validationErrors: [],
warnings: ['Fake assistant response'],
message: 'I fixed the current workflow.'
};
case 'explain':
return {
flow: null,
validationErrors: [],
warnings: [],
message: 'This workflow classifies incoming tickets and routes urgent cases to a human reviewer.'
};
}
}
}

View File

@ -28,45 +28,56 @@ describe('AssistantCallService', () => {
});
it('normalizes node families in assistant drafts and nested container subflows', async () => {
const request = firstValueFrom(service.draft({ userPrompt: 'Create a flow' }));
const request = firstValueFrom(service.getCall('call-1'));
httpMock.expectOne(`${environment.apiUrl}/assistant/flows/draft`).flush({
flow: {
name: 'Loop draft',
httpMock.expectOne(`${environment.apiUrl}/assistant/calls/call-1`).flush({
id: 'call-1',
sessionId: 'session-1',
status: 'COMPLETED',
phase: 'completed',
intent: 'DRAFT',
flowResult: {
flow: {
blocks: [{ id: 'root-block', typeName: 'LLMBlock', specificConfiguration: {} }],
containers: [{
id: 'loop-1',
typeName: 'LoopContainer',
specificConfiguration: {
subFlow: {
blocks: [{ id: 'body-block', typeName: 'LLMBlock', specificConfiguration: {} }],
containers: [{
id: 'nested-container',
typeName: 'GenericContainer',
specificConfiguration: {
subFlow: { blocks: [], containers: [], connections: [], dependencies: [] }
}
}],
connections: [],
dependencies: []
},
guardSubFlow: {
blocks: [{ id: 'guard-block', typeName: 'SwitchBlock', specificConfiguration: {} }],
containers: [],
connections: [],
dependencies: []
name: 'Loop draft',
flow: {
blocks: [{ id: 'root-block', typeName: 'LLMBlock', specificConfiguration: {} }],
containers: [{
id: 'loop-1',
typeName: 'LoopContainer',
specificConfiguration: {
subFlow: {
blocks: [{ id: 'body-block', typeName: 'LLMBlock', specificConfiguration: {} }],
containers: [{
id: 'nested-container',
typeName: 'GenericContainer',
specificConfiguration: {
subFlow: { blocks: [], containers: [], connections: [], dependencies: [] }
}
}],
connections: [],
dependencies: []
},
guardSubFlow: {
blocks: [{ id: 'guard-block', typeName: 'SwitchBlock', specificConfiguration: {} }],
containers: [],
connections: [],
dependencies: []
}
}
}
}],
connections: [],
dependencies: []
}
}],
connections: [],
dependencies: []
}
},
valid: true,
validationErrors: [],
warnings: [],
assistantRationale: 'Done'
}
});
const result = await request;
const flow = result.flow!.flow;
const flow = result.actionResult!.flow!.flow;
const loopConfiguration = flow.containers[0].specificConfiguration as Record<string, any>;
expect(flow.blocks[0].nodeFamily).toBe('block');
@ -94,7 +105,7 @@ describe('AssistantCallService', () => {
await expect(models).resolves.toEqual(['gpt-oss:20b']);
});
it('sends llmSelection only when supplied, for sessions and flow actions', async () => {
it('sends llmSelection only when supplied when creating a session', async () => {
const defaultSession = firstValueFrom(service.createSession({}));
const defaultSessionRequest = httpMock.expectOne(`${environment.apiUrl}/assistant/sessions`);
expect(defaultSessionRequest.request.body).toEqual({});
@ -107,20 +118,52 @@ describe('AssistantCallService', () => {
credentialId: 'credential-1',
phaseModels: { planningModel: 'planning-model' }
};
const actionRequests = [
['draft', service.draft({ userPrompt: 'Create a flow', llmSelection: selection })],
['refine', service.refine({ userPrompt: 'Refine it', flow: { name: 'Flow', flow: emptyFlow() }, llmSelection: selection })],
['fix', service.fix({ userPrompt: 'Fix it', flow: { name: 'Flow', flow: emptyFlow() }, llmSelection: selection })],
['explain', service.explain({ userPrompt: 'Explain it', flow: { name: 'Flow', flow: emptyFlow() }, llmSelection: selection })]
] as const;
const selectionSession = firstValueFrom(service.createSession({ llmSelection: selection }));
const selectionSessionRequest = httpMock.expectOne(`${environment.apiUrl}/assistant/sessions`);
expect(selectionSessionRequest.request.body).toEqual({ llmSelection: selection });
selectionSessionRequest.flush({ id: 'session-selection', messages: [] });
await selectionSession;
});
for (const [action, observable] of actionRequests) {
const result = firstValueFrom(observable);
const request = httpMock.expectOne(`${environment.apiUrl}/assistant/flows/${action}`);
expect(request.request.body.llmSelection).toEqual(selection);
request.flush({ message: 'Done' });
await result;
}
it('submits a session message and polls the call, then cancels it', async () => {
const accepted = firstValueFrom(service.submitMessage('session-1', {
message: 'Create a flow',
flow: { name: 'Flow', flow: emptyFlow() }
}));
const submitRequest = httpMock.expectOne(`${environment.apiUrl}/assistant/sessions/session-1/messages`);
expect(submitRequest.request.body).toEqual({
message: 'Create a flow',
flow: { name: 'Flow', flow: emptyFlow() }
});
submitRequest.flush({ sessionId: 'session-1', callId: 'call-1' });
await expect(accepted).resolves.toEqual({ sessionId: 'session-1', callId: 'call-1' });
const call = firstValueFrom(service.getCall('call-1'));
httpMock.expectOne(`${environment.apiUrl}/assistant/calls/call-1`).flush({
id: 'call-1',
sessionId: 'session-1',
status: 'RUNNING',
phase: 'planning',
progressMessage: 'Planning workflow blocks',
intent: 'DRAFT'
});
const runningCall = await call;
expect(runningCall.status).toBe('RUNNING');
expect(runningCall.phase).toBe('planning');
const cancelled = firstValueFrom(service.cancelCall('call-1'));
const cancelRequest = httpMock.expectOne(`${environment.apiUrl}/assistant/calls/call-1/cancel`);
expect(cancelRequest.request.method).toBe('PUT');
cancelRequest.flush({
id: 'call-1',
sessionId: 'session-1',
status: 'CANCELLED',
phase: 'cancelled',
progressMessage: 'Assistant request cancelled',
intent: 'DRAFT'
});
const cancelledCall = await cancelled;
expect(cancelledCall.status).toBe('CANCELLED');
});
});

View File

@ -1,11 +1,13 @@
import { HttpClient } from '@angular/common/http';
import { inject } from '@angular/core';
import {
AssistantCallAccepted,
AssistantCallState,
AssistantChatMessage,
AssistantConfig,
AssistantDraftPayload,
AssistantFlowActionResult,
AssistantFlowRequest,
AssistantSessionMessageRequest,
AssistantSessionRequest,
AssistantSessionState,
AssistantValidationIssue
@ -43,30 +45,56 @@ export class AssistantCallService extends AssistantCallServiceBase {
.pipe(map((raw) => mapAssistantSessionState(raw)));
}
override draft(request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
return this.runFlowAction('draft', request);
}
override refine(request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
return this.runFlowAction('refine', request);
}
override fix(request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
return this.runFlowAction('fix', request);
}
override explain(request: AssistantFlowRequest): Observable<AssistantFlowActionResult> {
return this.runFlowAction('explain', request);
}
private runFlowAction(
action: 'draft' | 'refine' | 'fix' | 'explain',
request: AssistantFlowRequest
): Observable<AssistantFlowActionResult> {
override submitMessage(sessionId: string, request: AssistantSessionMessageRequest): Observable<AssistantCallAccepted> {
return this.http
.post<unknown>(`${environment.apiUrl}/assistant/flows/${action}`, request)
.pipe(map((raw) => mapAssistantFlowActionResult(raw)));
.post<unknown>(`${environment.apiUrl}/assistant/sessions/${sessionId}/messages`, request)
.pipe(map((raw) => mapAssistantCallAccepted(raw)));
}
override getCall(callId: string): Observable<AssistantCallState> {
return this.http
.get<unknown>(`${environment.apiUrl}/assistant/calls/${callId}`)
.pipe(map((raw) => mapAssistantCallState(raw)));
}
override cancelCall(callId: string): Observable<AssistantCallState> {
return this.http
.put<unknown>(`${environment.apiUrl}/assistant/calls/${callId}/cancel`, {})
.pipe(map((raw) => mapAssistantCallState(raw)));
}
}
function mapAssistantCallAccepted(raw: unknown): AssistantCallAccepted {
const value = (raw ?? {}) as Record<string, unknown>;
return {
sessionId: String(value['sessionId'] ?? ''),
callId: String(value['callId'] ?? '')
};
}
function mapAssistantCallState(raw: unknown): AssistantCallState {
const value = (raw ?? {}) as Record<string, unknown>;
const intent = mapAssistantIntent(value['intent']);
const actionResult = mapAssistantFlowActionResult(value['flowResult'] ?? value['explainResult'] ?? {});
return {
id: String(value['id'] ?? ''),
sessionId: String(value['sessionId'] ?? ''),
status: String(value['status'] ?? 'QUEUED').toUpperCase() as AssistantCallState['status'],
phase: String(value['phase'] ?? 'queued') as AssistantCallState['phase'],
progressMessage: typeof value['progressMessage'] === 'string' ? value['progressMessage'] : undefined,
intent,
errorMessage: typeof value['errorMessage'] === 'string' ? value['errorMessage'] : undefined,
flowResult: actionResult.flow,
actionResult
};
}
function mapAssistantIntent(raw: unknown): AssistantCallState['intent'] {
if (typeof raw !== 'string') return null;
const normalized = raw.toLowerCase();
return normalized === 'draft' || normalized === 'refine' || normalized === 'fix' || normalized === 'explain'
? normalized
: null;
}
function mapAssistantConfig(raw: unknown): AssistantConfig {

View File

@ -1,7 +1,7 @@
import { Injectable } from '@angular/core';
import { environment } from '@environment';
import {
AssistantFlowRequest,
AssistantSessionMessageRequest,
AssistantSessionRequest
} from '@models/assistant';
import { AssistantCallServiceBase } from './assistant-call.base';
@ -28,20 +28,16 @@ export class AssistantService {
return this.assistantCall.createSession(request);
}
draft(request: AssistantFlowRequest) {
return this.assistantCall.draft(request);
submitMessage(sessionId: string, request: AssistantSessionMessageRequest) {
return this.assistantCall.submitMessage(sessionId, request);
}
refine(request: AssistantFlowRequest) {
return this.assistantCall.refine(request);
getCall(callId: string) {
return this.assistantCall.getCall(callId);
}
fix(request: AssistantFlowRequest) {
return this.assistantCall.fix(request);
}
explain(request: AssistantFlowRequest) {
return this.assistantCall.explain(request);
cancelCall(callId: string) {
return this.assistantCall.cancelCall(callId);
}
}

View File

@ -14,6 +14,7 @@ import {
AssistantDraftPayload,
AssistantFlowActionResult,
AssistantLlmSelection,
AssistantSessionMessageRequest,
AssistantSessionState,
VaultSecret
} from '@models/assistant';
@ -448,11 +449,12 @@ export class FlowAssistant implements OnInit, OnDestroy {
return;
}
this.runFlowAction(normalizedContent);
this.submitAssistantMessage(normalizedContent);
}
private runFlowAction(normalizedContent: string) {
const intent = this.resolveIntent(normalizedContent);
private submitAssistantMessage(normalizedContent: string) {
const sessionId = this.sessionState()?.id;
if (!sessionId) return;
if (this.isCreateModal()) this.createPromptSubmitted.set(true);
this.requestPending.set(true);
@ -468,30 +470,19 @@ export class FlowAssistant implements OnInit, OnDestroy {
]);
this.persistSnapshot();
const request = {
userPrompt: normalizedContent,
maxRepairAttempts: 2,
...(this.llmSelection() ? { llmSelection: this.llmSelection() } : {}),
...(intent === 'draft' ? {} : { flow: this.assistantFlowForRequest() }),
...(intent === 'fix' ? { validationErrors: this.sessionState()?.lastValidationErrors ?? [] } : {})
const flow = this.assistantFlowForRequest();
const request: AssistantSessionMessageRequest = {
message: normalizedContent,
...(flow ? { flow } : {})
};
const action = intent === 'draft'
? this.assistant.draft(request)
: intent === 'fix'
? this.assistant.fix(request)
: intent === 'explain'
? this.assistant.explain(request)
: this.assistant.refine(request);
action.pipe(take(1)).subscribe({
next: (result) => {
this.assistant.submitMessage(sessionId, request).pipe(take(1)).subscribe({
next: (accepted) => {
this.requestPending.set(false);
this.applyFlowActionResult(result);
this.createPromptSubmitted.set(false);
this.persistSnapshot();
this.beginPolling(accepted.callId);
},
error: (err) => {
console.error('Assistant flow action failed', err);
console.error('Assistant message submission failed', err);
this.requestPending.set(false);
this.discardFailedSession();
this.handleAssistantErrorWithRetry(normalizedContent, this.backendErrorMessage(err));
@ -499,6 +490,49 @@ export class FlowAssistant implements OnInit, OnDestroy {
});
}
private beginPolling(callId: string) {
this.stopPolling();
this.pollSubscription = interval(1000).pipe(
switchMap(() => this.assistant.getCall(callId))
).subscribe({
next: (call) => {
this.currentCall.set(call);
this.persistSnapshot();
if (call.status === 'COMPLETED') {
this.stopPolling();
this.applyFlowActionResult(call.actionResult ?? {
flow: null,
validationErrors: [],
warnings: [],
message: 'The assistant completed the request.'
});
this.createPromptSubmitted.set(false);
this.persistSnapshot();
return;
}
if (call.status === 'FAILED') {
this.stopPolling();
this.discardFailedSession();
this.handleAssistantErrorWithRetry(this.lastSubmittedPrompt(), call.errorMessage || undefined);
return;
}
if (call.status === 'CANCELLED') {
this.stopPolling();
this.createPromptSubmitted.set(false);
}
},
error: (err) => {
console.error('Assistant call polling failed', err);
this.stopPolling();
this.discardFailedSession();
this.handleAssistantErrorWithRetry(this.lastSubmittedPrompt(), this.backendErrorMessage(err));
}
});
}
private bootstrapAssistant() {
this.sessionLoading.set(true);
this.modelsError.set(null);
@ -598,13 +632,6 @@ export class FlowAssistant implements OnInit, OnDestroy {
return llmSelection ? { llmSelection } : {};
}
private resolveIntent(prompt: string): 'draft' | 'refine' | 'fix' | 'explain' {
const normalized = prompt.toLowerCase();
if (/\b(explain|what does|why|describe)\b/.test(normalized)) return 'explain';
if (/\b(fix|repair|invalid|error|broken)\b/.test(normalized)) return 'fix';
return this.canOfferCreate() ? 'draft' : 'refine';
}
private assistantFlowForRequest(): AssistantDraftPayload | undefined {
const draft = this.currentDraft();
if (draft) return draft;
@ -682,7 +709,7 @@ export class FlowAssistant implements OnInit, OnDestroy {
next: (session) => {
this.applySessionState(session);
this.persistSnapshot(flowKey);
if (promptToSend) this.runFlowAction(promptToSend);
if (promptToSend) this.submitAssistantMessage(promptToSend);
},
error: (err) => {
console.error('Assistant session creation failed', err);