rag pipeline
parent
3f52f491d7
commit
fa8ab4ea04
@ -1 +1,5 @@
|
|||||||
export * from './use-available-nodes-meta-data'
|
export * from './use-available-nodes-meta-data'
|
||||||
|
export * from './use-workflow-refresh-draft'
|
||||||
|
export * from './use-nodes-sync-draft'
|
||||||
|
export * from './use-workflow-run'
|
||||||
|
export * from './use-workflow-start-run'
|
||||||
|
|||||||
@ -0,0 +1,84 @@
|
|||||||
|
import { useCallback } from 'react'
|
||||||
|
import produce from 'immer'
|
||||||
|
import { useStoreApi } from 'reactflow'
|
||||||
|
import {
|
||||||
|
useWorkflowStore,
|
||||||
|
} from '@/app/components/workflow/store'
|
||||||
|
import {
|
||||||
|
useNodesReadOnly,
|
||||||
|
} from '@/app/components/workflow/hooks/use-workflow'
|
||||||
|
|
||||||
|
export const useNodesSyncDraft = () => {
|
||||||
|
const store = useStoreApi()
|
||||||
|
const workflowStore = useWorkflowStore()
|
||||||
|
const { getNodesReadOnly } = useNodesReadOnly()
|
||||||
|
|
||||||
|
const getPostParams = useCallback(() => {
|
||||||
|
const {
|
||||||
|
getNodes,
|
||||||
|
edges,
|
||||||
|
transform,
|
||||||
|
} = store.getState()
|
||||||
|
const [x, y, zoom] = transform
|
||||||
|
const {
|
||||||
|
pipelineId,
|
||||||
|
environmentVariables,
|
||||||
|
syncWorkflowDraftHash,
|
||||||
|
} = workflowStore.getState()
|
||||||
|
|
||||||
|
if (pipelineId) {
|
||||||
|
const nodes = getNodes()
|
||||||
|
|
||||||
|
const producedNodes = produce(nodes, (draft) => {
|
||||||
|
draft.forEach((node) => {
|
||||||
|
Object.keys(node.data).forEach((key) => {
|
||||||
|
if (key.startsWith('_'))
|
||||||
|
delete node.data[key]
|
||||||
|
})
|
||||||
|
})
|
||||||
|
})
|
||||||
|
const producedEdges = produce(edges, (draft) => {
|
||||||
|
draft.forEach((edge) => {
|
||||||
|
Object.keys(edge.data).forEach((key) => {
|
||||||
|
if (key.startsWith('_'))
|
||||||
|
delete edge.data[key]
|
||||||
|
})
|
||||||
|
})
|
||||||
|
})
|
||||||
|
return {
|
||||||
|
url: `/datasets/${pipelineId}/workflows/draft`,
|
||||||
|
params: {
|
||||||
|
graph: {
|
||||||
|
nodes: producedNodes,
|
||||||
|
edges: producedEdges,
|
||||||
|
viewport: {
|
||||||
|
x,
|
||||||
|
y,
|
||||||
|
zoom,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
environment_variables: environmentVariables,
|
||||||
|
hash: syncWorkflowDraftHash,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}, [store, workflowStore])
|
||||||
|
|
||||||
|
const syncWorkflowDraftWhenPageClose = useCallback(() => {
|
||||||
|
return true
|
||||||
|
}, [])
|
||||||
|
|
||||||
|
const doSyncWorkflowDraft = useCallback(async () => {
|
||||||
|
if (getNodesReadOnly())
|
||||||
|
return
|
||||||
|
const postParams = getPostParams()
|
||||||
|
|
||||||
|
if (postParams)
|
||||||
|
return true
|
||||||
|
}, [getPostParams, getNodesReadOnly])
|
||||||
|
|
||||||
|
return {
|
||||||
|
doSyncWorkflowDraft,
|
||||||
|
syncWorkflowDraftWhenPageClose,
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -0,0 +1,11 @@
|
|||||||
|
import { useCallback } from 'react'
|
||||||
|
|
||||||
|
export const useWorkflowRefreshDraft = () => {
|
||||||
|
const handleRefreshWorkflowDraft = useCallback(() => {
|
||||||
|
return true
|
||||||
|
}, [])
|
||||||
|
|
||||||
|
return {
|
||||||
|
handleRefreshWorkflowDraft,
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -0,0 +1,305 @@
|
|||||||
|
import { useCallback } from 'react'
|
||||||
|
import {
|
||||||
|
useReactFlow,
|
||||||
|
useStoreApi,
|
||||||
|
} from 'reactflow'
|
||||||
|
import produce from 'immer'
|
||||||
|
import { useWorkflowStore } from '@/app/components/workflow/store'
|
||||||
|
import { WorkflowRunningStatus } from '@/app/components/workflow/types'
|
||||||
|
import { useWorkflowUpdate } from '@/app/components/workflow/hooks/use-workflow-interactions'
|
||||||
|
import { useWorkflowRunEvent } from '@/app/components/workflow/hooks/use-workflow-run-event/use-workflow-run-event'
|
||||||
|
import type { IOtherOptions } from '@/service/base'
|
||||||
|
import { ssePost } from '@/service/base'
|
||||||
|
import { stopWorkflowRun } from '@/service/workflow'
|
||||||
|
import type { VersionHistory } from '@/types/workflow'
|
||||||
|
import { useNodesSyncDraft } from './use-nodes-sync-draft'
|
||||||
|
|
||||||
|
export const useWorkflowRun = () => {
|
||||||
|
const store = useStoreApi()
|
||||||
|
const workflowStore = useWorkflowStore()
|
||||||
|
const reactflow = useReactFlow()
|
||||||
|
const { doSyncWorkflowDraft } = useNodesSyncDraft()
|
||||||
|
const { handleUpdateWorkflowCanvas } = useWorkflowUpdate()
|
||||||
|
|
||||||
|
const {
|
||||||
|
handleWorkflowStarted,
|
||||||
|
handleWorkflowFinished,
|
||||||
|
handleWorkflowFailed,
|
||||||
|
handleWorkflowNodeStarted,
|
||||||
|
handleWorkflowNodeFinished,
|
||||||
|
handleWorkflowNodeIterationStarted,
|
||||||
|
handleWorkflowNodeIterationNext,
|
||||||
|
handleWorkflowNodeIterationFinished,
|
||||||
|
handleWorkflowNodeLoopStarted,
|
||||||
|
handleWorkflowNodeLoopNext,
|
||||||
|
handleWorkflowNodeLoopFinished,
|
||||||
|
handleWorkflowNodeRetry,
|
||||||
|
handleWorkflowAgentLog,
|
||||||
|
handleWorkflowTextChunk,
|
||||||
|
handleWorkflowTextReplace,
|
||||||
|
} = useWorkflowRunEvent()
|
||||||
|
|
||||||
|
const handleBackupDraft = useCallback(() => {
|
||||||
|
const {
|
||||||
|
getNodes,
|
||||||
|
edges,
|
||||||
|
} = store.getState()
|
||||||
|
const { getViewport } = reactflow
|
||||||
|
const {
|
||||||
|
backupDraft,
|
||||||
|
setBackupDraft,
|
||||||
|
environmentVariables,
|
||||||
|
} = workflowStore.getState()
|
||||||
|
|
||||||
|
if (!backupDraft) {
|
||||||
|
setBackupDraft({
|
||||||
|
nodes: getNodes(),
|
||||||
|
edges,
|
||||||
|
viewport: getViewport(),
|
||||||
|
environmentVariables,
|
||||||
|
})
|
||||||
|
doSyncWorkflowDraft()
|
||||||
|
}
|
||||||
|
}, [reactflow, workflowStore, store, doSyncWorkflowDraft])
|
||||||
|
|
||||||
|
const handleLoadBackupDraft = useCallback(() => {
|
||||||
|
const {
|
||||||
|
backupDraft,
|
||||||
|
setBackupDraft,
|
||||||
|
setEnvironmentVariables,
|
||||||
|
} = workflowStore.getState()
|
||||||
|
|
||||||
|
if (backupDraft) {
|
||||||
|
const {
|
||||||
|
nodes,
|
||||||
|
edges,
|
||||||
|
viewport,
|
||||||
|
environmentVariables,
|
||||||
|
} = backupDraft
|
||||||
|
handleUpdateWorkflowCanvas({
|
||||||
|
nodes,
|
||||||
|
edges,
|
||||||
|
viewport,
|
||||||
|
})
|
||||||
|
setEnvironmentVariables(environmentVariables)
|
||||||
|
setBackupDraft(undefined)
|
||||||
|
}
|
||||||
|
}, [handleUpdateWorkflowCanvas, workflowStore])
|
||||||
|
|
||||||
|
const handleRun = useCallback(async (
|
||||||
|
params: any,
|
||||||
|
callback?: IOtherOptions,
|
||||||
|
) => {
|
||||||
|
const {
|
||||||
|
getNodes,
|
||||||
|
setNodes,
|
||||||
|
} = store.getState()
|
||||||
|
const newNodes = produce(getNodes(), (draft) => {
|
||||||
|
draft.forEach((node) => {
|
||||||
|
node.data.selected = false
|
||||||
|
node.data._runningStatus = undefined
|
||||||
|
})
|
||||||
|
})
|
||||||
|
setNodes(newNodes)
|
||||||
|
await doSyncWorkflowDraft()
|
||||||
|
|
||||||
|
const {
|
||||||
|
onWorkflowStarted,
|
||||||
|
onWorkflowFinished,
|
||||||
|
onNodeStarted,
|
||||||
|
onNodeFinished,
|
||||||
|
onIterationStart,
|
||||||
|
onIterationNext,
|
||||||
|
onIterationFinish,
|
||||||
|
onLoopStart,
|
||||||
|
onLoopNext,
|
||||||
|
onLoopFinish,
|
||||||
|
onNodeRetry,
|
||||||
|
onAgentLog,
|
||||||
|
onError,
|
||||||
|
...restCallback
|
||||||
|
} = callback || {}
|
||||||
|
const { pipelineId } = workflowStore.getState()
|
||||||
|
workflowStore.setState({ historyWorkflowData: undefined })
|
||||||
|
const workflowContainer = document.getElementById('workflow-container')
|
||||||
|
|
||||||
|
const {
|
||||||
|
clientWidth,
|
||||||
|
clientHeight,
|
||||||
|
} = workflowContainer!
|
||||||
|
|
||||||
|
const url = `/rag/pipeline/${pipelineId}/workflows/draft/run`
|
||||||
|
|
||||||
|
const {
|
||||||
|
setWorkflowRunningData,
|
||||||
|
} = workflowStore.getState()
|
||||||
|
setWorkflowRunningData({
|
||||||
|
result: {
|
||||||
|
status: WorkflowRunningStatus.Running,
|
||||||
|
},
|
||||||
|
tracing: [],
|
||||||
|
resultText: '',
|
||||||
|
})
|
||||||
|
|
||||||
|
return true
|
||||||
|
|
||||||
|
ssePost(
|
||||||
|
url,
|
||||||
|
{
|
||||||
|
body: params,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
onWorkflowStarted: (params) => {
|
||||||
|
handleWorkflowStarted(params)
|
||||||
|
|
||||||
|
if (onWorkflowStarted)
|
||||||
|
onWorkflowStarted(params)
|
||||||
|
},
|
||||||
|
onWorkflowFinished: (params) => {
|
||||||
|
handleWorkflowFinished(params)
|
||||||
|
|
||||||
|
if (onWorkflowFinished)
|
||||||
|
onWorkflowFinished(params)
|
||||||
|
},
|
||||||
|
onError: (params) => {
|
||||||
|
handleWorkflowFailed()
|
||||||
|
|
||||||
|
if (onError)
|
||||||
|
onError(params)
|
||||||
|
},
|
||||||
|
onNodeStarted: (params) => {
|
||||||
|
handleWorkflowNodeStarted(
|
||||||
|
params,
|
||||||
|
{
|
||||||
|
clientWidth,
|
||||||
|
clientHeight,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
if (onNodeStarted)
|
||||||
|
onNodeStarted(params)
|
||||||
|
},
|
||||||
|
onNodeFinished: (params) => {
|
||||||
|
handleWorkflowNodeFinished(params)
|
||||||
|
|
||||||
|
if (onNodeFinished)
|
||||||
|
onNodeFinished(params)
|
||||||
|
},
|
||||||
|
onIterationStart: (params) => {
|
||||||
|
handleWorkflowNodeIterationStarted(
|
||||||
|
params,
|
||||||
|
{
|
||||||
|
clientWidth,
|
||||||
|
clientHeight,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
if (onIterationStart)
|
||||||
|
onIterationStart(params)
|
||||||
|
},
|
||||||
|
onIterationNext: (params) => {
|
||||||
|
handleWorkflowNodeIterationNext(params)
|
||||||
|
|
||||||
|
if (onIterationNext)
|
||||||
|
onIterationNext(params)
|
||||||
|
},
|
||||||
|
onIterationFinish: (params) => {
|
||||||
|
handleWorkflowNodeIterationFinished(params)
|
||||||
|
|
||||||
|
if (onIterationFinish)
|
||||||
|
onIterationFinish(params)
|
||||||
|
},
|
||||||
|
onLoopStart: (params) => {
|
||||||
|
handleWorkflowNodeLoopStarted(
|
||||||
|
params,
|
||||||
|
{
|
||||||
|
clientWidth,
|
||||||
|
clientHeight,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
if (onLoopStart)
|
||||||
|
onLoopStart(params)
|
||||||
|
},
|
||||||
|
onLoopNext: (params) => {
|
||||||
|
handleWorkflowNodeLoopNext(params)
|
||||||
|
|
||||||
|
if (onLoopNext)
|
||||||
|
onLoopNext(params)
|
||||||
|
},
|
||||||
|
onLoopFinish: (params) => {
|
||||||
|
handleWorkflowNodeLoopFinished(params)
|
||||||
|
|
||||||
|
if (onLoopFinish)
|
||||||
|
onLoopFinish(params)
|
||||||
|
},
|
||||||
|
onNodeRetry: (params) => {
|
||||||
|
handleWorkflowNodeRetry(params)
|
||||||
|
|
||||||
|
if (onNodeRetry)
|
||||||
|
onNodeRetry(params)
|
||||||
|
},
|
||||||
|
onAgentLog: (params) => {
|
||||||
|
handleWorkflowAgentLog(params)
|
||||||
|
|
||||||
|
if (onAgentLog)
|
||||||
|
onAgentLog(params)
|
||||||
|
},
|
||||||
|
onTextChunk: (params) => {
|
||||||
|
handleWorkflowTextChunk(params)
|
||||||
|
},
|
||||||
|
onTextReplace: (params) => {
|
||||||
|
handleWorkflowTextReplace(params)
|
||||||
|
},
|
||||||
|
...restCallback,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
}, [
|
||||||
|
store,
|
||||||
|
workflowStore,
|
||||||
|
doSyncWorkflowDraft,
|
||||||
|
handleWorkflowStarted,
|
||||||
|
handleWorkflowFinished,
|
||||||
|
handleWorkflowFailed,
|
||||||
|
handleWorkflowNodeStarted,
|
||||||
|
handleWorkflowNodeFinished,
|
||||||
|
handleWorkflowNodeIterationStarted,
|
||||||
|
handleWorkflowNodeIterationNext,
|
||||||
|
handleWorkflowNodeIterationFinished,
|
||||||
|
handleWorkflowNodeLoopStarted,
|
||||||
|
handleWorkflowNodeLoopNext,
|
||||||
|
handleWorkflowNodeLoopFinished,
|
||||||
|
handleWorkflowNodeRetry,
|
||||||
|
handleWorkflowTextChunk,
|
||||||
|
handleWorkflowTextReplace,
|
||||||
|
handleWorkflowAgentLog,
|
||||||
|
],
|
||||||
|
)
|
||||||
|
|
||||||
|
const handleStopRun = useCallback((taskId: string) => {
|
||||||
|
const { pipelineId } = workflowStore.getState()
|
||||||
|
|
||||||
|
stopWorkflowRun(`/rag/pipeline/${pipelineId}/workflow-runs/tasks/${taskId}/stop`)
|
||||||
|
}, [workflowStore])
|
||||||
|
|
||||||
|
const handleRestoreFromPublishedWorkflow = useCallback((publishedWorkflow: VersionHistory) => {
|
||||||
|
const nodes = publishedWorkflow.graph.nodes.map(node => ({ ...node, selected: false, data: { ...node.data, selected: false } }))
|
||||||
|
const edges = publishedWorkflow.graph.edges
|
||||||
|
const viewport = publishedWorkflow.graph.viewport!
|
||||||
|
handleUpdateWorkflowCanvas({
|
||||||
|
nodes,
|
||||||
|
edges,
|
||||||
|
viewport,
|
||||||
|
})
|
||||||
|
|
||||||
|
workflowStore.getState().setEnvironmentVariables(publishedWorkflow.environment_variables || [])
|
||||||
|
}, [handleUpdateWorkflowCanvas, workflowStore])
|
||||||
|
|
||||||
|
return {
|
||||||
|
handleBackupDraft,
|
||||||
|
handleLoadBackupDraft,
|
||||||
|
handleRun,
|
||||||
|
handleStopRun,
|
||||||
|
handleRestoreFromPublishedWorkflow,
|
||||||
|
}
|
||||||
|
}
|
||||||
@ -0,0 +1,68 @@
|
|||||||
|
import { useCallback } from 'react'
|
||||||
|
import { useStoreApi } from 'reactflow'
|
||||||
|
import { useWorkflowStore } from '@/app/components/workflow/store'
|
||||||
|
import {
|
||||||
|
BlockEnum,
|
||||||
|
WorkflowRunningStatus,
|
||||||
|
} from '@/app/components/workflow/types'
|
||||||
|
import { useWorkflowInteractions } from '@/app/components/workflow/hooks'
|
||||||
|
import {
|
||||||
|
useNodesSyncDraft,
|
||||||
|
useWorkflowRun,
|
||||||
|
} from '.'
|
||||||
|
|
||||||
|
export const useWorkflowStartRun = () => {
|
||||||
|
const store = useStoreApi()
|
||||||
|
const workflowStore = useWorkflowStore()
|
||||||
|
const { handleCancelDebugAndPreviewPanel } = useWorkflowInteractions()
|
||||||
|
const { handleRun } = useWorkflowRun()
|
||||||
|
const { doSyncWorkflowDraft } = useNodesSyncDraft()
|
||||||
|
|
||||||
|
const handleWorkflowStartRunInWorkflow = useCallback(async () => {
|
||||||
|
const {
|
||||||
|
workflowRunningData,
|
||||||
|
} = workflowStore.getState()
|
||||||
|
|
||||||
|
if (workflowRunningData?.result.status === WorkflowRunningStatus.Running)
|
||||||
|
return
|
||||||
|
|
||||||
|
const { getNodes } = store.getState()
|
||||||
|
const nodes = getNodes()
|
||||||
|
const startNode = nodes.find(node => node.data.type === BlockEnum.Start)
|
||||||
|
const startVariables = startNode?.data.variables || []
|
||||||
|
const {
|
||||||
|
showTestRunPanel,
|
||||||
|
setShowInputsPanel,
|
||||||
|
setShowEnvPanel,
|
||||||
|
setShowTestRunPanel,
|
||||||
|
} = workflowStore.getState()
|
||||||
|
|
||||||
|
setShowEnvPanel(false)
|
||||||
|
|
||||||
|
if (showTestRunPanel) {
|
||||||
|
setShowTestRunPanel?.(false)
|
||||||
|
handleCancelDebugAndPreviewPanel()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!startVariables.length) {
|
||||||
|
await doSyncWorkflowDraft()
|
||||||
|
handleRun({ inputs: {}, files: [] })
|
||||||
|
setShowTestRunPanel?.(true)
|
||||||
|
setShowInputsPanel(false)
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
setShowTestRunPanel?.(true)
|
||||||
|
setShowInputsPanel(true)
|
||||||
|
}
|
||||||
|
}, [store, workflowStore, handleCancelDebugAndPreviewPanel, handleRun, doSyncWorkflowDraft])
|
||||||
|
|
||||||
|
const handleStartWorkflowRun = useCallback(() => {
|
||||||
|
handleWorkflowStartRunInWorkflow()
|
||||||
|
}, [handleWorkflowStartRunInWorkflow])
|
||||||
|
|
||||||
|
return {
|
||||||
|
handleStartWorkflowRun,
|
||||||
|
handleWorkflowStartRunInWorkflow,
|
||||||
|
}
|
||||||
|
}
|
||||||
Loading…
Reference in New Issue