feat: implement SSE for job updates polling

- Add /getupdate_sse/:id endpoint for real-time job updates via Server-Sent Events
- SSE streams job status, logs, progress, and flow status updates
- Auto-stops streaming when job completes
- Frontend uses EventSource with graceful fallback to polling on errors
- Reduces server load and improves real-time responsiveness
- Update OpenAPI spec with new SSE endpoint definition

🤖 Generated with [Claude Code](https://claude.ai/code)

Co-authored-by: Ruben Fiszel <rubenfiszel@users.noreply.github.com>
This commit is contained in:
claude[bot]
2025-07-13 23:23:04 +00:00
parent ce442a7493
commit fba3145d07
3 changed files with 405 additions and 2 deletions
+30
View File
@@ -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
+209
View File
@@ -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<DB>,
Path((w_id, job_id)): Path<(String, Uuid)>,
Query(JobUpdateQuery { log_offset, get_progress, running }): Query<JobUpdateQuery>,
) -> 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<ApiAuthed>,
opt_tokened: OptTokened,
db: DB,
w_id: String,
job_id: Uuid,
initial_log_offset: i32,
get_progress: Option<bool>,
running: bool,
) -> impl futures::Stream<Item = String> {
let (tx, rx) = tokio::sync::mpsc::channel(32);
tokio::spawn(async move {
let mut log_offset = initial_log_offset;
let mut last_update: Option<String> = 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<ApiAuthed>,
opt_tokened: &OptTokened,
db: &DB,
w_id: &str,
job_id: &Uuid,
log_offset: i32,
get_progress: Option<bool>,
running: bool,
) -> Result<JobUpdate, Error> {
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<Box<RawValue>>\",
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<Box<RawValue>>| x.0),
})
}
pub fn filter_list_completed_query(
mut sqlb: SqlBuilder,
lq: &ListCompletedQuery,
@@ -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<boolean> {
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<void> {
if (currentId != id) {
dispatch('cancel', id)
@@ -326,6 +488,8 @@
onDestroy(async () => {
currentId = undefined
currentEventSource?.close()
currentEventSource = undefined
})
</script>