This commit is contained in:
Ruben Fiszel
2025-07-15 07:28:12 +00:00
parent 15403d1aca
commit 5b685af774
7 changed files with 283 additions and 229 deletions
+62 -55
View File
@@ -166,7 +166,13 @@
currentScript.language,
args,
currentScript.tag,
useLock ? currentScript.lock : undefined
useLock ? currentScript.lock : undefined,
undefined,
{
done(x) {
loadPastTests()
}
}
)
} else {
sendUserToast(`Bundle received ${lastCommandId} was obsolete, ignoring`, true)
@@ -221,56 +227,63 @@
})
async function testBundle(file: string, isTar: boolean) {
jobLoader?.abstractRun(async () => {
try {
const form = new FormData()
form.append(
'preview',
JSON.stringify({
content: currentScript?.content,
kind: isTar ? 'tarbundle' : 'bundle',
path: currentScript?.path,
args,
language: currentScript?.language,
tag: currentScript?.tag
})
)
// sendUserToast(JSON.stringify(file))
if (isTar) {
var array: number[] = []
file = atob(file)
for (var i = 0; i < file.length; i++) {
array.push(file.charCodeAt(i))
}
let blob = new Blob([new Uint8Array(array)], { type: 'application/octet-stream' })
form.append('file', blob)
} else {
form.append('file', file)
}
const url = '/api/w/' + workspace + '/jobs/run/preview_bundle'
const req = await fetch(url, {
method: 'POST',
body: form,
headers: {
Authorization: 'Bearer ' + token
}
})
if (req.status != 201) {
throw Error(
`Script snapshot creation was not successful: ${req.status} - ${
req.statusText
} - ${await req.text()}`
jobLoader?.abstractRun(
async () => {
try {
const form = new FormData()
form.append(
'preview',
JSON.stringify({
content: currentScript?.content,
kind: isTar ? 'tarbundle' : 'bundle',
path: currentScript?.path,
args,
language: currentScript?.language,
tag: currentScript?.tag
})
)
// sendUserToast(JSON.stringify(file))
if (isTar) {
var array: number[] = []
file = atob(file)
for (var i = 0; i < file.length; i++) {
array.push(file.charCodeAt(i))
}
let blob = new Blob([new Uint8Array(array)], { type: 'application/octet-stream' })
form.append('file', blob)
} else {
form.append('file', file)
}
const url = '/api/w/' + workspace + '/jobs/run/preview_bundle'
const req = await fetch(url, {
method: 'POST',
body: form,
headers: {
Authorization: 'Bearer ' + token
}
})
if (req.status != 201) {
throw Error(
`Script snapshot creation was not successful: ${req.status} - ${
req.statusText
} - ${await req.text()}`
)
}
return await req.text()
} catch (e) {
sendUserToast(`Failed to send bundle ${e}`, true)
throw Error(e)
}
},
{
done(x) {
loadPastTests()
}
return await req.text()
} catch (e) {
sendUserToast(`Failed to send bundle ${e}`, true)
throw Error(e)
}
})
)
loadingCodebaseButton = false
}
onDestroy(() => {
@@ -637,13 +650,7 @@
<svelte:window onkeydown={onKeyDown} />
<JobLoader
noCode={true}
on:done={loadPastTests}
bind:this={jobLoader}
bind:isLoading={testIsLoading}
bind:job={testJob}
/>
<JobLoader noCode={true} bind:this={jobLoader} bind:isLoading={testIsLoading} bind:job={testJob} />
<main class="h-screen w-full">
{#if mode == 'script'}
+169 -144
View File
@@ -1,15 +1,27 @@
<script lang="ts">
import { type Job, JobService, type FlowStatus, type Preview } from '$lib/gen'
import {
type Job,
JobService,
type FlowStatus,
type Preview,
type GetJobUpdatesResponse
} from '$lib/gen'
import { workspaceStore } from '$lib/stores'
import { onDestroy, tick, untrack } from 'svelte'
import { createEventDispatcher } from 'svelte'
import type { SupportedLanguage } from '$lib/common'
import { sendUserToast } from '$lib/toast'
import { isScriptPreview } from '$lib/utils'
// Will be set to number if job is not a flow
type Callbacks = { done: (x: any) => void; cancel: () => void; error: (err: Error) => void }
type Callbacks = {
done?: (x: Job & { result?: any }) => void
doneError?: (x: { id: string; error: Error }) => void
cancel?: () => void
error?: (err: Error) => void
started?: (id: string) => void
running?: (id: string) => void
}
interface Props {
isLoading?: boolean
@@ -17,6 +29,7 @@
noCode?: boolean
workspaceOverride?: string | undefined
notfound?: boolean
allowConcurentRequests?: boolean
jobUpdateLastFetch?: Date | undefined
toastError?: boolean
lazyLogs?: boolean
@@ -29,6 +42,7 @@
isLoading = $bindable(false),
job = $bindable(undefined),
noCode = false,
allowConcurentRequests = false,
workspaceOverride = undefined,
notfound = $bindable(false),
jobUpdateLastFetch = $bindable(undefined),
@@ -47,8 +61,6 @@
/// How often loader poll progress
const getProgressRate: number = 1000
const dispatch = createEventDispatcher()
let workspace = $derived(workspaceOverride ?? $workspaceStore)
let syncIteration: number = 0
@@ -57,6 +69,7 @@
let logOffset = 0
let lastCallbacks: Callbacks | undefined = undefined
let finished: string[] = []
let ITERATIONS_BEFORE_SLOW_REFRESH = 10
let ITERATIONS_BEFORE_SUPER_SLOW_REFRESH = 100
@@ -80,9 +93,10 @@
const startedAt = Date.now()
const testId = await fn()
if (lastStartedAt < startedAt) {
if (lastStartedAt < startedAt || allowConcurentRequests) {
lastStartedAt = startedAt
if (testId) {
callbacks?.started?.(testId)
try {
await watchJob(testId, callbacks)
} catch {
@@ -97,6 +111,7 @@
if (toastError) {
sendUserToast(err.body, true)
}
callbacks?.error?.(err)
// if error happens on submitting the job, reset UI state so the user can try again
isLoading = false
currentId = undefined
@@ -186,27 +201,28 @@
scriptProgress = undefined
lastTimeCheckedProgress = undefined
return abstractRun(() =>
JobService.runScriptPreview({
workspace: $workspaceStore!,
requestBody: {
path,
content: code,
args,
language: lang as Preview['language'],
tag,
lock,
script_hash: hash
}
})
return abstractRun(
() =>
JobService.runScriptPreview({
workspace: $workspaceStore!,
requestBody: {
path,
content: code,
args,
language: lang as Preview['language'],
tag,
lock,
script_hash: hash
}
}),
callbacks
)
}
export async function cancelJob() {
const id = currentId
if (id) {
dispatch('cancel', id)
lastCallbacks?.cancel()
lastCallbacks?.cancel?.()
lastCallbacks = undefined
currentId = undefined
// Clean up SSE connection
@@ -225,8 +241,10 @@
}
export async function clearCurrentJob() {
if (currentId) {
if (currentId && !allowConcurentRequests) {
job = undefined
lastCallbacks?.cancel?.()
lastCallbacks = undefined
await cancelJob()
}
}
@@ -243,42 +261,76 @@
currentEventSource = undefined
// Try SSE first, fall back to polling if needed
const isCompleted = await loadTestJobWithSSE(testId, callbacks)
if (!isCompleted && !currentEventSource) {
// If SSE didn't start (job might not be running yet), use polling
setTimeout(() => {
syncer(testId, callbacks)
}, 50)
}
await loadTestJobWithSSE(testId, 0, callbacks)
}
async function loadTestJob(id: string): Promise<boolean> {
function setJobProgress(job: Job) {
let getProgress: boolean | undefined = undefined
// We only pull individual job progress this way
// Flow's progress we are getting from FlowStatusModule of flow job
if (job.job_kind == 'script' || isScriptPreview(job.job_kind)) {
// First time, before running job, lastTimeCheckedProgress is always undefined
if (lastTimeCheckedProgress) {
const lastTimeCheckedMs = Date.now() - lastTimeCheckedProgress
// Ask for progress if the last time we asked is >5s OR the progress was once not undefined
if (
lastTimeCheckedMs > getProgressRetryRate ||
(scriptProgress != undefined && lastTimeCheckedMs > getProgressRate)
) {
lastTimeCheckedProgress = Date.now()
getProgress = true
}
} else {
// Make it think we asked for progress, but in reality we didnt. First 5s we want to wait without putting extra work on db
// 99.99% of the jobs won't have progress be set so we have to do a balance between having low-latency for jobs that use it and job that don't
// we would usually not care to have progress the first 5s and jobs that are less than 5s
lastTimeCheckedProgress = Date.now()
}
}
return getProgress
}
const clamp = (num: number, min: number, max: number) => Math.min(Math.max(num, min), max)
function updateJobFromProgress(
previewJobUpdates: GetJobUpdatesResponse,
offset: number,
job: Job
) {
// Clamp number between two values with the following line:
if (previewJobUpdates.progress) {
// Progress cannot go back and cannot be set to 100
scriptProgress = clamp(previewJobUpdates.progress, scriptProgress ?? 0, 99)
}
if (previewJobUpdates.new_logs) {
if (offset == 0) {
job.logs = previewJobUpdates.new_logs ?? ''
} else {
job.logs = (job?.logs ?? '').concat(previewJobUpdates.new_logs)
}
}
if (previewJobUpdates.log_offset) {
logOffset = previewJobUpdates.log_offset ?? 0
}
if (previewJobUpdates.flow_status) {
job.flow_status = previewJobUpdates.flow_status as FlowStatus
}
if (previewJobUpdates.mem_peak && job) {
job.mem_peak = previewJobUpdates.mem_peak
}
}
async function loadTestJob(id: string, callbacks?: Callbacks): Promise<boolean> {
let isCompleted = false
if (currentId === id) {
if (currentId === id || allowConcurentRequests) {
try {
if (job && `running` in job) {
let getProgress: boolean | undefined = undefined
// We only pull individual job progress this way
// Flow's progress we are getting from FlowStatusModule of flow job
if (job.job_kind == 'script' || isScriptPreview(job.job_kind)) {
// First time, before running job, lastTimeCheckedProgress is always undefined
if (lastTimeCheckedProgress) {
const lastTimeCheckedMs = Date.now() - lastTimeCheckedProgress
// Ask for progress if the last time we asked is >5s OR the progress was once not undefined
if (
lastTimeCheckedMs > getProgressRetryRate ||
(scriptProgress != undefined && lastTimeCheckedMs > getProgressRate)
) {
lastTimeCheckedProgress = Date.now()
getProgress = true
}
} else {
// Make it think we asked for progress, but in reality we didnt. First 5s we want to wait without putting extra work on db
// 99.99% of the jobs won't have progress be set so we have to do a balance between having low-latency for jobs that use it and job that don't
// we would usually not care to have progress the first 5s and jobs that are less than 5s
lastTimeCheckedProgress = Date.now()
}
}
callbacks?.running?.(id)
let getProgress: boolean | undefined = setJobProgress(job)
const offset = logOffset == 0 ? (job.logs?.length ? job.logs?.length + 1 : 0) : logOffset
@@ -290,35 +342,11 @@
getProgress: getProgress
})
// Clamp number between two values with the following line:
const clamp = (num, min, max) => Math.min(Math.max(num, min), max)
if (previewJobUpdates.progress) {
// Progress cannot go back and cannot be set to 100
scriptProgress = clamp(previewJobUpdates.progress, scriptProgress ?? 0, 99)
}
if (previewJobUpdates.new_logs) {
if (offset == 0) {
job.logs = previewJobUpdates.new_logs ?? ''
} else {
job.logs = (job?.logs ?? '').concat(previewJobUpdates.new_logs)
}
}
if (previewJobUpdates.log_offset) {
logOffset = previewJobUpdates.log_offset ?? 0
}
if (previewJobUpdates.flow_status) {
job.flow_status = previewJobUpdates.flow_status as FlowStatus
}
if (previewJobUpdates.mem_peak && job) {
job.mem_peak = previewJobUpdates.mem_peak
}
if ((previewJobUpdates.running ?? false) || (previewJobUpdates.completed ?? false)) {
job = await JobService.getJob({ workspace: workspace!, id, noCode })
}
updateJobFromProgress(previewJobUpdates, offset, job)
} else {
job = await JobService.getJob({ workspace: workspace!, id, noLogs: lazyLogs, noCode })
}
@@ -327,11 +355,7 @@
if (job?.type === 'CompletedJob') {
//only CompletedJob has success property
isCompleted = true
if (currentId === id) {
await tick()
dispatch('done', job)
currentId = undefined
}
onJobCompleted(id, job, callbacks)
}
notfound = false
} catch (err) {
@@ -350,7 +374,40 @@
let currentEventSource: EventSource | undefined = undefined
async function loadTestJobWithSSE(id: string, callbacks?: Callbacks): Promise<boolean> {
async function onJobCompleted(
id: string,
job: Job & { result?: any; success?: boolean },
callbacks?: Callbacks
) {
if (currentId === id || allowConcurentRequests) {
await tick()
if (
(callbacks?.error || callbacks?.doneError) &&
!job?.success &&
typeof job?.result == 'object' &&
'error' in (job?.result ?? {})
) {
callbacks?.error?.(job.result.error)
callbacks?.doneError?.({
id,
error: job.result.error
})
} else {
callbacks?.done?.(job)
}
if (!allowConcurentRequests) {
currentId = undefined
} else {
finished.push(id)
}
}
}
async function loadTestJobWithSSE(
id: string,
attempt: number,
callbacks?: Callbacks
): Promise<boolean> {
let isCompleted = false
if (currentId === id) {
try {
@@ -364,8 +421,7 @@
isCompleted = true
if (currentId === id) {
await tick()
dispatch('done', job)
callbacks?.done(job)
callbacks?.done?.(job)
currentId = undefined
}
return isCompleted
@@ -373,23 +429,9 @@
// Only start SSE if job is running and we haven't started it yet
if (job && `running` in job && !currentEventSource) {
let getProgress: boolean | undefined = undefined
callbacks?.running?.(id)
// Check if we should get progress updates
if (job.job_kind == 'script' || isScriptPreview(job.job_kind)) {
if (lastTimeCheckedProgress) {
const lastTimeCheckedMs = Date.now() - lastTimeCheckedProgress
if (
lastTimeCheckedMs > getProgressRetryRate ||
(scriptProgress != undefined && lastTimeCheckedMs > getProgressRate)
) {
lastTimeCheckedProgress = Date.now()
getProgress = true
}
} else {
lastTimeCheckedProgress = Date.now()
}
}
let getProgress: boolean | undefined = setJobProgress(job)
const offset = logOffset == 0 ? (job.logs?.length ? job.logs?.length + 1 : 0) : logOffset
@@ -417,31 +459,8 @@
const previewJobUpdates = JSON.parse(event.data)
jobUpdateLastFetch = new Date()
// Clamp number between two values with the following line:
const clamp = (num, min, max) => Math.min(Math.max(num, min), max)
if (previewJobUpdates.progress) {
// Progress cannot go back and cannot be set to 100
scriptProgress = clamp(previewJobUpdates.progress, scriptProgress ?? 0, 99)
}
if (previewJobUpdates.new_logs && job) {
if (offset == 0) {
job.logs = previewJobUpdates.new_logs ?? ''
} else {
job.logs = (job?.logs ?? '').concat(previewJobUpdates.new_logs)
}
}
if (previewJobUpdates.log_offset) {
logOffset = previewJobUpdates.log_offset ?? 0
}
if (previewJobUpdates.flow_status && job) {
job.flow_status = previewJobUpdates.flow_status as FlowStatus
}
if (previewJobUpdates.mem_peak && job) {
job.mem_peak = previewJobUpdates.mem_peak
if (job) {
updateJobFromProgress(previewJobUpdates, offset, job)
}
// Check if job is completed
@@ -452,15 +471,8 @@
currentEventSource?.close()
currentEventSource = undefined
if (job?.type === 'CompletedJob') {
isCompleted = true
if (currentId === id) {
await tick()
dispatch('done', job)
callbacks?.done(job)
currentId = undefined
}
}
isCompleted = true
onJobCompleted(id, job, callbacks)
}
} catch (parseErr) {
console.warn('Failed to parse SSE data:', parseErr)
@@ -471,8 +483,14 @@
console.warn('SSE error:', error)
currentEventSource?.close()
currentEventSource = undefined
// Fall back to polling on error
setTimeout(() => syncer(id), 1000)
if (attempt < 3) {
console.log(`SSE error )1), retrying ... attempt: ${attempt}/3`)
attempt++
setTimeout(() => loadTestJobWithSSE(id, attempt, callbacks), 1000)
} else {
// Fall back to polling on error
setTimeout(() => syncer(id), 1000)
}
}
currentEventSource.onopen = () => {
@@ -491,7 +509,15 @@
// Fall back to polling on error
currentEventSource?.close()
currentEventSource = undefined
setTimeout(() => syncer(id), 1000)
if (attempt < 3) {
console.log(`SSE error (2), retrying ... attempt: ${attempt}/3`)
attempt++
loadTestJobWithSSE(id, attempt, callbacks)
} else {
// Fall back to polling on error
setTimeout(() => syncer(id), 1000)
}
}
return isCompleted
} else {
@@ -503,13 +529,12 @@
}
async function syncer(id: string, callbacks?: Callbacks): Promise<void> {
if (currentId != id) {
dispatch('cancel', id)
callbacks?.cancel()
if ((currentId != id && !allowConcurentRequests) || finished.includes(id)) {
callbacks?.cancel?.()
return
}
syncIteration++
if (await loadTestJob(id)) {
if (await loadTestJob(id, callbacks)) {
return
}
let nextIteration = 50
+12 -4
View File
@@ -53,13 +53,21 @@
const val = mod.value
// let jobId: string | undefined = undefined
let callbacks = {
done(x) {
jobDone()
}
}
if (val.type == 'rawscript') {
await jobLoader?.runPreview(
val.path ?? ($pathStore ?? '') + '/' + mod.id,
val.content,
val.language,
mod.id === 'preprocessor' ? { _ENTRYPOINT_OVERRIDE: 'preprocessor', ...args } : args,
flowStore?.val?.tag ?? val.tag
flowStore?.val?.tag ?? val.tag,
undefined,
undefined,
callbacks
)
} else if (val.type == 'script') {
const script = val.hash
@@ -72,10 +80,11 @@
mod.id === 'preprocessor' ? { _ENTRYPOINT_OVERRIDE: 'preprocessor', ...args } : args,
flowStore?.val?.tag ?? (val.tag_override ? val.tag_override : script.tag),
script.lock,
val.hash ?? script.hash
val.hash ?? script.hash,
callbacks
)
} else if (val.type == 'flow') {
await jobLoader?.runFlowByPath(val.path, args)
await jobLoader?.runFlowByPath(val.path, args, callbacks)
} else {
throw Error('Not supported module type')
}
@@ -119,7 +128,6 @@
<JobLoader
noCode={true}
toastError={noEditor}
on:done={() => jobDone()}
bind:scriptProgress
bind:this={jobLoader}
bind:isLoading={
@@ -224,7 +224,14 @@
selectedTab === 'preprocessor' || kind === 'preprocessor'
? { _ENTRYPOINT_OVERRIDE: 'preprocessor', ...(args ?? {}) }
: (args ?? {}),
tag
tag,
undefined,
undefined,
{
done(_x) {
loadPastTests()
}
}
)
logPanel?.setFocusToLogs()
return job
@@ -446,7 +453,6 @@
<JobLoader
noCode={true}
on:done={loadPastTests}
bind:scriptProgress
bind:this={jobLoader}
bind:isLoading={testIsLoading}
@@ -44,7 +44,7 @@
const outputs = initOutput($worldStore, id, {
result: undefined,
loading: false,
jobId: undefined
jobId: undefined as string | undefined
})
initializing = false
@@ -61,13 +61,24 @@
outputs.loading.set(true)
const jobId = resolvedConfig?.['jobId']
if (jobId) {
jobLoader?.watchJob(jobId)
jobLoader?.watchJob(jobId, {
done(x) {
onDone(x)
}
})
}
})
}
})
let result: any = $state(undefined)
function onDone(job: Job & { result?: any }) {
outputs.loading.set(false)
outputs.jobId.set(job.id)
outputs.result.set(job.result)
result = job.result
}
</script>
{#each Object.keys(components['jobiddisplaycomponent'].initialData.configuration) as key (key)}
@@ -95,12 +106,6 @@
bind:this={jobLoader}
bind:isLoading={testIsLoading}
bind:job={testJob}
on:done={(e) => {
outputs.loading.set(false)
outputs.jobId.set(e.detail.id)
outputs.result.set(e.detail.result)
result = e.detail.result
}}
/>
<InitializeComponent {id} />
@@ -60,7 +60,14 @@
outputs.loading.set(true)
const jobId = resolvedConfig?.['jobId']
if (jobId) {
jobLoader?.watchJob(jobId)
let callbacks = {
done(x) {
outputs.loading.set(false)
outputs.jobId.set(x.id)
outputs.result.set(x.result)
}
}
jobLoader?.watchJob(jobId, callbacks)
}
})
}
@@ -92,11 +99,6 @@
bind:this={jobLoader}
bind:isLoading={testIsLoading}
bind:job={testJob}
on:done={(e) => {
outputs.loading.set(false)
outputs.jobId.set(e.detail.id)
outputs.result.set(e.detail.result)
}}
/>
<InitializeComponent {id} />
@@ -30,8 +30,7 @@
let result: any = $state()
function onDone(event: { detail: Job }) {
job = event.detail
function onDone(job: Job) {
result = job['result']
}
@@ -56,7 +55,15 @@
}
})
$effect(() => {
id && jobLoader && untrack(() => jobLoader?.watchJob(id))
id &&
jobLoader &&
untrack(() =>
jobLoader?.watchJob(id, {
done(x) {
onDone(x)
}
})
)
})
$effect(() => {
job?.logs == undefined && job && viewTab == 'logs' && untrack(() => jobLoader?.getLogs())
@@ -68,13 +75,7 @@
let jobLoader: JobLoader | undefined = $state(undefined)
</script>
<JobLoader
lazyLogs
workspaceOverride={workspace}
bind:job={currentJob}
bind:this={jobLoader}
on:done={onDone}
/>
<JobLoader lazyLogs workspaceOverride={workspace} bind:job={currentJob} bind:this={jobLoader} />
<div class="p-4 flex flex-col gap-2 items-start h-full">
{#if job}