From 3be1bb722b1d62f8c81f6cb7c47cf550237bbd19 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 14 Jul 2025 12:41:43 +0000 Subject: [PATCH] all --- backend/windmill-api/src/jobs.rs | 77 ++----------------- .../src/lib/components/TestJobLoader.svelte | 33 ++++---- 2 files changed, 25 insertions(+), 85 deletions(-) diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index bc9e655b9d..625a353f00 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -5742,77 +5742,16 @@ async fn get_job_update( Path((w_id, job_id)): Path<(String, Uuid)>, Query(JobUpdateQuery { log_offset, get_progress, running }): Query, ) -> JsonResult { - let record = sqlx::query!( - "SELECT - c.id IS NOT NULL AS completed, - CASE - WHEN q.id IS NOT NULL THEN (CASE WHEN NOT $5 AND q.running THEN true ELSE null END) - ELSE false - END AS running, - SUBSTR(logs, GREATEST($1 - log_offset, 0)) AS logs, - COALESCE(r.memory_peak, c.memory_peak) AS mem_peak, - CASE - -- flow step: - WHEN flow_step_id IS NOT NULL THEN NULL - -- completed: - WHEN c.id IS NOT NULL THEN COALESCE( - c.workflow_as_code_status || c.flow_status, - c.workflow_as_code_status, - c.flow_status - ) - -- not completed: - ELSE COALESCE( - f.workflow_as_code_status || f.flow_status, - f.workflow_as_code_status, - f.flow_status - ) - END AS \"flow_status: sqlx::types::Json>\", - job_logs.log_offset + CHAR_LENGTH(job_logs.logs) + 1 AS log_offset, - created_by AS \"created_by!\", - CASE WHEN $4::BOOLEAN THEN ( - SELECT scalar_int FROM job_stats WHERE job_id = $3 AND metric_id = 'progress_perc' - ) END AS progress - FROM v2_job j - LEFT JOIN v2_job_queue q USING (id) - LEFT JOIN v2_job_runtime r USING (id) - LEFT JOIN v2_job_status f USING (id) - LEFT JOIN v2_job_completed c USING (id) - LEFT JOIN job_logs ON job_logs.job_id = $3 - WHERE j.workspace_id = $2 AND j.id = $3", - log_offset, - &w_id, - job_id, - get_progress.unwrap_or(false), - running, - ) - .fetch_optional(&db) - .await? - .ok_or_else(|| Error::NotFound(format!("Job not found: {}", job_id)))?; - - if opt_authed.is_none() && record.created_by != "anonymous" { - return Err(Error::BadRequest( - "As a non logged in user, you can only see jobs ran by anonymous users".to_string(), - )); - } - log_job_view( + Ok(Json(get_job_update_data( + &opt_authed, + &opt_tokened, &db, - opt_authed.as_ref(), - opt_tokened.token.as_deref(), &w_id, &job_id, - ) - .await?; - Ok(Json(JobUpdate { - running: record.running, - completed: record.completed, - log_offset: record.log_offset, - new_logs: record.logs, - mem_peak: record.mem_peak, - progress: record.progress, - flow_status: record - .flow_status - .map(|x: sqlx::types::Json>| x.0), - })) + log_offset, + get_progress, + running, + ).await?)) } async fn get_job_update_sse( @@ -5947,7 +5886,7 @@ async fn get_job_update_data( log_offset: i32, get_progress: Option, running: bool, -) -> Result { +) -> error::Result { let record = sqlx::query!( "SELECT c.id IS NOT NULL AS completed, diff --git a/frontend/src/lib/components/TestJobLoader.svelte b/frontend/src/lib/components/TestJobLoader.svelte index dff1db0382..686df46ed5 100644 --- a/frontend/src/lib/components/TestJobLoader.svelte +++ b/frontend/src/lib/components/TestJobLoader.svelte @@ -205,13 +205,14 @@ errorIteration = 0 currentId = testId job = undefined - + // Clean up any existing SSE connection currentEventSource?.close() currentEventSource = undefined - + // Try SSE first, fall back to polling if needed const isCompleted = await loadTestJobWithSSE(testId) + console.log('isCompleted', isCompleted) if (!isCompleted && !currentEventSource) { // If SSE didn't start (job might not be running yet), use polling setTimeout(() => { @@ -326,7 +327,7 @@ if (!job) { job = await JobService.getJob({ workspace: workspace!, id, noLogs: lazyLogs }) } - + // If job is already completed, don't start SSE if (job?.type === 'CompletedJob') { isCompleted = true @@ -341,7 +342,7 @@ // 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 - + // Check if we should get progress updates if (job.job_kind == 'script' || isScriptPreview(job.job_kind)) { if (lastTimeCheckedProgress) { @@ -359,7 +360,7 @@ } const offset = logOffset == 0 ? (job.logs?.length ? job.logs?.length + 1 : 0) : logOffset - + // Build SSE URL with query parameters const params = new URLSearchParams({ running: job.running.toString(), @@ -368,22 +369,22 @@ if (getProgress !== undefined) { params.set('get_progress', getProgress.toString()) } - + const sseUrl = `/api/w/${workspace}/jobs_u/getupdate_sse/${id}?${params.toString()}` - + currentEventSource = new EventSource(sseUrl) - + currentEventSource.onmessage = async (event) => { if (currentId !== id) { currentEventSource?.close() currentEventSource = undefined return } - + try { 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) @@ -410,13 +411,13 @@ if (previewJobUpdates.mem_peak && job) { job.mem_peak = previewJobUpdates.mem_peak } - + // Check if job is completed if (previewJobUpdates.completed) { job = await JobService.getJob({ workspace: workspace!, id }) currentEventSource?.close() currentEventSource = undefined - + if (job?.type === 'CompletedJob') { isCompleted = true if (currentId === id) { @@ -425,14 +426,14 @@ currentId = undefined } } - } else if ((previewJobUpdates.running ?? false)) { + } else if (previewJobUpdates.running ?? false) { job = await JobService.getJob({ workspace: workspace!, id }) } } catch (parseErr) { console.warn('Failed to parse SSE data:', parseErr) } } - + currentEventSource.onerror = (error) => { console.warn('SSE error:', error) currentEventSource?.close() @@ -440,12 +441,12 @@ // Fall back to polling on error setTimeout(() => syncer(id), 1000) } - + currentEventSource.onopen = () => { console.log('SSE connection opened for job:', id) } } - + notfound = false } catch (err) { errorIteration += 1