add flow recording and offline replay (#8080)

Add the ability to record a flow test execution and replay it offline
without any API calls. This is useful for debugging, sharing, and
reviewing flow executions outside of a running Windmill instance.

Recording:
- "Test flow & record" option in the flow editor three-dots menu
  opens the test drawer in recording mode
- While in recording mode, running a test captures all job events
  (SSE streams, sub-job completions, flow status transitions) along
  with the flow definition into a downloadable JSON file
- Recording state module (flowRecording.svelte.ts) manages active
  recording/replay instances at the module level

Replay:
- Standalone /replay page where users upload a recording JSON file
  and watch the flow execute with real-time status transitions
- FlowRecordingReplay component handles timestamp rebasing, event
  ordering fixes, and drives FlowStatusViewer with recorded data
- JobLoader intercepts replay mode to feed recorded events via
  timed callbacks instead of real SSE/polling
- FlowStatusViewerInner and FlowLogViewer guard all API call sites
  to prevent network requests during replay
- Job links, log downloads, and resource lookups are suppressed
  in replay mode

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
hugocasa
2026-02-24 20:55:20 +01:00
committed by GitHub
parent bb03497967
commit 1f0e2b32b4
11 changed files with 794 additions and 90 deletions
+5
View File
@@ -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",
@@ -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()
}
]
}
@@ -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<Job>
@@ -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}
</div>
{#if flowInfo?.jobId}
{#if flowInfo?.jobId && !isReplay}
<a
href={getJobLink(flowInfo.jobId)}
class="text-xs text-gray-400 hover:text-primary pl-1"
@@ -615,7 +618,7 @@
{#if flowInfo?.logs}
<LogViewer
content={flowInfo.logs}
jobId={flowInfo.jobId}
jobId={isReplay ? undefined : flowInfo.jobId}
isLoading={false}
small={true}
download={false}
@@ -805,7 +808,7 @@
{/if}
</div>
{#if isLeafStep && jobId}
{#if isLeafStep && jobId && !isReplay}
<a
href={getJobLink(jobId ?? '')}
class="text-xs text-gray-400 hover:text-primary pl-1"
@@ -898,7 +901,7 @@
<div onclick={() => select(`${module.id}-logs`)}>
<LogViewer
content={logs}
jobId={jobId ?? ''}
jobId={isReplay ? undefined : (jobId ?? '')}
isLoading={false}
small={true}
download={false}
@@ -17,7 +17,16 @@
import FlowProgressBar from './flows/FlowProgressBar.svelte'
import FlowExecutionStatus from './runs/FlowExecutionStatus.svelte'
import JobDetailHeader from './runs/JobDetailHeader.svelte'
import { AlertTriangle, CornerDownLeft, Loader2, Play, RefreshCw, X } from 'lucide-svelte'
import {
AlertTriangle,
Circle,
CornerDownLeft,
Download,
Loader2,
Play,
RefreshCw,
X
} from 'lucide-svelte'
import { emptyString, sendUserToast, type StateStore } from '$lib/utils'
import { dfs } from './flows/dfs'
import { sliceModules } from './flows/flowStateUtils.svelte'
@@ -30,6 +39,8 @@
import FlowChat from './flows/conversations/FlowChat.svelte'
import { stateSnapshot } from '$lib/svelte5Utils.svelte'
import FlowRestartButton from './FlowRestartButton.svelte'
import { createFlowRecording, setActiveRecording } from './recording/flowRecording.svelte'
import type { FlowRecording } from './recording/types'
interface Props {
previewMode: 'upTo' | 'whole'
@@ -109,6 +120,16 @@
let currentJobId: string | undefined = $state(undefined)
let stepHistoryLoader = getStepHistoryLoaderContext()
let flowProgressBar: FlowProgressBar | undefined = $state(undefined)
// Recording & replay
let flowRecording = createFlowRecording()
let lastRecording: FlowRecording | undefined = $state(undefined)
let recordingMode: boolean = $state(false)
export function setRecordingMode(on: boolean) {
recordingMode = on
}
let loadingHistory = $state(false)
let shouldUseStreaming = $derived.by(() => {
@@ -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 @@
</div>
{#if isRunning}
<div class="mx-auto">
<div class="mx-auto flex items-center gap-2">
<Button
variant="accent"
destructive
@@ -322,9 +408,14 @@
>
Cancel
</Button>
{#if flowRecording.active}
<span class="text-xs text-red-500 font-medium flex items-center gap-1">
<Circle size={8} fill="currentColor" /> Recording
</span>
{/if}
</div>
{:else}
<div class="grow justify-center flex flex-row gap-2">
<div class="grow justify-center flex flex-row gap-2 items-center">
{#if jobId !== undefined && selectedJobStep !== undefined && selectedJobStepIsTopLevel}
<FlowRestartButton
{jobId}
@@ -346,7 +437,7 @@
startIcon={{ icon: isRunning ? RefreshCw : Play }}
size="sm"
btnClasses="w-full max-w-lg"
on:click={() => 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 @@
<Badge baseClass="ml-1" color="indigo">
{selectionManager.getSelectedId()}
</Badge>
{:else if recordingMode}
Test flow and record
{:else}
Test flow
{/if}
</Button>
{/if}
{#if lastRecording && recordingMode}
<Button
variant="subtle"
unifiedSize="sm"
title="Download recording"
on:click={downloadRecording}
startIcon={{ icon: Download }}
>
Download recording
</Button>
{/if}
</div>
{/if}
</div>
@@ -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
@@ -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<string, string | undefined> = $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}
<div class="w-full border rounded-sm bg-surface p-1 overflow-auto max-h-[90vh]">
<DisplayResult
workspaceId={job?.workspace_id}
{jobId}
workspaceId={isReplay ? undefined : job?.workspace_id}
jobId={isReplay ? undefined : jobId}
result_stream={job?.result_stream}
result={jobResults}
language={job?.language}
@@ -1337,9 +1351,9 @@
<div class="flex-1 min-h-0 overflow-auto rounded-md border bg-surface-tertiary p-4">
{#if job !== undefined && (job.result_stream || (job.type == 'CompletedJob' && 'result' in job && job.result !== undefined))}
<DisplayResult
workspaceId={job?.workspace_id}
workspaceId={isReplay ? undefined : job?.workspace_id}
result_stream={job.result_stream}
jobId={job?.id}
jobId={isReplay ? undefined : job?.id}
result={'result' in job ? job.result : undefined}
language={job?.language}
isTest={false}
@@ -1363,7 +1377,8 @@
<h3 class="shrink-0 text-xs font-semibold text-emphasis mb-1">Logs</h3>
<div class="flex-1 min-h-0 overflow-auto rounded-md border bg-surface-tertiary">
<LogViewer
jobId={job.id}
jobId={isReplay ? undefined : job.id}
download={!isReplay}
duration={job?.['duration_ms']}
mem={job?.['mem_peak']}
isLoading={job?.['running'] == false}
@@ -1379,9 +1394,9 @@
<div class="flex-1 overflow-auto rounded-md border bg-surface-tertiary p-4 max-h-screen">
{#if job !== undefined && (job.result_stream || (job.type == 'CompletedJob' && 'result' in job && job.result !== undefined))}
<DisplayResult
workspaceId={job?.workspace_id}
workspaceId={isReplay ? undefined : job?.workspace_id}
result_stream={job.result_stream}
jobId={job?.id}
jobId={isReplay ? undefined : job?.id}
result={'result' in job ? job.result : undefined}
language={job?.language}
isTest={false}
@@ -1445,7 +1460,7 @@
btnClasses="w-full flex justify-start"
on:click={async () => {
let storedJob = storedListJobs[j]
if (!storedJob) {
if (!storedJob && !isReplay) {
storedJob = await JobService.getJob({
workspace: workspaceId ?? $workspaceStore ?? '',
id: loopJobId,
@@ -1454,7 +1469,9 @@
})
storedListJobs[j] = storedJob
}
innerJobLoaded(storedJob, j, true, false)
if (storedJob) {
innerJobLoaded(storedJob, j, true, false)
}
}}
endIcon={{
icon: ChevronDown,
@@ -1947,21 +1964,21 @@
{#if selectedNode == 'end'}
<FlowJobResult
tagLabel={customUi?.tagLabel}
workspaceId={job?.workspace_id}
jobId={job?.id}
workspaceId={isReplay ? undefined : job?.workspace_id}
jobId={isReplay ? undefined : job?.id}
filename={job.id}
loading={job['running']}
tag={job?.tag}
col
result={job['result']}
logs={job.logs ?? ''}
downloadLogs={!hideDownloadLogs}
downloadLogs={!hideDownloadLogs && !isReplay}
/>
{:else if selectedNode == 'start'}
{#if job.args}
<JobArgs
id={job.id}
workspace={job.workspace_id ?? $workspaceStore ?? 'no_w'}
id={isReplay ? undefined : job.id}
workspace={isReplay ? undefined : (job.workspace_id ?? $workspaceStore ?? 'no_w')}
args={job.args}
/>
{:else}
@@ -1982,10 +1999,10 @@
>
<div class="overflow-auto max-h-[200px] p-2">
<DisplayResult
workspaceId={job?.workspace_id}
workspaceId={isReplay ? undefined : job?.workspace_id}
result={node.flow_jobs_results}
nodeId={selectedNode}
jobId={job?.id}
jobId={isReplay ? undefined : job?.id}
language={job?.language}
/>
</div>
@@ -2011,7 +2028,7 @@
{msToSec(node.duration_ms)} s
</Badge>
{/if}
{#if node.job_id}
{#if node.job_id && !isReplay}
<div class="grow w-full flex flex-row-reverse">
<a
class="text-right text-xs"
@@ -2030,8 +2047,8 @@
<div>
<div class="text-xs text-emphasis font-semibold mb-1">Inputs</div>
<JobArgs
id={node.job_id}
workspace={job.workspace_id ?? $workspaceStore ?? 'no_w'}
id={isReplay ? undefined : node.job_id}
workspace={isReplay ? undefined : (job.workspace_id ?? $workspaceStore ?? 'no_w')}
args={node.args}
/>
</div>
@@ -2039,8 +2056,8 @@
</div>
<FlowJobResult
tagLabel={customUi?.tagLabel}
workspaceId={job?.workspace_id}
jobId={node.job_id}
workspaceId={isReplay ? undefined : job?.workspace_id}
jobId={isReplay ? undefined : node.job_id}
loading={node.type != 'Success' && node.type != 'Failure'}
waitingForExecutor={node.type == 'WaitingForExecutor'}
refreshLog={node.type == 'InProgress'}
@@ -2049,7 +2066,7 @@
result={node.result}
tag={node.tag}
logs={node.logs}
downloadLogs={!hideDownloadLogs}
downloadLogs={!hideDownloadLogs && !isReplay}
aiAgentStatus={agentTools &&
node?.job_id &&
(node.type === 'Success' || node.type === 'Failure')
@@ -19,6 +19,11 @@
import type { SupportedLanguage } from '$lib/common'
import { sendUserToast } from '$lib/toast'
import { DynamicInput, isScriptPreview } from '$lib/utils'
import {
getActiveRecording,
getActiveReplay,
getReplayStartTime
} from './recording/flowRecording.svelte'
// Will be set to number if job is not a flow
@@ -348,6 +353,8 @@
return sseSupported
}
let replayTimeouts: ReturnType<typeof setTimeout>[] = []
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 = []
})
</script>
@@ -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<string | undefined> {
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
@@ -0,0 +1,222 @@
<script lang="ts">
import type { Job } from '$lib/gen'
import { workspaceStore } from '$lib/stores'
import FlowStatusViewer from '$lib/components/FlowStatusViewer.svelte'
import FlowViewer from '$lib/components/FlowViewer.svelte'
import FlowProgressBar from '$lib/components/flows/FlowProgressBar.svelte'
import FlowExecutionStatus from '$lib/components/runs/FlowExecutionStatus.svelte'
import { setActiveReplay } from './flowRecording.svelte'
import type { FlowRecording } from './types'
import { sendUserToast } from '$lib/toast'
import { Button } from '$lib/components/common'
import Tooltip from '$lib/components/meltComponents/Tooltip.svelte'
import { InfoIcon, LogOut, Play, Square } from 'lucide-svelte'
import { onDestroy } from 'svelte'
interface Props {
recording: FlowRecording
}
let { recording }: Props = $props()
type ReplayState = 'loaded' | 'playing'
let replayState: ReplayState = $state('loaded')
let rootJobId: string | undefined = $state(undefined)
let rootInitialJob: Job | undefined = $state(undefined)
let job: Job | undefined = $state(undefined)
let done = $derived((job as any)?.type === 'CompletedJob')
function stop() {
setActiveReplay(undefined)
job = undefined
initRecording()
}
function findRootJobId(data: FlowRecording): string | undefined {
for (const [id, recorded] of Object.entries(data.jobs)) {
const j = recorded.initial_job
if (
(j.job_kind === 'flow' || j.job_kind === 'flowpreview') &&
!j.parent_job
) {
return id
}
}
for (const [id, recorded] of Object.entries(data.jobs)) {
if (!recorded.initial_job.parent_job) return id
}
return Object.keys(data.jobs)[0]
}
/**
* Offset all absolute timestamps in the recording so they are relative to "now".
* TimelineCompute uses Date.now() for in-progress bars, so we need recorded
* started_at/created_at values to be near the current time.
*/
function rebaseTimestamps(data: FlowRecording, rootId: string): FlowRecording {
const rootJob = data.jobs[rootId]?.initial_job
const anchor = rootJob?.started_at ?? rootJob?.created_at
if (!anchor) return data
const earliest = new Date(anchor).getTime()
if (isNaN(earliest)) return data
const offset = Date.now() - earliest
function offsetDate(d: string | number | undefined): string | undefined {
if (!d) return d as undefined
const t = new Date(d).getTime()
if (isNaN(t)) return d as string
return new Date(t + offset).toISOString()
}
function offsetJobTimestamps(j: any) {
if (j.started_at) j.started_at = offsetDate(j.started_at)
if (j.created_at) j.created_at = offsetDate(j.created_at)
if (j.completed_at) j.completed_at = offsetDate(j.completed_at)
}
function offsetFlowStatus(fs: any) {
if (!fs?.modules) return
for (const mod of fs.modules) {
const durations = mod.flow_jobs_duration
if (durations?.started_at) {
durations.started_at = durations.started_at.map(
(d: string) => offsetDate(d) ?? d
)
}
}
}
for (const recorded of Object.values(data.jobs)) {
offsetJobTimestamps(recorded.initial_job)
for (const event of recorded.events) {
if (event.data?.job) offsetJobTimestamps(event.data.job)
if (event.data?.flow_status) offsetFlowStatus(event.data.flow_status)
}
}
return data
}
function buildInitialJob(data: FlowRecording, jobId: string): Job {
const initialJob = JSON.parse(JSON.stringify(data.jobs[jobId].initial_job))
if (data.flow?.value) {
initialJob.raw_flow = JSON.parse(JSON.stringify(data.flow.value))
}
return initialJob
}
function initRecording() {
const flowJobId = findRootJobId(recording)
if (!flowJobId) {
sendUserToast('Recording has no jobs', true)
return
}
rootJobId = flowJobId
rootInitialJob = buildInitialJob(recording, flowJobId)
replayState = 'loaded'
}
// Initialize on mount
initRecording()
/**
* Ensure the root flow's completed event fires after all sub-job events.
* During recording, addCompletedJob runs after the flow completes, so sub-job
* completed events can have a later `t` than the root's completed event.
*/
function fixEventOrdering(data: FlowRecording, rootId: string) {
const rootEvents = data.jobs[rootId]?.events
if (!rootEvents?.length) return
// Find the latest event.t across all sub-jobs
let maxSubJobT = 0
for (const [id, recorded] of Object.entries(data.jobs)) {
if (id === rootId) continue
for (const event of recorded.events) {
if (event.t > maxSubJobT) maxSubJobT = event.t
}
}
// Push the root's completed event to fire after all sub-job events
let completedIdx = -1
for (let i = rootEvents.length - 1; i >= 0; i--) {
if (rootEvents[i].data.completed) { completedIdx = i; break }
}
if (completedIdx >= 0 && rootEvents[completedIdx].t < maxSubJobT) {
rootEvents[completedIdx].t = maxSubJobT + 50
}
}
function startReplay() {
// JSON round-trip to unwrap reactive proxies and strip non-cloneable properties
const snapshot = JSON.parse(JSON.stringify(recording)) as FlowRecording
fixEventOrdering(snapshot, rootJobId!)
rebaseTimestamps(snapshot, rootJobId!)
setActiveReplay(snapshot)
rootInitialJob = buildInitialJob(snapshot, rootJobId!)
job = undefined
replayState = 'playing'
}
onDestroy(() => {
setActiveReplay(undefined)
})
</script>
{#if !recording?.flow}
<div class="flex flex-col items-center justify-center min-h-[60vh]">
<div class="border rounded-lg p-8 bg-surface-tertiary max-w-md w-full text-center">
<p class="text-xs text-secondary">
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.
</p>
</div>
</div>
{:else if replayState === 'loaded'}
<div class="flex flex-col gap-4">
<div class="flex items-center justify-between">
<div class="flex items-center gap-2">
<h2 class="text-lg font-semibold text-emphasis">{recording.flow_path}</h2>
<Tooltip placement="bottom">
<InfoIcon size={16} class="text-tertiary" />
<span class="text-2xs" slot="text">
Recorded {new Date(recording.recorded_at).toLocaleString()} &mdash;
{(recording.total_duration_ms / 1000).toFixed(1)}s
</span>
</Tooltip>
</div>
<Button variant="contained" color="blue" on:click={startReplay} startIcon={{ icon: Play }}>
Play
</Button>
</div>
<FlowViewer flow={recording.flow} noSummary />
</div>
{:else if replayState === 'playing' && rootJobId}
<div class="flex flex-col gap-4">
<div class="flex items-center justify-between">
<h2 class="text-lg font-semibold text-emphasis">Replaying: {recording.flow_path}</h2>
<Button variant="border" size="xs" on:click={stop} startIcon={{ icon: done ? LogOut : Square }}>
{done ? 'Exit' : 'Stop'}
</Button>
</div>
<FlowProgressBar {job} slim textPosition="bottom" showStepId />
{#if job}
<FlowExecutionStatus
{job}
workspaceId={$workspaceStore}
isOwner={false}
innerModules={job?.flow_status?.modules}
suspendStatus={{ val: {} }}
/>
{/if}
<FlowStatusViewer
jobId={rootJobId}
initialJob={rootInitialJob}
bind:job
workspaceId={$workspaceStore}
wideResults
showLogsWithResult
/>
</div>
{/if}
@@ -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<string, RecordedJob> = {}
let flow: OpenFlow | undefined = undefined
let watchedSubJobs = new Set<string>()
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<string, any>) {
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<string, any>
})
},
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<typeof createFlowRecording>
@@ -0,0 +1,20 @@
import type { Job, OpenFlow } from '$lib/gen'
export type RecordedEvent = {
t: number
data: Record<string, any>
}
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<string, RecordedJob>
flow?: OpenFlow
}
@@ -0,0 +1,59 @@
<script lang="ts">
import FlowRecordingReplay from '$lib/components/recording/FlowRecordingReplay.svelte'
import type { FlowRecording } from '$lib/components/recording/types'
import { sendUserToast } from '$lib/toast'
import { Button } from '$lib/components/common'
import FileInput from '$lib/components/common/fileInput/FileInput.svelte'
import { Upload } from 'lucide-svelte'
import { setActiveReplay } from '$lib/components/recording/flowRecording.svelte'
let recording: FlowRecording | undefined = $state(undefined)
function handleFileChange(event: CustomEvent<(string | ArrayBuffer | null)[]>) {
const content = event.detail?.[0]
if (!content || typeof content !== 'string') return
try {
const data = JSON.parse(content) as FlowRecording
if (data.version !== 1 || !data.jobs) {
sendUserToast('Invalid recording format', true)
return
}
recording = data
} catch (err) {
sendUserToast('Failed to load recording: ' + err, true)
}
}
function quit() {
setActiveReplay(undefined)
recording = undefined
}
</script>
<div class="max-w-7xl mx-auto px-4 py-8 w-full">
{#if recording}
<div class="flex justify-end mb-4">
<Button variant="border" size="xs" on:click={quit} startIcon={{ icon: Upload }}>
Load another recording
</Button>
</div>
<FlowRecordingReplay {recording} />
{:else}
<div class="flex flex-col items-center justify-center min-h-[60vh]">
<div class="flex flex-col items-center gap-2 max-w-md w-full">
<h2 class="text-lg font-semibold text-emphasis">Replay a flow recording</h2>
<p class="text-xs text-secondary mb-2">
Upload a recording JSON file to replay a flow execution offline.
</p>
<FileInput
accept=".json"
convertTo="text"
class="w-full"
on:change={handleFileChange}
>
Drag and drop a recording file
</FileInput>
</div>
</div>
{/if}
</div>