fix(task-execution): render execution graph topology
This commit is contained in:
parent
318531111d
commit
cc1f773c2d
|
|
@ -118,6 +118,8 @@ export type FlowNodeBase = {
|
|||
typeName: BlockTypeName;
|
||||
nodeFamily?: NodeFamily;
|
||||
laneId?: string | null;
|
||||
capabilities?: NodeTypeCapabilities;
|
||||
userInteractive?: boolean;
|
||||
};
|
||||
|
||||
export type BiasActivationMode =
|
||||
|
|
|
|||
|
|
@ -1,4 +1,4 @@
|
|||
import { FlowBlockConnection, FlowNode, FlowNodeDependency, FlowPort, LLMDescriptor } from './flow';
|
||||
import { FlowBlockConnection, FlowData, FlowNode, FlowNodeDependency, FlowPort, LLMDescriptor } from './flow';
|
||||
import { BiasExecutionContext } from './bias-impact';
|
||||
|
||||
export type TaskExecutionStatus = 'CREATED' | 'READY' | 'RUNNING' | 'WAITING' | 'SUSPENDED' | 'SUCCESS' | 'ERROR' | 'CANCELLED';
|
||||
|
|
@ -31,9 +31,10 @@ export type TaskExecution = {
|
|||
interactionSimulationEnabled?: boolean;
|
||||
simulationAvailable?: boolean;
|
||||
interactionSimulationDescriptor?: LLMDescriptor;
|
||||
flowSnapshot?: FlowData;
|
||||
stepConnections?: FlowBlockConnection[];
|
||||
stepDependencies?: FlowNodeDependency[];
|
||||
requiredAuthorizations?: Record<string, TaskExecutionAuthorizationRequirement>;
|
||||
requiredAuthorizations?: Record<string, TaskExecutionAuthorizationRequirement> | TaskExecutionAuthorizationRequirement[];
|
||||
providedAuthorizations?: Record<string, unknown>;
|
||||
missingAuthorizationKeys?: string[];
|
||||
missingGlobalInputKeys?: string[];
|
||||
|
|
@ -68,10 +69,15 @@ export type TaskExecutionContext = {
|
|||
startTime?: number | null;
|
||||
endTime?: number | null;
|
||||
errors: Record<string, string>;
|
||||
warnings: Record<string, string>;
|
||||
warnings: Record<string, string> | unknown[];
|
||||
steps: Record<string, TaskExecutionStep>;
|
||||
status: TaskExecutionStatus;
|
||||
waitingSteps: string[];
|
||||
authorizations?: Record<string, unknown>;
|
||||
executionVariables?: Record<string, unknown>;
|
||||
executionVariableDescriptors?: Record<string, unknown>;
|
||||
errorCodes?: Record<string, unknown>;
|
||||
outcomes?: unknown[];
|
||||
};
|
||||
|
||||
export type TaskExecutionGlobalInputDescriptor = {
|
||||
|
|
@ -86,11 +92,12 @@ export type TaskExecutionGlobalInputDescriptor = {
|
|||
export type TaskExecutionStep = {
|
||||
node?: FlowNode;
|
||||
id: string;
|
||||
inputs: TaskExecutionStepInput[];
|
||||
outputs: TaskExecutionStepOutput[];
|
||||
inputs?: TaskExecutionStepInput[];
|
||||
outputs?: TaskExecutionStepOutput[];
|
||||
result?: Record<string, unknown>;
|
||||
status: StepStatus;
|
||||
started: boolean;
|
||||
started?: boolean;
|
||||
skipReason?: string | null;
|
||||
simulated: boolean;
|
||||
};
|
||||
|
||||
|
|
|
|||
|
|
@ -177,6 +177,7 @@ export class BlocksCallService extends BlocksCallServiceBase {
|
|||
specificConfiguration,
|
||||
typeName,
|
||||
nodeFamily: 'block',
|
||||
...(value["capabilities"] == null ? {} : { capabilities: toNodeCapabilities(value["capabilities"]) }),
|
||||
biasAnnotations: Array.isArray(value["biasAnnotations"])
|
||||
? value["biasAnnotations"] as FlowBlock["biasAnnotations"]
|
||||
: []
|
||||
|
|
|
|||
|
|
@ -149,6 +149,7 @@ export class ContainersCallService extends ContainersCallServiceBase {
|
|||
specificConfiguration,
|
||||
typeName,
|
||||
nodeFamily: 'container',
|
||||
...(value["capabilities"] == null ? {} : { capabilities: toNodeCapabilities(value["capabilities"]) }),
|
||||
biasAnnotations: Array.isArray(value["biasAnnotations"])
|
||||
? value["biasAnnotations"] as FlowContainer["biasAnnotations"]
|
||||
: []
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ import { Injectable, signal } from '@angular/core';
|
|||
import { environment } from '@environment';
|
||||
import { Flow } from '@models/flow';
|
||||
import { FlowsCallServiceBase } from './flows-call.base';
|
||||
import { catchError, firstValueFrom, Observable, tap, throwError } from 'rxjs';
|
||||
import { catchError, firstValueFrom, Observable, of, tap, throwError } from 'rxjs';
|
||||
|
||||
@Injectable({
|
||||
providedIn: 'root',
|
||||
|
|
@ -34,6 +34,23 @@ export class FlowsService {
|
|||
return this.flows;
|
||||
}
|
||||
|
||||
getFlowById(flowId: string): Observable<Flow> {
|
||||
const cached = this._flows().find((flow) => flow.id === flowId);
|
||||
if (cached) return of(cached);
|
||||
|
||||
return this.flowsCallService.getFlowById(flowId).pipe(
|
||||
tap((flow) => {
|
||||
this._flows.update((flows) => {
|
||||
const index = flows.findIndex((candidate) => candidate.id === flow.id);
|
||||
if (index < 0) return [flow, ...flows];
|
||||
const next = [...flows];
|
||||
next[index] = flow;
|
||||
return next;
|
||||
});
|
||||
})
|
||||
);
|
||||
}
|
||||
|
||||
async refresh(force = false): Promise<void> {
|
||||
if (this.loadingPromise && !force) {
|
||||
return this.loadingPromise;
|
||||
|
|
|
|||
|
|
@ -45,6 +45,133 @@ describe('TaskExecutionsCallService bias APIs', () => {
|
|||
|
||||
afterEach(() => httpMock.verify());
|
||||
|
||||
it('maps the execution flow snapshot, branched topology, bias annotations and node capabilities', async () => {
|
||||
const result = firstValueFrom(service.retrieveAllTaskExecutions());
|
||||
const request = httpMock.expectOne(`${environment.apiUrl}/executions`);
|
||||
request.flush([{
|
||||
id: 'execution-1',
|
||||
name: 'Branched flow',
|
||||
creationTime: 1,
|
||||
flowId: 'flow-1',
|
||||
context: {
|
||||
inputs: {},
|
||||
result: {},
|
||||
errors: {},
|
||||
warnings: {},
|
||||
waitingSteps: [],
|
||||
status: 'READY',
|
||||
steps: {},
|
||||
connections: [
|
||||
{
|
||||
id: 'branch-a',
|
||||
sourceNodeId: 'decision',
|
||||
sourceOutput: 'accepted',
|
||||
targetNodeId: 'accepted-step',
|
||||
targetInput: 'input'
|
||||
}
|
||||
]
|
||||
},
|
||||
flowSnapshot: {
|
||||
blocks: [{
|
||||
id: 'decision',
|
||||
name: 'Decision',
|
||||
inputs: [{ name: 'input', type: 'ANY', multiple: false }],
|
||||
outputs: [{ name: 'accepted', type: 'ANY', multiple: false }],
|
||||
specificConfiguration: {},
|
||||
typeName: 'HumanDecisionBlock',
|
||||
biasAnnotations: [{ id: 'bias-1', category: 'SELECTION_BIAS' }],
|
||||
capabilities: {
|
||||
visualRole: 'DECISION',
|
||||
terminal: false,
|
||||
biasAnnotationsAllowed: true,
|
||||
allowsIncomingConnections: true,
|
||||
allowsOutgoingConnections: true,
|
||||
canDependOnOtherNodes: true,
|
||||
canHaveDependentNodes: true
|
||||
}
|
||||
}],
|
||||
containers: [],
|
||||
connections: [{
|
||||
id: 'branch-a',
|
||||
sourceId: 'decision',
|
||||
sourceName: 'accepted',
|
||||
targetId: 'accepted-step',
|
||||
targetName: 'input'
|
||||
}],
|
||||
dependencies: []
|
||||
}
|
||||
}]);
|
||||
|
||||
const execution = await result;
|
||||
expect(execution[0].flowSnapshot?.connections[0]).toEqual({
|
||||
id: 'branch-a',
|
||||
sourceId: 'decision',
|
||||
sourceName: 'accepted',
|
||||
targetId: 'accepted-step',
|
||||
targetName: 'input'
|
||||
});
|
||||
expect(execution[0].flowSnapshot?.blocks[0].biasAnnotations).toEqual([
|
||||
{ id: 'bias-1', category: 'SELECTION_BIAS' }
|
||||
]);
|
||||
expect(execution[0].flowSnapshot?.blocks[0].capabilities?.visualRole).toBe('DECISION');
|
||||
expect(execution[0].stepConnections?.[0].sourceName).toBe('accepted');
|
||||
});
|
||||
|
||||
it('keeps the documented root topology and execution node metadata', async () => {
|
||||
const result = firstValueFrom(service.retrieveAllTaskExecutions());
|
||||
const request = httpMock.expectOne(`${environment.apiUrl}/executions`);
|
||||
request.flush([{
|
||||
id: 'execution-id',
|
||||
name: 'test biased',
|
||||
sourceFlowId: 'flow-id',
|
||||
creationTime: 1784795504244,
|
||||
requiredAuthorizations: [],
|
||||
context: {
|
||||
inputs: {},
|
||||
result: {},
|
||||
errors: {},
|
||||
warnings: [],
|
||||
waitingSteps: [],
|
||||
status: 'CREATED',
|
||||
steps: {
|
||||
decision: {
|
||||
id: 'decision',
|
||||
status: 'WAITING_FOR_INPUT',
|
||||
skipReason: null,
|
||||
simulated: false,
|
||||
node: {
|
||||
id: 'decision',
|
||||
name: 'shortlist-decision',
|
||||
typeName: 'HumanDecisionBlock',
|
||||
userInteractive: true,
|
||||
position: { x: 600, y: 160 },
|
||||
inputs: [{ name: 'input', type: 'TEXT', multiple: false }],
|
||||
outputs: [
|
||||
{ name: 'approve', type: 'TEXT', multiple: false },
|
||||
{ name: 'reject', type: 'TEXT', multiple: false }
|
||||
],
|
||||
biasAnnotations: [{ id: 'selection-risk', category: 'SELECTION_BIAS' }],
|
||||
specificConfiguration: { name: 'shortlist-decision' }
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
stepConnections: [
|
||||
{ id: 'approve-edge', sourceId: 'decision', sourceName: 'approve', targetId: 'approved', targetName: 'input' },
|
||||
{ id: 'reject-edge', sourceId: 'decision', sourceName: 'reject', targetId: 'rejected', targetName: 'input' }
|
||||
],
|
||||
stepDependencies: []
|
||||
}]);
|
||||
|
||||
const [execution] = await result;
|
||||
const decision = execution.context.steps['decision'];
|
||||
expect(execution.stepConnections?.map((connection) => connection.sourceName)).toEqual(['approve', 'reject']);
|
||||
expect(execution.stepDependencies).toEqual([]);
|
||||
expect(decision.node?.position).toEqual({ x: 600, y: 160 });
|
||||
expect(decision.node?.userInteractive).toBe(true);
|
||||
expect(decision.node?.biasAnnotations?.[0].id).toBe('selection-risk');
|
||||
});
|
||||
|
||||
it('starts an asynchronous impact experiment and maps the job response', async () => {
|
||||
const result = firstValueFrom(service.runBiasImpactExperiment('execution-1', 'step-1', {
|
||||
annotationIds: ['annotation-1'],
|
||||
|
|
|
|||
|
|
@ -1,7 +1,7 @@
|
|||
import { HttpClient, HttpParams } from '@angular/common/http';
|
||||
import { inject } from '@angular/core';
|
||||
import { environment } from '@environment';
|
||||
import { LLMDescriptor } from '@models/flow';
|
||||
import { FlowBlockConnection, FlowData, FlowNodeDependency, LLMDescriptor } from '@models/flow';
|
||||
import {
|
||||
BiasDownstreamImpactEntry,
|
||||
BiasImpactExperimentRequest,
|
||||
|
|
@ -247,9 +247,37 @@ export class TaskExecutionsCallService extends TaskExecutionsCallServiceBase {
|
|||
context['globalInputDescriptors']
|
||||
?? execution['globalInputDescriptors']
|
||||
);
|
||||
const flowSnapshot = this.normalizeFlowSnapshot(
|
||||
execution['flowSnapshot']
|
||||
?? execution['flowData']
|
||||
?? context['flowSnapshot']
|
||||
?? context['flowData']
|
||||
?? (execution['flow'] && typeof execution['flow'] === 'object' ? execution['flow'] : null)
|
||||
);
|
||||
const rawStepConnections =
|
||||
execution['stepConnections']
|
||||
?? execution['connections']
|
||||
?? context['stepConnections']
|
||||
?? context['connections']
|
||||
?? flowSnapshot?.connections;
|
||||
const stepConnections = rawStepConnections === undefined
|
||||
? undefined
|
||||
: this.normalizeStepConnections(rawStepConnections);
|
||||
const rawStepDependencies =
|
||||
execution['stepDependencies']
|
||||
?? execution['dependencies']
|
||||
?? context['stepDependencies']
|
||||
?? context['dependencies']
|
||||
?? flowSnapshot?.dependencies;
|
||||
const stepDependencies = rawStepDependencies === undefined
|
||||
? undefined
|
||||
: this.normalizeStepDependencies(rawStepDependencies);
|
||||
|
||||
return {
|
||||
...execution,
|
||||
flowSnapshot,
|
||||
stepConnections,
|
||||
stepDependencies,
|
||||
context: {
|
||||
...(context as TaskExecution['context']),
|
||||
globalInputs,
|
||||
|
|
@ -258,6 +286,64 @@ export class TaskExecutionsCallService extends TaskExecutionsCallServiceBase {
|
|||
};
|
||||
}
|
||||
|
||||
private normalizeFlowSnapshot(raw: unknown): FlowData | undefined {
|
||||
if (!raw || typeof raw !== 'object' || Array.isArray(raw)) return undefined;
|
||||
const root = raw as Record<string, unknown>;
|
||||
const nested = root['flow'] && typeof root['flow'] === 'object' && !Array.isArray(root['flow'])
|
||||
? root['flow'] as Record<string, unknown>
|
||||
: root;
|
||||
if (!Array.isArray(nested['blocks']) && !Array.isArray(nested['containers'])) return undefined;
|
||||
const blocks = Array.isArray(nested['blocks'])
|
||||
? nested['blocks'].filter((node): node is FlowData['blocks'][number] => !!node && typeof node === 'object')
|
||||
.map((node) => ({ ...node, nodeFamily: 'block' as const }))
|
||||
: [];
|
||||
const containers = Array.isArray(nested['containers'])
|
||||
? nested['containers'].filter((node): node is FlowData['containers'][number] => !!node && typeof node === 'object')
|
||||
.map((node) => ({ ...node, nodeFamily: 'container' as const }))
|
||||
: [];
|
||||
|
||||
return {
|
||||
blocks,
|
||||
containers,
|
||||
connections: this.normalizeStepConnections(nested['connections']),
|
||||
dependencies: this.normalizeStepDependencies(nested['dependencies']),
|
||||
globalInputs: Array.isArray(nested['globalInputs']) ? nested['globalInputs'] as FlowData['globalInputs'] : [],
|
||||
lanes: Array.isArray(nested['lanes']) ? nested['lanes'] as FlowData['lanes'] : []
|
||||
};
|
||||
}
|
||||
|
||||
private normalizeStepConnections(raw: unknown): FlowBlockConnection[] {
|
||||
if (!Array.isArray(raw)) return [];
|
||||
return raw.flatMap((item, index) => {
|
||||
if (!item || typeof item !== 'object' || Array.isArray(item)) return [];
|
||||
const value = item as Record<string, unknown>;
|
||||
const sourceId = value['sourceId'] ?? value['sourceNodeId'] ?? value['sourceBlockId'];
|
||||
const sourceName = value['sourceName'] ?? value['sourceOutput'] ?? value['outputName'];
|
||||
const targetId = value['targetId'] ?? value['targetNodeId'] ?? value['targetBlockId'];
|
||||
const targetName = value['targetName'] ?? value['targetInput'] ?? value['inputName'];
|
||||
if ([sourceId, sourceName, targetId, targetName].some((part) => typeof part !== 'string' || !part)) return [];
|
||||
return [{
|
||||
id: String(value['id'] ?? `${sourceId}:${sourceName}->${targetId}:${targetName}:${index}`),
|
||||
sourceId: String(sourceId),
|
||||
sourceName: String(sourceName),
|
||||
targetId: String(targetId),
|
||||
targetName: String(targetName)
|
||||
}];
|
||||
});
|
||||
}
|
||||
|
||||
private normalizeStepDependencies(raw: unknown): FlowNodeDependency[] {
|
||||
if (!Array.isArray(raw)) return [];
|
||||
return raw.flatMap((item) => {
|
||||
if (!item || typeof item !== 'object' || Array.isArray(item)) return [];
|
||||
const value = item as Record<string, unknown>;
|
||||
const sourceId = value['sourceId'] ?? value['sourceNodeId'];
|
||||
const targetId = value['targetId'] ?? value['targetNodeId'];
|
||||
if (typeof sourceId !== 'string' || !sourceId || typeof targetId !== 'string' || !targetId) return [];
|
||||
return [{ sourceId, targetId }];
|
||||
});
|
||||
}
|
||||
|
||||
private biasImpactJobFromApi(raw: unknown): BiasImpactJob {
|
||||
const value = this.toRecord(raw);
|
||||
const status = this.toBiasImpactJobStatus(value['status']);
|
||||
|
|
|
|||
|
|
@ -71,6 +71,55 @@
|
|||
box-shadow: 0 0 0 3px rgba(180, 83, 9, 0.22), 0 10px 24px rgba(15, 23, 42, 0.12);
|
||||
}
|
||||
|
||||
.llm-node-metadata {
|
||||
display: flex;
|
||||
flex-wrap: wrap;
|
||||
align-items: center;
|
||||
gap: 4px;
|
||||
margin-top: 5px;
|
||||
}
|
||||
|
||||
.llm-node-capability-badge,
|
||||
.llm-node-bias-summary,
|
||||
.llm-node-bias-capability,
|
||||
.llm-node-skip-reason {
|
||||
display: inline-flex;
|
||||
align-items: center;
|
||||
gap: 3px;
|
||||
min-height: 18px;
|
||||
padding: 2px 6px;
|
||||
border: 1px solid rgba(148, 163, 184, 0.72);
|
||||
border-radius: 999px;
|
||||
background: rgba(255, 255, 255, 0.9);
|
||||
color: #475569;
|
||||
font-size: 9px;
|
||||
font-weight: 800;
|
||||
line-height: 1;
|
||||
white-space: nowrap;
|
||||
}
|
||||
|
||||
.llm-node-bias-summary,
|
||||
.llm-node-bias-capability {
|
||||
border-color: #c4b5fd;
|
||||
color: #6d28d9;
|
||||
background: #f5f3ff;
|
||||
}
|
||||
|
||||
.llm-node-bias-summary-active {
|
||||
border-color: #7c3aed;
|
||||
background: #7c3aed;
|
||||
color: #fff;
|
||||
}
|
||||
|
||||
.llm-node-skip-reason {
|
||||
max-width: 190px;
|
||||
overflow: hidden;
|
||||
border-color: #cbd5e1;
|
||||
background: #f1f5f9;
|
||||
color: #475569;
|
||||
text-overflow: ellipsis;
|
||||
}
|
||||
|
||||
.llm-bias-canvas-badge-wrap {
|
||||
position: absolute;
|
||||
top: -10px;
|
||||
|
|
|
|||
|
|
@ -74,6 +74,30 @@
|
|||
<div class="llm-subtitle-row">
|
||||
<span class="llm-subtitle">{{ name }}</span>
|
||||
</div>
|
||||
<div class="llm-node-metadata">
|
||||
<span class="llm-node-capability-badge" [title]="capabilitiesTooltip()">
|
||||
{{ visualRoleLabel() }}
|
||||
</span>
|
||||
@if (allBiasAnnotations().length) {
|
||||
<span
|
||||
class="llm-node-bias-summary"
|
||||
[class.llm-node-bias-summary-active]="activeBiasAnnotationCount() > 0"
|
||||
[title]="allBiasAnnotations().length + ' bias annotations; ' + activeBiasAnnotationCount() + ' active in this execution'">
|
||||
<i class="bi bi-clipboard2-pulse-fill"></i>
|
||||
{{ activeBiasAnnotationCount() }}/{{ allBiasAnnotations().length }}
|
||||
</span>
|
||||
}
|
||||
@if (biasCapabilities?.supported) {
|
||||
<span class="llm-node-bias-capability" title="Bias impact capabilities available">
|
||||
Bias capable
|
||||
</span>
|
||||
}
|
||||
@if (stepSkipReason(); as skipReason) {
|
||||
<span class="llm-node-skip-reason" [title]="skipReason">
|
||||
Skipped: {{ skipReason }}
|
||||
</span>
|
||||
}
|
||||
</div>
|
||||
</div>
|
||||
@if (hasMeasurableBiasAnnotations()) {
|
||||
<button
|
||||
|
|
|
|||
|
|
@ -101,6 +101,27 @@ describe('TaskStepNodeComponent bias canvas highlighting', () => {
|
|||
expect(component.isBiasRoutingChangeSource()).toBe(true);
|
||||
});
|
||||
|
||||
it('represents node capabilities and bias annotation counts in the execution node', () => {
|
||||
component.data.data.capabilities = {
|
||||
visualRole: 'DECISION',
|
||||
terminal: false,
|
||||
biasAnnotationsAllowed: true,
|
||||
allowsIncomingConnections: true,
|
||||
allowsOutgoingConnections: true,
|
||||
canDependOnOtherNodes: false,
|
||||
canHaveDependentNodes: false
|
||||
};
|
||||
component.data.data.biasAnnotations = [
|
||||
{ id: 'annotation-1', behavioralProbe: { activationMode: 'PROMPT_DIRECTIVE', instruction: 'Nudge it' } },
|
||||
{ id: 'annotation-2' }
|
||||
];
|
||||
|
||||
expect(component.visualRoleLabel()).toBe('Decision');
|
||||
expect(component.allBiasAnnotations()).toHaveLength(2);
|
||||
expect(component.activeBiasAnnotationCount()).toBe(1);
|
||||
expect(component.capabilitiesTooltip()).toContain('Bias annotations: allowed');
|
||||
});
|
||||
|
||||
describe('Measure bias impact availability', () => {
|
||||
beforeEach(() => {
|
||||
component.data.data.biasAnnotations = [
|
||||
|
|
|
|||
|
|
@ -2,7 +2,7 @@ import { CommonModule } from '@angular/common';
|
|||
import { ChangeDetectionStrategy, ChangeDetectorRef, Component, HostBinding, Input, inject } from '@angular/core';
|
||||
import { ClassicPreset } from 'rete';
|
||||
import { ReteModule } from 'rete-angular-plugin/21';
|
||||
import { BiasAnnotation, BlockInteractionContract, BlockType, FlowData, FlowPort, isProbeExecutable, FLOW_DEPENDANT_PORT_KEY, FLOW_DEPENDENCY_PORT_KEY } from '@models/flow';
|
||||
import { BiasAnnotation, BlockInteractionContract, BlockType, DEFAULT_NODE_CAPABILITIES, FlowData, FlowPort, isProbeExecutable, FLOW_DEPENDANT_PORT_KEY, FLOW_DEPENDENCY_PORT_KEY, NodeTypeCapabilities } from '@models/flow';
|
||||
import { BiasCapabilities } from '@models/bias-impact';
|
||||
import { BlocksService } from '@services/blocks/blocks';
|
||||
import { ContainersService } from '@services/containers/containers';
|
||||
|
|
@ -131,10 +131,6 @@ export class TaskStepNodeComponent {
|
|||
|
||||
@HostBinding('class.llm-node-readonly') readonlyClass = true;
|
||||
|
||||
outputs: { key: string; socket: ClassicPreset.Socket }[] = [];
|
||||
inputs: { key: string; socket: ClassicPreset.Socket }[] = [];
|
||||
dependantOutput: { key: string; socket: ClassicPreset.Socket } | null = null;
|
||||
dependencyInput: { key: string; socket: ClassicPreset.Socket } | null = null;
|
||||
parameterFields: DisplayField[] = [];
|
||||
parameterFieldGroups: DisplayFieldGroup[] = [];
|
||||
arrayFields: ArrayFieldView[] = [];
|
||||
|
|
@ -150,31 +146,37 @@ export class TaskStepNodeComponent {
|
|||
private arrayFieldDefinitions: ArrayFieldDefinition[] = [];
|
||||
private flowFieldDefinitions: SchemaFlowDataFieldDefinition[] = [];
|
||||
|
||||
get outputs(): { key: string; socket: ClassicPreset.Socket }[] {
|
||||
return Object.entries(this.data?.outputs ?? {})
|
||||
.filter(([key]) => key !== FLOW_DEPENDANT_PORT_KEY)
|
||||
.map(([key, output]) => ({ key, socket: (output as any).socket as ClassicPreset.Socket }));
|
||||
}
|
||||
|
||||
get inputs(): { key: string; socket: ClassicPreset.Socket }[] {
|
||||
return Object.entries(this.data?.inputs ?? {})
|
||||
.filter(([key]) => key !== FLOW_DEPENDENCY_PORT_KEY)
|
||||
.map(([key, input]) => ({ key, socket: (input as any).socket as ClassicPreset.Socket }));
|
||||
}
|
||||
|
||||
get dependantOutput(): { key: string; socket: ClassicPreset.Socket } | null {
|
||||
const output = this.data?.outputs?.[FLOW_DEPENDANT_PORT_KEY];
|
||||
return output
|
||||
? { key: FLOW_DEPENDANT_PORT_KEY, socket: (output as any).socket as ClassicPreset.Socket }
|
||||
: null;
|
||||
}
|
||||
|
||||
get dependencyInput(): { key: string; socket: ClassicPreset.Socket } | null {
|
||||
const input = this.data?.inputs?.[FLOW_DEPENDENCY_PORT_KEY];
|
||||
return input
|
||||
? { key: FLOW_DEPENDENCY_PORT_KEY, socket: (input as any).socket as ClassicPreset.Socket }
|
||||
: null;
|
||||
}
|
||||
|
||||
ngOnInit() {
|
||||
this.outputs = [];
|
||||
this.inputs = [];
|
||||
this.parameterFields = [];
|
||||
this.parameterFieldGroups = [];
|
||||
this.arrayFields = [];
|
||||
|
||||
Object.entries(this.data.outputs).forEach(([key, output]) => {
|
||||
const entry = { key, socket: (output as any).socket };
|
||||
if (key === FLOW_DEPENDANT_PORT_KEY) {
|
||||
this.dependantOutput = entry;
|
||||
return;
|
||||
}
|
||||
this.outputs.push(entry);
|
||||
});
|
||||
|
||||
Object.entries(this.data.inputs).forEach(([key, input]) => {
|
||||
const entry = { key, socket: (input as any).socket };
|
||||
if (key === FLOW_DEPENDENCY_PORT_KEY) {
|
||||
this.dependencyInput = entry;
|
||||
return;
|
||||
}
|
||||
this.inputs.push(entry);
|
||||
});
|
||||
|
||||
this.rebuildDisplayState();
|
||||
void this.loadSchemaContext();
|
||||
this.loadBiasCapabilities();
|
||||
|
|
@ -432,14 +434,48 @@ export class TaskStepNodeComponent {
|
|||
}
|
||||
|
||||
executableBiasAnnotations(): BiasAnnotation[] {
|
||||
return this.allBiasAnnotations()
|
||||
.filter((annotation) => isProbeExecutable(annotation.behavioralProbe));
|
||||
}
|
||||
|
||||
allBiasAnnotations(): BiasAnnotation[] {
|
||||
const annotations = this.data?.data?.biasAnnotations;
|
||||
return Array.isArray(annotations)
|
||||
? annotations.filter((annotation): annotation is BiasAnnotation => !!annotation && isProbeExecutable(annotation.behavioralProbe))
|
||||
? annotations.filter((annotation): annotation is BiasAnnotation => !!annotation)
|
||||
: [];
|
||||
}
|
||||
|
||||
activeBiasAnnotationCount(): number {
|
||||
const ids = this.blockConfiguration?.['__biasActiveAnnotationIds'];
|
||||
return Array.isArray(ids) ? ids.length : 0;
|
||||
}
|
||||
|
||||
typeCapabilities(): NodeTypeCapabilities {
|
||||
return this.data?.data?.capabilities
|
||||
?? this.blockDescriptor?.capabilities
|
||||
?? DEFAULT_NODE_CAPABILITIES;
|
||||
}
|
||||
|
||||
visualRoleLabel(): string {
|
||||
const role = this.typeCapabilities().visualRole.toLowerCase();
|
||||
return role.charAt(0).toUpperCase() + role.slice(1);
|
||||
}
|
||||
|
||||
capabilitiesTooltip(): string {
|
||||
const capabilities = this.typeCapabilities();
|
||||
return [
|
||||
`Role: ${this.visualRoleLabel()}`,
|
||||
`Terminal: ${capabilities.terminal ? 'yes' : 'no'}`,
|
||||
`Incoming connections: ${capabilities.allowsIncomingConnections ? 'allowed' : 'blocked'}`,
|
||||
`Outgoing connections: ${capabilities.allowsOutgoingConnections ? 'allowed' : 'blocked'}`,
|
||||
`Dependencies: ${capabilities.canDependOnOtherNodes || capabilities.canHaveDependentNodes ? 'supported' : 'blocked'}`,
|
||||
`Bias annotations: ${capabilities.biasAnnotationsAllowed ? 'allowed' : 'blocked'}`
|
||||
].join('\n');
|
||||
}
|
||||
|
||||
hasMeasurableBiasAnnotations(): boolean {
|
||||
return this.executableBiasAnnotations().length > 0
|
||||
return this.typeCapabilities().biasAnnotationsAllowed
|
||||
&& this.executableBiasAnnotations().length > 0
|
||||
&& this.biasCapabilities?.isolatedExperimentSupported === true;
|
||||
}
|
||||
|
||||
|
|
@ -497,6 +533,11 @@ export class TaskStepNodeComponent {
|
|||
return typeof status === 'string' ? status.toUpperCase() : '';
|
||||
}
|
||||
|
||||
stepSkipReason(): string | null {
|
||||
const reason = this.blockConfiguration?.['__stepSkipReason'];
|
||||
return typeof reason === 'string' && reason.trim().length > 0 ? reason.trim() : null;
|
||||
}
|
||||
|
||||
async openInteractionModal(event?: Event) {
|
||||
event?.preventDefault();
|
||||
event?.stopPropagation();
|
||||
|
|
@ -548,8 +589,11 @@ export class TaskStepNodeComponent {
|
|||
|
||||
private loadBiasCapabilities() {
|
||||
const blockType = this.blockType;
|
||||
if (!blockType || this.isContainerNode()) return;
|
||||
this.blocksService.retrieveBiasCapabilities(blockType).pipe(take(1)).subscribe({
|
||||
if (!blockType) return;
|
||||
const capabilities$ = this.isContainerNode()
|
||||
? this.containersService.retrieveBiasCapabilities(blockType)
|
||||
: this.blocksService.retrieveBiasCapabilities(blockType);
|
||||
capabilities$.pipe(take(1)).subscribe({
|
||||
next: (capabilities) => {
|
||||
this.biasCapabilities = capabilities;
|
||||
this.cdr.markForCheck();
|
||||
|
|
|
|||
|
|
@ -0,0 +1,147 @@
|
|||
import { BiasAnnotation, FlowBlock, FlowData, FlowNode } from '@models/flow';
|
||||
import { TaskExecution, TaskExecutionStep } from '@models/task-execution';
|
||||
import {
|
||||
mergeExecutionStepNode,
|
||||
resolveExecutionConnections,
|
||||
resolveExecutionDependencies
|
||||
} from './execution-graph';
|
||||
|
||||
function node(
|
||||
id: string,
|
||||
position: { x: number; y: number },
|
||||
biasAnnotations: BiasAnnotation[] = []
|
||||
): FlowBlock {
|
||||
return {
|
||||
id,
|
||||
name: id,
|
||||
position,
|
||||
inputs: [{ name: 'input', type: 'ANY', multiple: false }],
|
||||
outputs: [{ name: 'left', type: 'ANY', multiple: false }, { name: 'right', type: 'ANY', multiple: false }],
|
||||
specificConfiguration: { name: id },
|
||||
typeName: 'HumanDecisionBlock',
|
||||
nodeFamily: 'block',
|
||||
biasAnnotations
|
||||
};
|
||||
}
|
||||
|
||||
function step(flowNode: FlowNode): TaskExecutionStep {
|
||||
return {
|
||||
id: flowNode.id,
|
||||
node: { ...flowNode, position: undefined, biasAnnotations: undefined },
|
||||
inputs: [],
|
||||
outputs: [],
|
||||
status: 'READY',
|
||||
started: false,
|
||||
simulated: false
|
||||
};
|
||||
}
|
||||
|
||||
const sourceFlow: FlowData = {
|
||||
blocks: [
|
||||
node('decision', { x: 100, y: 200 }, [{ id: 'bias-1', category: 'SELECTION_BIAS' }]),
|
||||
node('left-target', { x: 500, y: 80 }),
|
||||
node('right-target', { x: 500, y: 320 })
|
||||
],
|
||||
containers: [],
|
||||
connections: [
|
||||
{ id: 'c-left', sourceId: 'decision', sourceName: 'left', targetId: 'left-target', targetName: 'input' },
|
||||
{ id: 'c-right', sourceId: 'decision', sourceName: 'right', targetId: 'right-target', targetName: 'input' }
|
||||
],
|
||||
dependencies: [{ sourceId: 'left-target', targetId: 'right-target' }]
|
||||
};
|
||||
|
||||
describe('execution graph topology', () => {
|
||||
const steps = sourceFlow.blocks.map(step);
|
||||
const execution: TaskExecution = {
|
||||
id: 'execution-1',
|
||||
name: 'Execution',
|
||||
creationTime: 1,
|
||||
stepConnections: sourceFlow.connections,
|
||||
stepDependencies: sourceFlow.dependencies,
|
||||
context: {
|
||||
inputs: {},
|
||||
result: {},
|
||||
errors: {},
|
||||
warnings: {},
|
||||
steps: {},
|
||||
status: 'READY',
|
||||
waitingSteps: []
|
||||
}
|
||||
};
|
||||
|
||||
it('uses the explicit execution branches instead of a linear inferred fallback', () => {
|
||||
const inferred = [
|
||||
{ id: 'linear-1', sourceId: 'decision', sourceName: 'left', targetId: 'left-target', targetName: 'input' },
|
||||
{ id: 'linear-2', sourceId: 'left-target', sourceName: 'left', targetId: 'right-target', targetName: 'input' }
|
||||
];
|
||||
|
||||
expect(resolveExecutionConnections(execution, steps, sourceFlow, inferred)).toEqual(sourceFlow.connections);
|
||||
});
|
||||
|
||||
it('uses source metadata only when it is missing from the execution step snapshot', () => {
|
||||
const merged = mergeExecutionStepNode(steps[0], sourceFlow);
|
||||
|
||||
expect(merged?.position).toEqual({ x: 100, y: 200 });
|
||||
expect(merged?.biasAnnotations).toEqual([{ id: 'bias-1', category: 'SELECTION_BIAS' }]);
|
||||
expect(merged?.outputs.map((port) => port.name)).toEqual(['left', 'right']);
|
||||
});
|
||||
|
||||
it('uses explicit execution dependencies', () => {
|
||||
expect(resolveExecutionDependencies(execution, steps, sourceFlow)).toEqual(sourceFlow.dependencies);
|
||||
});
|
||||
|
||||
it('uses a source-flow fallback for legacy executions with no explicit topology', () => {
|
||||
const partialSteps = steps.filter((item) => item.id !== 'right-target');
|
||||
const legacyExecution = { ...execution, stepConnections: undefined };
|
||||
|
||||
expect(resolveExecutionConnections(legacyExecution, partialSteps, sourceFlow, [])).toEqual(sourceFlow.connections);
|
||||
});
|
||||
|
||||
it('does not let a current source flow override the execution snapshot topology', () => {
|
||||
const explicitExecution = {
|
||||
...execution,
|
||||
stepConnections: [
|
||||
{ id: 'executed', sourceId: 'decision', sourceName: 'right', targetId: 'right-target', targetName: 'input' }
|
||||
]
|
||||
};
|
||||
|
||||
expect(resolveExecutionConnections(explicitExecution, steps, sourceFlow, [])).toEqual(explicitExecution.stepConnections);
|
||||
});
|
||||
|
||||
it('preserves positions supplied by the execution even if the source flow has moved', () => {
|
||||
const executionPosition = { x: 640, y: 180 };
|
||||
const executionStep = {
|
||||
...steps[0],
|
||||
node: { ...steps[0].node!, position: executionPosition }
|
||||
};
|
||||
|
||||
expect(mergeExecutionStepNode(executionStep, sourceFlow)?.position).toEqual(executionPosition);
|
||||
});
|
||||
|
||||
it('reproduces the documented five-edge decision graph, including both decision outputs', () => {
|
||||
const ids = Array.from({ length: 6 }, (_, index) =>
|
||||
`b1a50000-0000-4000-8000-00000000000${index + 1}`
|
||||
);
|
||||
const documentedConnections = [
|
||||
{ id: 'c1', sourceId: ids[0], sourceName: 'output', targetId: ids[1], targetName: 'candidateProfile' },
|
||||
{ id: 'c2', sourceId: ids[1], sourceName: 'response', targetId: ids[2], targetName: 'input' },
|
||||
{ id: 'c3', sourceId: ids[2], sourceName: 'approve', targetId: ids[3], targetName: 'input' },
|
||||
{ id: 'c4', sourceId: ids[2], sourceName: 'reject', targetId: ids[4], targetName: 'input' },
|
||||
{ id: 'c5', sourceId: ids[3], sourceName: 'output', targetId: ids[5], targetName: 'input' }
|
||||
];
|
||||
const documentedSteps = ids.map((id, index) => step(node(id, {
|
||||
x: index < 3 ? index * 300 : index === 4 ? 900 : index === 5 ? 1200 : 900,
|
||||
y: index === 4 ? 300 : index > 2 ? 40 : 160
|
||||
})));
|
||||
const documentedExecution = {
|
||||
...execution,
|
||||
stepConnections: documentedConnections,
|
||||
stepDependencies: []
|
||||
};
|
||||
|
||||
const resolved = resolveExecutionConnections(documentedExecution, documentedSteps, null, []);
|
||||
expect(resolved).toEqual(documentedConnections);
|
||||
expect(resolved.filter((connection) => connection.sourceId === ids[2]).map((connection) => connection.sourceName))
|
||||
.toEqual(['approve', 'reject']);
|
||||
});
|
||||
});
|
||||
|
|
@ -0,0 +1,109 @@
|
|||
import { FlowBlockConnection, FlowData, FlowNode, FlowNodeDependency } from '@models/flow';
|
||||
import { getTaskExecutionStepNode, TaskExecution, TaskExecutionStep } from '@models/task-execution';
|
||||
|
||||
export function executionStepNodeId(step: TaskExecutionStep): string {
|
||||
return String(getTaskExecutionStepNode(step)?.id ?? step.id);
|
||||
}
|
||||
|
||||
export function mergeExecutionStepNode(
|
||||
step: TaskExecutionStep,
|
||||
sourceFlow: FlowData | null | undefined
|
||||
): FlowNode | null {
|
||||
const executionNode = getTaskExecutionStepNode(step);
|
||||
if (!executionNode) return null;
|
||||
|
||||
const sourceNode = findFlowNode(sourceFlow, executionNode.id)
|
||||
?? findFlowNode(sourceFlow, step.id);
|
||||
if (!sourceNode) return executionNode;
|
||||
|
||||
return {
|
||||
...sourceNode,
|
||||
...executionNode,
|
||||
id: executionNode.id,
|
||||
position: executionNode.position ?? sourceNode.position,
|
||||
inputs: executionNode.inputs?.length ? executionNode.inputs : sourceNode.inputs,
|
||||
outputs: executionNode.outputs?.length ? executionNode.outputs : sourceNode.outputs,
|
||||
specificConfiguration: {
|
||||
...(sourceNode.specificConfiguration ?? {}),
|
||||
...(executionNode.specificConfiguration ?? {})
|
||||
},
|
||||
biasAnnotations: executionNode.biasAnnotations ?? sourceNode.biasAnnotations,
|
||||
capabilities: executionNode.capabilities ?? sourceNode.capabilities,
|
||||
nodeFamily: sourceNode.nodeFamily ?? executionNode.nodeFamily
|
||||
} as FlowNode;
|
||||
}
|
||||
|
||||
export function resolveExecutionConnections(
|
||||
execution: TaskExecution | null | undefined,
|
||||
steps: TaskExecutionStep[],
|
||||
sourceFlow: FlowData | null | undefined,
|
||||
inferredConnections: FlowBlockConnection[]
|
||||
): FlowBlockConnection[] {
|
||||
const explicit = execution?.stepConnections;
|
||||
const source = sourceFlow?.connections ?? [];
|
||||
const selected = Array.isArray(explicit)
|
||||
? explicit
|
||||
: source.length
|
||||
? source
|
||||
: inferredConnections;
|
||||
return normalizeConnectionsForSteps(selected, steps, sourceFlow);
|
||||
}
|
||||
|
||||
export function resolveExecutionDependencies(
|
||||
execution: TaskExecution | null | undefined,
|
||||
steps: TaskExecutionStep[],
|
||||
sourceFlow: FlowData | null | undefined
|
||||
): FlowNodeDependency[] {
|
||||
const explicit = execution?.stepDependencies;
|
||||
const source = sourceFlow?.dependencies ?? [];
|
||||
const selected = Array.isArray(explicit) ? explicit : source;
|
||||
const ids = executionNodeIds(steps, sourceFlow);
|
||||
|
||||
return selected.flatMap((dependency) => {
|
||||
const sourceId = ids.get(String(dependency.sourceId));
|
||||
const targetId = ids.get(String(dependency.targetId));
|
||||
return sourceId && targetId ? [{ sourceId, targetId }] : [];
|
||||
});
|
||||
}
|
||||
|
||||
function findFlowNode(flow: FlowData | null | undefined, id: string): FlowNode | null {
|
||||
if (!flow) return null;
|
||||
return [...(flow.blocks ?? []), ...(flow.containers ?? [])].find((node) => node.id === id) ?? null;
|
||||
}
|
||||
|
||||
function normalizeConnectionsForSteps(
|
||||
connections: FlowBlockConnection[],
|
||||
steps: TaskExecutionStep[],
|
||||
sourceFlow: FlowData | null | undefined
|
||||
): FlowBlockConnection[] {
|
||||
const ids = executionNodeIds(steps, sourceFlow);
|
||||
return connections.flatMap((connection, index) => {
|
||||
const sourceId = ids.get(String(connection.sourceId));
|
||||
const targetId = ids.get(String(connection.targetId));
|
||||
if (!sourceId || !targetId) return [];
|
||||
|
||||
return [{
|
||||
id: String(connection.id || `${sourceId}:${connection.sourceName}->${targetId}:${connection.targetName}:${index}`),
|
||||
sourceId,
|
||||
sourceName: String(connection.sourceName),
|
||||
targetId,
|
||||
targetName: String(connection.targetName)
|
||||
}];
|
||||
});
|
||||
}
|
||||
|
||||
function executionNodeIds(
|
||||
steps: TaskExecutionStep[],
|
||||
sourceFlow?: FlowData | null
|
||||
): Map<string, string> {
|
||||
const ids = new Map<string, string>();
|
||||
for (const step of steps) {
|
||||
const nodeId = executionStepNodeId(step);
|
||||
ids.set(String(step.id), nodeId);
|
||||
ids.set(nodeId, nodeId);
|
||||
}
|
||||
for (const node of [...(sourceFlow?.blocks ?? []), ...(sourceFlow?.containers ?? [])]) {
|
||||
ids.set(node.id, node.id);
|
||||
}
|
||||
return ids;
|
||||
}
|
||||
|
|
@ -321,7 +321,24 @@ export function getExecutionErrors(stepId: string, contextErrors: Record<string,
|
|||
return [];
|
||||
}
|
||||
|
||||
export function getExecutionWarnings(stepId: string, contextWarnings: Record<string, string>): string[] {
|
||||
const raw = contextWarnings[stepId];
|
||||
return raw && raw.trim().length > 0 ? [raw] : [];
|
||||
export function getExecutionWarnings(stepId: string, contextWarnings: unknown): string[] {
|
||||
if (!contextWarnings || typeof contextWarnings !== 'object') return [];
|
||||
if (Array.isArray(contextWarnings)) {
|
||||
return contextWarnings.flatMap((warning) => {
|
||||
if (typeof warning === 'string') return [warning];
|
||||
if (!warning || typeof warning !== 'object') return [];
|
||||
const value = warning as Record<string, unknown>;
|
||||
const warningStepId = value['stepId'] ?? value['nodeId'];
|
||||
const message = value['message'] ?? value['warning'];
|
||||
return String(warningStepId ?? '') === stepId && typeof message === 'string' && message.trim()
|
||||
? [message]
|
||||
: [];
|
||||
});
|
||||
}
|
||||
const raw = (contextWarnings as Record<string, unknown>)[stepId];
|
||||
if (typeof raw === 'string' && raw.trim().length > 0) return [raw];
|
||||
if (Array.isArray(raw)) {
|
||||
return raw.filter((value): value is string => typeof value === 'string' && value.trim().length > 0);
|
||||
}
|
||||
return [];
|
||||
}
|
||||
|
|
|
|||
|
|
@ -35,13 +35,14 @@ import {
|
|||
import { NodeSettingField, NodeSettingsDialogService } from '@services/dialogs/node-settings-dialog';
|
||||
import { FieldRetriever } from '@services/retriever/field-retriever';
|
||||
import { TaskExecutionsService } from '@services/task-executions/task-executions';
|
||||
import { FlowsService } from '@services/flows/flows';
|
||||
import { ContainersService } from '@services/containers/containers';
|
||||
import { BlocksService } from '@services/blocks/blocks';
|
||||
import { BiasRerunDialogService, BiasRerunCandidate } from '@services/dialogs/bias-rerun-dialog';
|
||||
import { BiasCompareDialogService } from '@services/dialogs/bias-compare-dialog';
|
||||
import { BiasComparisonViewStateService } from '@services/bias/bias-comparison-view-state';
|
||||
import { BiasImpactReportListComponent } from '@shared/bias-impact-report-list/bias-impact-report-list';
|
||||
import { firstValueFrom } from 'rxjs';
|
||||
import { firstValueFrom, take } from 'rxjs';
|
||||
import {
|
||||
ExecutionOutputEntry,
|
||||
ExecutionOutputGroup,
|
||||
|
|
@ -70,6 +71,11 @@ import {
|
|||
getExecutionErrors,
|
||||
getExecutionWarnings,
|
||||
} from './execution-viewer.utils';
|
||||
import {
|
||||
mergeExecutionStepNode,
|
||||
resolveExecutionConnections,
|
||||
resolveExecutionDependencies
|
||||
} from './execution-graph';
|
||||
|
||||
@Component({
|
||||
selector: 'app-task-execution-viewer',
|
||||
|
|
@ -81,6 +87,7 @@ import {
|
|||
export class TaskExecutionViewerComponent implements OnDestroy {
|
||||
private static readonly EVENTS_POLL_INTERVAL_MS = 5000;
|
||||
private taskExecutionsService = inject(TaskExecutionsService);
|
||||
private flowsService = inject(FlowsService);
|
||||
private humanInteractionDialog = inject(HumanInteractionDialogService);
|
||||
private settingsDialog = inject(NodeSettingsDialogService);
|
||||
private fieldRetriever = inject(FieldRetriever);
|
||||
|
|
@ -112,9 +119,13 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
readonly outputPreviewModal = signal<ExecutionOutputEntry | null>(null);
|
||||
readonly intermediateInputPreviewModal = signal<ExecutionIntermediateInputEntry | null>(null);
|
||||
readonly executionLogs = signal<ExecutionEventLogEntry[]>([]);
|
||||
readonly sourceFlowData = signal<FlowData | null>(null);
|
||||
readonly sourceFlowLoading = signal(false);
|
||||
readonly logsLoading = signal(false);
|
||||
readonly logsError = signal<string | null>(null);
|
||||
private readonly logsScrollViewport = viewChild<ElementRef<HTMLDivElement>>('logsScrollViewport');
|
||||
private sourceFlowRequestVersion = 0;
|
||||
private readonly sourceFlowCache = new Map<string, FlowData>();
|
||||
|
||||
constructor() {
|
||||
effect(() => {
|
||||
|
|
@ -133,6 +144,53 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
this.logsLoading.set(false);
|
||||
});
|
||||
|
||||
effect(() => {
|
||||
const execution = this.execution();
|
||||
const requestVersion = ++this.sourceFlowRequestVersion;
|
||||
const embeddedFlow = execution?.flowSnapshot ?? null;
|
||||
if (embeddedFlow) {
|
||||
this.sourceFlowData.set(embeddedFlow);
|
||||
this.sourceFlowLoading.set(false);
|
||||
return;
|
||||
}
|
||||
|
||||
if (Array.isArray(execution?.stepConnections)) {
|
||||
this.sourceFlowData.set(null);
|
||||
this.sourceFlowLoading.set(false);
|
||||
return;
|
||||
}
|
||||
|
||||
const flowId = String(execution?.sourceFlowId ?? execution?.flowId ?? '').trim();
|
||||
if (!flowId) {
|
||||
this.sourceFlowData.set(null);
|
||||
this.sourceFlowLoading.set(false);
|
||||
return;
|
||||
}
|
||||
|
||||
const cached = this.sourceFlowCache.get(flowId);
|
||||
if (cached) {
|
||||
this.sourceFlowData.set(cached);
|
||||
this.sourceFlowLoading.set(false);
|
||||
return;
|
||||
}
|
||||
|
||||
this.sourceFlowData.set(null);
|
||||
this.sourceFlowLoading.set(true);
|
||||
this.flowsService.getFlowById(flowId).pipe(take(1)).subscribe({
|
||||
next: (flow) => {
|
||||
if (requestVersion !== this.sourceFlowRequestVersion) return;
|
||||
this.sourceFlowCache.set(flowId, flow.data);
|
||||
this.sourceFlowData.set(flow.data);
|
||||
this.sourceFlowLoading.set(false);
|
||||
},
|
||||
error: () => {
|
||||
if (requestVersion !== this.sourceFlowRequestVersion) return;
|
||||
this.sourceFlowData.set(null);
|
||||
this.sourceFlowLoading.set(false);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
effect(() => {
|
||||
if (!this.executionOutputTabEnabled() && this.activeAsideTab() === 'output') {
|
||||
this.activeAsideTab.set('inputs');
|
||||
|
|
@ -271,12 +329,30 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
const waitingSteps = this.execution()?.context.waitingSteps ?? [];
|
||||
const activeAnnotationIdsByNode = this.execution()?.biasExecutionContext?.activeAnnotationIdsByNode ?? {};
|
||||
const steps = this.stepsArray();
|
||||
const execution = this.execution();
|
||||
const sourceFlow = execution?.flowSnapshot ?? this.sourceFlowData();
|
||||
const useSourceGraphFallback = !Array.isArray(execution?.stepConnections);
|
||||
const connections = this.getExecutionConnections(steps, sourceFlow);
|
||||
const dependencies = this.getExecutionDependencies(steps, sourceFlow);
|
||||
const blocks: FlowBlock[] = [];
|
||||
const containers: FlowContainer[] = [];
|
||||
const renderedNodeIds = new Set<string>();
|
||||
|
||||
for (const [index, step] of steps.entries()) {
|
||||
const stepNode = getTaskExecutionStepNode(step);
|
||||
const stepNode = mergeExecutionStepNode(step, sourceFlow);
|
||||
if (!stepNode) continue;
|
||||
const connectedInputs = Array.from(new Set([
|
||||
...getConnectedInputs(step),
|
||||
...connections
|
||||
.filter((connection) => connection.targetId === stepNode.id)
|
||||
.map((connection) => connection.targetName)
|
||||
]));
|
||||
const connectedOutputs = Array.from(new Set([
|
||||
...getConnectedOutputs(step),
|
||||
...connections
|
||||
.filter((connection) => connection.sourceId === stepNode.id)
|
||||
.map((connection) => connection.sourceName)
|
||||
]));
|
||||
|
||||
const executionNode: FlowNode = {
|
||||
...stepNode,
|
||||
|
|
@ -288,14 +364,21 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
__executionStatus: this.execution()?.context.status ?? null,
|
||||
__interactionSimulationEnabled: this.execution()?.interactionSimulationEnabled === true,
|
||||
__stepStatus: step.status,
|
||||
__stepSkipReason: step.skipReason ?? null,
|
||||
__stepSimulated: step.simulated === true,
|
||||
__stepUserInteractive: stepNode.userInteractive === true,
|
||||
__executionStatusGroup: executionStatusGroup,
|
||||
__isWaitingStep: waitingSteps.includes(step.id),
|
||||
__executionInputs: getExecutionInputValues(step, contextInputs),
|
||||
__connectedInputs: getConnectedInputs(step),
|
||||
__connectedInputs: connectedInputs,
|
||||
__executionOutputs: getExecutionOutputValues(step, contextResults),
|
||||
__connectedOutputs: getConnectedOutputs(step),
|
||||
__hasDependencyInputConnection: this.hasIncomingDependency(step.id),
|
||||
__hasDependantOutputConnection: this.hasOutgoingDependency(step.id),
|
||||
__connectedOutputs: connectedOutputs,
|
||||
__hasDependencyInputConnection: dependencies.some(
|
||||
(dependency) => dependency.targetId === stepNode.id
|
||||
),
|
||||
__hasDependantOutputConnection: dependencies.some(
|
||||
(dependency) => dependency.sourceId === stepNode.id
|
||||
),
|
||||
__executionErrors: getExecutionErrors(step.id, contextErrors),
|
||||
__executionWarnings: getExecutionWarnings(step.id, contextWarnings),
|
||||
__stepResultData: step.result ?? null,
|
||||
|
|
@ -307,6 +390,49 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
y: 100 + Math.floor(index / 3) * 220
|
||||
}
|
||||
};
|
||||
renderedNodeIds.add(executionNode.id);
|
||||
|
||||
if (executionNode.nodeFamily === 'container') {
|
||||
containers.push(executionNode);
|
||||
} else {
|
||||
blocks.push(executionNode);
|
||||
}
|
||||
}
|
||||
|
||||
for (const sourceNode of useSourceGraphFallback
|
||||
? [...(sourceFlow?.blocks ?? []), ...(sourceFlow?.containers ?? [])]
|
||||
: []) {
|
||||
if (renderedNodeIds.has(sourceNode.id)) continue;
|
||||
|
||||
const connectedInputs = (sourceFlow?.connections ?? [])
|
||||
.filter((connection) => connection.targetId === sourceNode.id)
|
||||
.map((connection) => connection.targetName);
|
||||
const connectedOutputs = (sourceFlow?.connections ?? [])
|
||||
.filter((connection) => connection.sourceId === sourceNode.id)
|
||||
.map((connection) => connection.sourceName);
|
||||
const executionNode: FlowNode = {
|
||||
...sourceNode,
|
||||
specificConfiguration: {
|
||||
...(sourceNode.specificConfiguration ?? {}),
|
||||
__executionId: this.execution()?.id ?? null,
|
||||
__executionNodeId: sourceNode.id,
|
||||
__executionStatus: this.execution()?.context.status ?? null,
|
||||
__executionStatusGroup: executionStatusGroup,
|
||||
__stepStatus: 'SKIPPED',
|
||||
__isWaitingStep: false,
|
||||
__executionInputs: {},
|
||||
__connectedInputs: connectedInputs,
|
||||
__executionOutputs: {},
|
||||
__connectedOutputs: connectedOutputs,
|
||||
__hasDependencyInputConnection: this.hasIncomingDependency(sourceNode.id),
|
||||
__hasDependantOutputConnection: this.hasOutgoingDependency(sourceNode.id),
|
||||
__executionErrors: [],
|
||||
__executionWarnings: [],
|
||||
__stepResultData: null,
|
||||
__executionPartialResult: this.execution()?.context.partialResult ?? null,
|
||||
__biasActiveAnnotationIds: activeAnnotationIdsByNode[sourceNode.id] ?? []
|
||||
}
|
||||
};
|
||||
|
||||
if (executionNode.nodeFamily === 'container') {
|
||||
containers.push(executionNode);
|
||||
|
|
@ -315,14 +441,13 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
}
|
||||
}
|
||||
|
||||
const connections = this.getExecutionConnections(steps);
|
||||
const dependencies = this.getExecutionDependencies();
|
||||
return {
|
||||
blocks,
|
||||
containers,
|
||||
connections,
|
||||
dependencies,
|
||||
globalInputs: []
|
||||
globalInputs: sourceFlow?.globalInputs ?? [],
|
||||
lanes: sourceFlow?.lanes ?? []
|
||||
};
|
||||
});
|
||||
|
||||
|
|
@ -1000,30 +1125,27 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
return connections;
|
||||
}
|
||||
|
||||
private getExecutionConnections(steps: TaskExecutionStep[]): FlowBlockConnection[] {
|
||||
const explicitConnections = this.execution()?.stepConnections;
|
||||
if (explicitConnections?.length) {
|
||||
return explicitConnections.map((connection) => ({
|
||||
id: String(connection.id),
|
||||
sourceId: String(connection.sourceId),
|
||||
sourceName: String(connection.sourceName),
|
||||
targetId: String(connection.targetId),
|
||||
targetName: String(connection.targetName)
|
||||
}));
|
||||
}
|
||||
|
||||
return this.inferConnections(steps);
|
||||
private getExecutionConnections(
|
||||
steps: TaskExecutionStep[],
|
||||
sourceFlow: FlowData | null
|
||||
): FlowBlockConnection[] {
|
||||
return resolveExecutionConnections(
|
||||
this.execution(),
|
||||
steps,
|
||||
sourceFlow,
|
||||
this.inferConnections(steps)
|
||||
);
|
||||
}
|
||||
|
||||
private getExecutionDependencies(): FlowNodeDependency[] {
|
||||
return (this.execution()?.stepDependencies ?? []).map((dependency) => ({
|
||||
sourceId: String(dependency.sourceId),
|
||||
targetId: String(dependency.targetId)
|
||||
}));
|
||||
private getExecutionDependencies(
|
||||
steps = this.stepsArray(),
|
||||
sourceFlow = this.execution()?.flowSnapshot ?? this.sourceFlowData()
|
||||
): FlowNodeDependency[] {
|
||||
return resolveExecutionDependencies(this.execution(), steps, sourceFlow);
|
||||
}
|
||||
|
||||
private hasIncomingDependency(stepId: string): boolean {
|
||||
return (this.execution()?.stepDependencies ?? []).some((dependency) => String(dependency.targetId) === stepId);
|
||||
return this.getExecutionDependencies().some((dependency) => String(dependency.targetId) === stepId);
|
||||
}
|
||||
|
||||
private async biasRerunCandidates(): Promise<BiasRerunCandidate[]> {
|
||||
|
|
@ -1052,7 +1174,7 @@ export class TaskExecutionViewerComponent implements OnDestroy {
|
|||
}
|
||||
|
||||
private hasOutgoingDependency(stepId: string): boolean {
|
||||
return (this.execution()?.stepDependencies ?? []).some((dependency) => String(dependency.sourceId) === stepId);
|
||||
return this.getExecutionDependencies().some((dependency) => String(dependency.sourceId) === stepId);
|
||||
}
|
||||
|
||||
private pickBestConnectionCandidate(
|
||||
|
|
|
|||
|
|
@ -181,6 +181,7 @@ export function exportGraph(editor: NodeEditor<HFSchemes>) {
|
|||
inputs,
|
||||
outputs,
|
||||
specificConfiguration: cloneValue(blockData?.specificConfiguration ?? {}),
|
||||
capabilities: cloneValue(blockData?.capabilities),
|
||||
[biasAnnotationsProperty]: cloneValue(blockRecord?.[biasAnnotationsProperty] ?? []),
|
||||
typeName: blockData?.typeName ?? "LLMBlock",
|
||||
nodeFamily: blockData?.nodeFamily === 'container' ? 'container' : 'block',
|
||||
|
|
@ -638,6 +639,8 @@ function getSocket(editor: NodeEditor<HFSchemes>, type: string) {
|
|||
}
|
||||
|
||||
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;
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue