diff --git a/backend/windmill-api/openapi.yaml b/backend/windmill-api/openapi.yaml index 3564d13541..e955dc37c4 100644 --- a/backend/windmill-api/openapi.yaml +++ b/backend/windmill-api/openapi.yaml @@ -7591,6 +7591,36 @@ paths: flow_status: $ref: "#/components/schemas/WorkflowStatusRecord" + /w/{workspace}/jobs_u/getupdate_sse/{id}: + get: + summary: get job updates via server-sent events + operationId: getJobUpdatesSSE + tags: + - job + parameters: + - $ref: "#/components/parameters/WorkspaceId" + - $ref: "#/components/parameters/JobId" + - name: running + in: query + schema: + type: boolean + - name: log_offset + in: query + schema: + type: integer + - name: get_progress + in: query + schema: + type: boolean + + responses: + "200": + description: server-sent events stream of job updates + content: + text/event-stream: + schema: + type: string + /w/{workspace}/jobs_u/get_log_file/{path}: get: summary: get log file from object store diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index b178b20934..bc9e655b9d 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -272,6 +272,7 @@ pub fn workspace_unauthed_service() -> Router { get(get_completed_job_result_maybe), ) .route("/getupdate/:id", get(get_job_update)) + .route("/getupdate_sse/:id", get(get_job_update_sse)) .route("/get_log_file/*file_path", get(get_log_file)) .route("/queue/cancel/:id", post(cancel_job_api)) .route( @@ -5814,6 +5815,214 @@ async fn get_job_update( })) } +async fn get_job_update_sse( + OptAuthed(opt_authed): OptAuthed, + opt_tokened: OptTokened, + Extension(db): Extension, + Path((w_id, job_id)): Path<(String, Uuid)>, + Query(JobUpdateQuery { log_offset, get_progress, running }): Query, +) -> Response { + let stream = get_job_update_sse_stream( + opt_authed, + opt_tokened, + db, + w_id, + job_id, + log_offset, + get_progress, + running, + ); + + let body = axum::body::Body::from_stream(stream.map(Result::<_, std::convert::Infallible>::Ok)); + + Response::builder() + .status(200) + .header("Content-Type", "text/event-stream") + .header("Cache-Control", "no-cache") + .header("Connection", "keep-alive") + .body(body) + .unwrap() +} + +fn get_job_update_sse_stream( + opt_authed: Option, + opt_tokened: OptTokened, + db: DB, + w_id: String, + job_id: Uuid, + initial_log_offset: i32, + get_progress: Option, + running: bool, +) -> impl futures::Stream { + let (tx, rx) = tokio::sync::mpsc::channel(32); + + tokio::spawn(async move { + let mut log_offset = initial_log_offset; + let mut last_update: Option = None; + let mut completion_sent = false; + + // Send initial update immediately + if let Ok(update) = get_job_update_data( + &opt_authed, + &opt_tokened, + &db, + &w_id, + &job_id, + log_offset, + get_progress, + running, + ).await { + if let Ok(serialized) = serde_json::to_string(&update) { + let event_data = format!("data: {}\n\n", serialized); + if tx.send(event_data.clone()).await.is_err() { + return; + } + last_update = Some(serialized); + if let Some(new_offset) = update.log_offset { + log_offset = new_offset; + } + completion_sent = update.completed.unwrap_or(false); + } + } + + // If job is already completed, no need to poll + if completion_sent { + return; + } + + // Poll for updates every 1 second + let mut interval = tokio::time::interval(std::time::Duration::from_secs(1)); + loop { + interval.tick().await; + + match get_job_update_data( + &opt_authed, + &opt_tokened, + &db, + &w_id, + &job_id, + log_offset, + get_progress, + running, + ).await { + Ok(update) => { + if let Ok(serialized) = serde_json::to_string(&update) { + // Only send if the update has changed + if last_update.as_ref() != Some(&serialized) { + let event_data = format!("data: {}\n\n", serialized); + if tx.send(event_data).await.is_err() { + break; + } + last_update = Some(serialized); + + // Update log offset if available + if let Some(new_offset) = update.log_offset { + log_offset = new_offset; + } + + // Check if job is completed + if update.completed.unwrap_or(false) { + break; + } + } + } + } + Err(_) => { + // Job might have been deleted or access denied, break the loop + break; + } + } + } + }); + + tokio_stream::wrappers::ReceiverStream::new(rx) +} + +async fn get_job_update_data( + opt_authed: &Option, + opt_tokened: &OptTokened, + db: &DB, + w_id: &str, + job_id: &Uuid, + log_offset: i32, + get_progress: Option, + running: bool, +) -> Result { + 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( + db, + opt_authed.as_ref(), + opt_tokened.token.as_deref(), + w_id, + job_id, + ) + .await?; + + Ok(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), + }) +} + pub fn filter_list_completed_query( mut sqlb: SqlBuilder, lq: &ListCompletedQuery, diff --git a/frontend/src/lib/components/TestJobLoader.svelte b/frontend/src/lib/components/TestJobLoader.svelte index 28990e7a53..dff1db0382 100644 --- a/frontend/src/lib/components/TestJobLoader.svelte +++ b/frontend/src/lib/components/TestJobLoader.svelte @@ -177,6 +177,9 @@ if (id) { dispatch('cancel', id) currentId = undefined + // Clean up SSE connection + currentEventSource?.close() + currentEventSource = undefined try { await JobService.cancelQueuedJob({ workspace: $workspaceStore ?? '', @@ -202,8 +205,15 @@ errorIteration = 0 currentId = testId job = undefined - const isCompleted = await loadTestJob(testId) - if (!isCompleted) { + + // Clean up any existing SSE connection + currentEventSource?.close() + currentEventSource = undefined + + // Try SSE first, fall back to polling if needed + const isCompleted = await loadTestJobWithSSE(testId) + if (!isCompleted && !currentEventSource) { + // If SSE didn't start (job might not be running yet), use polling setTimeout(() => { syncer(testId) }, 50) @@ -306,6 +316,158 @@ } } + let currentEventSource: EventSource | undefined = undefined + + async function loadTestJobWithSSE(id: string): Promise { + let isCompleted = false + if (currentId === id) { + try { + // First load the job to get initial state + 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 + if (currentId === id) { + await tick() + dispatch('done', job) + currentId = undefined + } + return isCompleted + } + + // 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) { + const lastTimeCheckedMs = Date.now() - lastTimeCheckedProgress + if ( + lastTimeCheckedMs > getProgressRetryRate || + (scriptProgress != undefined && lastTimeCheckedMs > getProgressRate) + ) { + lastTimeCheckedProgress = Date.now() + getProgress = true + } + } else { + lastTimeCheckedProgress = Date.now() + } + } + + 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(), + log_offset: offset.toString() + }) + 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) + + 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 + } + + // 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) { + await tick() + dispatch('done', job) + currentId = undefined + } + } + } 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() + currentEventSource = undefined + // 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 + if (errorIteration == 5) { + notfound = true + job = undefined + } + console.warn(err) + // Fall back to polling on error + currentEventSource?.close() + currentEventSource = undefined + setTimeout(() => syncer(id), 1000) + } + return isCompleted + } else { + // Clean up SSE connection if current ID changed + currentEventSource?.close() + currentEventSource = undefined + return true + } + } + async function syncer(id: string): Promise { if (currentId != id) { dispatch('cancel', id) @@ -326,6 +488,8 @@ onDestroy(async () => { currentId = undefined + currentEventSource?.close() + currentEventSource = undefined })