fix: restore set_progress feature with sse
This commit is contained in:
@@ -6702,7 +6702,7 @@ async fn get_job_update(
|
||||
&job_id,
|
||||
log_offset,
|
||||
stream_offset,
|
||||
get_progress,
|
||||
get_progress.unwrap_or(false),
|
||||
running,
|
||||
true,
|
||||
false,
|
||||
@@ -6806,6 +6806,7 @@ fn start_job_update_sse_stream(
|
||||
// Send initial update immediately
|
||||
let mut running = running;
|
||||
let mut mem_peak = 0;
|
||||
|
||||
match get_job_update_data(
|
||||
&opt_authed,
|
||||
&opt_tokened,
|
||||
@@ -6814,7 +6815,7 @@ fn start_job_update_sse_stream(
|
||||
&job_id,
|
||||
log_offset,
|
||||
stream_offset,
|
||||
get_progress,
|
||||
false,
|
||||
running,
|
||||
true,
|
||||
true,
|
||||
@@ -6871,11 +6872,12 @@ fn start_job_update_sse_stream(
|
||||
}
|
||||
}
|
||||
|
||||
let mut get_progress_m: bool = false;
|
||||
// Poll for updates every 1 second
|
||||
let mut i = 0;
|
||||
let start = Instant::now();
|
||||
let mut last_ping = Instant::now();
|
||||
|
||||
let mut last_progress_check = Instant::now();
|
||||
loop {
|
||||
i += 1;
|
||||
|
||||
@@ -6918,6 +6920,10 @@ fn start_job_update_sse_stream(
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(ms_duration)).await;
|
||||
|
||||
// Check progress if the user requested it, and check periodically if the job has progress
|
||||
// Once it has progress, we always check progress
|
||||
let check_progress = get_progress.unwrap_or(false)
|
||||
&& (get_progress_m || last_progress_check.elapsed().as_secs() > 5);
|
||||
match get_job_update_data(
|
||||
&opt_authed,
|
||||
&opt_tokened,
|
||||
@@ -6926,7 +6932,7 @@ fn start_job_update_sse_stream(
|
||||
&job_id,
|
||||
log_offset,
|
||||
stream_offset,
|
||||
get_progress,
|
||||
check_progress,
|
||||
running,
|
||||
false,
|
||||
true,
|
||||
@@ -6947,6 +6953,13 @@ fn start_job_update_sse_stream(
|
||||
if update.new_logs.as_ref().is_some_and(|x| x.is_empty()) {
|
||||
update.new_logs = None;
|
||||
}
|
||||
if check_progress {
|
||||
if update.progress.is_some() {
|
||||
get_progress_m = true;
|
||||
} else {
|
||||
last_progress_check = Instant::now();
|
||||
}
|
||||
}
|
||||
|
||||
// if !only_result.unwrap_or(false) {
|
||||
// tracing::error!("update {:?}", update);
|
||||
@@ -7044,7 +7057,7 @@ async fn get_job_update_data(
|
||||
job_id: &Uuid,
|
||||
log_offset: Option<i32>,
|
||||
stream_offset: Option<i32>,
|
||||
get_progress: Option<bool>,
|
||||
get_progress: bool,
|
||||
running: Option<bool>,
|
||||
log_view: bool,
|
||||
get_full_job_on_completion: bool,
|
||||
@@ -7281,7 +7294,7 @@ async fn get_job_update_data(
|
||||
log_offset,
|
||||
w_id,
|
||||
job_id,
|
||||
get_progress.unwrap_or(false),
|
||||
get_progress,
|
||||
running,
|
||||
tags.as_ref().map(|v| v.as_slice()) as Option<&[&str]>,
|
||||
no_logs.unwrap_or(false),
|
||||
|
||||
@@ -14,14 +14,18 @@
|
||||
import { deepEqual } from 'fast-equals'
|
||||
import { isWindmillTooBigObject } from './job_args'
|
||||
|
||||
export let id: string | undefined = undefined
|
||||
export let args: any
|
||||
export let argLabel: string | undefined = undefined
|
||||
export let workspace: string | undefined = undefined
|
||||
interface Props {
|
||||
id?: string | undefined
|
||||
args: any
|
||||
argLabel?: string | undefined
|
||||
workspace?: string | undefined
|
||||
}
|
||||
|
||||
let jsonViewer: Drawer
|
||||
let runLocally: Drawer
|
||||
let jsonStr = ''
|
||||
let { id = undefined, args, argLabel = undefined, workspace = undefined }: Props = $props()
|
||||
|
||||
let jsonViewer: Drawer | undefined = $state()
|
||||
let runLocally: Drawer | undefined = $state()
|
||||
let jsonStr = $state('')
|
||||
|
||||
function pythonCode() {
|
||||
return `
|
||||
@@ -53,9 +57,9 @@ ${Object.entries(args)
|
||||
}
|
||||
</script>
|
||||
|
||||
{#if args && typeof args === 'object' && deepEqual( Object.keys(args), ['reason'] ) && args['reason'] == 'PREPROCESSOR_ARGS_ARE_DISCARDED'}
|
||||
{#if args && typeof args === 'object' && deepEqual( Object.keys(args ?? {}), ['reason'] ) && args['reason'] == 'PREPROCESSOR_ARGS_ARE_DISCARDED'}
|
||||
Preprocessor args are discarded
|
||||
{:else if id && workspace && args && typeof args === 'object' && deepEqual( Object.keys(args), ['reason'] ) && args['reason'] == 'WINDMILL_TOO_BIG'}
|
||||
{:else if id && workspace && args && typeof args === 'object' && deepEqual( Object.keys(args ?? {}), ['reason'] ) && args['reason'] == 'WINDMILL_TOO_BIG'}
|
||||
The args are too big in size to be able to fetch alongside job. Please <a
|
||||
href="/api/w/{workspace}/jobs_u/get_args/{id}"
|
||||
target="_blank">download the JSON file to view them</a
|
||||
@@ -68,21 +72,21 @@ ${Object.entries(args)
|
||||
<Cell head first>{argLabel ?? 'Arg'}</Cell>
|
||||
<Cell head last>Value</Cell>
|
||||
</tr>
|
||||
<svelte:fragment slot="headerAction">
|
||||
{#snippet headerAction()}
|
||||
<button
|
||||
on:click={() => {
|
||||
onclick={() => {
|
||||
jsonStr = JSON.stringify(args, null, 4)
|
||||
jsonViewer.openDrawer()
|
||||
jsonViewer?.openDrawer()
|
||||
}}
|
||||
>
|
||||
<Expand size={18} />
|
||||
</button>
|
||||
</svelte:fragment>
|
||||
{/snippet}
|
||||
</Head>
|
||||
|
||||
<tbody class="divide-y w-full">
|
||||
{#if args && typeof args === 'object' && Object.keys(args).length > 0}
|
||||
{#each Object.entries(args).sort((a, b) => a[0].localeCompare(b[0])) as [arg, value]}
|
||||
{#if args && typeof args === 'object' && Object.keys(args ?? {}).length > 0}
|
||||
{#each Object.entries(args ?? {}).sort( (a, b) => a?.[0]?.localeCompare(b?.[0]) ) as [arg, value]}
|
||||
<Row>
|
||||
<Cell first>{arg}</Cell>
|
||||
<Cell><ArgInfo {value} /></Cell>
|
||||
@@ -124,7 +128,7 @@ ${Object.entries(args)
|
||||
Download
|
||||
</Button>
|
||||
<Button
|
||||
on:click={runLocally.openDrawer}
|
||||
on:click={() => runLocally?.openDrawer()}
|
||||
color="light"
|
||||
size="xs"
|
||||
startIcon={{ icon: ChevronRightSquare }}
|
||||
|
||||
@@ -68,15 +68,6 @@
|
||||
children
|
||||
}: Props = $props()
|
||||
|
||||
/// Last time asked for job progress
|
||||
let lastTimeCheckedProgress: number | undefined = undefined
|
||||
|
||||
/// Will try to poll progress every 5s and if once progress returned was not undefined, will be ignored
|
||||
/// and getProgressRate will be used instead
|
||||
const getProgressRetryRate: number = 5000
|
||||
/// How often loader poll progress
|
||||
const getProgressRate: number = 1000
|
||||
|
||||
let workspace = $derived(workspaceOverride ?? $workspaceStore)
|
||||
|
||||
let syncIteration: number = 0
|
||||
@@ -146,6 +137,7 @@
|
||||
export async function abstractRun(fn: () => Promise<string>, callbacks?: Callbacks) {
|
||||
try {
|
||||
isLoading = true
|
||||
scriptProgress = undefined
|
||||
lastCompletedJobId = undefined
|
||||
clearCurrentJob()
|
||||
lastCallbacks = callbacks
|
||||
@@ -298,10 +290,6 @@
|
||||
hash?: string,
|
||||
callbacks?: Callbacks
|
||||
): Promise<string> {
|
||||
// Reset in case we rerun job without reloading
|
||||
scriptProgress = undefined
|
||||
lastTimeCheckedProgress = undefined
|
||||
|
||||
return abstractRun(
|
||||
() =>
|
||||
JobService.runScriptPreview({
|
||||
@@ -363,6 +351,7 @@
|
||||
syncIteration = 0
|
||||
errorIteration = 0
|
||||
currentId = testId
|
||||
scriptProgress = undefined
|
||||
if (loadPlaceholderJobOnStart) {
|
||||
job = structuredClone(loadPlaceholderJobOnStart)
|
||||
} else {
|
||||
@@ -382,33 +371,6 @@
|
||||
}
|
||||
}
|
||||
|
||||
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(
|
||||
@@ -424,6 +386,7 @@
|
||||
}
|
||||
}
|
||||
if (previewJobUpdates.progress) {
|
||||
console.log('progress', previewJobUpdates.progress)
|
||||
// Progress cannot go back and cannot be set to 100
|
||||
scriptProgress = clamp(previewJobUpdates.progress, scriptProgress ?? 0, 99)
|
||||
}
|
||||
@@ -484,7 +447,6 @@
|
||||
try {
|
||||
if (job && `running` in job) {
|
||||
callbacks?.running?.({ id })
|
||||
let getProgress: boolean | undefined = setJobProgress(job)
|
||||
|
||||
refreshLogOffset()
|
||||
|
||||
@@ -494,7 +456,7 @@
|
||||
running: job.running,
|
||||
logOffset: logOffset,
|
||||
streamOffset: resultStreamOffset,
|
||||
getProgress: getProgress
|
||||
getProgress: false
|
||||
})
|
||||
|
||||
if ((previewJobUpdates.running ?? false) || (previewJobUpdates.completed ?? false)) {
|
||||
@@ -640,7 +602,9 @@
|
||||
}
|
||||
|
||||
let getProgress: boolean | undefined =
|
||||
onlyResult || !job ? undefined : setJobProgress(job)
|
||||
onlyResult || !job
|
||||
? undefined
|
||||
: job.job_kind == 'script' || isScriptPreview(job.job_kind)
|
||||
|
||||
refreshLogOffset()
|
||||
// Build SSE URL with query parameters
|
||||
|
||||
@@ -168,7 +168,6 @@
|
||||
reloadError = undefined
|
||||
try {
|
||||
const { input_transforms, schema } = await loadSchemaFromModule(flowModule)
|
||||
console.log('reload', schema)
|
||||
validCode = true
|
||||
|
||||
if (inputTransformSchemaForm) {
|
||||
|
||||
@@ -14,10 +14,10 @@
|
||||
|
||||
let {
|
||||
job = undefined,
|
||||
compact = $bindable(false),
|
||||
scriptProgress = $bindable(undefined),
|
||||
hideStepTitle = $bindable(false),
|
||||
class: className = $bindable('')
|
||||
compact = false,
|
||||
scriptProgress = undefined,
|
||||
hideStepTitle = false,
|
||||
class: className = ''
|
||||
}: Props = $props()
|
||||
|
||||
let error: number | undefined = $state(undefined)
|
||||
@@ -28,6 +28,8 @@
|
||||
let nextInProgress = false
|
||||
|
||||
let progressBar: ProgressBar | undefined = $state(undefined)
|
||||
let lastJobId = $state()
|
||||
|
||||
function updateJobProgress(job: Job) {
|
||||
if (!job['running'] && !job['success']) {
|
||||
error = 0
|
||||
@@ -49,6 +51,13 @@
|
||||
scriptProgress = undefined
|
||||
}
|
||||
|
||||
$effect(() => {
|
||||
if (lastJobId && job && job.id !== lastJobId) {
|
||||
lastJobId = job.id
|
||||
reset()
|
||||
}
|
||||
})
|
||||
|
||||
$effect(() => {
|
||||
if (job) updateJobProgress(job)
|
||||
})
|
||||
|
||||
@@ -222,7 +222,11 @@
|
||||
async function onJobLoaded() {
|
||||
// We want to set up scriptProgress once job is loaded
|
||||
// We need this to show progress bar if job has progress and is finished
|
||||
if (job && job.type == 'CompletedJob') {
|
||||
if (
|
||||
job &&
|
||||
job.type == 'CompletedJob' &&
|
||||
(job.job_kind == 'script' || isScriptPreview(job.job_kind))
|
||||
) {
|
||||
// If error occured and job is completed
|
||||
// than we fetch progress from server to display on what progress did it fail
|
||||
// Could be displayed after run or as a historical page
|
||||
|
||||
Reference in New Issue
Block a user