diff --git a/frontend/package.json b/frontend/package.json index e463138510..a5df2674ec 100644 --- a/frontend/package.json +++ b/frontend/package.json @@ -281,6 +281,11 @@ "svelte": "./package/components/FlowStatusViewer.svelte", "default": "./package/components/FlowStatusViewer.svelte" }, + "./components/FlowRecordingReplay.svelte": { + "types": "./package/components/recording/FlowRecordingReplay.svelte.d.ts", + "svelte": "./package/components/recording/FlowRecordingReplay.svelte", + "default": "./package/components/recording/FlowRecordingReplay.svelte" + }, "./components/FlowWrapper.svelte": { "types": "./package/components/FlowWrapper.svelte.d.ts", "svelte": "./package/components/FlowWrapper.svelte", diff --git a/frontend/src/lib/components/FlowBuilder.svelte b/frontend/src/lib/components/FlowBuilder.svelte index 35dda2ed84..b90d91158a 100644 --- a/frontend/src/lib/components/FlowBuilder.svelte +++ b/frontend/src/lib/components/FlowBuilder.svelte @@ -57,7 +57,8 @@ Circle, CheckCircle, RefreshCw, - CheckCheck + CheckCheck, + Focus } from 'lucide-svelte' import Awareness from './Awareness.svelte' import { getAllModules } from './flows/flowExplorer' @@ -913,6 +914,11 @@ icon: CheckCheck } ] + }, + { + displayName: 'Test flow & record', + icon: Focus, + action: () => flowPreviewButtons?.openRecordingPreview() } ] } diff --git a/frontend/src/lib/components/FlowLogViewer.svelte b/frontend/src/lib/components/FlowLogViewer.svelte index 14a19c676e..465275be9d 100644 --- a/frontend/src/lib/components/FlowLogViewer.svelte +++ b/frontend/src/lib/components/FlowLogViewer.svelte @@ -25,6 +25,7 @@ import FlowLogRow from './FlowLogRow.svelte' import { Tooltip } from './meltComponents' import FlowTimelineBar from './FlowTimelineBar.svelte' + import { getActiveReplay } from './recording/flowRecording.svelte' type RootJobData = Partial @@ -93,6 +94,8 @@ showTimeline = true }: Props = $props() + let isReplay = $derived(!!getActiveReplay()) + function getJobLink(jobId: string | undefined): string { if (!jobId) return '' return `${base}/run/${jobId}?workspace=${workspaceId ?? $workspaceStore}` @@ -595,7 +598,7 @@ {/if} - {#if flowInfo?.jobId} + {#if flowInfo?.jobId && !isReplay} - {#if isLeafStep && jobId} + {#if isLeafStep && jobId && !isReplay} select(`${module.id}-logs`)}> { @@ -230,6 +251,53 @@ renderCount++ } + async function recordAndTest() { + lastRecording = undefined + flowRecording.start($pathStore) + flowRecording.setFlow(extractFlow(previewMode)) + setActiveRecording(flowRecording) + await runPreview(previewArgs.val, undefined) + } + + function collectSubJobIds(flowStatus: Job['flow_status']): string[] { + if (!flowStatus) return [] + const ids: string[] = [] + for (const mod of flowStatus.modules ?? []) { + if (mod.job) ids.push(mod.job) + if (mod.flow_jobs) ids.push(...mod.flow_jobs) + } + if (flowStatus.failure_module?.job) ids.push(flowStatus.failure_module.job) + if (flowStatus.preprocessor_module?.job) ids.push(flowStatus.preprocessor_module.job) + return ids + } + + async function recordSubJobs(completedJob: Job) { + const subJobIds = collectSubJobIds(completedJob.flow_status) + await Promise.all( + subJobIds.map(async (subId) => { + try { + const subJob = await JobService.getJob({ + workspace: $workspaceStore!, + id: subId + }) + flowRecording.addCompletedJob(subId, subJob) + // Recurse into nested flows (flow-within-flow) + if (subJob.flow_status) { + await recordSubJobs(subJob) + } + } catch (e) { + console.warn('[recording] failed to fetch sub-job', subId, e) + } + }) + ) + } + + function downloadRecording() { + if (lastRecording) { + flowRecording.download(lastRecording) + } + } + let scrollableDiv: HTMLDivElement | undefined = $state(undefined) function handleScroll() { let newScroll = scrollableDiv?.scrollTop ?? 0 @@ -250,6 +318,24 @@ scrollableDiv && render && untrack(() => onScrollableDivChange()) }) + // During recording, watch sub-job SSE streams to capture incremental logs + $effect(() => { + // Only watch sub-jobs for the current job (jobId is set after runPreview starts) + if (flowRecording.active && job?.flow_status && job?.id === jobId) { + const modules = job.flow_status.modules ?? [] + untrack(() => { + for (const mod of modules) { + if (mod.job) { + flowRecording.watchSubJob(mod.job, $workspaceStore!) + } + } + if (job?.flow_status?.failure_module?.job) { + flowRecording.watchSubJob(job.flow_status.failure_module.job, $workspaceStore!) + } + }) + } + }) + export async function cancelTest() { isRunning = false try { @@ -308,7 +394,7 @@ {#if isRunning} -
+
+ {#if flowRecording.active} + + Recording + + {/if}
{:else} -
+
{#if jobId !== undefined && selectedJobStep !== undefined && selectedJobStepIsTopLevel} runPreview(previewArgs.val, undefined)} + on:click={() => recordingMode ? recordAndTest() : runPreview(previewArgs.val, undefined)} id="flow-editor-test-flow-drawer" shortCut={{ Icon: CornerDownLeft }} > @@ -355,11 +446,24 @@ {selectionManager.getSelectedId()} + {:else if recordingMode} + Test flow and record {:else} Test flow {/if} {/if} + {#if lastRecording && recordingMode} + + {/if}
{/if}
@@ -542,9 +646,17 @@ wideResults bind:flowState={flowStateStore.val} {jobId} - onDone={() => { + onDone={async ({ job: completedJob }) => { isRunning = false $executionCount = $executionCount + 1 + if (flowRecording.active) { + lastRecording = flowRecording.stop() + setActiveRecording(undefined) + // Fetch sub-job data and add to recording (they load after the root completes) + await recordSubJobs(completedJob) + // Trigger reactivity update for lastRecording + lastRecording = lastRecording + } onJobDone?.() }} bind:selectedJobStep diff --git a/frontend/src/lib/components/FlowStatusViewerInner.svelte b/frontend/src/lib/components/FlowStatusViewerInner.svelte index 006016caa2..410ba822ab 100644 --- a/frontend/src/lib/components/FlowStatusViewerInner.svelte +++ b/frontend/src/lib/components/FlowStatusViewerInner.svelte @@ -53,6 +53,7 @@ import { SelectionManager } from './graph/selectionUtils.svelte' import { useThrottle } from 'runed' import { Splitpanes, Pane } from 'svelte-splitpanes' + import { getActiveReplay } from './recording/flowRecording.svelte' let { flowState: flowStateStore, @@ -175,6 +176,9 @@ let getTopModuleStates = $derived(topModuleStates ?? localModuleStates) + // Cache replay check to avoid repeated function calls in the template + let isReplay = $derived(!!getActiveReplay()) + let resultStreams: Record = $state({}) if (onResultStreamUpdate == undefined) { @@ -213,12 +217,14 @@ for (const asset of inputAssets) { if (asset.kind === 'resource' && !(asset.path in resourceMetadataCache)) { resourceMetadataCache[asset.path] = undefined - ResourceService.getResource({ - workspace: workspace ?? $workspaceStore!, - path: asset.path - }) - .then((r) => (resourceMetadataCache[asset.path] = r)) - .catch((err) => {}) + if (!isReplay) { + ResourceService.getResource({ + workspace: workspace ?? $workspaceStore!, + path: asset.path + }) + .then((r) => (resourceMetadataCache[asset.path] = r)) + .catch((err) => {}) + } } } extendedFlowGraphAssetsCtx.val = clone(_flowGraphAssetsCtx?.val) @@ -534,28 +540,30 @@ mod.type === 'WaitingForExecutor' && localModuleStates[mod.id ?? '']?.scheduled_for == undefined ) { - JobService.getJob({ - workspace: workspaceId ?? $workspaceStore ?? '', - id: mod.job ?? '', - noLogs: true, - noCode: true - }) - .then((job) => { - const newState = { - type: mod.type, - scheduled_for: job?.['scheduled_for'], - job_id: job?.id, - parent_module: mod['parent_module'], - args: job?.args, - tag: job?.tag, - script_hash: job?.script_hash - } + if (!isReplay) { + JobService.getJob({ + workspace: workspaceId ?? $workspaceStore ?? '', + id: mod.job ?? '', + noLogs: true, + noCode: true + }) + .then((job) => { + const newState = { + type: mod.type, + scheduled_for: job?.['scheduled_for'], + job_id: job?.id, + parent_module: mod['parent_module'], + args: job?.args, + tag: job?.tag, + script_hash: job?.script_hash + } - setModuleState(mod.id ?? '', newState) - }) - .catch((e) => { - console.error(`Could not load inner module for job ${mod.job}`, e) - }) + setModuleState(mod.id ?? '', newState) + }) + .catch((e) => { + console.error(`Could not load inner module for job ${mod.job}`, e) + }) + } } else if ( (mod.flow_jobs || mod.branch_chosen) && (mod.type == 'Success' || mod.type == 'Failure') && @@ -623,44 +631,50 @@ updateDurationStatuses(key, durationStatuses) if (missingStartedAtIds.length > 0) { + // Mark as pending to prevent duplicate fetches missingStartedAtIds.forEach((id) => { jobMissingStartedAt[id] = 'P' }) - JobService.getStartedAtByIds({ - workspace: workspaceId ?? $workspaceStore ?? '', - requestBody: missingStartedAtIds - }) - .then((jobs) => { - let lastStarted: string | undefined = undefined - let anySet = false - let nDurationStatuses = localDurationStatuses[key]?.byJob - missingStartedAtIds.forEach((id, idx) => { - const startedAt = jobs[idx] - const time = startedAt ? new Date(startedAt).getTime() : undefined - if (time) { - jobMissingStartedAt[id] = time - } else { - delete jobMissingStartedAt[id] - } - if (nDurationStatuses && time) { - if (!nDurationStatuses[id]?.duration_ms) { - anySet = true - lastStarted = id - nDurationStatuses[id] = { - created_at: time, - started_at: time + if (!isReplay) { + JobService.getStartedAtByIds({ + workspace: workspaceId ?? $workspaceStore ?? '', + requestBody: missingStartedAtIds + }) + .then((jobs) => { + let lastStarted: string | undefined = undefined + let anySet = false + let nDurationStatuses = localDurationStatuses[key]?.byJob + missingStartedAtIds.forEach((id, idx) => { + const startedAt = jobs[idx] + const time = startedAt ? new Date(startedAt).getTime() : undefined + if (time) { + jobMissingStartedAt[id] = time + } else { + delete jobMissingStartedAt[id] + } + if (nDurationStatuses && time) { + if (!nDurationStatuses[id]?.duration_ms) { + anySet = true + lastStarted = id + nDurationStatuses[id] = { + created_at: time, + started_at: time + } } } + }) + if (anySet) { + updateDurationStatuses(key, nDurationStatuses) + if (lastStarted) setSelectedLoopSwitch(lastStarted, mod) } }) - if (anySet) { - updateDurationStatuses(key, nDurationStatuses) - if (lastStarted) setSelectedLoopSwitch(lastStarted, mod) - } - }) - .catch((e) => { - console.error(`Could not load inner module duration status for job ${mod.job}`, e) - }) + .catch((e) => { + console.error( + `Could not load inner module duration status for job ${mod.job}`, + e + ) + }) + } } else { setIteration(0, mod.flow_jobs?.[0] ?? '', false, mod.id ?? '', true) } @@ -1313,8 +1327,8 @@ {#if render}
{#if job !== undefined && (job.result_stream || (job.type == 'CompletedJob' && 'result' in job && job.result !== undefined))} Logs [] = [] + let startedWatchingJob: number | undefined = undefined export async function watchJob(testId: string, callbacks?: Callbacks) { logOffset = 0 @@ -356,6 +363,52 @@ errorIteration = 0 currentId = testId scriptProgress = undefined + + // Replay mode: feed recorded events instead of real SSE + const replay = getActiveReplay() + if (replay) { + const recorded = replay.jobs[testId] + if (recorded) { + job = structuredClone(recorded.initial_job) + callbacks?.change?.(job) + // Compute delays relative to replay start so sub-jobs (discovered + // later by FlowStatusViewerInner) stay in sync with the root job. + const elapsed = Date.now() - getReplayStartTime() + for (const event of recorded.events) { + const delay = Math.max(0, event.t - elapsed) + const timeout = setTimeout(() => { + if (job) { + updateJobFromProgress(event.data, job, callbacks) + } + if (event.data.completed) { + const njob = (event.data as any).job as Job & { result_stream?: string } + if (njob) { + // Use whichever logs are more complete (longer) + const streamedLogs = job?.logs ?? '' + const completedLogs = njob.logs ?? '' + njob.logs = + streamedLogs.length >= completedLogs.length ? streamedLogs : completedLogs + const streamedResult = job?.result_stream ?? '' + const completedResult = njob.result_stream ?? '' + njob.result_stream = + streamedResult.length >= completedResult.length ? streamedResult : completedResult + job = njob + onJobCompleted(testId, job, callbacks) + } + } + }, delay) + replayTimeouts.push(timeout) + } + return + } + // Job not in recording — stub it to prevent API calls + const stubJob = { id: testId, type: 'CompletedJob', success: true } as Job + job = stubJob + callbacks?.change?.(stubJob) + clearCurrentId() + return + } + if (loadPlaceholderJobOnStart) { job = structuredClone(loadPlaceholderJobOnStart) } else { @@ -581,6 +634,7 @@ }) callbacks?.change?.(job) + getActiveRecording()?.recordInitialJob(id, job) } if (!onlyResult) { @@ -683,6 +737,7 @@ throw new Error('Not found') } jobUpdateLastFetch = new Date() + getActiveRecording()?.recordEvent(id, previewJobUpdates) if (job) { updateJobFromProgress(previewJobUpdates, job, callbacks) @@ -806,6 +861,8 @@ clearCurrentId() currentEventSource?.close() currentEventSource = undefined + replayTimeouts.forEach(clearTimeout) + replayTimeouts = [] }) diff --git a/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte b/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte index db23a25160..206c9f199b 100644 --- a/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte +++ b/frontend/src/lib/components/flows/header/FlowPreviewButtons.svelte @@ -53,6 +53,16 @@ flowPreviewContent?.test() } + export async function openRecordingPreview() { + if (!previewOpen) { + previewOpen = true + await tick() + flowPreviewContent?.refresh() + } + previewMode = 'whole' + flowPreviewContent?.setRecordingMode(true) + } + export async function runPreview(conversationId?: string): Promise { if (!previewOpen) { deferContent = true @@ -172,6 +182,7 @@ // keep the data in the preview content deferContent = true previewOpen = false + flowPreviewContent?.setRecordingMode(false) }} on:openTriggers={(e) => { previewOpen = false diff --git a/frontend/src/lib/components/recording/FlowRecordingReplay.svelte b/frontend/src/lib/components/recording/FlowRecordingReplay.svelte new file mode 100644 index 0000000000..cf6c3b2bdf --- /dev/null +++ b/frontend/src/lib/components/recording/FlowRecordingReplay.svelte @@ -0,0 +1,222 @@ + + +{#if !recording?.flow} +
+
+

+ This recording does not include a flow definition. It was likely recorded with an older + version. Re-record the flow to include the flow definition. +

+
+
+{:else if replayState === 'loaded'} +
+
+
+

{recording.flow_path}

+ + + + Recorded {new Date(recording.recorded_at).toLocaleString()} — + {(recording.total_duration_ms / 1000).toFixed(1)}s + + +
+ +
+ +
+{:else if replayState === 'playing' && rootJobId} +
+
+

Replaying: {recording.flow_path}

+ +
+ + {#if job} + + {/if} + +
+{/if} diff --git a/frontend/src/lib/components/recording/flowRecording.svelte.ts b/frontend/src/lib/components/recording/flowRecording.svelte.ts new file mode 100644 index 0000000000..2e715c7082 --- /dev/null +++ b/frontend/src/lib/components/recording/flowRecording.svelte.ts @@ -0,0 +1,192 @@ +import type { Job, OpenFlow } from '$lib/gen' +import type { FlowRecording, RecordedJob } from './types' + +// Module-level active instances (bypasses context/portal issues) +let activeRecording: FlowRecordingStore | undefined = undefined +let activeReplay: FlowRecording | undefined = undefined +let replayStartTime: number = 0 + +export function getActiveRecording() { + return activeRecording +} +export function setActiveRecording(r: FlowRecordingStore | undefined) { + activeRecording = r +} +export function getActiveReplay() { + return activeReplay +} +export function setActiveReplay(r: FlowRecording | undefined) { + activeReplay = r + replayStartTime = r ? Date.now() : 0 +} +export function getReplayStartTime() { + return replayStartTime +} + +export function createFlowRecording() { + let active = $state(false) + let startTime = 0 + let flowPath = '' + let jobs: Record = {} + let flow: OpenFlow | undefined = undefined + let watchedSubJobs = new Set() + let subJobSources: EventSource[] = [] + + return { + get active() { + return active + }, + start(path: string) { + subJobSources.forEach((es) => es.close()) + subJobSources = [] + watchedSubJobs.clear() + active = true + startTime = Date.now() + flowPath = path + jobs = {} + flow = undefined + }, + setFlow(f: OpenFlow) { + // JSON round-trip to strip non-serializable properties (event handlers, etc.) + flow = JSON.parse(JSON.stringify(f)) as OpenFlow + }, + recordInitialJob(jobId: string, job: Job) { + if (!active) return + // $state.snapshot unwraps Svelte 5 reactive proxies (structuredClone fails on them) + jobs[jobId] = { + initial_job: $state.snapshot(job) as Job, + events: [] + } + }, + recordEvent(jobId: string, data: Record) { + if (!active) return + if (!jobs[jobId]) { + jobs[jobId] = { + initial_job: (data as any).job + ? ($state.snapshot((data as any).job) as Job) + : ({} as Job), + events: [] + } + } + jobs[jobId].events.push({ + t: Date.now() - startTime, + data: $state.snapshot(data) as Record + }) + }, + stop(): FlowRecording { + active = false + subJobSources.forEach((es) => es.close()) + subJobSources = [] + const recording: FlowRecording = { + version: 1, + recorded_at: new Date().toISOString(), + flow_path: flowPath, + total_duration_ms: Date.now() - startTime, + jobs, + flow + } + return recording + }, + /** Add a completed sub-job to the recording (called after stop for sub-jobs fetched post-completion) */ + addCompletedJob(jobId: string, completedJob: Job) { + const snapshotJob = $state.snapshot(completedJob) as Job + if (!jobs[jobId]) { + jobs[jobId] = { + initial_job: snapshotJob, + events: [] + } + } else if (!jobs[jobId].initial_job?.id) { + // Replace placeholder initial_job created by recordEvent + jobs[jobId].initial_job = snapshotJob + } + // Only add synthetic completed event if SSE didn't already capture one + const hasCompleted = jobs[jobId].events.some((e) => e.data.completed) + if (!hasCompleted) { + jobs[jobId].events.push({ + t: Date.now() - startTime, + data: { completed: true, job: snapshotJob } + }) + } + }, + /** Watch a sub-job's SSE stream during recording to capture incremental log events */ + watchSubJob(jobId: string, workspace: string) { + if (!active || watchedSubJobs.has(jobId)) return + watchedSubJobs.add(jobId) + + let subJobLogOffset = 0 + + const params = new URLSearchParams({ + log_offset: '0', + running: 'true', + fast: 'true' + }) + const url = `/api/w/${workspace}/jobs_u/getupdate_sse/${jobId}?${params}` + const es = new EventSource(url) + subJobSources.push(es) + + es.onmessage = (event) => { + try { + const data = JSON.parse(event.data) + if (data.type === 'ping' || data.type === 'timeout') return + if (data.type === 'error' || data.type === 'not_found') { + es.close() + return + } + if (!active) { + es.close() + return + } + + // Deduplicate log data: SSE may resend full log dumps on reconnect. + // Drop log data if the offset hasn't advanced past what we already recorded. + if (data.new_logs != null && data.log_offset != null) { + if (subJobLogOffset > 0 && data.log_offset <= subJobLogOffset) { + delete data.new_logs + delete data.log_offset + } else { + subJobLogOffset = data.log_offset + } + } else if (data.log_offset != null && data.log_offset > subJobLogOffset) { + subJobLogOffset = data.log_offset + } + + // Record initial job from the first event that carries job data + if (!jobs[jobId] && data.job) { + jobs[jobId] = { + initial_job: data.job as Job, + events: [] + } + } else if (!jobs[jobId]) { + jobs[jobId] = { + initial_job: {} as Job, + events: [] + } + } + jobs[jobId].events.push({ + t: Date.now() - startTime, + data + }) + if (data.completed) { + es.close() + } + } catch { + // Ignore parse errors + } + } + es.onerror = () => { + es.close() + } + }, + download(recording: FlowRecording) { + const blob = new Blob([JSON.stringify(recording, null, 2)], { type: 'application/json' }) + const url = URL.createObjectURL(blob) + const a = document.createElement('a') + a.href = url + a.download = `flow-recording-${recording.flow_path.replace(/\//g, '-')}-${Date.now()}.json` + a.click() + URL.revokeObjectURL(url) + } + } +} + +export type FlowRecordingStore = ReturnType diff --git a/frontend/src/lib/components/recording/types.ts b/frontend/src/lib/components/recording/types.ts new file mode 100644 index 0000000000..86e7aac754 --- /dev/null +++ b/frontend/src/lib/components/recording/types.ts @@ -0,0 +1,20 @@ +import type { Job, OpenFlow } from '$lib/gen' + +export type RecordedEvent = { + t: number + data: Record +} + +export type RecordedJob = { + initial_job: Job + events: RecordedEvent[] +} + +export type FlowRecording = { + version: 1 + recorded_at: string + flow_path: string + total_duration_ms: number + jobs: Record + flow?: OpenFlow +} diff --git a/frontend/src/routes/(root)/(logged)/replay/+page.svelte b/frontend/src/routes/(root)/(logged)/replay/+page.svelte new file mode 100644 index 0000000000..4309a2078f --- /dev/null +++ b/frontend/src/routes/(root)/(logged)/replay/+page.svelte @@ -0,0 +1,59 @@ + + +
+ {#if recording} +
+ +
+ + {:else} +
+
+

Replay a flow recording

+

+ Upload a recording JSON file to replay a flow execution offline. +

+ + Drag and drop a recording file + +
+
+ {/if} +